From 2f614eb6c7c2b460df3f5f766da79e8800d14ec9 Mon Sep 17 00:00:00 2001 From: Jim Huang Date: Sun, 4 Oct 2026 06:47:01 +0800 Subject: [PATCH] Hold the room until the report is received A completed interview could end with no report: the agent published the report packet, slept 250 ms and left, and a successful publish only queues the packet, so leaving could drop it and the page fell back to "the interviewer never returned a report". Every report the agent sends, provisional, regenerated or final, now republishes the same bytes up to three times, five seconds each, until the page answers with a receipt naming their SHA-256 digest, and the agent leaves only after that. No retry calls the report model again, and the Live session closes beside the first delivery rather than after it. A candidate who drops before acknowledging the provisional report gets it again on rejoining, and a recovery wait that begins with them already gone starts the rejoin grace instead of holding the room for the whole window. The page renders the first copy at once, ignores retransmissions, and keeps the room, Done and its fallback exits until its receipt has had its chance to leave. Its escape wait grows to 155 s and its wait for a regenerated report to 145 s, so both cover delivery. Close #133 --- docs/provider-cost-and-degradation.md | 32 +- scripts/gen-wire-fixtures.mjs | 13 + src/gemini.rs | 2 +- src/livekit.rs | 4 +- src/livekit/report.rs | 285 ++++++++++-- tests/browser/account.test.js | 10 +- tests/browser/dom-contract.test.js | 15 +- tests/browser/lib.test.js | 237 ++++++++++ tests/browser/report-recovery.test.js | 16 +- tests/fixtures/report-receipt.json | 12 + tests/fixtures/report-recovery.json | 2 +- tests/unit/livekit.rs | 46 +- tests/unit/livekit/report.rs | 629 +++++++++++++++++++++++++- web/interview.js | 104 +++-- web/lib.js | 74 ++- web/report-recovery.js | 6 +- 16 files changed, 1388 insertions(+), 99 deletions(-) create mode 100644 tests/fixtures/report-receipt.json diff --git a/docs/provider-cost-and-degradation.md b/docs/provider-cost-and-degradation.md index baf40d1c..4a9d5663 100644 --- a/docs/provider-cost-and-degradation.md +++ b/docs/provider-cost-and-degradation.md @@ -105,20 +105,32 @@ exhausted quota rotation also qualifies, including when another interview exhausted it. The offer expires after five minutes and the agent leaving ends it. A candidate who leaves has 30 seconds to rejoin under the same identity, which is how a full LiveKit reconnect looks, and a regeneration under way keeps -running through it and is published once they are back; a candidate already -gone when the interview ends gets the failure at once with no window. The -Live session is closed after the report is published, or right after the -provisional report when a window opens. While the offer is open the agent +running through it and is published once they are back. One who rejoins +before acknowledging the provisional report gets it again, since the drop may +have lost it; a candidate already gone when the interview ends gets the failure +at once with no window. The +Live session is closed while the report, or the provisional report when a +window opens, is being delivered. While the offer is open the agent answers each request on the control topic, `report_retry` with status `accepted` or `early` (with the seconds still to wait, which does not spend the regeneration), and announces `closed` at expiry or when no key can ever answer; a duplicate request during -regeneration gets no answer. The browser waits at most 140 seconds from its +regeneration gets no answer. The browser waits at most 145 seconds from its request or the acceptance, whichever came last, so neither a reconnect nor a lost answer cuts off a regeneration still in progress. A transient burst followed by a schema failure offers no regeneration: the terminal failure determines eligibility. Reloads and process restarts cannot recover the inputs. +Every report packet, provisional, regenerated or final, is delivered the same +way. The agent publishes those same bytes at most three times, each with a +five-second publication and receipt window, without another Gemini call. The +Live session closes beside the first delivery, so retries spend no Live time. +Each delivery keeps the room open for at most fifteen extra seconds and stops +waiting when the candidate leaves or the room disconnects. A candidate receipt +names the SHA-256 digest of the packet bytes. Missing receipts mean delivery is +unconfirmed, not that the candidate received no report; only publication +attempts that all fail or time out produce `report_delivery_failed`. + Quiet-pause interim reviews use that same report model and quota. A review is eligible after 8 seconds of candidate quiet and 150 seconds from interview start, no more often than every 150 seconds, and only after six new candidate turns; @@ -190,10 +202,12 @@ Watch `codetrial dispatch_refused ... reason=at_capacity` and `reason=finalizing`, `livekit quota:` transitions, token HTTP 429 with `Retry-After`, `gemini report transport_failed call=... retry=...`, `gemini report retry_unavailable` (a key rotation lost during backoff, without another -HTTP call), `codetrial report_recovery_notice_failed`, `codetrial live_usage -... outcome=billing`, and the bounded incomplete-report categories. Each -recovery window also keeps the agent and the candidate connected to LiveKit for -up to about seven minutes after the interview, which counts against +HTTP call), `codetrial report_recovery_notice_failed`, `codetrial +report_publish_failed`, `codetrial report_receipt_missing`, `codetrial +report_delivery ... outcome=acknowledged|unconfirmed|failed`, `codetrial +live_usage ... outcome=billing`, and the bounded incomplete-report categories. +Each recovery window also keeps the agent and the candidate connected to +LiveKit for up to about seven minutes after the interview, which counts against connection-minute quota like interview time. Raise concurrency only after checking provider minutes, Gemini limits, CPU and diff --git a/scripts/gen-wire-fixtures.mjs b/scripts/gen-wire-fixtures.mjs index f8e339fd..22829541 100755 --- a/scripts/gen-wire-fixtures.mjs +++ b/scripts/gen-wire-fixtures.mjs @@ -384,6 +384,19 @@ const files = { cases: codeUpdateCases(languages), }, "control.json": { topic: lib.topics.control, cases: controlCases() }, + "report-receipt.json": { + topic: lib.topics.control, + cases: [ + { + name: "received report", + // The page's own digest, so the agent test that rehashes these + // bytes checks the browser's hashing and not a copy of it. + payload: await lib.reportReceipt( + new TextEncoder().encode('{"codingScore":80}'), + ), + }, + ], + }, "test-results.json": { topic: lib.topics.tests, cases: testResultsCases() }, "integrity-chain.json": await integrityChain(), }; diff --git a/src/gemini.rs b/src/gemini.rs index 1f7077b2..abb01910 100644 --- a/src/gemini.rs +++ b/src/gemini.rs @@ -44,7 +44,7 @@ const WRITE_TIMEOUT: Duration = Duration::from_secs(5); /// Bounds the close handshake the same way. The socket being closed is most /// often the one that stopped answering, and a close that waits on it held up /// the reconnect and the agent's exit behind a peer that was already gone. -const CLOSE_TIMEOUT: Duration = Duration::from_secs(2); +pub(crate) const CLOSE_TIMEOUT: Duration = Duration::from_secs(2); /// How often the room loop pings. A socket nobody is writing to cannot fail a /// write, and a paused interview writes nothing, so without a ping a peer that /// went away there would not be noticed until something was finally said. diff --git a/src/livekit.rs b/src/livekit.rs index 4798e12d..b21fc496 100644 --- a/src/livekit.rs +++ b/src/livekit.rs @@ -3014,8 +3014,8 @@ async fn handle_data_packet( ) .await?; - // Give the report packet a moment to leave before the agent goes. - tokio::time::sleep(Duration::from_millis(250)).await; + // Every report went out through a receipt wait, so leaving now drops + // nothing the candidate has not acknowledged or stopped listening for. leave_room(room).await; Ok(ControlFlow::Break(())) } diff --git a/src/livekit/report.rs b/src/livekit/report.rs index a4736879..f59842d9 100644 --- a/src/livekit/report.rs +++ b/src/livekit/report.rs @@ -10,7 +10,8 @@ //! //! Split out of `livekit.rs` along the line its module doc already drew. -use ::livekit::prelude::{DataPacket, Room}; +use ::livekit::prelude::{DataPacket, Room, RoomEvent}; +use std::time::Duration; use crate::agent::{ ModelInputKind, ReportPromptInput, RuntimeState, final_report, format_test_run, @@ -91,6 +92,177 @@ pub(super) async fn generate_report_bounded( .await } +pub(super) const DELIVERY_ATTEMPTS: usize = 3; +pub(super) const DELIVERY_WAIT: Duration = Duration::from_secs(5); + +fn is_report_receipt( + topic: Option<&str>, + sender: Option<&str>, + candidate: &str, + payload: &[u8], + id: &str, + candidate_gone: bool, +) -> bool { + // LiveKit resolves the sender against its current roster, so a receipt that + // lands after the candidate's departure arrives with no participant. The + // report is broadcast, so any peer could echo its digest; an unattributed + // receipt counts only once the candidate is gone and no retry could reach + // them anyway. + topic == Some(crate::runtime::TOPIC_CONTROL) + && (sender.is_none() && candidate_gone + || super::is_interview_participant(sender, candidate)) + && serde_json::from_slice::(payload) + .is_ok_and(|value| value["type"] == "report_received" && value["deliveryId"] == id) +} + +pub(super) trait DeliveryEvent { + fn acknowledges(&self, candidate: &str, id: &str, candidate_gone: bool) -> bool; + fn disconnected(&self) -> bool { + false + } +} + +impl DeliveryEvent for RoomEvent { + fn acknowledges(&self, candidate: &str, id: &str, candidate_gone: bool) -> bool { + if let Self::DataReceived { + payload, + topic, + participant, + .. + } = self + { + let sender = participant + .as_ref() + .map(|participant| participant.identity().0); + is_report_receipt( + topic.as_deref(), + sender.as_deref(), + candidate, + payload, + id, + candidate_gone, + ) + } else { + false + } + } + fn disconnected(&self) -> bool { + matches!(self, Self::Disconnected { .. }) + } +} + +/// A successful SDK publish only queues the packet. Keep the room alive until +/// the candidate acknowledges it, retransmitting the immutable bytes on loss. +/// Neither retries nor receipt timeouts call the report model again. +pub(super) async fn deliver_report( + mut publish: F, + events: &mut tokio::sync::mpsc::UnboundedReceiver, + candidate: &str, + room_name: &str, + candidate_present: impl Fn() -> bool, + packet: DataPacket, +) -> Result> +where + F: FnMut(DataPacket) -> Fut, + Fut: std::future::Future>, + Event: DeliveryEvent, +{ + let id = crate::sha256_hex(&[&packet.payload]); + let mut published = false; + for attempt in 1..=DELIVERY_ATTEMPTS { + let deadline = tokio::time::Instant::now() + DELIVERY_WAIT; + let failure = match tokio::time::timeout_at(deadline, publish(packet.clone())).await { + Ok(Ok(())) => { + published = true; + None + } + Ok(Err(_)) => Some("publish_error"), + Err(_) => Some("timeout"), + }; + if let Some(cause) = failure { + eprintln!( + "codetrial report_publish_failed room={room_name} attempt={attempt} cause={cause}" + ); + } + + // A failed publish waits out the rest of its window. The wait also + // accepts a late receipt for an earlier attempt and observes + // departures. + let gone = |events: &mut tokio::sync::mpsc::UnboundedReceiver| { + if queued_receipt(events, candidate, &id) { + Wait::Acknowledged + } else { + Wait::Gone + } + }; + let wait = if candidate_present() { + tokio::time::timeout_at(deadline, async { + while let Some(event) = events.recv().await { + // Presence first: the receipt that arrives unattributed + // because its sender just left is this very event. + let candidate_gone = event.disconnected() || !candidate_present(); + if event.acknowledges(candidate, &id, candidate_gone) { + return Wait::Acknowledged; + } + if candidate_gone { + return gone(events); + } + } + Wait::Gone + }) + .await + .unwrap_or(Wait::TimedOut) + } else { + gone(events) + }; + match wait { + Wait::Acknowledged => { + eprintln!( + "codetrial report_delivery room={room_name} attempt={attempt} outcome=acknowledged" + ); + return Ok(true); + } + Wait::Gone => break, + Wait::TimedOut => (), + } + if published { + eprintln!("codetrial report_receipt_missing room={room_name} attempt={attempt}"); + } + } + if published { + eprintln!("codetrial report_delivery room={room_name} outcome=unconfirmed"); + Ok(false) + } else { + eprintln!("codetrial report_delivery room={room_name} outcome=failed"); + Err("report_delivery_failed: all publication attempts failed or timed out".into()) + } +} + +/// How one attempt's receipt wait ended. `Gone` covers the candidate leaving, +/// the room disconnecting and the event stream closing: nobody is left to +/// retransmit to. +enum Wait { + Acknowledged, + Gone, + TimedOut, +} + +/// The roster drops the candidate when LiveKit handles the departure, not when +/// this loop reads the queue, so a receipt the page sent just before leaving +/// can still be waiting behind events the loop has not reached. +fn queued_receipt( + events: &mut tokio::sync::mpsc::UnboundedReceiver, + candidate: &str, + id: &str, +) -> bool { + while let Ok(event) = events.try_recv() { + if event.acknowledges(candidate, id, true) { + return true; + } + } + false +} + const REPORT_RECOVERY_WINDOW: std::time::Duration = std::time::Duration::from_secs(300); const REPORT_RETRY_COOLDOWN: std::time::Duration = std::time::Duration::from_secs(30); @@ -204,14 +376,17 @@ trait RecoveryRoom { fn candidate_present(&self) -> bool; fn next(&mut self) -> impl std::future::Future + Send; fn notify(&mut self, notice: RecoveryNotice) -> impl std::future::Future + Send; + /// `Ok(true)` once the page acknowledged the packet, `Ok(false)` when it + /// went out unconfirmed. fn publish( &mut self, packet: DataPacket, - ) -> impl std::future::Future>> + Send; + ) -> impl std::future::Future>> + Send; } struct LiveRecoveryRoom<'a> { room: &'a Room, + room_name: &'a str, candidate: &'a str, events: &'a mut tokio::sync::mpsc::UnboundedReceiver<::livekit::RoomEvent>, } @@ -258,8 +433,9 @@ impl RecoveryRoom for LiveRecoveryRoom<'_> { fn candidate_present(&self) -> bool { self.room .remote_participants() - .values() - .any(|participant| participant.identity().0 == self.candidate) + .contains_key(&::livekit::id::ParticipantIdentity( + self.candidate.to_string(), + )) } async fn next(&mut self) -> RecoveryEvent { @@ -268,22 +444,47 @@ impl RecoveryRoom for LiveRecoveryRoom<'_> { async fn notify(&mut self, notice: RecoveryNotice) { // A notice that cannot be sent leaves the page on its own fallback - // timers, which is no worse than before notices existed. - let sent = match browser_packet(crate::runtime::TOPIC_CONTROL, &recovery_notice(notice)) { - Ok(packet) => self.publish(packet).await, - Err(error) => Err(error.into()), - }; + // timers, which is no worse than before notices existed, so it is + // published once and never waits for a receipt. + let sent: Result<(), Box> = + match browser_packet(crate::runtime::TOPIC_CONTROL, &recovery_notice(notice)) { + Ok(packet) => self + .room + .local_participant() + .publish_data(packet) + .await + .map_err(Into::into), + Err(error) => Err(error.into()), + }; if let Err(error) = sent { eprintln!("codetrial report_recovery_notice_failed notice={notice:?} error={error}"); } } + /// Reports only. Retry requests and departures that arrive during the + /// receipt wait are consumed by it; `recover_report` reads presence from + /// the room again when its own wait begins and after each republish. async fn publish( &mut self, packet: DataPacket, - ) -> Result<(), Box> { - self.room.local_participant().publish_data(packet).await?; - Ok(()) + ) -> Result> { + let Self { + room, + room_name, + candidate, + events, + } = self; + let participant = room.local_participant(); + let identity = ::livekit::id::ParticipantIdentity(candidate.to_string()); + deliver_report( + |packet| participant.publish_data(packet), + events, + candidate, + room_name, + || room.remote_participants().contains_key(&identity), + packet, + ) + .await } } @@ -316,15 +517,25 @@ fn track_presence( std::ops::ControlFlow::Continue(()) } +/// `unconfirmed` is the provisional report when no receipt came back for it. +/// The candidate may have dropped before it arrived, so a rejoin republishes +/// it; without that the rejoined page never learns a retry is on offer. async fn recover_report( events: &mut impl RecoveryRoom, + mut unconfirmed: Option, generation: impl std::future::Future, window: std::time::Duration, cooldown: std::time::Duration, started: tokio::time::Instant, readiness: impl Fn() -> Readiness, ) -> Option { + // A departure the provisional report's receipt wait consumed is not + // replayed, so absence starts the rejoin grace here instead of the wait + // holding the slot for a retry from nobody. let mut presence = super::CandidatePresence::default(); + if !events.candidate_present() { + presence.left(tokio::time::Instant::now().into_std()); + } loop { let event = tokio::select! { biased; @@ -338,9 +549,28 @@ async fn recover_report( if track_presence(event, &mut presence).is_break() { return None; } + if event == RecoveryEvent::Back + && let Some(packet) = unconfirmed.take() + { + match events.publish(packet.clone()).await { + Ok(true) => {} + Ok(false) => unconfirmed = Some(packet), + Err(error) => { + eprintln!("codetrial report_republish_failed error={error}"); + unconfirmed = Some(packet); + } + } + // The delivery read the events a departure would have arrived on. + if !events.candidate_present() { + presence.left(tokio::time::Instant::now().into_std()); + } + continue; + } let RecoveryEvent::Retry = event else { continue; }; + // Only a page holding the provisional report can ask for a retry. + unconfirmed = None; let waited = started.elapsed(); if waited < cooldown { events @@ -411,6 +641,7 @@ pub(super) async fn publish_with_recovery( let behavioral_round_opened = crate::agent::BehavioralRound::of(&state).opened(); let mut room = LiveRecoveryRoom { room, + room_name: boot.room_name, candidate, events, }; @@ -458,29 +689,31 @@ async fn run_recovery( keys, } = report; - // The Live session is closed once, after the report it would otherwise - // delay by up to its close timeout, and before a recovery wait that must - // not hold it open for minutes. + // The Live session is closed once, beside the first report's delivery: + // waiting for it would delay the report by up to its close timeout, and + // closing after it would hold Live open through the receipt retries and a + // recovery wait of minutes. let Some(cooldown) = regeneration_cooldown(&generated, keys).filter(|_| room.candidate_present()) else { - let published: Result<(), Box> = async { - room.publish(report_packet(boot, state, reason, keys, generated)?) - .await - } - .await; - close_live.await; - return published; + let (published, ()) = tokio::join!( + async { + room.publish(report_packet(boot, state, reason, keys, generated)?) + .await + }, + close_live + ); + return published.map(|_| ()); }; let mut provisional = report_value(boot, state, reason, keys, generated); provisional["reportRecovery"] = recovery_metadata(cooldown); + let provisional = report_data_packet(provisional)?; let started = clock(); - let published: Result<(), Box> = - async { room.publish(report_data_packet(provisional)?).await }.await; - close_live.await; - published?; + let (acknowledged, ()) = tokio::join!(room.publish(provisional.clone()), close_live); + let unconfirmed = (!acknowledged?).then_some(provisional); if let Some(generated) = recover_report( room, + unconfirmed, regenerate, REPORT_RECOVERY_WINDOW, cooldown, diff --git a/tests/browser/account.test.js b/tests/browser/account.test.js index afce4b2a..e8319e52 100644 --- a/tests/browser/account.test.js +++ b/tests/browser/account.test.js @@ -75,8 +75,14 @@ test("interview history routes through the shared persistence helper", () => { for (const name of ["receiveReport", "showReport"]) { const body = functionBody(script, name); assert.match(body, /renderReport\(\)/); - assert.match(body, /renderReportSaveStatus\(await saving\)/); - assert.ok(body.indexOf("renderReport()") < body.indexOf("await saving")); + // An agent report also waits on its receipt before Done may navigate. + const awaited = body.search(/await (saving|Promise\.all\(\[saving\b)/); + assert.ok(awaited !== -1, `${name} awaits the history save`); + assert.match( + body, + /renderReportSaveStatus\(await saving\)|await Promise\.all\(\[saving, flushed\]\)[\s\S]*renderReportSaveStatus\(/, + ); + assert.ok(body.indexOf("renderReport()") < awaited); } assert.match(functionBody(script, "renderReport"), /saveResult: null/); assert.match( diff --git a/tests/browser/dom-contract.test.js b/tests/browser/dom-contract.test.js index fa57776f..cc69ebc4 100644 --- a/tests/browser/dom-contract.test.js +++ b/tests/browser/dom-contract.test.js @@ -938,14 +938,15 @@ test("runner progress statuses stay wired to each execution path", () => { test("a report that never lands still offers the offline summary", () => { const ending = functionBody(interviewSource(), "endInterview"); - // Ordering, not a byte window: both reveals live in the escape timeout, and - // the only thing worth pinning is that the offline summary is revealed there - // too rather than left hidden behind "leave the room". - const escape = ending.indexOf("REPORT_ESCAPE_WAIT_MS"); - assert.ok(escape !== -1, "the escape timeout left endInterview"); - assert.ok(ending.indexOf("nodes.leaveRoom.hidden = false") > escape); + // Inspect the timeout callback rather than a mention in its comment. + const deadline = ending.search(/\},\s*REPORT_ESCAPE_WAIT_MS\s*\)/); + assert.ok(deadline !== -1, "the escape timeout left endInterview"); + const callback = ending.lastIndexOf("setTimeout(() => {", deadline); + assert.ok(callback !== -1 && callback < deadline); + const escape = ending.slice(callback, deadline); + assert.ok(escape.includes("nodes.leaveRoom.hidden = false")); assert.ok( - ending.indexOf("nodes.forceReport.hidden = false") > escape, + escape.includes("nodes.forceReport.hidden = false"), "past the escape deadline the offline summary is offered beside leaving, not instead of it", ); }); diff --git a/tests/browser/lib.test.js b/tests/browser/lib.test.js index 5728f4b4..46adbb31 100644 --- a/tests/browser/lib.test.js +++ b/tests/browser/lib.test.js @@ -3088,3 +3088,240 @@ test("requested thinking time marks response windows until a release or intervie [true, true, false], ); }); + +test("report retries are claimed once and every copy receives a receipt", async () => { + const { claimReport, reportReceipt, reportReceiptPayload } = + await import("../../web/lib.js"); + const received = new Set(); + const packet = new TextEncoder().encode(JSON.stringify({ codingScore: 80 })); + // Synchronous, so the page decides before its first await. + assert.equal(claimReport(packet, received), true); + assert.equal(claimReport(packet, received), false); + const receipt = await reportReceipt(packet); + assert.deepEqual(receipt, reportReceiptPayload(receipt.deliveryId)); + assert.deepEqual(await reportReceipt(packet), receipt); + const { createHash } = await import("node:crypto"); + assert.equal( + receipt.deliveryId, + createHash("sha256").update(packet).digest("hex"), + ); +}); + +test("concurrent report retries publish receipts and render once even when a receipt throws", async (t) => { + const { claimReport, receiveReportDelivery } = + await import("../../web/lib.js"); + const delays = []; + const originalTimeout = globalThis.setTimeout; + t.mock.method(globalThis, "setTimeout", (callback, delay, ...args) => { + delays.push(delay); + return originalTimeout(callback, delay, ...args); + }); + const packet = new TextEncoder().encode('{"codingScore":80}'); + const received = new Set(); + let receipts = 0; + let renders = 0; + const publish = () => { + receipts++; + throw new Error("disconnected"); + }; + const render = (bytes) => { + assert.equal(bytes, packet); + renders++; + }; + await Promise.all([ + receiveReportDelivery( + packet, + claimReport(packet, received), + publish, + render, + ), + receiveReportDelivery( + packet, + claimReport(packet, received), + publish, + render, + ), + ]); + assert.equal(receipts, 2); + assert.equal(renders, 1); + assert.deepEqual(delays, [1000, 1000]); +}); + +test( + "a stalled receipt cannot keep a delivered report off screen", + { timeout: 3000 }, + async () => { + const { receiveReportDelivery } = await import("../../web/lib.js"); + let displayed = false; + let published = false; + let flushed = null; + await receiveReportDelivery( + new TextEncoder().encode('{"codingScore":80}'), + true, + () => { + published = true; + return new Promise(() => {}); + }, + (_bytes, receipt) => { + // Drawn before the receipt has even been handed to the room. + assert.equal(published, false); + displayed = true; + flushed = receipt; + }, + ); + assert.equal(displayed, true); + assert.equal(published, true); + // The disconnect waits on this, and a stalled publish must not hold it. + assert.equal(await flushed, undefined); + }, +); + +test( + "a hash that never finishes holds neither the report nor the disconnect", + { timeout: 3000 }, + async (t) => { + const { receiveReportDelivery } = await import("../../web/lib.js"); + t.mock.method(crypto.subtle, "digest", () => new Promise(() => {})); + let flushed = null; + let published = false; + await receiveReportDelivery( + new TextEncoder().encode('{"codingScore":80}'), + true, + () => { + published = true; + }, + (_bytes, receipt) => { + flushed = receipt; + }, + ); + assert.notEqual(flushed, null); + assert.equal(await flushed, undefined); + assert.equal(published, false); + }, +); + +test("an agent report leaves the room only after its receipt had its chance", () => { + const receive = functionBody(read("web/interview.js"), "receiveReport"); + // Disconnecting, Done and the render-failure exits all drop a receipt + // still leaving. + assert.doesNotMatch(receive, /void room\??\.disconnect\(/); + assert.equal( + receive.match(/flushed\.then\(\(\) => room\??\.disconnect\(\)\)/g)?.length, + 2, + ); + // Each pair is found before it is ordered, since a missing string's -1 + // would otherwise sort first and pass. + for (const [first, then] of [ + ["Promise.all([saving, flushed])", "renderReportSaveStatus("], + ["await flushed;", "reportRenderFailed("], + ]) { + const at = receive.indexOf(first); + assert.ok(at !== -1, first); + assert.ok(receive.indexOf(then, at) !== -1, `${then} after ${first}`); + } +}); + +test("a retry is recognised whether or not each copy could be hashed", async (t) => { + const { claimReport, receiveReportDelivery } = + await import("../../web/lib.js"); + const digest = crypto.subtle.digest.bind(crypto.subtle); + // Hashing never works, or fails for the first copy only. + for (const hashes of [() => false, (call) => call > 0]) { + let calls = 0; + const mock = t.mock.method(crypto.subtle, "digest", async (...args) => { + if (!hashes(calls++)) throw new Error("unavailable"); + return digest(...args); + }); + const packet = new TextEncoder().encode('{"codingScore":80}'); + const received = new Set(); + let receipts = 0; + let renders = 0; + for (let copy = 0; copy < 2; copy++) { + await receiveReportDelivery( + packet, + claimReport(packet, received), + () => receipts++, + () => renders++, + ); + } + mock.mock.restore(); + assert.equal(calls, 2); + assert.equal(renders, 1); + assert.equal(receipts, hashes(0) + hashes(1)); + } +}); + +test("only the first report copy claims the page, before its receipt is handed off", () => { + const connect = functionBody(read("web/interview.js"), "connectLiveKit"); + const accepted = connect.indexOf( + "if (!acceptsReport(topic, participant)) return;", + ); + const claim = connect.indexOf( + "const first = claimReport(payload, receivedReports);", + accepted, + ); + const guarded = connect.indexOf("if (first) {", claim); + const handoff = connect.indexOf("await receiveReportDelivery(", guarded); + assert.ok(accepted >= 0 && claim > accepted && guarded > claim); + // Everything that hides the ways out sits inside the first-copy branch. + const branch = connect.slice( + guarded, + connect.indexOf("\n }\n", guarded), + ); + for (const step of [ + "state.reportReceiving = true;", + "nodes.forceReport.hidden = true;", + "nodes.leaveRoom.hidden = true;", + "stopEndingEscape();", + ]) { + assert.ok(branch.includes(step), step); + assert.equal( + connect.indexOf(step, accepted), + connect.indexOf(step, guarded), + ); + } + assert.ok(branch.length < handoff - guarded); + assert.match( + connect.slice(handoff), + /^await receiveReportDelivery\(\s*payload,\s*first,/, + ); +}); + +test("an arriving report holds the offline summary and leaving until it fails to draw", async () => { + const { runInNewContext } = await import("node:vm"); + const page = read("web/interview.js"); + const state = { phase: "ending", reportReceiving: true }; + let navigated = false; + const context = { + state, + window: { + location: { + set href(_value) { + navigated = true; + }, + }, + }, + }; + runInNewContext( + `${functionBody(page, "showReport")}\n}\n` + + `${functionBody(page, "leaveRoom")}\n}\n` + + "void showReport(); leaveRoom();", + context, + ); + assert.equal(state.phase, "ending"); + assert.equal(navigated, false); + const failed = functionBody(page, "reportRenderFailed"); + assert.match(failed, /state\.reportReceiving = false;/); +}); + +test("an accepted report prevents a later End click from changing replay provenance", async () => { + const { runInNewContext } = await import("node:vm"); + const page = read("web/interview.js"); + const state = { phase: "live", reportReceiving: true }; + runInNewContext( + functionBody(page, "endInterview") + + '\n}\nendInterview("candidate_ended");', + { state }, + ); + assert.equal(state.phase, "live"); +}); diff --git a/tests/browser/report-recovery.test.js b/tests/browser/report-recovery.test.js index 1ab7a9f0..f4922998 100644 --- a/tests/browser/report-recovery.test.js +++ b/tests/browser/report-recovery.test.js @@ -1,8 +1,14 @@ import { test } from "node:test"; import assert from "node:assert/strict"; -import { createReportRecovery } from "../../web/report-recovery.js"; +import { + createReportRecovery, + reportRecoveryLimits, +} from "../../web/report-recovery.js"; import { retryReportPayload, sanitizeReport, topics } from "../../web/lib.js"; +// Counted from the agent's answer, or from the click when no answer comes. +const retryWaitMs = reportRecoveryLimits.retryWaitSeconds * 1000; + function setup() { let now = 0; let next = 0; @@ -63,7 +69,7 @@ test("report recovery spends one retry after cooldown and times it from acceptan advance(100_000); assert.equal(finalized(events).length, 0); assert.equal(recovery.notice(accepted), true); - advance(139_000); + advance(retryWaitMs - 1000); assert.equal(finalized(events).length, 0); advance(1000); assert.equal(events.filter((e) => e?.type === "retry_report").length, 1); @@ -80,9 +86,9 @@ test("an accepted retry outlives the offer's own expiry", () => { recovery.notice(accepted); // A closing notice cannot cut short a generation the agent already began. assert.equal(recovery.notice(closed), false); - advance(100_000); + advance(retryWaitMs - 1000); assert.equal(finalized(events).length, 0); - advance(40_000); + advance(1000); assert.equal(finalized(events).length, 1); }); @@ -179,7 +185,7 @@ test("a retry whose answer is lost is bounded by its own wait, not the offer", ( // The offer would have run out here; the agent may still be generating. advance(30_000); assert.equal(finalized(events).length, 0); - advance(109_000); + advance(retryWaitMs - 31_000); assert.equal(finalized(events).length, 0); advance(1000); assert.equal(finalized(events).length, 1); diff --git a/tests/fixtures/report-receipt.json b/tests/fixtures/report-receipt.json new file mode 100644 index 00000000..e807d589 --- /dev/null +++ b/tests/fixtures/report-receipt.json @@ -0,0 +1,12 @@ +{ + "topic": "control", + "cases": [ + { + "name": "received report", + "payload": { + "type": "report_received", + "deliveryId": "a6eb5f83ac3b057d6ca883e4c844bb990c1047e8b42bd05fae824424c042a93f" + } + } + ] +} diff --git a/tests/fixtures/report-recovery.json b/tests/fixtures/report-recovery.json index 25f3aee6..39f37a8a 100644 --- a/tests/fixtures/report-recovery.json +++ b/tests/fixtures/report-recovery.json @@ -2,7 +2,7 @@ "expiresInSeconds": 300, "retryAfterSeconds": 30, "quotaRetryAfterSeconds": 60, - "retryWaitSeconds": 140, + "retryWaitSeconds": 145, "closeGraceSeconds": 15, "retryStatuses": [ "accepted", diff --git a/tests/unit/livekit.rs b/tests/unit/livekit.rs index 6d195cec..c068f139 100644 --- a/tests/unit/livekit.rs +++ b/tests/unit/livekit.rs @@ -1211,13 +1211,53 @@ fn the_browser_escape_hatch_outlasts_the_report_deadline() { .expect("REPORT_ESCAPE_WAIT_MS is a number"), ); + // The Gemini close sits between the frozen report and its first publish. + let close = crate::gemini::CLOSE_TIMEOUT; + let delivery = report::DELIVERY_WAIT * report::DELIVERY_ATTEMPTS as u32; assert!( - wait >= REPORT_TIMEOUT + WRAP_UP_WAIT, - "a report bounded at {REPORT_TIMEOUT:?} after a {WRAP_UP_WAIT:?} wrap-up cannot land \ - before the page offers to leave at {wait:?}" + wait > REPORT_TIMEOUT + WRAP_UP_WAIT + close + delivery, + "a report bounded at {REPORT_TIMEOUT:?} after a {WRAP_UP_WAIT:?} wrap-up and a \ + {close:?} Gemini close cannot land with up to {delivery:?} delivery time before \ + the page offers to leave at {wait:?}" ); } +/// A publish only queues the packet, so leaving right behind it can drop the +/// report. The ending leaves once delivery settles, with no fixed sleep +/// standing in for it, and a report goes out only through the receipt wait. +#[test] +fn the_ending_leaves_only_after_the_report_delivery_settles() { + let source = include_str!("../../src/livekit.rs"); + // To the next item at column zero, as the source tests above read a body. + let ending = source + .split("async fn handle_data_packet(") + .nth(1) + .expect("handle_data_packet is still defined here") + .split("\n}\n") + .next() + .unwrap_or_default(); + let publish = ending + .find("publish_with_recovery(") + .expect("the ending publishes through recovery"); + let leave = ending.find("leave_room(").expect("the ending leaves"); + assert!(publish < leave); + assert!(!ending[publish..leave].contains("sleep(")); + + let report = include_str!("../../src/livekit/report.rs"); + let live = report + .split("impl RecoveryRoom for LiveRecoveryRoom") + .nth(1) + .expect("the live recovery room is still defined here") + .split("\n}\n") + .next() + .unwrap_or_default(); + let publish = live + .split("async fn publish(") + .nth(1) + .expect("the live room publishes reports"); + assert!(publish.contains("deliver_report(")); +} + /// Each pause reads the stretch since the last one, and never that stretch /// twice. /// diff --git a/tests/unit/livekit/report.rs b/tests/unit/livekit/report.rs index 5bdd169d..5ee10787 100644 --- a/tests/unit/livekit/report.rs +++ b/tests/unit/livekit/report.rs @@ -962,16 +962,21 @@ fn report_metadata_uses_the_frozen_assessment() { ); } +/// Events in; notices and reports out; whether the candidate is present; +/// whether a published report is acknowledged; and, when set, the presence a +/// publish leaves behind, for a candidate who leaves during its receipt wait. struct RecoveryFixture( tokio::sync::mpsc::UnboundedReceiver, Vec, Vec, bool, + bool, + Option, ); impl RecoveryFixture { fn new(events: tokio::sync::mpsc::UnboundedReceiver) -> Self { - Self(events, Vec::new(), Vec::new(), true) + Self(events, Vec::new(), Vec::new(), true, true, None) } } @@ -991,10 +996,13 @@ impl RecoveryRoom for RecoveryFixture { async fn publish( &mut self, packet: DataPacket, - ) -> Result<(), Box> { + ) -> Result> { assert_eq!(packet.topic.as_deref(), Some(TOPIC_REPORT)); self.2.push(serde_json::from_slice(&packet.payload)?); - Ok(()) + if let Some(present) = self.5 { + self.3 = present; + } + Ok(self.4) } } @@ -1021,6 +1029,7 @@ async fn recovery_wait_never_generates_before_a_valid_request() { }; let result = recover_report( &mut fixture, + None, generation, std::time::Duration::from_millis(2), REPORT_RETRY_COOLDOWN, @@ -1049,6 +1058,7 @@ async fn recovery_generates_once_despite_duplicate_requests() { }; let result = recover_report( &mut fixture, + None, generation, REPORT_RECOVERY_WINDOW, std::time::Duration::ZERO, @@ -1069,6 +1079,7 @@ async fn leaving_during_regeneration_cancels_the_call() { let mut fixture = RecoveryFixture::new(rx); let result = recover_report( &mut fixture, + None, std::future::pending(), REPORT_RECOVERY_WINDOW, std::time::Duration::ZERO, @@ -1096,6 +1107,7 @@ async fn an_early_retry_is_told_how_long_to_wait_and_keeps_its_turn() { let (result, _tx) = tokio::join!( recover_report( &mut fixture, + None, async { Ok(Ok(serde_json::json!({"result": "complete"}))) }, REPORT_RECOVERY_WINDOW, REPORT_RETRY_COOLDOWN, @@ -1131,6 +1143,7 @@ async fn a_retry_after_expiry_is_not_accepted() { let (result, _tx) = tokio::join!( recover_report( &mut fixture, + None, async { Ok(Ok(serde_json::json!({}))) }, REPORT_RECOVERY_WINDOW, REPORT_RETRY_COOLDOWN, @@ -1171,7 +1184,13 @@ fn recovery_offer_matches_browser_limits_and_generation_deadline() { recovery_metadata(crate::gemini::QUOTA_COOLDOWN)["retryAfterSeconds"], limits["quotaRetryAfterSeconds"] ); - assert!(limits["retryWaitSeconds"].as_u64().unwrap() > REPORT_TIMEOUT.as_secs() + 3); + + // The regenerated report goes through the same receipt wait, so a page that + // stops waiting before every attempt could land drops one in flight. + let delivery = DELIVERY_WAIT * DELIVERY_ATTEMPTS as u32; + assert!( + limits["retryWaitSeconds"].as_u64().unwrap() > (REPORT_TIMEOUT + delivery).as_secs() + 3 + ); } #[tokio::test] @@ -1400,6 +1419,7 @@ async fn scripted_recovery( let (result, _tx) = tokio::join!( recover_report( &mut fixture, + None, generation, REPORT_RECOVERY_WINDOW, std::time::Duration::ZERO, @@ -1491,3 +1511,604 @@ async fn a_retry_waits_for_keys_that_went_out_after_the_offer() { assert_eq!(notices, [RecoveryNotice::Closed]); assert!(started.elapsed() < std::time::Duration::from_secs(1)); } + +#[test] +fn report_receipt_requires_the_candidate_and_matching_delivery() { + let fixture: serde_json::Value = + serde_json::from_str(include_str!("../../fixtures/report-receipt.json")).unwrap(); + let cases = fixture["cases"].as_array().unwrap(); + assert_eq!(cases.len(), 1); + let payload = serde_json::to_vec(&cases[0]["payload"]).unwrap(); + let id = cases[0]["payload"]["deliveryId"].as_str().unwrap(); + // The fixture holds the page's digest; the agent must compute the same. + assert_eq!(id, crate::sha256_hex(&[br#"{"codingScore":80}"#])); + let valid = |topic, sender, id, payload: &[u8]| { + is_report_receipt(topic, sender, "candidate", payload, id, false) + }; + assert!(valid(Some("control"), Some("candidate"), id, &payload)); + assert!(!valid(Some("report"), Some("candidate"), id, &payload)); + assert!(!valid(Some("control"), Some("peer"), id, &payload)); + + // LiveKit drops the sender of a packet that lands after its departure. Any + // peer could echo the digest, so that counts only once the candidate is + // gone. + assert!(!valid(Some("control"), None, id, &payload)); + assert!(is_report_receipt( + Some("control"), + None, + "candidate", + &payload, + id, + true + )); + assert!(!valid( + Some("control"), + Some("candidate"), + "old-report", + &payload + )); + assert!(!valid( + Some("control"), + Some("candidate"), + "malformed", + b"invalid" + )); +} + +#[tokio::test(start_paused = true)] +async fn report_delivery_retries_identical_bytes_after_a_publish_failure() { + let (_sender, mut events) = tokio::sync::mpsc::unbounded_channel::(); + let payload = br#"{"codingScore":80}"#.to_vec(); + let mut published = Vec::new(); + let started = tokio::time::Instant::now(); + let result = deliver_report( + |packet| { + assert_eq!(packet.topic.as_deref(), Some(TOPIC_REPORT)); + assert!(packet.reliable); + published.push(packet.payload); + std::future::ready(if published.len() == 1 { + Err("transport") + } else { + Ok(()) + }) + }, + &mut events, + "candidate", + "room", + || true, + delivery_test_packet(payload.clone()), + ) + .await; + assert!(!result.unwrap()); + assert_eq!(published, vec![payload; DELIVERY_ATTEMPTS]); + assert_eq!(started.elapsed(), DELIVERY_WAIT * DELIVERY_ATTEMPTS as u32); +} + +#[tokio::test(start_paused = true)] +async fn a_hung_report_publish_is_bounded_and_retried() { + let (_sender, mut events) = tokio::sync::mpsc::unbounded_channel::(); + let mut attempts = 0; + let started = tokio::time::Instant::now(); + let result = deliver_report( + |_| { + attempts += 1; + std::future::pending::>() + }, + &mut events, + "candidate", + "room", + || true, + delivery_test_packet(br#"{"codingScore":80}"#.to_vec()), + ) + .await; + assert!(result.is_err()); + assert_eq!(attempts, DELIVERY_ATTEMPTS); + assert_eq!(started.elapsed(), DELIVERY_WAIT * DELIVERY_ATTEMPTS as u32); +} + +struct ReceiptTestEvent { + sender: &'static str, + payload: Vec, +} + +impl DeliveryEvent for ReceiptTestEvent { + fn acknowledges(&self, candidate: &str, id: &str, candidate_gone: bool) -> bool { + is_report_receipt( + Some("control"), + Some(self.sender), + candidate, + &self.payload, + id, + candidate_gone, + ) + } +} + +#[tokio::test(start_paused = true)] +async fn delivery_stops_on_the_candidate_receipt_after_retry() { + let (sender, mut events) = tokio::sync::mpsc::unbounded_channel(); + let fixture: serde_json::Value = + serde_json::from_str(include_str!("../../fixtures/report-receipt.json")).unwrap(); + let receipt = serde_json::to_vec(&fixture["cases"][0]["payload"]).unwrap(); + sender + .send(ReceiptTestEvent { + sender: "peer", + payload: receipt.clone(), + }) + .unwrap(); + sender + .send(ReceiptTestEvent { + sender: "candidate", + payload: receipt_bytes("old"), + }) + .unwrap(); + let mut attempts = 0; + let result = deliver_report( + |_| { + attempts += 1; + if attempts == 2 { + sender + .send(ReceiptTestEvent { + sender: "candidate", + payload: receipt.clone(), + }) + .unwrap(); + } + std::future::ready(if attempts == 1 { + Err("transport") + } else { + Ok(()) + }) + }, + &mut events, + "candidate", + "room", + || true, + delivery_test_packet(br#"{"codingScore":80}"#.to_vec()), + ) + .await; + assert!(result.unwrap()); + assert_eq!(attempts, 2); +} + +#[tokio::test(start_paused = true)] +async fn a_lost_packet_is_retransmitted_until_receipt() { + for acknowledged_on in [1, 2] { + let (sender, mut events) = tokio::sync::mpsc::unbounded_channel(); + let payload = br#"{"codingScore":80}"#.to_vec(); + let id = crate::sha256_hex(&[&payload]); + let mut published = Vec::new(); + let started = tokio::time::Instant::now(); + let result = deliver_report( + |packet| { + published.push(packet.payload); + if published.len() == acknowledged_on { + sender + .send(ReceiptTestEvent { + sender: "candidate", + payload: receipt_bytes(&id), + }) + .unwrap(); + } + std::future::ready(Ok::<(), std::io::Error>(())) + }, + &mut events, + "candidate", + "room", + || true, + delivery_test_packet(payload.clone()), + ) + .await; + assert!(result.unwrap()); + assert_eq!(published, vec![payload; acknowledged_on]); + assert_eq!( + started.elapsed(), + DELIVERY_WAIT * (acknowledged_on - 1) as u32 + ); + } +} + +#[tokio::test(start_paused = true)] +async fn a_disconnected_room_stops_waiting_without_claiming_delivery_failure() { + let (sender, mut events) = tokio::sync::mpsc::unbounded_channel(); + sender + .send(RoomEvent::Disconnected { + reason: ::livekit::DisconnectReason::ClientInitiated, + }) + .unwrap(); + let mut attempts = 0; + let started = tokio::time::Instant::now(); + let result = deliver_report( + |_| { + attempts += 1; + std::future::ready(Ok::<(), std::io::Error>(())) + }, + &mut events, + "candidate", + "room", + || true, + delivery_test_packet(b"{}".to_vec()), + ) + .await; + assert!(!result.unwrap()); + assert_eq!(attempts, 1); + assert_eq!(started.elapsed(), Duration::ZERO); +} + +#[tokio::test(start_paused = true)] +async fn failed_report_publishes_are_retried_and_reported_as_failure() { + let (_sender, mut events) = tokio::sync::mpsc::unbounded_channel::(); + let mut attempts = 0; + let result = deliver_report( + |_| { + attempts += 1; + std::future::ready(Err::<(), _>("transport")) + }, + &mut events, + "candidate", + "room", + || true, + delivery_test_packet(b"{}".to_vec()), + ) + .await; + assert!( + result + .unwrap_err() + .to_string() + .contains("report_delivery_failed") + ); + assert_eq!(attempts, DELIVERY_ATTEMPTS); +} + +/// What the page sends back for a report whose bytes hash to `id`. +fn receipt_bytes(id: &str) -> Vec { + serde_json::to_vec(&serde_json::json!({ + "type": "report_received", "deliveryId": id, + })) + .unwrap() +} + +/// Built by the same `report_data_packet` the agent sends, so the topic and +/// reliability under test cannot drift from the real packet's. +fn delivery_test_packet(payload: Vec) -> DataPacket { + let packet = report_data_packet(serde_json::from_slice(&payload).unwrap()).unwrap(); + assert_eq!(packet.payload, payload, "the test bytes are the bytes sent"); + packet +} + +#[tokio::test(start_paused = true)] +async fn a_candidate_already_absent_gets_one_publish_without_a_receipt_wait() { + let (_sender, mut events) = tokio::sync::mpsc::unbounded_channel::(); + let started = tokio::time::Instant::now(); + let mut attempts = 0; + let result = deliver_report( + |_| { + attempts += 1; + std::future::ready(Ok::<(), std::io::Error>(())) + }, + &mut events, + "candidate", + "room", + || false, + delivery_test_packet(b"{}".to_vec()), + ) + .await; + assert!(!result.unwrap()); + assert_eq!(attempts, 1); + assert_eq!(started.elapsed(), Duration::ZERO); +} + +#[tokio::test(start_paused = true)] +async fn failed_publication_followed_by_disconnect_is_a_delivery_failure() { + let (sender, mut events) = tokio::sync::mpsc::unbounded_channel(); + sender + .send(RoomEvent::Disconnected { + reason: ::livekit::DisconnectReason::ClientInitiated, + }) + .unwrap(); + let result = deliver_report( + |_| std::future::ready(Err::<(), _>("transport")), + &mut events, + "candidate", + "room", + || true, + delivery_test_packet(b"{}".to_vec()), + ) + .await; + assert!( + result + .unwrap_err() + .to_string() + .contains("report_delivery_failed") + ); +} + +// LiveKit constructs RemoteParticipant internally, so the real room's +// participant map is supplied through the presence probe rather than faked. +#[tokio::test(start_paused = true)] +async fn a_candidate_leaving_during_the_receipt_wait_stops_retries() { + let (sender, mut events) = tokio::sync::mpsc::unbounded_channel(); + sender + .send(ReceiptTestEvent { + sender: "peer", + payload: b"{}".to_vec(), + }) + .unwrap(); + let presence_reads = std::cell::Cell::new(0); + let mut attempts = 0; + let started = tokio::time::Instant::now(); + let result = deliver_report( + |_| { + attempts += 1; + std::future::ready(Ok::<(), std::io::Error>(())) + }, + &mut events, + "candidate", + "room", + || { + let reads = presence_reads.get(); + presence_reads.set(reads + 1); + + // Present after publication, gone when the queued event is + // consumed. + reads == 0 + }, + delivery_test_packet(b"{}".to_vec()), + ) + .await; + assert!(!result.unwrap()); + assert_eq!(attempts, 1); + assert_eq!(presence_reads.get(), 2); + assert_eq!(started.elapsed(), Duration::ZERO); +} + +#[test] +fn an_unattributed_receipt_counts_only_once_the_candidate_is_gone() { + let payload = br#"{"codingScore":80}"#; + let id = crate::sha256_hex(&[payload]); + let receipt = |topic: &str| RoomEvent::DataReceived { + payload: std::sync::Arc::new(receipt_bytes(&id)), + topic: Some(topic.to_string()), + kind: ::livekit::prelude::DataPacketKind::Reliable, + participant: None, + }; + assert!(receipt("control").acknowledges("candidate", &id, true)); + assert!(!receipt("control").acknowledges("candidate", &id, false)); + assert!(!receipt("report").acknowledges("candidate", &id, true)); + assert!(!receipt("control").disconnected()); + assert!( + RoomEvent::Disconnected { + reason: ::livekit::DisconnectReason::ClientInitiated, + } + .disconnected() + ); +} + +/// The page sends its receipt and then disconnects, and the roster can lose the +/// candidate before this loop reads the receipt queued ahead of that. +#[tokio::test(start_paused = true)] +async fn a_receipt_queued_before_the_candidate_left_is_still_acknowledged() { + let payload = br#"{"codingScore":80}"#.to_vec(); + let id = crate::sha256_hex(&[&payload]); + let receipt = || ReceiptTestEvent { + sender: "candidate", + payload: receipt_bytes(&id), + }; + // Gone before the wait starts, and gone while an unrelated event is read. + for unrelated_first in [false, true] { + let (sender, mut events) = tokio::sync::mpsc::unbounded_channel(); + if unrelated_first { + sender + .send(ReceiptTestEvent { + sender: "peer", + payload: b"{}".to_vec(), + }) + .unwrap(); + } + sender.send(receipt()).unwrap(); + let presence_reads = std::cell::Cell::new(0); + let result = deliver_report( + |_| std::future::ready(Ok::<(), std::io::Error>(())), + &mut events, + "candidate", + "room", + || { + let reads = presence_reads.get(); + presence_reads.set(reads + 1); + unrelated_first && reads == 0 + }, + delivery_test_packet(payload.clone()), + ) + .await; + assert!(result.unwrap(), "unrelated_first={unrelated_first}"); + } +} + +/// The provisional report's receipt wait reads the same event stream, so a +/// departure during it never reaches `recover_report`. Absence at the start +/// begins the rejoin grace instead of a five-minute wait for nobody, and a +/// rejoin within it keeps the offer. +#[tokio::test(start_paused = true)] +async fn a_candidate_gone_when_the_wait_begins_gets_only_the_rejoin_grace() { + for rejoins in [false, true] { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + if rejoins { + tx.send(RecoveryEvent::Back).unwrap(); + } + let mut fixture = RecoveryFixture::new(rx); + fixture.3 = false; + let started = tokio::time::Instant::now(); + let result = recover_report( + &mut fixture, + None, + async { Ok(Ok(serde_json::json!({}))) }, + REPORT_RECOVERY_WINDOW, + REPORT_RETRY_COOLDOWN, + started, + || Readiness::Ready, + ) + .await; + assert!(result.is_none()); + let waited = started.elapsed(); + if rejoins { + // Back in time, so only the window's own expiry ends the wait. + assert_eq!(waited, REPORT_RECOVERY_WINDOW); + assert_eq!(fixture.1, vec![RecoveryNotice::Closed]); + } else { + assert_eq!(waited, REJOIN_GRACE); + assert!(fixture.1.is_empty()); + } + drop(tx); + } +} + +/// The receipt that arrives unattributed because its sender just left is +/// itself the event that reveals the departure, so presence is read before the +/// receipt is judged rather than only for the events behind it. +#[tokio::test(start_paused = true)] +async fn an_unattributed_receipt_that_reveals_the_departure_is_acknowledged() { + let payload = br#"{"codingScore":80}"#.to_vec(); + let id = crate::sha256_hex(&[&payload]); + let (sender, mut events) = tokio::sync::mpsc::unbounded_channel(); + sender + .send(RoomEvent::DataReceived { + payload: std::sync::Arc::new(receipt_bytes(&id)), + topic: Some("control".to_string()), + kind: ::livekit::prelude::DataPacketKind::Reliable, + participant: None, + }) + .unwrap(); + let presence_reads = std::cell::Cell::new(0); + let result = deliver_report( + |_| std::future::ready(Ok::<(), std::io::Error>(())), + &mut events, + "candidate", + "room", + || { + let reads = presence_reads.get(); + presence_reads.set(reads + 1); + // Present when the wait starts, gone once the receipt is read. + reads == 0 + }, + delivery_test_packet(payload), + ) + .await; + assert!(result.unwrap()); +} + +/// A candidate who dropped before the provisional report's receipt came back +/// may never have seen it. Their rejoin republishes it, once acknowledged it is +/// never sent again, and a retry request proves the page already holds it. +#[tokio::test(start_paused = true)] +async fn a_rejoin_republishes_a_provisional_report_never_acknowledged() { + let provisional = || delivery_test_packet(br#"{"incomplete":true}"#.to_vec()); + for (events, unconfirmed, republished) in [ + // Away and back twice: one republish, acknowledged the first time. + ( + vec![ + RecoveryEvent::Back, + RecoveryEvent::Away, + RecoveryEvent::Back, + ], + true, + 1, + ), + // Already acknowledged: a rejoin sends nothing. + (vec![RecoveryEvent::Back], false, 0), + // A retry first: the page has it, so a later rejoin sends nothing. + (vec![RecoveryEvent::Retry, RecoveryEvent::Back], true, 0), + ] { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + for event in events { + tx.send(event).unwrap(); + } + let mut fixture = RecoveryFixture::new(rx); + let started = tokio::time::Instant::now(); + let result = recover_report( + &mut fixture, + unconfirmed.then(provisional), + async { Ok(Ok(serde_json::json!({}))) }, + std::time::Duration::from_secs(1), + REPORT_RETRY_COOLDOWN, + started, + || Readiness::Ready, + ) + .await; + assert!(result.is_none()); + assert_eq!(fixture.2.len(), republished); + if republished > 0 { + assert_eq!(fixture.2[0], serde_json::json!({ "incomplete": true })); + } + drop(tx); + } +} + +/// Only a provisional report that went out unacknowledged is handed to the +/// wait for republishing; an acknowledged one is never sent twice. +#[tokio::test(start_paused = true)] +async fn only_an_unacknowledged_provisional_report_is_republished() { + for (acknowledged, reports) in [(false, 2), (true, 1)] { + let config = report_test_config(); + let boot = bootstrap(&config, "interview-fixed", Some("two-sum"), 45); + let mut live = RuntimeState::default(); + let mut frozen = freeze_assessment(&boot, &mut live, 12.0); + let keys = GeminiKeys::single("republish"); + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + tx.send(RecoveryEvent::Back).unwrap(); + let mut room = RecoveryFixture::new(rx); + room.4 = acknowledged; + run_recovery( + &mut room, + RecoveryReport { + boot: &boot, + state: &mut frozen.state, + reason: "time_up", + keys: &keys, + }, + Err(elapsed().await), + async { Ok(Ok(serde_json::json!({}))) }, + tokio::time::Instant::now, + async {}, + ) + .await + .unwrap(); + assert_eq!(room.2.len(), reports, "acknowledged={acknowledged}"); + assert!(room.2.iter().all(|report| report["incomplete"] == true)); + drop(tx); + } +} + +/// A republish waits for its own receipt on the event stream a second +/// departure would arrive on, so presence is read again after it: gone means +/// the rejoin grace, present means the offer runs its course. +#[tokio::test(start_paused = true)] +async fn presence_is_read_again_after_a_republish() { + for present in [false, true] { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + tx.send(RecoveryEvent::Back).unwrap(); + let mut fixture = RecoveryFixture::new(rx); + fixture.4 = false; + fixture.5 = Some(present); + let started = tokio::time::Instant::now(); + let result = recover_report( + &mut fixture, + Some(delivery_test_packet(br#"{"incomplete":true}"#.to_vec())), + async { Ok(Ok(serde_json::json!({}))) }, + REPORT_RECOVERY_WINDOW, + REPORT_RETRY_COOLDOWN, + started, + || Readiness::Ready, + ) + .await; + assert!(result.is_none()); + assert_eq!(fixture.2.len(), 1); + if present { + assert_eq!(started.elapsed(), REPORT_RECOVERY_WINDOW); + assert_eq!(fixture.1, vec![RecoveryNotice::Closed]); + } else { + assert_eq!(started.elapsed(), REJOIN_GRACE); + assert!(fixture.1.is_empty()); + } + drop(tx); + } +} diff --git a/web/interview.js b/web/interview.js index f57aa3b1..be20510c 100644 --- a/web/interview.js +++ b/web/interview.js @@ -36,6 +36,7 @@ import { import { acceptsReport, CANDIDATE_CASE_LIMIT, + claimReport, clamp, codeUpdatePayload, codingLoop, @@ -49,6 +50,7 @@ import { integrityEventPayload, isAgent, providerUiState, + receiveReportDelivery, roomInterviewer, sanitizeReport, sessionReport, @@ -307,6 +309,7 @@ const state = { /// worked examples are hidden: it is a worked case in its own right. candidateCaseHint: "", report: null, + reportReceiving: false, /// Set when the interviewer's report reached this page and could not be /// rendered. The offline summary that follows is written from what this page /// holds either way; what this changes is the sentence explaining why the @@ -1180,15 +1183,38 @@ async function connectLiveKit(connection, preflight, presenting = false) { if (!livekit?.Room) throw new Error("LiveKit browser SDK is unavailable."); const room = new livekit.Room({ adaptiveStream: true, dynacast: true }); + const receivedReports = new Set(); room.on( livekit.RoomEvent.DataReceived, - (payload, participant, kind, topic) => { + async (payload, participant, kind, topic) => { if (topic === topics.control && isAgent(participant)) { receiveControl(payload); return; } if (!acceptsReport(topic, participant)) return; - void receiveReport(room, payload); + // The report has arrived. Nothing may offer the offline summary or a way + // out while it is drawn, or the candidate could save an unscored summary + // beside it; `reportRenderFailed` gives the ways out back. Only the copy + // that renders claims this: a retransmission landing after a failed draw + // would otherwise hide the ways out again with nothing to restore them. + const first = claimReport(payload, receivedReports); + if (first) { + state.reportReceiving = true; + nodes.end.disabled = true; + nodes.forceReport.hidden = true; + nodes.leaveRoom.hidden = true; + stopEndingEscape(); + } + await receiveReportDelivery( + payload, + first, + (receipt) => + room.localParticipant.publishData( + new TextEncoder().encode(JSON.stringify(receipt)), + { reliable: true, topic: topics.control }, + ), + (report, flushed) => receiveReport(room, report, flushed), + ); }, ); room.on(livekit.RoomEvent.ParticipantAttributesChanged, updateAgentState); @@ -1480,7 +1506,10 @@ function finalizeRecoveryOnPageHide() { window.addEventListener("pagehide", finalizeRecoveryOnPageHide); -async function receiveReport(room, payload) { +/// `flushed` settles once the report receipt has had its chance to leave; the +/// room stays up until then so the agent can stop retransmitting. Recovery's +/// own finalize has no receipt to wait for. +async function receiveReport(room, payload, flushed = Promise.resolve()) { // The wait ended when this packet arrived, whatever becomes of it below. stopEndingEscape(); let rendered = false; @@ -1493,6 +1522,9 @@ async function receiveReport(room, payload) { // into an incomplete report rather than a throw. if (raw?.incomplete && raw?.reportRecovery) { const live = state.phase === "live"; + // Recovery owns the ways out from here, its retry and its "save and + // leave", so the claim that held them for a report being drawn ends. + state.reportReceiving = false; if (reportRecovery.start(raw)) { // Only once the offer is taken: one that is refused falls through to // the ordinary path below, which writes this frame itself. @@ -1540,10 +1572,12 @@ async function receiveReport(room, payload) { renderAttempted = true; renderReport(); rendered = true; - void room?.disconnect().catch(() => {}); + void flushed.then(() => room?.disconnect()).catch(() => {}); state.room = null; state.connected = false; - renderReportSaveStatus(await saving); + // Done navigates away, which would drop a receipt still leaving. + const [saved] = await Promise.all([saving, flushed]); + renderReportSaveStatus(saved); } catch (error) { reportRecovery.stop(); restoreEndingOverlay(); @@ -1561,7 +1595,7 @@ async function receiveReport(room, payload) { // shown, until the candidate pressed something. stopAvatar(); stopLocalMedia(); - void room?.disconnect().catch(() => {}); + void flushed.then(() => room?.disconnect()).catch(() => {}); state.room = null; state.connected = false; // The interviewer closed this one itself, and the `ended` row for it sits @@ -1585,6 +1619,10 @@ async function receiveReport(room, payload) { // It draws through the same `renderReport`, so a throw from inside that // one is a throw the click reproduces, and the guard around the second // attempt can do no more than say so again. + // + // Both ways out it offers disconnect, so they wait for the receipt that + // stops the agent's retries. + await flushed; reportRenderFailed(error, { offlineSummary: !renderAttempted }); } } @@ -2432,8 +2470,9 @@ function flushPendingLanguagePublish() { /// the agent's own deadline by /// the_browser_escape_hatch_outlasts_the_report_deadline, which reads this /// declaration: the value lives here, and Rust checks that it clears -/// REPORT_TIMEOUT plus WRAP_UP_WAIT rather than keeping a copy of it. -const REPORT_ESCAPE_WAIT_MS = 135000; +/// the generation, wrap-up, Gemini close and delivery bounds rather than +/// keeping copies. +const REPORT_ESCAPE_WAIT_MS = 155000; /// The two timers that speak for a report nobody has seen yet, held so that a /// report which arrives can take them back. Both say a wait is still running, @@ -2442,14 +2481,10 @@ const REPORT_ESCAPE_WAIT_MS = 135000; /// "preparing your report" and then with an offer to retry the provider. let endingEscape = []; -/// How long a report published just before the interviewer left has to arrive. -/// -/// The agent publishes and then leaves, so the moment it goes is also the -/// moment its report may be one packet away: `publish_report` in -/// src/livekit.rs is followed by a sleep and then `leave_room`. Offering the -/// offline summary inside that window offers to replace an evaluation that -/// exists with an unscored one, and the click disconnects the room that was -/// about to deliver it. +/// A final packet can arrive after the participant's departure notification, +/// particularly from an older agent that leaves without waiting for a receipt. +/// Offering the offline summary immediately lets the candidate disconnect +/// before an evaluation already in transit reaches the page. const REPORT_DELIVERY_GRACE_MS = 3000; /// Armed when the interviewer leaves while the page is still waiting, so a @@ -2464,7 +2499,7 @@ function stopEndingEscape() { } function endInterview(reason) { - if (state.phase !== "live") return; + if (state.phase !== "live" || state.reportReceiving) return; state.phase = "ending"; // The hint outranks the ending overlay in the stacking order, so a candidate // who ends while it is still up would read the report status through it. @@ -2503,12 +2538,10 @@ function endInterview(reason) { nodes.forceReport.hidden = reportComing; startEndingClock(); if (reportComing) { - // Longer than the agent's worst case, not shorter: the report is bounded - // by REPORT_TIMEOUT in src/livekit.rs and a timer-driven end spends - // WRAP_UP_WAIT ahead of it. Offering "leave the room" before that elapses - // invites the candidate to walk out on a report that is still coming, and - // leaving never saves it. REPORT_ESCAPE_WAIT_MS has to clear both, and - // says where that is checked. + // The wait clears report generation, the timer-driven wrap-up, the Gemini + // close, and all delivery retries, summed as if none overlapped. An + // earlier escape can disconnect a packet still in transit; the Rust + // deadline test holds this number against those bounds. endingEscape = [ setTimeout(() => { if (state.phase === "ending") @@ -2566,6 +2599,7 @@ function leaveRoom() { reportRecovery.finish(); return; } + if (state.reportReceiving) return; stopEndingClock(); void state.room?.disconnect?.(); stopAvatar(); @@ -2574,13 +2608,15 @@ function leaveRoom() { } async function showReport() { - if (state.phase === "report") return; + // An arrived report outranks the offline summary, including from the timer + // `endInterview` sets when no report is expected. + if (state.phase === "report" || state.reportReceiving) return; state.phase = "report"; // The offline summary is now offered while a room is still up -- it used to // be hidden for the whole life of one -- so this is the one report path that - // can leave the agent in an interview nobody is attending. `receiveReport` - // disconnects because the agent published and left; here nothing has, and a - // candidate reading a local summary is still paying for a Gemini session. + // can leave the agent in an interview nobody is attending. Receiving an + // agent report ends assessment; a local summary cannot do that, and the + // candidate would otherwise still be paying for a Gemini session. void state.room?.disconnect?.().catch?.(() => {}); state.room = null; state.connected = false; @@ -2644,6 +2680,9 @@ async function showReport() { function reportRenderFailed(error, { offlineSummary }) { console.warn("codetrial report_render_failed", error); state.reportUnreadable = true; + // The ways out below are the fallback for the report that failed, so they + // have to work again. + state.reportReceiving = false; state.phase = "ending"; stopEndingClock(); // Both of these speak for a report still on its way, and nothing is on its @@ -2834,15 +2873,14 @@ function updateAgentState() { "The interviewer disconnected. Nothing you typed is lost; end the interview to get your report.", ); } - // `endInterview` decides once whether a report is coming, and the - // interviewer can leave a moment later: the agent leaves the room right - // after publishing, and the failure behind issue 77 is the one where it - // leaves without publishing at all. The overlay covers the banner above, - // so a candidate already waiting learns none of this and sits out the - // whole escape wait for a report with nobody left to send it. + // An interviewer can leave after endInterview starts waiting, including + // after exhausting delivery retries. The overlay covers the banner above, + // so it must offer an exit once no sender remains, after allowing packets + // already in transit to arrive. if ( state.sawAgent && state.phase === "ending" && + !state.reportReceiving && !state.reportUnreadable && !deliveryGrace ) { diff --git a/web/lib.js b/web/lib.js index 127e550c..a7aaf312 100644 --- a/web/lib.js +++ b/web/lib.js @@ -293,6 +293,70 @@ export function acceptsReport(topic, participant) { return topic === topics.report && isAgent(participant); } +/// Claims a report copy, synchronously, so the caller can decide before its +/// first await whether this copy is the one that renders. Keyed on the bytes +/// rather than the digest, so a copy whose digest failed and one whose digest +/// worked are still the same delivery. +export function claimReport(payload, received) { + const key = new TextDecoder().decode(payload); + const first = !received.has(key); + received.add(key); + return first; +} + +/// Receipts name the exact packet bytes, leaving the report schema unchanged. +/// Null where Web Crypto is unavailable; display does not depend on it. +export async function reportReceipt(payload) { + try { + return reportReceiptPayload(await sha256Hex(payload)); + } catch { + return null; + } +} + +/// Renders a first copy without waiting on its receipt, not even on hashing +/// it. `receive` is handed a promise that settles once the receipt has had its +/// chance to leave, which is the earliest the page may disconnect or navigate; +/// a retransmitted copy only needs its receipt. +export async function receiveReportDelivery( + payload, + first, + publishReceipt, + receive, +) { + const flushed = sendReportReceipt(payload, publishReceipt); + if (first) await receive(payload, flushed); + await flushed; +} + +/// Hashing and publication share one bound, so neither can hold the page. +async function sendReportReceipt(payload, publishReceipt) { + let timer; + let queued; + try { + queued = await Promise.race([ + reportReceipt(payload) + .then(async (receipt) => { + if (!receipt) return false; + await publishReceipt(receipt); + return true; + }) + .catch(() => false), + new Promise((resolve) => { + timer = setTimeout(() => resolve(false), 1000); + }), + ]); + } finally { + clearTimeout(timer); + } + // Publication queues the receipt; let it leave before the page disconnects. + if (queued) await new Promise((resolve) => setTimeout(resolve, 250)); +} + +export function reportReceiptPayload(deliveryId) { + return { type: "report_received", deliveryId }; +} + // Data-channel payload builders. The agent decodes these by key, so they are // a wire contract; web/tests/lib.test.js pins the exact shapes. @@ -465,12 +529,16 @@ export async function integrityEventPayload( detail: event.detail, }; if (event.sourceEventIds.length) body.sourceEventIds = event.sourceEventIds; - const bytes = new TextEncoder().encode(canonicalJson(body)); + event.hash = await sha256Hex(new TextEncoder().encode(canonicalJson(body))); + return event; +} + +/// Lowercase hex, the spelling `sha256_hex` produces on the agent side. +async function sha256Hex(bytes) { const digest = await crypto.subtle.digest("SHA-256", bytes); - event.hash = [...new Uint8Array(digest)] + return [...new Uint8Array(digest)] .map((byte) => byte.toString(16).padStart(2, "0")) .join(""); - return event; } /// The two frameworks, kept apart on purpose. diff --git a/web/report-recovery.js b/web/report-recovery.js index ac77055f..443c5469 100644 --- a/web/report-recovery.js +++ b/web/report-recovery.js @@ -2,7 +2,7 @@ export const reportRecoveryLimits = Object.freeze({ expiresInSeconds: 300, retryAfterSeconds: 30, quotaRetryAfterSeconds: 60, - retryWaitSeconds: 140, + retryWaitSeconds: 145, // The agent announces the end of its window. This page's clock starts when // the offer arrives, later than the agent's, so it waits a little past it // for that word and falls back on its own only when the word never comes. @@ -57,8 +57,8 @@ export function createReportRecovery({ report = null; ready = false; } - // The agent gives regeneration its own 125-second deadline. Delivery gets - // the same grace as the ordinary report path. + // The agent gives regeneration its own 125-second deadline, and the report + // it produces may take every delivery attempt to arrive. function waitForReport() { if (wait) timers.clearTimeout(wait); wait = timers.setTimeout(