diff --git a/.cargo/mutants.toml b/.cargo/mutants.toml index a3db85c5..1f0205b2 100644 --- a/.cargo/mutants.toml +++ b/.cargo/mutants.toml @@ -1,7 +1,7 @@ # Functions the mutation gate cannot judge, because `cargo test` cannot reach # them. Matched against the mutant names that `cargo mutants --list` prints. # -# EXCLUSIONS: 60 +# EXCLUSIONS: 62 # # That number is checked by `scripts/test.sh`, so adding an entry means editing # this line too. The point is not the count, it is that the list only ever grows @@ -148,12 +148,19 @@ # `parse_and_validate_report`, which is tested; the wait in between is the part # that needs the network. # -# `publish_report` is the two lines around `report_packet`: it takes `&Room` and -# publishes what that returns, so `Ok(())` looks the same from outside as -# publishing, exactly as `publish_interviewer_state` above does. It belonged -# here as soon as `report_packet` did and was missed because both lived in -# `livekit.rs`, where `report_packet`'s own entry already covered the pair; -# splitting `src/livekit/report.rs` out scored the wrapper on its own. +# `publish_with_recovery` and `LiveRecoveryRoom` are the report's way into a +# live `&Room`: the first builds the room and the report call and hands both to +# `run_recovery`, the second reads the room's events and roster and publishes +# packets on it, so `Ok(())`, `()`, `true` or `false` look the same from +# outside as the real thing, exactly as `publish_interviewer_state` above does. +# Everything they carry is decided and tested beside them: `run_recovery` and +# `recover_report` against a scripted room, `recovery_event`, `candidate_only`, +# `recovery_request` and `recovery_notice` directly. Two of `recovery_event`'s +# arms are named on their own for the same reason: deleting the arm for a +# participant leaving or joining or a packet arriving only shows on an event +# that carries a `RemoteParticipant`, which nothing outside the SDK can +# construct. Their decisions are `candidate_only` and `recovery_request`, and +# the arms for the room itself closing stay judged. # # The media helpers are the same shape one layer down. `handle_media_event` # and `attach_audio` take a LiveKit participant, `next_audio_frame` and @@ -308,7 +315,9 @@ exclude_re = [ "publish_transcript", "generate_report", "report_packet", - "publish_report", + "publish_with_recovery", + "LiveRecoveryRoom", + "delete match arm Some\\(::livekit::RoomEvent::(ParticipantDisconnected|ParticipantConnected|DataReceived)\\b.* in recovery_event", "replace_gemini_session", "open_cold_session", "open_live_session_at", diff --git a/docs/provider-cost-and-degradation.md b/docs/provider-cost-and-degradation.md index d69d79cd..baf40d1c 100644 --- a/docs/provider-cost-and-degradation.md +++ b/docs/provider-cost-and-degradation.md @@ -83,12 +83,41 @@ Candidate video, off by default, therefore sends one frame in five seconds and asks for the low media resolution; the camera is there for presence, and the code reaches the model as text. -Final reporting has its own hard budget of six Gemini HTTP calls: initial -generation plus one semantic repair, each generation allowing its first call and -at most two transient retries. The counter is consumed immediately before the -network request, so no future loop change can exceed the budget by accident. -Authentication failures, bad models, malformed responses, and other permanent -failures get no transport retry. +Final reporting has a shared hard budget of five Gemini HTTP calls across +initial generation, up to two semantic repairs, and transient retries. The +counter is consumed immediately before the network request. Authentication +failures, bad models, malformed responses, and other permanent failures get no +transport retry. + +After transient transport exhaustion or the overall report deadline, the +connected candidate may request one regeneration after a 30-second cooldown, +or for quota exhaustion until the first configured key leaves its cooldown, at +most 60 seconds; a sole key waits the whole 60, and a usable backup only the 30. +A missed deadline asks the key rotation the same question, and so does an +accepted request, so keys another interview sent back to quota during the wait +answer `early` rather than spend the regeneration on a call that cannot start. It +reuses the frozen assessment, has its own five-call pool and 125-second +deadline, and therefore bounds final reporting at ten HTTP calls per interview. +Successful or salvaged reports cannot be regenerated. A rejected credential +never qualifies on its own; an available or quota-limited backup may still +offer recovery, regardless of which credential failed last. An already +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 +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 +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. 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, @@ -106,9 +135,16 @@ limit: the final report still receives the complete transcript and editor state. The server admits at most `CODETRIAL_MAX_CONCURRENT_INTERVIEWS` live local agents, 16 by default. A reload of an already-live room reuses its slot; -completion and panic release it. Provider projects are rotated and their LiveKit -connection-minute quota is refreshed in the background, and known-exhausted -projects are skipped. +completion and panic release it. A report recovery window keeps its slot and +LiveKit connection until expiry or 30 seconds after candidate departure, at most +five minutes plus a 125-second regeneration and a 30-second rejoin. Provider outages can therefore fill the slots +with completed interviews awaiting recovery; admission still refuses excess +starts rather than exceeding this bound. A fixed local room refuses a new +start while its existing agent is finalizing or awaiting report recovery, as +`dispatch_refused ... reason=finalizing` rather than `at_capacity`; it cannot +reuse that agent as an interviewer for another candidate. Provider projects are +rotated and their LiveKit connection-minute quota is refreshed in the +background, and known-exhausted projects are skipped. The token endpoint allows 30 starts per 60-second bucket. Signed-in buckets are keyed by account and anonymous buckets by client address, so anonymous traffic @@ -126,7 +162,14 @@ and nothing else. Never candidate prompts, transcript, code, profile, job description or resume grounding, provider output, repair output, framework evidence, or personalized feedback. Report requests are built from the current session and each retry uses that same session's immutable prompt, so there is no -cross-session response cache. +cross-session response cache. A failed report may keep its frozen prompt and +assessment in the same agent's process memory for the five-minute recovery +window and one bounded regeneration. They are discarded on final outcome, +expiry, disconnection or process exit, never persisted or reused by another +session. The provisional failure report itself is not an input: the browser +saves it on arrival like any report, and the final outcome replaces it under +the same report id, so a tab lost during the window keeps the failure rather +than nothing. ## What the candidate sees @@ -143,10 +186,15 @@ never presents canned feedback as an agent evaluation. ## Operating it -Watch `codetrial dispatch_refused ... reason=at_capacity`, `livekit quota:` -transitions, token HTTP 429 with `Retry-After`, `gemini report -transport_failed call=... retry=...`, `codetrial live_usage ... outcome=billing`, -and the bounded incomplete-report categories. +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 +connection-minute quota like interview time. Raise concurrency only after checking provider minutes, Gemini limits, CPU and audio capacity, and the token burst policy. The deterministic dispatcher, diff --git a/scripts/gen-wire-fixtures.mjs b/scripts/gen-wire-fixtures.mjs index 74f9f6b9..f8e339fd 100755 --- a/scripts/gen-wire-fixtures.mjs +++ b/scripts/gen-wire-fixtures.mjs @@ -20,6 +20,9 @@ import { fileURLToPath } from "node:url"; const ROOT = join(dirname(fileURLToPath(import.meta.url)), ".."); const FIXTURES = join(ROOT, "tests", "fixtures"); +const { reportRecoveryLimits } = await import( + join(ROOT, "web", "report-recovery.js") +); const lib = await import(join(ROOT, "web", "lib.js")); const { ALL_LANGUAGES } = await import( join(ROOT, "web", "compiler-explorer.js") @@ -58,6 +61,7 @@ function codeUpdateCases(languages) { function controlCases() { return [ + { name: "retry report", payload: lib.retryReportPayload() }, { name: "thinking start", payload: lib.thinkingPayload(true) }, { name: "thinking end", payload: lib.thinkingPayload(false) }, { name: "yield turn", payload: lib.yieldTurnPayload() }, @@ -374,6 +378,7 @@ async function integrityChain() { // to recognize. const languages = ALL_LANGUAGES; const files = { + "report-recovery.json": reportRecoveryLimits, "code-update.json": { topic: lib.topics.code, cases: codeUpdateCases(languages), diff --git a/src/dispatch.rs b/src/dispatch.rs index 7665d5c2..df164815 100644 --- a/src/dispatch.rs +++ b/src/dispatch.rs @@ -5,11 +5,11 @@ //! an interview, long after the binary finished deciding what it is. Here it //! also gets tests that do not need a process. -use std::collections::HashSet; +use std::collections::HashMap; use std::sync::{Arc, Mutex}; use crate::config::{AgentConfig, Provider}; -use crate::web::RoomDispatcher; +use crate::web::{DispatchRefusal, RoomDispatcher}; /// Runs the interviewer for rooms this process just named, in this process. /// @@ -23,22 +23,29 @@ pub struct LocalDispatcher { /// that arrives with the room, from the same lookup that minted the token. pub config: AgentConfig, pub runtime: tokio::runtime::Handle, - pub live: Arc>>, + + // True while assessing; false while finalizing the report. Both hold + // capacity. + pub live: Arc>>, /// From `config::max_concurrent_interviews`, so the ceiling an operator set /// is the ceiling this enforces. pub max_concurrent: usize, } impl RoomDispatcher for LocalDispatcher { - fn ensure_agent(&self, room_name: &str, provider: &Provider) -> bool { + fn ensure_agent(&self, room_name: &str, provider: &Provider) -> Result<(), DispatchRefusal> { let slot = match self.reserve(room_name) { - Reservation::Existing => return true, + Reservation::Existing => return Ok(()), Reservation::Full => { eprintln!( "codetrial dispatch_refused room={room_name} reason=at_capacity limit={}", self.max_concurrent ); - return false; + return Err(DispatchRefusal::AtCapacity); + } + Reservation::Finalizing => { + eprintln!("codetrial dispatch_refused room={room_name} reason=finalizing"); + return Err(DispatchRefusal::Finalizing); } Reservation::New(slot) => slot, }; @@ -48,17 +55,18 @@ impl RoomDispatcher for LocalDispatcher { // outlives the response by the length of the interview. self.runtime.spawn(async move { eprintln!("codetrial dispatch room={}", slot.room_name); - if let Err(error) = crate::livekit::run_room( + if let Err(error) = crate::livekit::run_room_with_slot( &config, &slot.room_name, crate::web::current_epoch_seconds(), + Some(&slot), ) .await { eprintln!("codetrial agent_failed room={}: {error}", slot.room_name); } }); - true + Ok(()) } } @@ -66,6 +74,7 @@ enum Reservation { Existing, New(Slot), Full, + Finalizing, } impl LocalDispatcher { @@ -74,13 +83,20 @@ impl LocalDispatcher { // A reload mints a token for the same fixed room in local mode, and two // agents in one room evict each other. - if live.contains(room_name) { - return Reservation::Existing; + if let Some(assessing) = live.get(room_name) { + // A fixed local room cannot start a new interview while its old + // agent is finalizing. Reusing it would mint a token for an agent + // that only accepts the original candidate's report retry. + return if *assessing { + Reservation::Existing + } else { + Reservation::Finalizing + }; } if live.len() >= self.max_concurrent { return Reservation::Full; } - live.insert(room_name.to_string()); + live.insert(room_name.to_string(), true); Reservation::New(Slot { live: Arc::clone(&self.live), room_name: room_name.to_string(), @@ -91,10 +107,23 @@ impl LocalDispatcher { /// One interview's claim on this process's capacity, held for as long as the /// task that owns it. pub struct Slot { - pub live: Arc>>, + pub live: Arc>>, pub room_name: String, } +impl Slot { + pub(crate) fn assessment_finished(&self) { + if let Some(assessing) = self + .live + .lock() + .unwrap_or_else(|error| error.into_inner()) + .get_mut(&self.room_name) + { + *assessing = false; + } + } +} + impl Drop for Slot { fn drop(&mut self) { self.live diff --git a/src/gemini.rs b/src/gemini.rs index 06759bd7..9915de5c 100644 --- a/src/gemini.rs +++ b/src/gemini.rs @@ -20,10 +20,10 @@ use crate::runtime::{ mod credentials; pub use credentials::GeminiKeys; -pub(crate) use credentials::exhausted_until; use credentials::{ ApiFailure, ApiSurface, CredentialFailure, credential_failure, failure_from_reason, }; +pub(crate) use credentials::{QUOTA_COOLDOWN, exhausted_until}; const LIVE_WEBSOCKET_ENDPOINT: &str = "wss://generativelanguage.googleapis.com/ws/google.ai.generativelanguage.v1beta.GenerativeService.BidiGenerateContent"; const SETUP_TIMEOUT: Duration = Duration::from_secs(15); @@ -697,11 +697,34 @@ pub(crate) async fn generate_report_with_keys( prompt: &str, problem: &crate::agent::Problem, scope: &str, +) -> Result> { + generate_report_with_keys_at( + keys, + &gemini_generate_content_url(model), + REPORT_RETRY_BACKOFF, + prompt, + problem, + scope, + ) + .await +} + +/// The same report against another endpoint and first backoff. Private: the +/// tests reach it through `gemini::tests::generate_report_at` to drive the +/// whole path from a local server, and nothing else targets another URL. +async fn generate_report_with_keys_at( + keys: &GeminiKeys, + url: &str, + backoff: Duration, + prompt: &str, + problem: &crate::agent::Problem, + scope: &str, ) -> Result> { let mut calls = ReportCalls { keys, - url: gemini_generate_content_url(model), + url: url.to_string(), budget: ReportCallBudget::new(), + backoff, scope, }; let (report, salvaged) = report_attempts(prompt, problem, &mut calls).await?; @@ -724,6 +747,7 @@ struct ReportCalls<'a> { keys: &'a GeminiKeys, url: String, budget: ReportCallBudget, + backoff: Duration, scope: &'a str, } @@ -737,7 +761,7 @@ impl ReportTransport for ReportCalls<'_> { &self.url, prompt, &mut self.budget, - REPORT_RETRY_BACKOFF, + self.backoff, self.scope, ) } @@ -948,7 +972,25 @@ async fn generate_report_transport( eprintln!( "gemini report transport_failed room={scope} call={call} final=true error={detail}" ); - return Err(io::Error::other(detail).into()); + let retry_after = match keys.select_report() { + Err(error) => report_regeneration_retry_after(&error), + Ok(_) => match failure { + None => is_retryable(error.as_ref()).then_some(Duration::ZERO), + + // A backup that selects now is usable now. A sole key is + // never taken out of rotation, so it selects again at once + // and has to sit out the cooldown itself. + Some(CredentialFailure::Quota) if keys.has_backups() => Some(Duration::ZERO), + Some(CredentialFailure::Quota) => Some(QUOTA_COOLDOWN), + Some(_) if retryable && budget.is_exhausted() => Some(Duration::ZERO), + Some(_) => None, + }, + }; + return Err(ReportTransportFailure { + detail, + retry_after, + } + .into()); }; failures += 1; @@ -963,17 +1005,31 @@ async fn generate_report_transport( "gemini report transport_failed room={scope} call={call} backoff_s={} error={detail}", backoff.as_secs() ); - tokio::time::sleep(backoff).await; - - // Chosen again after the wait, which another interview may have spent - // ruling this key out. - let Ok(next) = keys.select_report() else { - return Err(io::Error::other(detail).into()); - }; - api_key = next; + api_key = report_key_after_backoff(keys, backoff, detail, scope, call).await?; } } +/// The failed HTTP call was already counted before the wait. A rotation that +/// becomes unavailable during it must preserve that call's diagnosis without +/// logging a second transport failure for the same request. +async fn report_key_after_backoff( + keys: &GeminiKeys, + backoff: Duration, + detail: String, + scope: &str, + call: usize, +) -> Result> { + tokio::time::sleep(backoff).await; + keys.select_report().map_err(|error| { + eprintln!("gemini report retry_unavailable room={scope} call={call}"); + ReportTransportFailure { + detail, + retry_after: report_regeneration_retry_after(&error), + } + .into() + }) +} + /// The wait after the `failure`th transport failure of one call, counting from /// one. fn report_retry_backoff(first: Duration, failure: u32) -> Duration { @@ -1014,6 +1070,33 @@ fn repair_prompt(original: &str, invalid: &str, errors: &[String]) -> String { ) } +#[derive(Debug)] +struct ReportTransportFailure { + detail: String, + retry_after: Option, +} + +impl std::fmt::Display for ReportTransportFailure { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str(&self.detail) + } +} + +impl std::error::Error for ReportTransportFailure {} + +pub(crate) fn report_regeneration_retry_after( + error: &(dyn std::error::Error + 'static), +) -> Option { + if let Some(failure) = error.downcast_ref::() { + return failure.retry_after; + } + + // The rotation knows when its first key is back, which is usually sooner + // than a whole cooldown from now. + credentials::exhausted_until(error) + .map(|at| at.saturating_duration_since(std::time::Instant::now())) +} + /// Transient upstream conditions only. A bad key or a bad model is answered the /// same way every time, so retrying it just makes the candidate wait longer. /// @@ -1936,6 +2019,7 @@ fn parse_server_message(text: &str) -> ServerMessage { } } +// Visible to the crate so other modules' tests can share its report fixtures. #[cfg(test)] #[path = "../tests/unit/gemini.rs"] -mod tests; +pub(crate) mod tests; diff --git a/src/gemini/credentials.rs b/src/gemini/credentials.rs index 468a0b1b..71a25506 100644 --- a/src/gemini/credentials.rs +++ b/src/gemini/credentials.rs @@ -7,7 +7,7 @@ use serde_json::Value; use crate::config::AgentConfig; -const QUOTA_COOLDOWN: Duration = Duration::from_secs(60); +pub(crate) const QUOTA_COOLDOWN: Duration = Duration::from_secs(60); // As long as the web pool trusts a refused LiveKit credential verdict. A // refusal that was about the project rather than the key, and so marked every @@ -63,9 +63,10 @@ struct Cooldown { live: Option, report: Option, - // Read only to decide whether an exhausted Live rotation is worth waiting - // out. It never outlives `live`, so pruning and selection ignore it. + // Rejections cannot recover just by waiting. Each marker expires no later + // than its surface deadline, so pruning needs only those deadlines. live_rejected: Option, + report_rejected: Option, // Read only to say why a rotation ran dry, in the same way. billing: Option, @@ -146,7 +147,33 @@ impl GeminiKeys { self.select_for(ApiSurface::Report) } + /// A rotation whose every key is out on report quota, for callers outside + /// this module that decide from the rotation's answer. + #[cfg(test)] + pub(crate) fn report_quota_exhausted(keys: &[&str]) -> Self { + let rotation = Self::new(keys.iter().map(|key| key.to_string()).collect()); + for key in keys { + rotation.failed(key, CredentialFailure::Quota, ApiSurface::Report); + } + rotation + } + + #[cfg(test)] + pub(super) fn select_report_at(&self, now: Instant) -> Result { + // A simulated future must not expire another test's credentials. + self.select_for_at(ApiSurface::Report, now, false) + } + fn select_for(&self, surface: ApiSurface) -> Result { + self.select_for_at(surface, Instant::now(), true) + } + + fn select_for_at( + &self, + surface: ApiSurface, + now: Instant, + prune: bool, + ) -> Result { // Short of the shared map as well, which another interview's list // holding the same key string would otherwise write for it. if !self.has_backups() { @@ -160,10 +187,11 @@ impl GeminiKeys { .get_or_init(Mutex::default) .lock() .unwrap_or_else(|error| error.into_inner()); - let now = Instant::now(); - cooldowns.retain(|_, cooldown| { - cooling_down(cooldown.live, now) || cooling_down(cooldown.report, now) - }); + if prune { + cooldowns.retain(|_, cooldown| { + cooling_down(cooldown.live, now) || cooling_down(cooldown.report, now) + }); + } let mut current = self .current .lock() @@ -189,8 +217,8 @@ impl GeminiKeys { }) }) .ok_or_else(|| { - // The earliest Live key to come back that was out on quota - // alone. A refused key would only be refused again. + // A quota deadline can recover either surface. A rejected + // credential must not masquerade as a temporary rate limit. let retry_at = match surface { ApiSurface::Live => self .keys @@ -199,7 +227,13 @@ impl GeminiKeys { .filter(|cooldown| !cooling_down(cooldown.live_rejected, now)) .filter_map(|cooldown| cooldown.live) .min(), - ApiSurface::Report => None, + ApiSurface::Report => self + .keys + .iter() + .filter_map(|key| cooldowns.get(key)) + .filter(|cooldown| !cooling_down(cooldown.report_rejected, now)) + .filter_map(|cooldown| cooldown.report) + .min(), }; let billing = self.keys.iter().any(|key| { cooldowns @@ -233,6 +267,7 @@ impl GeminiKeys { extend(&mut entry.live, rejected); extend(&mut entry.report, rejected); extend(&mut entry.live_rejected, rejected); + extend(&mut entry.report_rejected, rejected); if failure == CredentialFailure::Billing { extend(&mut entry.billing, rejected); } @@ -242,7 +277,10 @@ impl GeminiKeys { extend(&mut entry.live, rejected); extend(&mut entry.live_rejected, rejected); } - ApiSurface::Report => extend(&mut entry.report, rejected), + ApiSurface::Report => { + extend(&mut entry.report, rejected); + extend(&mut entry.report_rejected, rejected); + } }, CredentialFailure::Quota => extend( match surface { @@ -292,7 +330,7 @@ pub(super) fn exhausted_by_billing(error: &(dyn std::error::Error + 'static)) -> as_exhausted(error).is_some_and(|exhausted| exhausted.billing) } -/// When an exhausted Live rotation has a key back from its quota cooldown. +/// When an exhausted rotation has a key back from its quota cooldown. pub(crate) fn exhausted_until(error: &(dyn std::error::Error + 'static)) -> Option { as_exhausted(error)?.retry_at } diff --git a/src/livekit.rs b/src/livekit.rs index 69e8e378..35440674 100644 --- a/src/livekit.rs +++ b/src/livekit.rs @@ -108,7 +108,7 @@ use session::{ send_wrap_up_and_wait, set_agent_state, }; -use report::{freeze_report_prompt, generate_report_bounded, publish_report}; +use report::{freeze_assessment, generate_report_bounded}; use rooms::{evict_duplicate_agent, isolate_local_agent}; // Re-exported rather than merely used: `web::setup` calls this to check the @@ -192,8 +192,14 @@ impl CandidatePresence { /// Present, or absent for less than the grace, both mean carry on. fn gave_up(&self, now: Instant) -> bool { - self.left_at - .is_some_and(|left| now.duration_since(left) >= CANDIDATE_ABSENCE_LIMIT) + self.deadline(CANDIDATE_ABSENCE_LIMIT) + .is_some_and(|deadline| now >= deadline) + } + + /// When a grace of `limit` runs out, if the candidate is away. The report + /// recovery wait holds the same rule with a shorter limit. + fn deadline(&self, limit: Duration) -> Option { + self.left_at.map(|left| left + limit) } } @@ -1143,7 +1149,7 @@ fn take_interim_review_window(state: &mut RuntimeState, boot: &RuntimeBootstrap< /// caller's, carried only by `boot`. struct OpenSession<'a> { room: Room, - events: tokio::sync::mpsc::UnboundedReceiver, + events: tokio::sync::Mutex>, agent_identity: String, candidate_identity: String, boot: RuntimeBootstrap<'a>, @@ -1185,6 +1191,7 @@ async fn open_session<'a>( eprintln!("no candidate joined room={room_name}; leaving it"); return Ok(None); }; + let events = tokio::sync::Mutex::new(events); let boot = candidate_bootstrap(config, room_name, Some(&candidate_metadata)); eprintln!( "starting interview: room={} problem={} duration={}min", @@ -1213,6 +1220,9 @@ async fn open_session<'a>( keys, boot: &boot, started_at: setup_began, + candidate_identity: &candidate_identity, + events: &events, + slot: None, }; crate::gemini::first_open_within( crate::gemini::FIRST_OPEN_LIMIT, @@ -1967,6 +1977,15 @@ pub async fn run_room( config: &AgentConfig, room_name: &str, now_seconds: u64, +) -> Result<(), Box> { + run_room_with_slot(config, room_name, now_seconds, None).await +} + +pub(crate) async fn run_room_with_slot( + config: &AgentConfig, + room_name: &str, + now_seconds: u64, + slot: Option<&crate::dispatch::Slot>, ) -> Result<(), Box> { // Names this run on every usage line. Milliseconds rather than the // dispatch's seconds, so a room retried within the same second is not @@ -1979,7 +1998,7 @@ pub async fn run_room( let keys = Arc::new(GeminiKeys::from_config(config)); let Some(OpenSession { room, - mut events, + events, agent_identity, candidate_identity, boot, @@ -2019,6 +2038,9 @@ pub async fn run_room( keys: &keys, boot: &boot, started_at, + candidate_identity: &candidate_identity, + events: &events, + slot, }; let ids = RoomIdentities { room_name, @@ -2063,7 +2085,7 @@ pub async fn run_room( turn.context(&mut output_audio, &mut gemini, &mut media); on_watch_tick(&room, &mut context, &mut loops, interview).await? } - event = events.recv() => { + event = async { events.lock().await.recv().await } => { let Some(event) = event else { eprintln!( "LiveKit event stream ended for room={room_name}; ending with no report, \ @@ -2507,6 +2529,9 @@ struct InterviewContext<'a> { keys: &'a Arc, boot: &'a RuntimeBootstrap<'a>, started_at: Instant, + candidate_identity: &'a str, + events: &'a tokio::sync::Mutex>, + slot: Option<&'a crate::dispatch::Slot>, } /// The interview's starting state, from the plan the token was minted for. @@ -2624,7 +2649,7 @@ async fn settle_hold( if result.thinking_changed == Some(true) { if let Err(error) = begin_button_hold( context.gemini, - context.candidate_audio, + &mut context.media.audio_bytes, context.activity, at, grace, @@ -2663,7 +2688,7 @@ async fn settle_hold( } let finalize_prompt = reply.clone(); let finalize = async { - flush_audio(context.gemini, context.candidate_audio).await?; + flush_audio(context.gemini, &mut context.media.audio_bytes).await?; if context.turns_candidate_open() && let Some(prompt) = reply.as_deref() { @@ -2910,6 +2935,9 @@ async fn handle_data_packet( let Some(reason) = result.finish_interview else { return Ok(ControlFlow::Continue(())); }; + if let Some(slot) = interview.slot { + slot.assessment_finished(); + } session::flush_thinking_notice(room, context.state).await; // The assessment ends here, before the goodbye: the reducer has closed the @@ -2920,18 +2948,20 @@ async fn handle_data_packet( // can add is the skips a started behavioral round leaves to it, and the // report scores those steps as unassessed with or without them. // - // Not `?`, and neither is anything else between here and `publish_report`. - // The turns go out as a LiveKit text stream and the report as a data - // packet, so the one failing says nothing about the other, and a final - // segment the panel never saw is cosmetic where a missing report is not. + // Not `?`, and neither is anything else between here and + // `publish_with_recovery`. The turns go out as a LiveKit text stream and + // the report as a data packet, so the one failing says nothing about the + // other, and a final segment the panel never saw is cosmetic where a + // missing report is not. if let Err(error) = close_turns(room, context).await { eprintln!("closing the last turns failed ({error}); writing the report anyway"); } - let prompt = freeze_report_prompt( + let assessment = freeze_assessment( interview.boot, context.state, interview.started_at.elapsed().as_secs_f64() / 60.0, ); + let mut recovery_events = interview.events.lock().await; let api_key = &**interview.keys; let farewell = async { // The goodbye is the only part of the ending that needs the Live @@ -2953,24 +2983,32 @@ async fn handle_data_packet( } }; let (generated, ()) = tokio::join!( - generate_report_bounded(interview.boot, &prompt, api_key), + generate_report_bounded(interview.boot, &assessment.prompt, api_key), farewell ); - publish_report( - room, - interview.boot, - context.state, - &reason, - api_key, + context.media.audio = None; + context.media.video = None; + context.media.audio_bytes.clear(); + let close_live = async { + if let Err(error) = context.gemini.shutdown().await { + eprintln!("Gemini close failed ({error}); leaving anyway"); + } + }; + report::publish_with_recovery( + report::ReportRecovery { + boot: interview.boot, + assessment, + reason: &reason, + keys: api_key, + candidate: interview.candidate_identity, + }, generated, + &mut recovery_events, + room, + close_live, ) .await?; - // The report is out; a close that fails now changes nothing but whether the - // agent leaves, and it has to. - if let Err(error) = context.gemini.shutdown().await { - eprintln!("Gemini close failed ({error}); leaving anyway"); - } // Give the report packet a moment to leave before the agent goes. tokio::time::sleep(Duration::from_millis(250)).await; leave_room(room).await; diff --git a/src/livekit/report.rs b/src/livekit/report.rs index ea143344..5da31814 100644 --- a/src/livekit/report.rs +++ b/src/livekit/report.rs @@ -1,8 +1,9 @@ //! Building the report packet the browser receives when an interview ends. //! -//! One region because it is one output. The interview loop freezes the prompt -//! with `freeze_report_prompt`, runs `generate_report_bounded` beside the -//! farewell, and hands what came back to `publish_report`; everything below +//! One region because it is one output. The interview loop freezes the +//! assessment with `freeze_assessment`, runs `generate_report_bounded` beside +//! the farewell, and hands what came back to `publish_with_recovery`, which +//! publishes it or offers one regeneration first; everything below //! is how the packet is assembled, and the pieces are separated so that a //! failure in one of them is a note in the report rather than no report at //! all. @@ -29,7 +30,7 @@ pub(super) type GeneratedReport = Result< /// The report prompt, built and counted once the interview's assessment is /// over and before the farewell is spoken, so the call can run while it plays. -pub(super) fn freeze_report_prompt( +fn freeze_report_prompt( boot: &RuntimeBootstrap<'_>, state: &mut RuntimeState, elapsed_min: f64, @@ -45,6 +46,23 @@ pub(super) fn freeze_report_prompt( prompt } +pub(super) struct FrozenAssessment { + pub prompt: String, + state: RuntimeState, +} + +pub(super) fn freeze_assessment( + boot: &RuntimeBootstrap<'_>, + state: &mut RuntimeState, + elapsed_min: f64, +) -> FrozenAssessment { + let prompt = freeze_report_prompt(boot, state, elapsed_min); + FrozenAssessment { + prompt, + state: state.clone(), + } +} + /// The report call under `REPORT_TIMEOUT`. Borrows nothing of the interview /// state, which is what lets it run beside the farewell that still needs it. pub(super) async fn generate_report_bounded( @@ -65,27 +83,435 @@ pub(super) async fn generate_report_bounded( .await } -pub(super) async fn publish_report( +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); + +/// How long the candidate waits before a regeneration may start, or `None` +/// when the failure is not one waiting can fix. A deadline says nothing about +/// the keys, so it asks the rotation itself: a retry offered while every key is +/// still out on quota would fail before its first call and spend the one retry. +fn regeneration_cooldown( + generated: &GeneratedReport, + keys: &GeminiKeys, +) -> Option { + let retry_after = match generated { + Err(_) => match report_readiness(keys) { + Readiness::Ready => Some(std::time::Duration::ZERO), + Readiness::Wait(delay) => Some(delay), + Readiness::Never => None, + }, + Ok(Err(error)) => crate::gemini::report_regeneration_retry_after(error.as_ref()), + Ok(Ok(_)) => None, + }; + retry_after.map(|delay| delay.max(REPORT_RETRY_COOLDOWN)) +} + +/// Whether the key rotation can make a report call now, and if not, whether +/// waiting would help. Another interview sharing the keys can put them back on +/// quota at any moment, so this is asked again when a retry arrives rather than +/// trusted from when the offer went out. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum Readiness { + Ready, + Wait(std::time::Duration), + Never, +} + +fn report_readiness(keys: &GeminiKeys) -> Readiness { + match keys.select_report() { + Ok(_) => Readiness::Ready, + Err(error) => match crate::gemini::report_regeneration_retry_after(&error) { + Some(delay) => Readiness::Wait(delay), + None => Readiness::Never, + }, + } +} + +fn recovery_request( + topic: Option<&str>, + sender: Option<&str>, + candidate: &str, + payload: &[u8], +) -> bool { + super::interview_packet(topic, sender, candidate, payload).is_some_and(|(topic, payload)| { + topic == crate::runtime::TOPIC_CONTROL && payload["type"] == "retry_report" + }) +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum RecoveryEvent { + Retry, + /// The room itself is gone, so nothing published can arrive. + Left, + /// The candidate left, which a full LiveKit rejoin under the same identity + /// also looks like until `Back`. + Away, + Back, + Ignore, +} + +/// How long a candidate who left may take to rejoin before the wait gives up. +/// A full LiveKit rejoin after a network drop leaves and returns under the same +/// identity, and treating the leave as final spent the one retry on a candidate +/// who never went anywhere. +const REJOIN_GRACE: std::time::Duration = std::time::Duration::from_secs(30); + +/// The agent's answer to a retry, and its word that the window closed. Without +/// them the page had to guess both from its own clock, which starts later than +/// this one and stops for a reconnect this one never sees. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum RecoveryNotice { + Accepted, + Early(std::time::Duration), + Closed, +} + +/// A wait as the page is told it, rounded up, so a page that waits exactly +/// this long is not early. +fn whole_seconds(wait: std::time::Duration) -> u64 { + wait.as_millis().div_ceil(1000) as u64 +} + +fn recovery_notice(notice: RecoveryNotice) -> serde_json::Value { + match notice { + RecoveryNotice::Accepted => { + serde_json::json!({ "type": "report_retry", "status": "accepted" }) + } + RecoveryNotice::Early(wait) => serde_json::json!({ + "type": "report_retry", + "status": "early", + + "retryAfterSeconds": whole_seconds(wait).max(1), + }), + RecoveryNotice::Closed => serde_json::json!({ "type": "report_retry", "status": "closed" }), + } +} + +/// The room as recovery uses it: what arrives, and what goes back out. A +/// separate seam from the room itself so the wait can be driven by a script. +trait RecoveryRoom { + /// Whether the candidate is in the room now. Their departure is an event + /// the interview loop may already have consumed, and a wait for a retry + /// from nobody held the slot for the whole window. + 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; + fn publish( + &mut self, + packet: DataPacket, + ) -> impl std::future::Future>> + Send; +} + +struct LiveRecoveryRoom<'a> { + room: &'a Room, + candidate: &'a str, + events: &'a mut tokio::sync::mpsc::UnboundedReceiver<::livekit::RoomEvent>, +} + +/// What one room event means to a recovery wait. `None` is the event stream +/// ending, which only happens once the room is gone. +fn recovery_event(event: Option<::livekit::RoomEvent>, candidate: &str) -> RecoveryEvent { + match event { + None | Some(::livekit::RoomEvent::Disconnected { .. }) => RecoveryEvent::Left, + Some(::livekit::RoomEvent::ParticipantDisconnected(p)) => { + candidate_only(&p.identity().0, candidate, RecoveryEvent::Away) + } + Some(::livekit::RoomEvent::ParticipantConnected(p)) => { + candidate_only(&p.identity().0, candidate, RecoveryEvent::Back) + } + Some(::livekit::RoomEvent::DataReceived { + topic, + payload, + participant, + .. + }) => { + let sender = participant.as_ref().map(|p| p.identity().0); + if recovery_request(topic.as_deref(), sender.as_deref(), candidate, &payload) { + RecoveryEvent::Retry + } else { + RecoveryEvent::Ignore + } + } + _ => RecoveryEvent::Ignore, + } +} + +/// Only the candidate coming and going matters. Anyone else, an observer or +/// the recording egress, says nothing about whether a retry can come. +fn candidate_only(identity: &str, candidate: &str, event: RecoveryEvent) -> RecoveryEvent { + if identity == candidate { + event + } else { + RecoveryEvent::Ignore + } +} + +impl RecoveryRoom for LiveRecoveryRoom<'_> { + fn candidate_present(&self) -> bool { + self.room + .remote_participants() + .values() + .any(|participant| participant.identity().0 == self.candidate) + } + + async fn next(&mut self) -> RecoveryEvent { + recovery_event(self.events.recv().await, self.candidate) + } + + 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()), + }; + if let Err(error) = sent { + eprintln!("codetrial report_recovery_notice_failed notice={notice:?} error={error}"); + } + } + + async fn publish( + &mut self, + packet: DataPacket, + ) -> Result<(), Box> { + self.room.local_participant().publish_data(packet).await?; + Ok(()) + } +} + +/// Sleeps until an absent candidate's rejoin grace runs out, or forever while +/// they are present. +async fn rejoin_expired(presence: &super::CandidatePresence) { + match presence.deadline(REJOIN_GRACE) { + Some(deadline) => { + tokio::time::sleep_until(tokio::time::Instant::from_std(deadline)).await; + } + None => std::future::pending().await, + } +} + +/// Follows the candidate's comings and goings, `Break` once the room itself is +/// gone. A retry can only come from the candidate, so it shows them present. +fn track_presence( + event: RecoveryEvent, + presence: &mut super::CandidatePresence, +) -> std::ops::ControlFlow<()> { + match event { + RecoveryEvent::Left => return std::ops::ControlFlow::Break(()), + + // The tokio clock, read as a std instant, so a paused test clock drives + // the grace the same way the real one does. + RecoveryEvent::Away => presence.left(tokio::time::Instant::now().into_std()), + RecoveryEvent::Back | RecoveryEvent::Retry => presence.returned(), + RecoveryEvent::Ignore => {} + } + std::ops::ControlFlow::Continue(()) +} + +async fn recover_report( + events: &mut impl RecoveryRoom, + generation: impl std::future::Future, + window: std::time::Duration, + cooldown: std::time::Duration, + started: tokio::time::Instant, + readiness: impl Fn() -> Readiness, +) -> Option { + let mut presence = super::CandidatePresence::default(); + loop { + let event = tokio::select! { + biased; + _ = tokio::time::sleep_until(started + window) => { + events.notify(RecoveryNotice::Closed).await; + return None; + } + _ = rejoin_expired(&presence) => return None, + event = events.next() => event, + }; + if track_presence(event, &mut presence).is_break() { + return None; + } + let RecoveryEvent::Retry = event else { + continue; + }; + let waited = started.elapsed(); + if waited < cooldown { + events + .notify(RecoveryNotice::Early(cooldown - waited)) + .await; + continue; + } + match readiness() { + Readiness::Ready => { + events.notify(RecoveryNotice::Accepted).await; + break; + } + Readiness::Wait(delay) => events.notify(RecoveryNotice::Early(delay)).await, + Readiness::Never => { + events.notify(RecoveryNotice::Closed).await; + return None; + } + } + } + + // The generation keeps running while the candidate is away within the + // grace, and what it returns waits for them to be back before it goes out. + tokio::pin!(generation); + let mut generated = None; + loop { + if presence.deadline(REJOIN_GRACE).is_none() + && let Some(result) = generated.take() + { + return Some(result); + } + tokio::select! { + _ = rejoin_expired(&presence) => return None, + event = events.next() => { + if track_presence(event, &mut presence).is_break() { + return None; + } + } + result = &mut generation, if generated.is_none() => generated = Some(result), + } + } +} + +pub(super) struct ReportRecovery<'a> { + pub boot: &'a RuntimeBootstrap<'a>, + pub assessment: FrozenAssessment, + pub reason: &'a str, + pub keys: &'a GeminiKeys, + pub candidate: &'a str, +} + +/// Recovery runs a separate event loop: the ended interview must never pump +/// new media into Gemini or interpret its deliberate shutdown as a reconnect. +pub(super) async fn publish_with_recovery( + recovery: ReportRecovery<'_>, + generated: GeneratedReport, + events: &mut tokio::sync::mpsc::UnboundedReceiver<::livekit::RoomEvent>, room: &Room, - boot: &RuntimeBootstrap<'_>, - state: &mut RuntimeState, - reason: &str, - api_key: &GeminiKeys, + close_live: impl std::future::Future, +) -> Result<(), Box> { + let ReportRecovery { + boot, + assessment, + reason, + keys, + candidate, + } = recovery; + let FrozenAssessment { prompt, mut state } = assessment; + let mut room = LiveRecoveryRoom { + room, + candidate, + events, + }; + run_recovery( + &mut room, + RecoveryReport { + boot, + state: &mut state, + reason, + keys, + }, + generated, + generate_report_bounded(boot, &prompt, keys), + tokio::time::Instant::now, + close_live, + ) + .await +} + +/// What each published report is built from. The frozen state, never the live +/// one, so a regenerated report carries the same evidence as the failure it +/// replaces. +struct RecoveryReport<'a> { + boot: &'a RuntimeBootstrap<'a>, + state: &'a mut RuntimeState, + reason: &'a str, + keys: &'a GeminiKeys, +} + +/// `clock` names when the offer went out. The agent reads the real clock; a +/// test reads one already past the cooldown rather than pausing time under a +/// real HTTP exchange. `close_live` runs exactly once, on every path. +async fn run_recovery( + room: &mut impl RecoveryRoom, + report: RecoveryReport<'_>, generated: GeneratedReport, + regenerate: impl std::future::Future, + clock: fn() -> tokio::time::Instant, + close_live: impl std::future::Future, ) -> Result<(), Box> { - room.local_participant() - .publish_data(report_packet(boot, state, reason, api_key, generated)?) - .await?; + let RecoveryReport { + boot, + state, + reason, + 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. + 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 mut provisional = report_value(boot, state, reason, keys, generated); + provisional["reportRecovery"] = recovery_metadata(cooldown); + let started = clock(); + let published: Result<(), Box> = + async { room.publish(report_data_packet(provisional)?).await }.await; + close_live.await; + published?; + if let Some(generated) = recover_report( + room, + regenerate, + REPORT_RECOVERY_WINDOW, + cooldown, + started, + || report_readiness(keys), + ) + .await + { + room.publish(report_packet(boot, state, reason, keys, generated)?) + .await?; + } Ok(()) } +fn recovery_metadata(cooldown: std::time::Duration) -> serde_json::Value { + serde_json::json!({ + "expiresInSeconds": REPORT_RECOVERY_WINDOW.as_secs(), + "retryAfterSeconds": whole_seconds(cooldown), + }) +} + fn report_packet( boot: &RuntimeBootstrap<'_>, state: &mut RuntimeState, reason: &str, - api_key: &GeminiKeys, + keys: &GeminiKeys, generated: GeneratedReport, ) -> Result> { + Ok(report_data_packet(report_value( + boot, state, reason, keys, generated, + ))?) +} + +fn report_value( + boot: &RuntimeBootstrap<'_>, + state: &mut RuntimeState, + reason: &str, + api_key: &GeminiKeys, + generated: GeneratedReport, +) -> serde_json::Value { let mut report = match generated { Ok(Ok(raw)) => final_report(Some(&raw), state.hints_used, None, boot.problem), Ok(Err(error)) => final_report( @@ -126,9 +552,7 @@ fn report_packet( }; stamp_report_debrief(&mut report, boot, state); stamp_report_contract(&mut report); - Ok(report_data_packet(report_with_integrity_events( - report, state, reason, - ))?) + report_with_integrity_events(report, state, reason) } /// The teaching material that becomes useful only after an interview ends. diff --git a/src/livekit/session.rs b/src/livekit/session.rs index 70063302..a80f0f71 100644 --- a/src/livekit/session.rs +++ b/src/livekit/session.rs @@ -110,8 +110,7 @@ pub(super) struct GeminiEventContext<'a> { pub(super) agent_state: &'a mut String, pub(super) activity: &'a mut RuntimeActivity, pub(super) turns: &'a mut SpeakerTurns, - candidate_identity: Option<&'a str>, - pub(super) candidate_audio: &'a mut Vec, + pub(super) media: &'a mut CandidateMedia, } impl GeminiEventContext<'_> { @@ -565,7 +564,8 @@ async fn on_input_transcript( text: &str, interruptible: Interruptible, ) -> Result<(), Box> { - if let (Some(_), Some(identity)) = (transcript_text(text), context.candidate_identity) { + let candidate_identity = context.media.identity.clone(); + if let (Some(_), Some(identity)) = (transcript_text(text), candidate_identity.as_deref()) { // The candidate is talking over audio Gemini finished producing a while // ago. Gemini will not call this an interruption, because as far as it // is concerned that turn ended when it stopped generating; only this @@ -1361,7 +1361,7 @@ pub(super) async fn close_turns( room: &Room, context: &mut GeminiEventContext<'_>, ) -> Result<(), Box> { - let candidate_identity = context.candidate_identity.map(str::to_string); + let candidate_identity = context.media.identity.clone(); let order = closing_order( context.turns.interviewer.transcript_line(), @@ -1431,8 +1431,7 @@ impl TurnState { agent_state: &mut self.agent_state, activity: &mut self.activity, turns: &mut self.turns, - candidate_identity: media.identity.as_deref(), - candidate_audio: &mut media.audio_bytes, + media, } } } diff --git a/src/web/mod.rs b/src/web/mod.rs index 3a823254..9089b5aa 100644 --- a/src/web/mod.rs +++ b/src/web/mod.rs @@ -528,13 +528,26 @@ impl AppState { /// then holds a token for one LiveKit project while the interviewer waits in /// another. There is one lookup, and this is its result. /// -/// Returns `false` when no interviewer will come, so the caller can refuse +/// Returns the refusal when no interviewer will come, so the caller can refuse /// instead of handing out a token for a room nobody will ever join. /// /// Implementations must be idempotent per room and must not block: this is /// called on the request path. pub trait RoomDispatcher: Send + Sync + 'static { - fn ensure_agent(&self, room_name: &str, provider: &crate::config::Provider) -> bool; + fn ensure_agent( + &self, + room_name: &str, + provider: &crate::config::Provider, + ) -> Result<(), DispatchRefusal>; +} + +/// Why no interviewer will come. The two answer the candidate differently: a +/// full server frees up as any interview ends, a finalizing room only when its +/// own report is out. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DispatchRefusal { + AtCapacity, + Finalizing, } pub fn web_service(config: WebServerConfig) -> IntoMakeServiceWithConnectInfo { diff --git a/src/web/token.rs b/src/web/token.rs index aaacc115..11274b35 100644 --- a/src/web/token.rs +++ b/src/web/token.rs @@ -441,9 +441,9 @@ async fn dispatch_or_release_consent( let Some(dispatcher) = &state.dispatcher else { return None; }; - if dispatcher.ensure_agent(room_name, provider) { + let Err(refusal) = dispatcher.ensure_agent(room_name, provider) else { return None; - } + }; if let Some(interview) = interview { let accounts = accounts.clone(); let interview = interview.clone(); @@ -459,11 +459,17 @@ async fn dispatch_or_release_consent( // Refusing is the honest failure. Handing out the token anyway would put // the candidate in an empty room reading "Waiting" with nothing, on screen // or in any log they can see, saying why. + let error = match refusal { + super::DispatchRefusal::AtCapacity => { + "The server is running as many interviews as it can right now. Try again in a few minutes." + } + super::DispatchRefusal::Finalizing => { + "This room is still finishing the previous interview's report. Try again in a few minutes." + } + }; Some(json_response( StatusCode::SERVICE_UNAVAILABLE, - json!({ - "error": "The server is running as many interviews as it can right now. Try again in a few minutes." - }), + json!({ "error": error }), )) } diff --git a/tests/agent/wire.rs b/tests/agent/wire.rs index 7472ef8b..4b760c09 100644 --- a/tests/agent/wire.rs +++ b/tests/agent/wire.rs @@ -946,6 +946,30 @@ fn cited_source_ids_are_bounded_the_way_the_browser_bounds_them() { ); } +/// A report retry is answered only by the recovery wait after an interview +/// ends. Reaching the live reducer, it must be nothing: not an end, not a +/// reply, and not evidence, or a stale button could steer an interview. +#[test] +fn a_report_retry_during_a_live_interview_changes_nothing() { + let (topic, cases) = wire_fixture(include_str!("../fixtures/control.json")); + let mut state = RuntimeState { + code: "return 1".into(), + ..RuntimeState::default() + }; + let before = format!("{state:?}"); + let result = apply_data_event( + &mut state, + &topic, + wire_case(&cases, "retry report"), + TEST_REACTION_COOLDOWN_S, + ); + assert_eq!(result, DataEventResult::default()); + // Counted as received, as every packet is, and nothing else. + assert_eq!(state.evidence_ledger.metrics.raw_events_received, 1); + state.evidence_ledger.metrics.raw_events_received = 0; + assert_eq!(format!("{state:?}"), before); +} + #[test] fn browser_thinking_controls_keep_evidence_live_and_yield_the_floor() { let (topic, cases) = wire_fixture(include_str!("../fixtures/control.json")); diff --git a/tests/browser/account.test.js b/tests/browser/account.test.js index e052bb00..afce4b2a 100644 --- a/tests/browser/account.test.js +++ b/tests/browser/account.test.js @@ -186,7 +186,7 @@ test("the lobby offers one interview and carries no mode to the room", () => { interview, "JSON.stringify({ problemId: problem.page, durationMin, interviewId, interviewLoop, interviewProfile, ...(interviewGrounding", ); - assertIncludesCompact(interview, "interviewLoop, report: state.report"); + assertIncludesCompact(interview, "durationMin, interviewLoop, report, };"); }); test("interview loop is explicit, budgeted, gated, and carried into artifacts", () => { @@ -337,6 +337,67 @@ test("report history writes local storage before account sync", async () => { assert.equal(posts[0].options.method, "POST"); }); +test("a report saved again under its id replaces the earlier copy", async () => { + const storage = memoryStorage(); + const fetcher = async () => response({ signedIn: false }); + await saveReportHistory( + { id: "older", problemId: "two-sum", report: { decision: "HIRE" } }, + { storage, fetcher }, + ); + for (const summary of ["provisional failure", "regenerated"]) { + await saveReportHistory( + { + id: "interview-report", + problemId: "two-sum", + report: { incomplete: summary !== "regenerated", summary }, + }, + { storage, fetcher }, + ); + } + for (const rows of [readLocalHistory(storage), readReviewHistory(storage)]) { + assert.deepEqual( + rows.map((row) => row.id), + ["interview-report", "older"], + ); + assert.equal(rows[0].report.summary, "regenerated"); + } +}); + +test("account saves of one report id land in the order they were made", async () => { + const posted = []; + let releaseFirst; + const firstSession = new Promise((resolve) => { + releaseFirst = () => resolve(response({ signedIn: true })); + }); + let sessions = 0; + const fetcher = async (url, options) => { + if (url === "/api/session") + return ++sessions === 1 ? firstSession : response({ signedIn: true }); + posted.push(JSON.parse(options.body).report.summary); + return response({ id: "same" }); + }; + const entry = (summary) => ({ + id: "same", + problemId: "two-sum", + report: { incomplete: true, summary }, + }); + const provisional = saveReportHistory(entry("provisional"), { + storage: memoryStorage(), + fetcher, + }); + const final = saveReportHistory(entry("final"), { + storage: memoryStorage(), + fetcher, + }); + // The second save's own session check would answer at once; it still waits. + await new Promise((resolve) => setImmediate(resolve)); + assert.deepEqual(posted, []); + releaseFirst(); + assert.deepEqual(await provisional, { local: "saved", account: "saved" }); + assert.deepEqual(await final, { local: "saved", account: "saved" }); + assert.deepEqual(posted, ["provisional", "final"]); +}); + test("review inputs survive the full-report cap", async () => { const storage = memoryStorage(); for (let index = 0; index < 21; index += 1) { diff --git a/tests/browser/report-recovery.test.js b/tests/browser/report-recovery.test.js new file mode 100644 index 00000000..1ab7a9f0 --- /dev/null +++ b/tests/browser/report-recovery.test.js @@ -0,0 +1,541 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { createReportRecovery } from "../../web/report-recovery.js"; +import { retryReportPayload, sanitizeReport, topics } from "../../web/lib.js"; + +function setup() { + let now = 0; + let next = 0; + const scheduled = new Map(); + const events = []; + const recovery = createReportRecovery({ + timers: { + setTimeout(fn, ms) { + scheduled.set(++next, { fn, at: now + ms }); + return next; + }, + clearTimeout(id) { + scheduled.delete(id); + }, + }, + send: () => events.push(retryReportPayload()), + offer: (ready) => events.push(ready), + waiting: () => events.push("waiting"), + finalize: (report) => events.push(report), + keep: (report) => events.push({ kept: report }), + }); + function advance(ms) { + now += ms; + for (const [id, timer] of [...scheduled]) { + if (timer.at <= now && scheduled.has(id)) { + scheduled.delete(id); + timer.fn(); + } + } + } + return { recovery, advance, events, scheduled }; +} +const raw = { + incomplete: true, + summary: "503", + reportRecovery: { expiresInSeconds: 300, retryAfterSeconds: 30 }, +}; + +const accepted = { type: "report_retry", status: "accepted" }; +const closed = { type: "report_retry", status: "closed" }; +const finalized = (events) => events.filter((e) => e?.summary); + +test("report recovery keeps the provisional failure the moment it arrives", () => { + const { recovery, events } = setup(); + assert.equal(recovery.start(raw), true); + assert.deepEqual(events[0], { kept: { incomplete: true, summary: "503" } }); + assert.equal(finalized(events).length, 0); +}); + +test("report recovery spends one retry after cooldown and times it from acceptance", () => { + const { recovery, events, advance } = setup(); + recovery.start(raw); + assert.equal(recovery.retry(), false); + advance(30_000); + assert.equal(recovery.retry(), true); + assert.equal(recovery.retry(), false); + // Queued behind a reconnect: the agent's answer restarts the wait. + advance(100_000); + assert.equal(finalized(events).length, 0); + assert.equal(recovery.notice(accepted), true); + advance(139_000); + assert.equal(finalized(events).length, 0); + advance(1000); + assert.equal(events.filter((e) => e?.type === "retry_report").length, 1); + assert.deepEqual(events.at(-1), { incomplete: true, summary: "503" }); + recovery.finish(); + assert.equal(finalized(events).length, 1); +}); + +test("an accepted retry outlives the offer's own expiry", () => { + const { recovery, events, advance } = setup(); + recovery.start(raw); + advance(290_000); + recovery.retry(); + recovery.notice(accepted); + // A closing notice cannot cut short a generation the agent already began. + assert.equal(recovery.notice(closed), false); + advance(100_000); + assert.equal(finalized(events).length, 0); + advance(40_000); + assert.equal(finalized(events).length, 1); +}); + +test("an early answer arriving after the offer ran out finalizes instead", () => { + const { recovery, advance, events } = setup(); + recovery.start(raw); + advance(290_000); + recovery.retry(); + // The offer's own expiry passes while the request is in flight. + advance(30_000); + assert.equal(finalized(events).length, 0); + assert.equal( + recovery.notice({ + type: "report_retry", + status: "early", + retryAfterSeconds: 2, + }), + true, + ); + assert.equal(finalized(events).length, 1); + advance(2000); + assert.equal(recovery.retry(), false); + assert.notEqual(events.at(-1), true); +}); + +test("an early answer re-offers the retry after the agent's wait", () => { + const { recovery, events, advance } = setup(); + recovery.start(raw); + advance(30_000); + recovery.retry(); + assert.equal( + recovery.notice({ + type: "report_retry", + status: "early", + retryAfterSeconds: 2, + }), + true, + ); + assert.equal(events.at(-1), false); + assert.equal(recovery.retry(), false); + advance(2000); + assert.equal(events.at(-1), true); + assert.equal(recovery.retry(), true); + assert.equal(events.filter((e) => e?.type === "retry_report").length, 2); +}); + +test("answers to a request never sent, or malformed ones, change nothing", () => { + const { recovery, events, advance, scheduled } = setup(); + recovery.start(raw); + const timers = scheduled.size; + for (const message of [ + accepted, + { type: "report_retry", status: "early", retryAfterSeconds: 2 }, + { type: "report_retry", status: "nonsense" }, + { type: "thinking_state", thinking: true }, + ]) + assert.equal(recovery.notice(message), false); + advance(30_000); + recovery.retry(); + assert.equal( + recovery.notice({ type: "report_retry", status: "early" }), + false, + ); + // The offer's expiry and the retry's own wait; no early offer re-armed. + assert.equal(scheduled.size, timers); + assert.equal(finalized(events).length, 0); +}); + +test("the agent's closing notice, or its absence, finalizes exactly once", () => { + for (const close of ["notice", "fallback", "departure"]) { + const { recovery, advance, events } = setup(); + recovery.start(raw); + advance(30_000); + // A retry sent too late is answered by the close, not left spinning. + if (close === "notice") { + recovery.retry(); + assert.equal(recovery.notice(closed), true); + } + if (close === "departure") recovery.finish(); + advance(284_000); + assert.equal(finalized(events).length, close === "fallback" ? 0 : 1); + advance(1000); + recovery.finish(); + recovery.notice(closed); + assert.equal(finalized(events).length, 1); + } +}); + +test("a retry whose answer is lost is bounded by its own wait, not the offer", () => { + const { recovery, advance, events } = setup(); + recovery.start(raw); + advance(290_000); + assert.equal(recovery.retry(), true); + // 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); + assert.equal(finalized(events).length, 0); + advance(1000); + assert.equal(finalized(events).length, 1); +}); +test("a delivered final report cancels recovery without saving provisional failure", () => { + const { recovery, advance, events, scheduled } = setup(); + recovery.start(raw); + advance(30_000); + recovery.retry(); + recovery.notice(accepted); + recovery.stop(); + advance(500_000); + assert.equal(events.filter((e) => e?.summary).length, 0); + assert.equal(scheduled.size, 0); +}); + +test("duplicate provisional reports cannot extend the window or allow rerolls", () => { + const { recovery, advance, events } = setup(); + recovery.start(raw); + advance(20_000); + assert.equal(recovery.start(raw), false); + advance(10_000); + recovery.retry(); + recovery.stop(); + assert.equal(recovery.start(raw), false); + assert.equal(events.filter((e) => e?.type).length, 1); +}); + +test("untrusted recovery metadata cannot offer unbounded retention", () => { + for (const data of [ + { ...raw, incomplete: false }, + { + ...raw, + reportRecovery: { expiresInSeconds: 301, retryAfterSeconds: 30 }, + }, + { ...raw, reportRecovery: { expiresInSeconds: 300, retryAfterSeconds: 0 } }, + ]) { + const { recovery, scheduled } = setup(); + assert.equal(recovery.start(data), false); + assert.equal(scheduled.size, 0); + } +}); + +// Exercise the real page functions with a small DOM and LiveKit boundary. +// The recovery controller uses fake time; rendering and saving are observed. +import { functionBody, read } from "./source.js"; + +function pageHarness(phase = "live") { + const source = read("web/interview.js"); + let now = 0; + let next = 0; + const scheduled = new Map(); + const events = []; + const saves = []; + const timers = { + setTimeout(fn, ms) { + scheduled.set(++next, { fn, at: now + ms }); + return next; + }, + clearTimeout(id) { + scheduled.delete(id); + }, + }; + const room = { disconnect: async () => events.push("disconnect") }; + const nodes = Object.fromEntries( + [ + "endingTitle", + "frameworkHint", + "retryReport", + "editor", + "ending", + "forceReport", + "leaveRoom", + "endingDetail", + ].map((key) => [key, {}]), + ); + const spinner = {}; + nodes.ending.querySelector = () => spinner; + nodes.endingTitle.textContent = "Jim is writing up your evaluation..."; + nodes.leaveRoom.textContent = "Taking too long? Leave the room"; + const scope = { + state: { phase, room, connected: true }, + nodes, + topics, + sanitizeReport, + createReportRecovery: (args) => createReportRecovery({ ...args, timers }), + globalThis: { ...timers }, + setTimeout: timers.setTimeout, + clearTimeout: timers.clearTimeout, + frameworkHintTimer: 0, + codePublishTimer: null, + pendingLanguagePublish: null, + endingEscape: [], + REPORT_DELIVERY_GRACE_MS: 3000, + deliveryGrace: 0, + stopEndingEscape() { + for (const id of scope.endingEscape) timers.clearTimeout(id); + scope.endingEscape = []; + timers.clearTimeout(scope.deliveryGrace); + scope.deliveryGrace = 0; + }, + stopAvatar: () => events.push("stop-avatar"), + stopLocalMedia: () => events.push("stop-media"), + startEndingClock: () => {}, + stopEndingClock: () => {}, + publish: (topic, payload) => { + events.push(payload); + return Promise.resolve(); + }, + retryReportPayload, + recordReplay: (_, row) => events.push(row), + flushReplay: async () => {}, + setBanner: () => {}, + providerUiState: () => ({ message: "failure" }), + setLocalAudioEnabled: () => {}, + saveHistory: async (report = scope.state.report) => { + events.push("save"); + saves.push(report); + }, + renderReport: () => events.push("render"), + renderReportSaveStatus: () => {}, + interviewLoop: "coding_only", + roomParticipants: () => [], + roomInterviewer: () => null, + setAgentStateLabel: () => {}, + window: { location: {}, addEventListener() {} }, + }; + const recovery = source.slice( + source.indexOf("const reportRecovery ="), + source.indexOf("async function receiveReport"), + ); + const names = [ + "receiveReport", + "leaveRoom", + "updateAgentState", + "receiveControl", + ]; + const bodies = names + .map((name) => `${functionBody(source, name)}\n}`) + .join("\n"); + const page = new Function( + "scope", + `with(scope) { ${recovery}\n${bodies}\nreturn {reportRecovery, finalizeRecoveryOnPageHide, ${names.join(",")}}; }`, + )(new Proxy(scope, { has: (target, key) => key in target })); + function advance(ms) { + now += ms; + for (const [id, timer] of [...scheduled]) { + if (timer.at <= now && scheduled.has(id)) { + scheduled.delete(id); + timer.fn(); + } + } + } + return { ...page, room, scope, events, saves, advance }; +} +const packet = (value) => new TextEncoder().encode(JSON.stringify(value)); +const settle = () => new Promise((resolve) => setImmediate(resolve)); + +const saveCount = (page) => + page.events.filter((event) => event === "save").length; +const control = (value) => packet(value); + +test("page saves the transient failure on arrival, then once more on leaving", async () => { + for (const phase of ["live", "ending"]) { + const page = pageHarness(phase); + await page.receiveReport(page.room, packet(raw)); + assert.equal(page.scope.state.phase, "report_recovery"); + assert.ok(page.events.includes("stop-media")); + assert.equal(saveCount(page), 1); + assert.equal(page.saves[0].summary, "503"); + assert.equal(page.saves[0].incomplete, true); + assert.ok(!page.events.includes("disconnect")); + page.leaveRoom(); + await settle(); + assert.equal(page.scope.state.phase, "report"); + assert.equal(saveCount(page), 2); + assert.equal( + page.events.filter((event) => event === "disconnect").length, + 1, + ); + } +}); + +test("page finalizes a regeneration once, cancelling provisional timers", async () => { + const page = pageHarness(); + await page.receiveReport(page.room, packet(raw)); + page.advance(30_000); + assert.equal(page.reportRecovery.retry(), true); + page.receiveControl(control({ type: "report_retry", status: "accepted" })); + const final = { incomplete: true, summary: "final outcome" }; + await page.receiveReport(page.room, packet(final)); + await page.receiveReport(page.room, packet(final)); + page.advance(500_000); + await settle(); + assert.equal(saveCount(page), 2); + assert.equal(page.saves.at(-1).summary, "final outcome"); + assert.equal( + page.events.filter((event) => event?.state === "ended").length, + 1, + ); + assert.equal( + page.events.filter((event) => event?.state === "rounds_final").length, + 1, + ); +}); + +test("the agent's closing notice finalizes the page through its control topic", async () => { + const page = pageHarness(); + await page.receiveReport(page.room, packet(raw)); + page.advance(30_000); + page.reportRecovery.retry(); + page.receiveControl(control({ type: "report_retry", status: "closed" })); + await settle(); + assert.equal(page.scope.state.phase, "report"); + assert.equal(saveCount(page), 2); +}); + +test("a report packet that parses to null is still delivered", async () => { + const page = pageHarness(); + await page.receiveReport(page.room, packet(null)); + assert.equal(page.scope.state.phase, "report"); + assert.ok(page.events.includes("render")); + assert.equal(saveCount(page), 1); +}); + +test("a regenerated report that fails to render does not keep the recovery overlay", async () => { + const page = pageHarness(); + await page.receiveReport(page.room, packet(raw)); + page.advance(30_000); + page.reportRecovery.retry(); + assert.equal( + page.scope.nodes.endingTitle.textContent, + "Retrying your evaluation", + ); + page.scope.renderReport = () => { + throw new Error("undrawable"); + }; + page.scope.reportRenderFailed = () => page.events.push("render-failed"); + page.scope.console = { warn() {} }; + await page.receiveReport( + page.room, + packet({ incomplete: true, summary: "final outcome" }), + ); + assert.ok(page.events.includes("render-failed")); + const { nodes } = page.scope; + assert.equal( + nodes.endingTitle.textContent, + "Jim is writing up your evaluation...", + ); + assert.equal(nodes.leaveRoom.textContent, "Taking too long? Leave the room"); + assert.equal(nodes.ending.querySelector(".spinner").hidden, false); + assert.equal(nodes.retryReport.hidden, true); +}); + +test("a refused offer in a live interview writes its end frame once", async () => { + const page = pageHarness(); + await page.receiveReport( + page.room, + packet({ + ...raw, + endReason: "time_up", + reportRecovery: { expiresInSeconds: 999, retryAfterSeconds: 30 }, + }), + ); + const ended = page.events.filter((event) => event?.state === "ended"); + assert.equal(page.scope.state.phase, "report"); + assert.equal(ended.length, 1); + assert.equal(ended[0].reason, "time_up"); +}); + +test("agent departure at cooldown boundary still finalizes the original failure", async () => { + const page = pageHarness(); + await page.receiveReport(page.room, packet(raw)); + page.advance(28_000); + page.updateAgentState(); + page.advance(2000); + assert.equal(page.scope.nodes.retryReport.disabled, false); + page.advance(1000); + await settle(); + assert.equal(page.scope.state.phase, "report"); + assert.equal(saveCount(page), 2); +}); + +test("terminal room disconnect finalizes recovery before clearing the room", async () => { + const page = pageHarness(); + await page.receiveReport(page.room, packet(raw)); + const source = read("web/interview.js"); + const start = + source.indexOf("room.on(livekit.RoomEvent.Disconnected, () => {") + + "room.on(livekit.RoomEvent.Disconnected, () => {".length; + const body = source.slice(start, source.indexOf("\n });", start)); + page.scope.reportRecovery = page.reportRecovery; + page.scope.console = { warn() {} }; + const disconnect = new Function( + "scope", + `with(scope) { return () => { ${body} }; }`, + )(new Proxy(page.scope, { has: (target, key) => key in target })); + disconnect(); + await settle(); + assert.equal(page.scope.state.room, null); + assert.equal(page.scope.state.connected, false); + assert.equal(page.scope.state.phase, "report"); + assert.equal(saveCount(page), 2); + page.advance(500_000); + assert.equal(saveCount(page), 2); +}); + +test("a finalization with no room left still renders and saves", async () => { + const page = pageHarness(); + await page.receiveReport(page.room, packet(raw)); + page.scope.state.room = null; + page.reportRecovery.finish(); + await settle(); + assert.equal(page.scope.state.phase, "report"); + assert.ok(page.events.includes("render")); + assert.equal(saveCount(page), 2); +}); + +test("page exit finalizes once, after the arrival save already landed", async () => { + const page = pageHarness(); + await page.receiveReport(page.room, packet(raw)); + assert.equal(saveCount(page), 1); + page.finalizeRecoveryOnPageHide(); + assert.equal(saveCount(page), 2); + page.finalizeRecoveryOnPageHide(); + page.advance(500_000); + await settle(); + assert.equal(saveCount(page), 2); +}); + +test("a returning agent cancels failure finalization during delivery grace", async () => { + const page = pageHarness(); + await page.receiveReport(page.room, packet(raw)); + page.advance(28_000); + page.updateAgentState(); + page.updateAgentState(); + page.scope.roomInterviewer = () => ({}); + page.advance(3000); + await settle(); + assert.equal(page.scope.state.phase, "report_recovery"); + assert.equal(saveCount(page), 1); +}); + +test("quota recovery waits for the provider cooldown and hides framework hints", async () => { + const page = pageHarness(); + await page.receiveReport( + page.room, + packet({ + ...raw, + reportRecovery: { expiresInSeconds: 300, retryAfterSeconds: 60 }, + }), + ); + assert.equal(page.scope.nodes.frameworkHint.hidden, true); + assert.match(page.scope.nodes.endingDetail.textContent, /60 seconds/); + page.advance(30_000); + assert.equal(page.reportRecovery.retry(), false); + page.advance(30_000); + assert.equal(page.reportRecovery.retry(), true); +}); diff --git a/tests/fixtures/README.md b/tests/fixtures/README.md index b8d61f26..1ba94c45 100644 --- a/tests/fixtures/README.md +++ b/tests/fixtures/README.md @@ -13,7 +13,8 @@ only thing in the repo that crosses that boundary. | File | Producer | Consumer | |---|---|---| | `code-update.json` | `codeUpdatePayload` | `apply_code_update` | -| `control.json` | `thinkingPayload`, `yieldTurnPayload`, `timeWarningPayload`, `endInterviewPayload` | `apply_control` | +| `report-recovery.json` | `reportRecoveryLimits` in `web/report-recovery.js` | `recovery_metadata`, `recovery_notice`, `REPORT_TIMEOUT` | +| `control.json` | `thinkingPayload`, `yieldTurnPayload`, `timeWarningPayload`, `endInterviewPayload`, `retryReportPayload` | `apply_control`, `recovery_request` | | `test-results.json` | `testPayload` | `apply_test_results` | | `integrity-chain.json` | `integrityEventPayload` | `apply_integrity` | diff --git a/tests/fixtures/control.json b/tests/fixtures/control.json index 09e6f5b7..dc8c36ed 100644 --- a/tests/fixtures/control.json +++ b/tests/fixtures/control.json @@ -1,6 +1,12 @@ { "topic": "control", "cases": [ + { + "name": "retry report", + "payload": { + "type": "retry_report" + } + }, { "name": "thinking start", "payload": { diff --git a/tests/fixtures/report-recovery.json b/tests/fixtures/report-recovery.json new file mode 100644 index 00000000..25f3aee6 --- /dev/null +++ b/tests/fixtures/report-recovery.json @@ -0,0 +1,12 @@ +{ + "expiresInSeconds": 300, + "retryAfterSeconds": 30, + "quotaRetryAfterSeconds": 60, + "retryWaitSeconds": 140, + "closeGraceSeconds": 15, + "retryStatuses": [ + "accepted", + "early", + "closed" + ] +} diff --git a/tests/test_analyze_gemini_usage.py b/tests/test_analyze_gemini_usage.py index be7b0262..1d6421c0 100644 --- a/tests/test_analyze_gemini_usage.py +++ b/tests/test_analyze_gemini_usage.py @@ -297,6 +297,7 @@ def test_http_failures_and_context_refreshes_are_counted(self): "gemini report transport_failed room=a call=1 backoff_s=2 " "error=secret prompt_tokens=5\n", "gemini report transport_failed room=a call=2 backoff_s=4 error=x\n", + "gemini report retry_unavailable room=a call=2\n", "interim review skipped room=a: quota\n", "codetrial context_refresh room=a session=1 bytes=300\n", "codetrial context_refresh room=a session=1 bytes=200\n", diff --git a/tests/unit/dispatch.rs b/tests/unit/dispatch.rs index ab43599b..de536847 100644 --- a/tests/unit/dispatch.rs +++ b/tests/unit/dispatch.rs @@ -10,8 +10,10 @@ use super::*; /// would refuse every interview after them for the life of the process. #[tokio::test] async fn a_panicking_interview_gives_its_capacity_back() { - let live = Arc::new(Mutex::new(HashSet::new())); - live.lock().unwrap().insert("interview-boom".to_string()); + let live = Arc::new(Mutex::new(HashMap::new())); + live.lock() + .unwrap() + .insert("interview-boom".to_string(), true); let slot = Slot { live: Arc::clone(&live), room_name: "interview-boom".to_string(), @@ -33,8 +35,10 @@ async fn a_panicking_interview_gives_its_capacity_back() { /// one gets one, and the refusal names the number they chose. #[tokio::test] async fn the_configured_cap_is_the_one_enforced() { - let live = Arc::new(Mutex::new(HashSet::new())); - live.lock().unwrap().insert("interview-first".to_string()); + let live = Arc::new(Mutex::new(HashMap::new())); + live.lock() + .unwrap() + .insert("interview-first".to_string(), true); let dispatcher = LocalDispatcher { config: crate::config::load_from_pairs([ ("LIVEKIT_URL", "wss://primary.example"), @@ -56,18 +60,22 @@ async fn the_configured_cap_is_the_one_enforced() { }; assert!( - !dispatcher.ensure_agent("interview-second", &provider), + dispatcher.ensure_agent("interview-second", &provider) == Err(DispatchRefusal::AtCapacity), "a second room must be refused at a cap of one" ); // The room already running is still admitted, because a reload asks for the // same room and refusing it would break the page that is open. - assert!(dispatcher.ensure_agent("interview-first", &provider)); + assert!( + dispatcher + .ensure_agent("interview-first", &provider) + .is_ok() + ); } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn a_concurrent_burst_never_overbooks_and_released_capacity_returns() { - let live = Arc::new(Mutex::new(HashSet::new())); + let live = Arc::new(Mutex::new(HashMap::new())); let dispatcher = LocalDispatcher { config: crate::config::load_from_pairs([ ("LIVEKIT_URL", "wss://primary.example"), @@ -96,7 +104,9 @@ async fn a_concurrent_burst_never_overbooks_and_released_capacity_returns() { match dispatcher.reserve(&format!("interview-burst-{index}")) { Reservation::New(slot) => Some(slot), Reservation::Full => None, - Reservation::Existing => panic!("burst ids are unique"), + Reservation::Existing | Reservation::Finalizing => { + panic!("burst ids are unique") + } } })); } @@ -173,3 +183,44 @@ fn provider_credentials_replace_the_base_without_blanking_the_key() { assert_eq!(selected.livekit_api_key, inherited.livekit_api_key); assert_eq!(selected.livekit_api_secret, inherited.livekit_api_secret); } + +#[tokio::test] +async fn a_room_finalizing_its_report_cannot_be_reused_for_a_new_interview() { + let live = Arc::new(Mutex::new(HashMap::new())); + let dispatcher = LocalDispatcher { + config: crate::config::load_from_pairs([ + ("LIVEKIT_URL", "wss://primary.example"), + ("LIVEKIT_API_KEY", "primary-key"), + ("LIVEKIT_API_SECRET", "primary-secret"), + ("GOOGLE_API_KEY", "primary-google"), + ]) + .unwrap(), + runtime: tokio::runtime::Handle::current(), + live: Arc::clone(&live), + max_concurrent: 2, + }; + let Reservation::New(slot) = dispatcher.reserve("local-room") else { + panic!("fresh slot"); + }; + assert!(matches!( + dispatcher.reserve("local-room"), + Reservation::Existing + )); + slot.assessment_finished(); + assert!(matches!( + dispatcher.reserve("local-room"), + Reservation::Finalizing + )); + assert_eq!(live.lock().unwrap().len(), 1); + // Other rooms can still use the unoccupied capacity. + let Reservation::New(other) = dispatcher.reserve("another-room") else { + panic!("spare slot"); + }; + drop(slot); + assert!(matches!( + dispatcher.reserve("local-room"), + Reservation::New(_) + )); + drop(other); + assert!(live.lock().unwrap().is_empty()); +} diff --git a/tests/unit/gemini.rs b/tests/unit/gemini.rs index e9f36e45..2b13096f 100644 --- a/tests/unit/gemini.rs +++ b/tests/unit/gemini.rs @@ -241,11 +241,30 @@ async fn report_failover_fixture( } async fn report_transport_fixture( + prefix: &str, + statuses: Vec, + single_key: bool, + budget: ReportCallBudget, + backoff: Duration, +) -> TransportFixture { + report_transport_fixture_with_body( + prefix, + statuses, + single_key, + budget, + backoff, + json!({"error":{"message":"do not expose upstream credentials"}}), + ) + .await +} + +async fn report_transport_fixture_with_body( prefix: &str, statuses: Vec, single_key: bool, mut budget: ReportCallBudget, backoff: Duration, + error_body: Value, ) -> TransportFixture { let config = live_config(&[( "GOOGLE_API_KEYS", @@ -263,6 +282,7 @@ async fn report_transport_fixture( axum::routing::post(move |headers: axum::http::HeaderMap| { let seen = Arc::clone(&seen); let statuses = statuses.clone(); + let error_body = error_body.clone(); async move { let mut seen = seen.lock().unwrap(); let status = statuses[seen.len()]; @@ -272,7 +292,7 @@ async fn report_transport_fixture( axum::Json(if status == 200 { json!({"candidates":[{"content":{"parts":[{"text":"report"}]}}]}) } else { - json!({"error":{"message":"do not expose upstream credentials"}}) + error_body }), ) } @@ -929,7 +949,21 @@ fn report_requests_are_session_local_and_never_reuse_personalized_output() { type ReportResult = Result>; -fn valid_report() -> Value { +/// The whole report path against a local server, for other modules' tests. +/// Lives here so the endpoint and backoff overrides stay out of the crate's +/// own API. +pub(crate) async fn generate_report_at( + keys: &GeminiKeys, + url: &str, + backoff: Duration, + prompt: &str, + problem: &crate::agent::Problem, + scope: &str, +) -> Result> { + generate_report_with_keys_at(keys, url, backoff, prompt, problem, scope).await +} + +pub(crate) fn valid_report() -> Value { let improvements = [ ("Algorithm", "Explain complexity"), ("Test", "Test boundaries"), @@ -3104,3 +3138,260 @@ fn an_interruption_precedes_candidate_speech_in_the_same_frame() { assert_eq!(events, expected); } } + +#[tokio::test] +async fn exhausted_transient_reports_can_be_regenerated_but_permanent_failures_cannot() { + for (statuses, allowed) in [ + (vec![503; MAX_REPORT_HTTP_ATTEMPTS], true), + (vec![429; MAX_REPORT_HTTP_ATTEMPTS], true), + (vec![400], false), + (vec![401], false), + (vec![402], false), + ] { + let (result, _, _) = report_transport_fixture( + "regeneration", + statuses, + true, + ReportCallBudget::new(), + Duration::ZERO, + ) + .await; + assert_eq!( + report_regeneration_retry_after(result.unwrap_err().as_ref()).is_some(), + allowed + ); + } +} + +/// The wait an exhausted rotation offers is until its first key is back: just +/// under a whole cooldown when the keys went out moments ago, and never more. +/// `since` is taken before the keys went out, so the wait plus the time since +/// covers a whole cooldown however slowly the test ran. +fn assert_quota_wait(delay: Option, since: std::time::Instant) { + let delay = delay.expect("a quota-exhausted rotation offers a regeneration"); + assert!(delay <= credentials::QUOTA_COOLDOWN, "{delay:?}"); + assert!( + delay + since.elapsed() >= credentials::QUOTA_COOLDOWN, + "{delay:?}" + ); +} + +#[tokio::test] +async fn quota_regeneration_waits_until_the_whole_key_rotation_can_make_a_call() { + let since = std::time::Instant::now(); + use std::sync::atomic::{AtomicUsize, Ordering}; + let prefix = "quota-regeneration-ready"; + let (failed, seen, _) = report_failover_fixture(prefix, vec![429, 429], false).await; + assert_eq!(seen.len(), 2); + let delay = report_regeneration_retry_after(failed.unwrap_err().as_ref()).unwrap(); + let config = live_config(&[( + "GOOGLE_API_KEYS", + &format!("{prefix}-first,{prefix}-second"), + )]); + let keys = GeminiKeys::from_config(&config); + assert!(keys.select_report().is_err()); + let selected = keys + .select_report_at(std::time::Instant::now() + delay) + .expect("retry must wait until a key returns"); + assert_quota_wait(Some(delay), since); + + let calls = Arc::new(AtomicUsize::new(0)); + let observed = Arc::clone(&calls); + let app = axum::Router::new().route( + "/", + axum::routing::post(move || { + let observed = Arc::clone(&observed); + async move { + observed.fetch_add(1, Ordering::SeqCst); + axum::Json(json!({"candidates":[{"content":{"parts":[{"text":"report"}]}}]})) + } + }), + ); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}/", listener.local_addr().unwrap()); + let server = tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + assert_eq!( + generate_report_once(&selected, &url, "same frozen prompt", "quota-regeneration") + .await + .unwrap(), + "report" + ); + assert_eq!(calls.load(Ordering::SeqCst), 1); + server.abort(); +} + +#[tokio::test] +async fn quota_body_failures_offer_the_same_recovery_as_http_429() { + let since = std::time::Instant::now(); + for status in [400, 403] { + let (result, seen, _) = report_transport_fixture_with_body( + &format!("quota-body-{status}"), + vec![status; MAX_REPORT_HTTP_ATTEMPTS], + true, + ReportCallBudget::new(), + Duration::ZERO, + json!({"error":{"status":"RESOURCE_EXHAUSTED"}}), + ) + .await; + assert_eq!(seen.len(), MAX_REPORT_HTTP_ATTEMPTS); + assert_quota_wait( + report_regeneration_retry_after(result.unwrap_err().as_ref()), + since, + ); + } +} + +#[tokio::test] +async fn report_recovery_follows_quota_availability_not_the_last_key_failure() { + for (statuses, allowed) in [ + (vec![429, 401], true), + (vec![401, 429], true), + (vec![429, 402], true), + (vec![401, 403], false), + (vec![402, 402], false), + ] { + let prefix = format!("recovery-mixed-{}-{}", statuses[0], statuses[1]); + let (result, seen, _) = report_failover_fixture(&prefix, statuses, false).await; + assert_eq!(seen.len(), 2); + assert_eq!( + report_regeneration_retry_after(result.unwrap_err().as_ref()).is_some(), + allowed + ); + + // A second interview encounters the same shared rotation at entry, + // before any request. It must get the same recovery eligibility. + let (result, seen, _) = report_failover_fixture(&prefix, vec![], false).await; + assert!(seen.is_empty()); + assert_eq!( + report_regeneration_retry_after(result.unwrap_err().as_ref()).is_some(), + allowed + ); + } +} + +#[tokio::test] +async fn concurrent_quota_exhaustion_keeps_its_delay_after_a_report_503() { + let since = std::time::Instant::now(); + let config = live_config(&[( + "GOOGLE_API_KEYS", + "concurrent-report-first,concurrent-report-second", + )]); + let keys = Arc::new(GeminiKeys::from_config(&config)); + let shared = Arc::clone(&keys); + let app = axum::Router::new().route( + "/", + axum::routing::post(move || { + let shared = Arc::clone(&shared); + async move { + for key in ["concurrent-report-first", "concurrent-report-second"] { + shared.failed(key, CredentialFailure::Quota, ApiSurface::Report); + } + ( + axum::http::StatusCode::SERVICE_UNAVAILABLE, + axum::Json(json!({})), + ) + } + }), + ); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}/", listener.local_addr().unwrap()); + let server = tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + let mut budget = ReportCallBudget::new(); + let error = generate_report_transport( + &keys, + &url, + "frozen prompt", + &mut budget, + Duration::ZERO, + "concurrent-quota", + ) + .await + .unwrap_err(); + assert_quota_wait(report_regeneration_retry_after(error.as_ref()), since); + assert_eq!(budget.remaining, MAX_REPORT_HTTP_ATTEMPTS - 1); + server.abort(); +} + +#[tokio::test] +async fn exhaustion_during_backoff_keeps_the_original_failure_and_recovery_delay() { + let since = std::time::Instant::now(); + for (failure, expected) in [ + (CredentialFailure::Quota, true), + (CredentialFailure::Invalid, false), + ] { + let prefix = format!("backoff-preserves-{failure:?}"); + let config = live_config(&[( + "GOOGLE_API_KEYS", + &format!("{prefix}-first,{prefix}-second"), + )]); + let keys = GeminiKeys::from_config(&config); + let original = "Gemini HTTP 503 Service Unavailable"; + let next = report_key_after_backoff( + &keys, + Duration::from_millis(20), + original.into(), + "backoff-room", + 1, + ); + tokio::pin!(next); + std::future::poll_fn(|cx| { + assert!(std::future::Future::poll(next.as_mut(), cx).is_pending()); + std::task::Poll::Ready(()) + }) + .await; + for suffix in ["first", "second"] { + keys.failed(&format!("{prefix}-{suffix}"), failure, ApiSurface::Report); + } + let error = next.await.unwrap_err(); + assert_eq!(error.to_string(), original); + let delay = report_regeneration_retry_after(error.as_ref()); + if expected { + assert_quota_wait(delay, since); + } else { + assert_eq!(delay, None); + } + } +} + +#[tokio::test] +async fn a_healthy_backup_can_regenerate_after_the_last_call_rejects_a_key() { + let prefix = "healthy-backup-regeneration"; + let (result, seen, remaining) = + report_failover_fixture(prefix, vec![503, 503, 503, 503, 401], false).await; + assert_eq!(seen.len(), MAX_REPORT_HTTP_ATTEMPTS); + assert_eq!(remaining, 0); + assert_eq!( + report_regeneration_retry_after(result.unwrap_err().as_ref()), + Some(Duration::ZERO) + ); + let (result, seen, _) = report_failover_fixture(prefix, vec![200], false).await; + assert_eq!(result.unwrap(), "report"); + assert_eq!(seen, [format!("{prefix}-second")]); +} + +/// A quota failure on the last call waits only as long as the rotation needs: +/// a backup that selects now is usable now, so nothing beyond the regeneration +/// floor, while a sole key is never ruled out and has to sit out the cooldown. +#[tokio::test] +async fn a_quota_failure_waits_only_when_no_other_key_can_answer() { + let (result, seen, remaining) = + report_failover_fixture("quota-last-call", vec![503, 503, 503, 503, 429], false).await; + assert_eq!(seen.len(), MAX_REPORT_HTTP_ATTEMPTS); + assert_eq!(remaining, 0); + assert_eq!( + report_regeneration_retry_after(result.unwrap_err().as_ref()), + Some(Duration::ZERO) + ); + + let (result, seen, _) = + report_failover_fixture("quota-sole-key", vec![429; MAX_REPORT_HTTP_ATTEMPTS], true).await; + assert_eq!(seen.len(), MAX_REPORT_HTTP_ATTEMPTS); + assert_eq!( + report_regeneration_retry_after(result.unwrap_err().as_ref()), + Some(credentials::QUOTA_COOLDOWN) + ); +} diff --git a/tests/unit/gemini/credentials.rs b/tests/unit/gemini/credentials.rs index b5760be9..c63f888c 100644 --- a/tests/unit/gemini/credentials.rs +++ b/tests/unit/gemini/credentials.rs @@ -506,8 +506,8 @@ fn an_unexplained_refusal_stays_on_its_own_surface() { assert_eq!(keys.select().unwrap(), "report-refused-a"); } -/// An exhausted Live rotation is worth waiting for only while some key is out -/// on quota alone; a Report rotation never waits. +/// Live waits for a quota key; Report exposes that deadline for the candidate +/// to request recovery. A rejection on one surface must not poison the other. #[test] fn an_exhausted_rotation_names_its_first_quota_key_back() { let keys = GeminiKeys::new(vec!["exhausted-wait-a".into(), "exhausted-wait-b".into()]); @@ -528,13 +528,28 @@ fn an_exhausted_rotation_names_its_first_quota_key_back() { CredentialFailure::Quota, ApiSurface::Report, ); - assert_eq!(exhausted_until(&keys.select_report().unwrap_err()), None); + let report_back = COOLDOWNS.get().unwrap().lock().unwrap()["exhausted-wait-b"].report; + assert_eq!( + exhausted_until(&keys.select_report().unwrap_err()), + report_back + ); + keys.failed( "exhausted-wait-b", CredentialFailure::Refused, ApiSurface::Live, ); assert_eq!(exhausted_until(&keys.select().unwrap_err()), None); + assert_eq!( + exhausted_until(&keys.select_report().unwrap_err()), + report_back + ); + keys.failed( + "exhausted-wait-b", + CredentialFailure::Refused, + ApiSurface::Report, + ); + assert_eq!(exhausted_until(&keys.select_report().unwrap_err()), None); } /// Slow quota rejections from enough keys outlast the first key's cooldown. diff --git a/tests/unit/livekit/report.rs b/tests/unit/livekit/report.rs index 74cd260c..f24bc57a 100644 --- a/tests/unit/livekit/report.rs +++ b/tests/unit/livekit/report.rs @@ -752,3 +752,650 @@ fn a_frozen_report_prompt_is_counted_and_a_missed_deadline_still_reports() { "{summary}" ); } + +#[test] +fn only_the_original_candidate_can_request_report_regeneration() { + let fixture: serde_json::Value = + serde_json::from_str(include_str!("../../fixtures/control.json")).unwrap(); + let cases = fixture + .as_array() + .unwrap_or_else(|| fixture["cases"].as_array().unwrap()); + let retry = cases + .iter() + .find(|case| case["payload"]["type"] == "retry_report") + .expect("generated retry fixture"); + let payload = serde_json::to_vec(&retry["payload"]).unwrap(); + assert!(recovery_request( + Some(crate::runtime::TOPIC_CONTROL), + Some("candidate-original"), + "candidate-original", + &payload + )); + for sender in [None, Some("candidate-other"), Some("interviewer-room")] { + assert!(!recovery_request( + Some(crate::runtime::TOPIC_CONTROL), + sender, + "candidate-original", + &payload + )); + } + assert!(!recovery_request( + Some(crate::runtime::TOPIC_CODE_UPDATE), + Some("candidate-original"), + "candidate-original", + &payload + )); + assert!(!recovery_request( + Some(crate::runtime::TOPIC_CONTROL), + Some("candidate-original"), + "candidate-original", + b"not JSON" + )); +} + +#[tokio::test] +async fn only_transient_failure_and_deadline_can_offer_regeneration() { + let keys = GeminiKeys::single("test"); + assert!(regeneration_cooldown(&Ok(Ok(serde_json::json!({}))), &keys).is_none()); + let schema_failure = Ok(Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "schema failure", + ) + .into())); + assert!(regeneration_cooldown(&schema_failure, &keys).is_none()); + assert_eq!( + regeneration_cooldown(&Err(elapsed().await), &keys), + Some(REPORT_RETRY_COOLDOWN) + ); +} + +async fn elapsed() -> tokio::time::error::Elapsed { + tokio::time::timeout(std::time::Duration::ZERO, std::future::pending::<()>()) + .await + .unwrap_err() +} + +/// A deadline missed while every key sat out on quota: the 30-second offer +/// would start the one regeneration before any key is back, and it would fail +/// at selection without a single HTTP call. +#[tokio::test] +async fn a_deadline_during_quota_exhaustion_waits_for_the_keys() { + let keys = + GeminiKeys::report_quota_exhausted(&["deadline-quota-first", "deadline-quota-second"]); + assert!(keys.select_report().is_err()); + let cooldown = regeneration_cooldown(&Err(elapsed().await), &keys).unwrap(); + + // Until the first key is back, which here is just under a whole cooldown, + // and the page is told the whole second it must wait. + assert!(cooldown > REPORT_RETRY_COOLDOWN, "{cooldown:?}"); + assert!(cooldown <= crate::gemini::QUOTA_COOLDOWN, "{cooldown:?}"); + assert_eq!( + recovery_metadata(cooldown)["retryAfterSeconds"], + crate::gemini::QUOTA_COOLDOWN.as_secs() + ); +} + +#[test] +fn report_metadata_uses_the_frozen_assessment() { + 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); + live.code = "post-interview edits".into(); + live.hints_used = 5; + assert!(!frozen.prompt.contains("post-interview edits")); + assert_eq!(live.evidence_ledger.metrics.final_report_prompt_count, 1); + assert_eq!( + frozen + .state + .evidence_ledger + .metrics + .final_report_prompt_count, + 1 + ); + live.integrity_events.push(serde_json::json!({"seq": 100})); + let packet = report_packet( + &boot, + &mut frozen.state, + "interview_complete", + &GeminiKeys::single("test"), + Ok(Err(std::io::Error::other("503").into())), + ) + .unwrap(); + let report: serde_json::Value = serde_json::from_slice(&packet.payload).unwrap(); + assert!(report["integrityEvents"].as_array().unwrap().is_empty()); + assert_ne!( + report["integrityEvents"], + report_with_integrity_events(serde_json::json!({}), &live, "interview_complete")["integrityEvents"] + ); +} + +struct RecoveryFixture( + tokio::sync::mpsc::UnboundedReceiver, + Vec, + Vec, + bool, +); + +impl RecoveryFixture { + fn new(events: tokio::sync::mpsc::UnboundedReceiver) -> Self { + Self(events, Vec::new(), Vec::new(), true) + } +} + +impl RecoveryRoom for RecoveryFixture { + fn candidate_present(&self) -> bool { + self.3 + } + + async fn next(&mut self) -> RecoveryEvent { + self.0.recv().await.unwrap_or(RecoveryEvent::Left) + } + + async fn notify(&mut self, notice: RecoveryNotice) { + self.1.push(notice); + } + + async fn publish( + &mut self, + packet: DataPacket, + ) -> Result<(), Box> { + assert_eq!(packet.topic.as_deref(), Some(TOPIC_REPORT)); + self.2.push(serde_json::from_slice(&packet.payload)?); + Ok(()) + } +} + +#[tokio::test(start_paused = true)] +async fn recovery_wait_never_generates_before_a_valid_request() { + for (event, notices) in [ + (RecoveryEvent::Left, vec![]), + ( + RecoveryEvent::Retry, + vec![ + RecoveryNotice::Early(REPORT_RETRY_COOLDOWN), + RecoveryNotice::Closed, + ], + ), + (RecoveryEvent::Ignore, vec![RecoveryNotice::Closed]), + ] { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + tx.send(event).unwrap(); + let mut fixture = RecoveryFixture::new(rx); + let calls = std::sync::atomic::AtomicUsize::new(0); + let generation = async { + calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + Ok(Ok(serde_json::json!({}))) + }; + let result = recover_report( + &mut fixture, + generation, + std::time::Duration::from_millis(2), + REPORT_RETRY_COOLDOWN, + tokio::time::Instant::now(), + || Readiness::Ready, + ) + .await; + assert!(result.is_none()); + assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 0); + assert_eq!(fixture.1, notices); + // Held open, so expiry is what ended the wait rather than the channel. + drop(tx); + } +} + +#[tokio::test] +async fn recovery_generates_once_despite_duplicate_requests() { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + tx.send(RecoveryEvent::Retry).unwrap(); + tx.send(RecoveryEvent::Retry).unwrap(); + let mut fixture = RecoveryFixture::new(rx); + let calls = std::sync::atomic::AtomicUsize::new(0); + let generation = async { + calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + Ok(Ok(serde_json::json!({"result": "complete"}))) + }; + let result = recover_report( + &mut fixture, + generation, + REPORT_RECOVERY_WINDOW, + std::time::Duration::ZERO, + tokio::time::Instant::now(), + || Readiness::Ready, + ) + .await; + assert_eq!(result.unwrap().unwrap().unwrap()["result"], "complete"); + assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 1); + assert_eq!(fixture.1, [RecoveryNotice::Accepted]); +} + +#[tokio::test] +async fn leaving_during_regeneration_cancels_the_call() { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + tx.send(RecoveryEvent::Retry).unwrap(); + tx.send(RecoveryEvent::Left).unwrap(); + let mut fixture = RecoveryFixture::new(rx); + let result = recover_report( + &mut fixture, + std::future::pending(), + REPORT_RECOVERY_WINDOW, + std::time::Duration::ZERO, + tokio::time::Instant::now(), + || Readiness::Ready, + ) + .await; + assert!(result.is_none()); +} + +/// An early request is answered with the time still to wait, and does not +/// spend the regeneration: the same page asking again after it is accepted. +#[tokio::test(start_paused = true)] +async fn an_early_retry_is_told_how_long_to_wait_and_keeps_its_turn() { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + let mut fixture = RecoveryFixture::new(rx); + let started = tokio::time::Instant::now(); + let requests = async { + tokio::time::sleep(std::time::Duration::from_millis(10_500)).await; + tx.send(RecoveryEvent::Retry).unwrap(); + tokio::time::sleep(std::time::Duration::from_secs(20)).await; + tx.send(RecoveryEvent::Retry).unwrap(); + tx + }; + let (result, _tx) = tokio::join!( + recover_report( + &mut fixture, + async { Ok(Ok(serde_json::json!({"result": "complete"}))) }, + REPORT_RECOVERY_WINDOW, + REPORT_RETRY_COOLDOWN, + started, + || Readiness::Ready, + ), + requests, + ); + assert!(result.is_some()); + let [RecoveryNotice::Early(wait), RecoveryNotice::Accepted] = fixture.1[..] else { + panic!("unexpected notices {:?}", fixture.1); + }; + assert_eq!(wait, std::time::Duration::from_millis(19_500)); + assert_eq!( + recovery_notice(RecoveryNotice::Early(wait))["retryAfterSeconds"], + 20 + ); +} + +/// The page's clock starts when the provisional report arrives, later than +/// this one. A retry it still thinks is in time but that lands after expiry +/// gets the closing notice rather than a spinner nobody will answer. +#[tokio::test(start_paused = true)] +async fn a_retry_after_expiry_is_not_accepted() { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + let mut fixture = RecoveryFixture::new(rx); + let started = tokio::time::Instant::now(); + let late = async { + tokio::time::sleep(REPORT_RECOVERY_WINDOW + std::time::Duration::from_secs(1)).await; + let _ = tx.send(RecoveryEvent::Retry); + tx + }; + let (result, _tx) = tokio::join!( + recover_report( + &mut fixture, + async { Ok(Ok(serde_json::json!({}))) }, + REPORT_RECOVERY_WINDOW, + REPORT_RETRY_COOLDOWN, + started, + || Readiness::Ready, + ), + late, + ); + assert!(result.is_none()); + assert_eq!(fixture.1, [RecoveryNotice::Closed]); +} + +#[test] +fn recovery_notices_use_the_statuses_the_page_reads() { + let limits: serde_json::Value = + serde_json::from_str(include_str!("../../fixtures/report-recovery.json")).unwrap(); + let statuses = limits["retryStatuses"].as_array().unwrap(); + for notice in [ + RecoveryNotice::Accepted, + RecoveryNotice::Early(std::time::Duration::from_millis(1)), + RecoveryNotice::Closed, + ] { + let value = recovery_notice(notice); + assert_eq!(value["type"], "report_retry"); + assert!(statuses.contains(&value["status"]), "{value}"); + } + assert_eq!(statuses.len(), 3); +} + +#[test] +fn recovery_offer_matches_browser_limits_and_generation_deadline() { + let limits: serde_json::Value = + serde_json::from_str(include_str!("../../fixtures/report-recovery.json")).unwrap(); + let offer = recovery_metadata(REPORT_RETRY_COOLDOWN); + assert_eq!(offer["expiresInSeconds"], limits["expiresInSeconds"]); + assert_eq!(offer["retryAfterSeconds"], limits["retryAfterSeconds"]); + assert_eq!( + recovery_metadata(crate::gemini::QUOTA_COOLDOWN)["retryAfterSeconds"], + limits["quotaRetryAfterSeconds"] + ); + assert!(limits["retryWaitSeconds"].as_u64().unwrap() > REPORT_TIMEOUT.as_secs() + 3); +} + +#[tokio::test] +async fn recovery_room_events_ignore_unrelated_activity_and_stop_on_disconnect() { + let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel(); + tx.send(::livekit::RoomEvent::Reconnecting).unwrap(); + tx.send(::livekit::RoomEvent::DataReceived { + payload: std::sync::Arc::new(br#"{"type":"retry_report"}"#.to_vec()), + topic: Some(crate::runtime::TOPIC_CONTROL.into()), + kind: ::livekit::DataPacketKind::Reliable, + participant: None, + }) + .unwrap(); + tx.send(::livekit::RoomEvent::Disconnected { + reason: ::livekit::DisconnectReason::ClientInitiated, + }) + .unwrap(); + let candidate = "candidate-original"; + assert_eq!( + recovery_event(rx.recv().await, candidate), + RecoveryEvent::Ignore + ); + assert_eq!( + recovery_event(rx.recv().await, candidate), + RecoveryEvent::Ignore + ); + assert_eq!( + recovery_event(rx.recv().await, candidate), + RecoveryEvent::Left + ); + drop(tx); + assert_eq!( + recovery_event(rx.recv().await, candidate), + RecoveryEvent::Left + ); +} + +fn past_cooldown() -> tokio::time::Instant { + tokio::time::Instant::now() - REPORT_RETRY_COOLDOWN +} + +/// The recovery the issue asked for, end to end over HTTP: an interview whose +/// report call ran out of retries on 503s is offered one regeneration, and the +/// candidate's retry gets a complete report from the same frozen interview, +/// with no new Live session and inside the ten-call ceiling. +#[tokio::test] +async fn a_report_lost_to_503s_is_regenerated_from_the_frozen_interview() { + let bodies = std::sync::Arc::new(std::sync::Mutex::new(Vec::::new())); + let seen = std::sync::Arc::clone(&bodies); + let app = axum::Router::new().route( + "/", + axum::routing::post(move |body: String| { + let seen = std::sync::Arc::clone(&seen); + async move { + let mut seen = seen.lock().unwrap(); + seen.push(body); + + // Every call of the first generation fails, then the service + // comes back for the regeneration. + if seen.len() <= 5 { + ( + axum::http::StatusCode::SERVICE_UNAVAILABLE, + axum::Json(serde_json::json!({"error": {"code": 503, "message": "overloaded"}})), + ) + } else { + let text = crate::gemini::tests::valid_report().to_string(); + ( + axum::http::StatusCode::OK, + axum::Json(serde_json::json!({"candidates": [{"content": {"parts": [{"text": text}]}}]})), + ) + } + } + }), + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}/", listener.local_addr().unwrap()); + let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + + 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("recovery-e2e"); + let generate = |prompt: &str| { + let (keys, url, problem) = (&keys, url.clone(), boot.problem); + let prompt = prompt.to_string(); + async move { + tokio::time::timeout( + REPORT_TIMEOUT, + crate::gemini::tests::generate_report_at( + keys, + &url, + std::time::Duration::ZERO, + &prompt, + problem, + "interview-fixed", + ), + ) + .await + } + }; + let first = generate(&frozen.prompt).await; + let regenerate = generate(&frozen.prompt); + + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + tx.send(RecoveryEvent::Retry).unwrap(); + let mut room = RecoveryFixture::new(rx); + let closed = std::sync::atomic::AtomicBool::new(false); + run_recovery( + &mut room, + RecoveryReport { + boot: &boot, + state: &mut frozen.state, + reason: "interview_complete", + keys: &keys, + }, + first, + regenerate, + past_cooldown, + async { closed.store(true, std::sync::atomic::Ordering::SeqCst) }, + ) + .await + .unwrap(); + server.abort(); + assert!( + closed.load(std::sync::atomic::Ordering::SeqCst), + "the Live session outlived the offer" + ); + + let [provisional, report] = &room.2[..] else { + panic!("expected two reports, got {:?}", room.2); + }; + assert_eq!(provisional["incomplete"], true); + assert_eq!( + provisional["reportRecovery"]["retryAfterSeconds"], + REPORT_RETRY_COOLDOWN.as_secs() + ); + assert!(report.get("reportRecovery").is_none()); + assert_ne!(report["incomplete"], true, "{report}"); + assert_eq!(report["codingScore"], 82); + assert_eq!(room.1, [RecoveryNotice::Accepted]); + + let bodies = bodies.lock().unwrap(); + assert_eq!(bodies.len(), 6); + assert!(bodies.len() <= 10); + assert!(bodies.iter().all(|body| body == &bodies[0])); + drop(tx); +} + +/// A candidate who left before the interview ended cannot ask for a retry, and +/// their departure was consumed by the interview loop: the final failure goes +/// out at once instead of a five-minute wait in an empty room. +#[tokio::test] +async fn an_absent_candidate_gets_the_failure_without_a_recovery_window() { + 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("absent"); + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + let mut room = RecoveryFixture::new(rx); + room.3 = false; + let regenerated = std::sync::atomic::AtomicBool::new(false); + let closed = std::sync::atomic::AtomicBool::new(false); + run_recovery( + &mut room, + RecoveryReport { + boot: &boot, + state: &mut frozen.state, + reason: "time_up", + keys: &keys, + }, + Err(elapsed().await), + async { + regenerated.store(true, std::sync::atomic::Ordering::SeqCst); + Ok(Ok(serde_json::json!({}))) + }, + tokio::time::Instant::now, + async { closed.store(true, std::sync::atomic::Ordering::SeqCst) }, + ) + .await + .unwrap(); + assert!( + closed.load(std::sync::atomic::Ordering::SeqCst), + "the Live session stayed open" + ); + let [report] = &room.2[..] else { + panic!("expected one report, got {:?}", room.2); + }; + assert_eq!(report["incomplete"], true); + assert!(report.get("reportRecovery").is_none()); + assert!(room.1.is_empty()); + assert!(!regenerated.load(std::sync::atomic::Ordering::SeqCst)); + drop(tx); +} + +#[test] +fn only_the_candidate_coming_and_going_moves_a_recovery_wait() { + for event in [RecoveryEvent::Away, RecoveryEvent::Back] { + assert_eq!( + candidate_only("candidate-original", "candidate-original", event), + event + ); + assert_eq!( + candidate_only("observer-1", "candidate-original", event), + RecoveryEvent::Ignore + ); + } +} + +/// Drives a recovery wait from a script of events, each sent after its delay. +async fn scripted_recovery( + script: Vec<(u64, RecoveryEvent)>, + generation: impl std::future::Future, + readiness: Readiness, +) -> (Option, Vec) { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + let mut fixture = RecoveryFixture::new(rx); + let sender = async move { + for (delay, event) in script { + tokio::time::sleep(std::time::Duration::from_secs(delay)).await; + let _ = tx.send(event); + } + tx + }; + let (result, _tx) = tokio::join!( + recover_report( + &mut fixture, + generation, + REPORT_RECOVERY_WINDOW, + std::time::Duration::ZERO, + tokio::time::Instant::now(), + move || readiness, + ), + sender, + ); + (result, fixture.1) +} + +/// A full LiveKit rejoin looks like the candidate leaving and coming back. An +/// accepted regeneration keeps running through it, and the report it returns +/// while they are away waits for them instead of going into an empty room. +#[tokio::test(start_paused = true)] +async fn a_rejoin_within_the_grace_keeps_an_accepted_regeneration() { + let started = tokio::time::Instant::now(); + let generation = async { + tokio::time::sleep(std::time::Duration::from_secs(5)).await; + Ok(Ok(serde_json::json!({"result": "complete"}))) + }; + let (result, notices) = scripted_recovery( + vec![ + (0, RecoveryEvent::Retry), + (1, RecoveryEvent::Away), + (19, RecoveryEvent::Back), + ], + generation, + Readiness::Ready, + ) + .await; + assert_eq!(result.unwrap().unwrap().unwrap()["result"], "complete"); + assert_eq!(notices, [RecoveryNotice::Accepted]); + assert!(started.elapsed() >= std::time::Duration::from_secs(20)); +} + +/// Gone past the grace is gone, before a retry or during the regeneration. +#[tokio::test(start_paused = true)] +async fn a_candidate_gone_past_the_grace_ends_the_wait() { + for script in [ + vec![(0, RecoveryEvent::Away)], + vec![(0, RecoveryEvent::Retry), (1, RecoveryEvent::Away)], + ] { + let started = tokio::time::Instant::now(); + let (result, _) = scripted_recovery(script, std::future::pending(), Readiness::Ready).await; + assert!(result.is_none()); + let waited = started.elapsed(); + assert!(waited >= REJOIN_GRACE, "{waited:?}"); + assert!( + waited < REJOIN_GRACE + std::time::Duration::from_secs(2), + "{waited:?}" + ); + } +} + +/// The keys are asked again when a retry arrives: another interview can have +/// put them back on quota since the offer went out. Accepting then would fail +/// before the first call and spend the one retry. +#[tokio::test(start_paused = true)] +async fn a_retry_waits_for_keys_that_went_out_after_the_offer() { + let calls = std::sync::atomic::AtomicUsize::new(0); + let generation = async { + calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + Ok(Ok(serde_json::json!({}))) + }; + let wait = std::time::Duration::from_secs(42); + let (result, notices) = scripted_recovery( + vec![(0, RecoveryEvent::Retry)], + generation, + Readiness::Wait(wait), + ) + .await; + assert!(result.is_none()); + assert_eq!( + notices, + [RecoveryNotice::Early(wait), RecoveryNotice::Closed] + ); + assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 0); + + // Keys that can never answer close the window rather than spend it. + let started = tokio::time::Instant::now(); + let (result, notices) = scripted_recovery( + vec![(0, RecoveryEvent::Retry)], + std::future::pending(), + Readiness::Never, + ) + .await; + assert!(result.is_none()); + assert_eq!(notices, [RecoveryNotice::Closed]); + assert!(started.elapsed() < std::time::Duration::from_secs(1)); +} diff --git a/tests/web.rs b/tests/web.rs index 8703fd97..be645177 100644 --- a/tests/web.rs +++ b/tests/web.rs @@ -24,9 +24,9 @@ use codetrial::token::{ livekit_token, }; use codetrial::web::{ - MAX_BODY_BYTES, MAX_REPORT_BYTES, READ_RATE_LIMIT, REPLAY_RATE_LIMIT, RoomDispatcher, - TOKEN_RATE_LIMIT, TokenConfig, WebServerConfig, initialize_account_database, login_config, - static_file_meta, token_response, + DispatchRefusal, MAX_BODY_BYTES, MAX_REPORT_BYTES, READ_RATE_LIMIT, REPLAY_RATE_LIMIT, + RoomDispatcher, TOKEN_RATE_LIMIT, TokenConfig, WebServerConfig, initialize_account_database, + login_config, static_file_meta, token_response, }; #[path = "common/http.rs"] @@ -322,18 +322,27 @@ struct RecordingDispatcher { staffed: std::sync::Mutex>, /// Set to refuse, standing in for a process already at capacity. at_capacity: std::sync::atomic::AtomicBool, + /// Set to refuse, standing in for a fixed room still finishing a report. + finalizing: std::sync::atomic::AtomicBool, } impl RoomDispatcher for RecordingDispatcher { - fn ensure_agent(&self, room_name: &str, provider: &codetrial::config::Provider) -> bool { + fn ensure_agent( + &self, + room_name: &str, + provider: &codetrial::config::Provider, + ) -> Result<(), DispatchRefusal> { if self.at_capacity.load(std::sync::atomic::Ordering::Relaxed) { - return false; + return Err(DispatchRefusal::AtCapacity); + } + if self.finalizing.load(std::sync::atomic::Ordering::Relaxed) { + return Err(DispatchRefusal::Finalizing); } self.staffed .lock() .unwrap() .push((room_name.to_string(), provider.url.clone())); - true + Ok(()) } } diff --git a/tests/web/contract.rs b/tests/web/contract.rs index 4291d67f..23b3403b 100644 --- a/tests/web/contract.rs +++ b/tests/web/contract.rs @@ -191,16 +191,18 @@ fn static_interview_script_keeps_transcript_and_report_contract() { // The report card and markdown export are asserted behaviorally in // tests/browser/render.test.js. What only Rust can check is that the // browser still routes a report through the sanitizer before rendering it. - assert_eq!( - call_arguments(report, "sanitizeReport("), - ["JSON.parse(new TextDecoder().decode(payload))"] - ); + assert_eq!(call_arguments(report, "sanitizeReport("), ["raw"]); let report = compact(report); + + // The first call above is the provisional branch, which returns early; the + // report that is rendered goes through its own. + assert!(report.contains(&compact("state.report = sanitizeReport(raw);"))); + assert!(report.contains("constraw=JSON.parse(newTextDecoder().decode(payload))")); for snippet in [ "setLocalAudioEnabled(false)", "saveHistory()", "renderReport()", - "room.disconnect()", + "room?.disconnect()", "state.room = null", ] { assert!( diff --git a/tests/web/token.rs b/tests/web/token.rs index 26d12d6a..844cfa00 100644 --- a/tests/web/token.rs +++ b/tests/web/token.rs @@ -946,6 +946,39 @@ async fn a_room_that_cannot_be_staffed_is_refused_rather_than_sold() { remove_database(db_path).await; } +/// A room still finishing its last report is not a full server, and saying it +/// was sent operators hunting for capacity that was there all along. +#[tokio::test] +async fn a_room_still_finishing_a_report_says_so() { + let (mut config, cookie, db_path) = signed_in_web_config("dispatch-finalizing"); + config.production = true; + config.fixed_room_name = None; + let dispatcher = std::sync::Arc::::default(); + dispatcher + .finalizing + .store(true, std::sync::atomic::Ordering::Relaxed); + let (base, server) = + spawn_web_server_with_dispatcher(config, std::sync::Arc::clone(&dispatcher)).await; + + let response = http_client() + .post(format!("{base}/api/token")) + .header("cookie", &cookie) + .json(&json!({"problemId":"two-sum","durationMin":45})) + .send() + .await + .unwrap(); + + assert_eq!(response.status(), 503); + let body: Value = response.json().await.unwrap(); + let error = body["error"].as_str().unwrap(); + assert!(error.contains("finishing"), "{body}"); + assert!(!error.contains("as many interviews"), "{body}"); + assert!(body.get("token").is_none(), "a refused room has no token"); + + server.shutdown().await; + remove_database(db_path).await; +} + #[tokio::test] async fn token_api_rejects_oversize_body() { let (config, cookie, db_path) = signed_in_web_config("oversize"); diff --git a/web/history.js b/web/history.js index daad3847..0143911c 100644 --- a/web/history.js +++ b/web/history.js @@ -244,10 +244,29 @@ export async function saveReportHistory( { fetcher = fetch, storage } = {}, ) { const local = saveLocalReport(entry, storage); - const account = await saveAccountReport(entry, fetcher); + const account = await inOrder(entry?.id, () => + saveAccountReport(entry, fetcher), + ); return { local, account }; } +// Account saves of one report id, in the order they were asked for. A report +// saved again under its id -- an outcome over the provisional failure it +// replaces -- otherwise raced the first save's session check, and whichever +// POST landed last was the copy the account kept. +const accountSaves = new Map(); + +function inOrder(id, save) { + if (id == null) return save(); + const next = (accountSaves.get(id) ?? Promise.resolve()).then(save, save); + accountSaves.set(id, next); + const forget = () => { + if (accountSaves.get(id) === next) accountSaves.delete(id); + }; + next.then(forget, forget); + return next; +} + /// `account` is what the page knows: whether the history it is showing came /// from an account. The session is still rechecked here rather than trusted /// from page load, because signing in from another tab has to reach the @@ -352,9 +371,15 @@ function saveLocalReport(entry, storage) { // Reading it can rebuild it, and that write is the 164 KB one: spending the // remaining quota on it here refused the save of a report the device had // room for, and the candidate was told it was gone. + // A save under an id already held replaces that row: an interview's final + // outcome is saved over the provisional failure it kept while the report + // could still be regenerated. + const others = previous.filter( + (row) => entry?.id == null || row?.id !== entry.id, + ); storage.setItem( historyKey, - JSON.stringify([entry, ...previous].slice(0, 20)), + JSON.stringify([entry, ...others].slice(0, 20)), ); // The short history is what the lobby draws this report from, so a review // store that refuses the write costs the reopen twenty attempts from now, diff --git a/web/interview.html b/web/interview.html index fb7ecaa9..caf800aa 100644 --- a/web/interview.html +++ b/web/interview.html @@ -400,8 +400,8 @@

Media preflight