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(