From 6dd08c898799b17d56182b2c4b7f22ebb36de241 Mon Sep 17 00:00:00 2001 From: Jim Huang Date: Fri, 2 Oct 2026 12:18:06 +0800 Subject: [PATCH 1/2] Rebuild a resumed session that keeps failing Gemini Live can close a socket with 1011 and then close every socket resumed from its checkpoint 10 to 25 seconds later, so the interview spent all eight restarts resuming the same failing conversation and ended. A resumed socket that dies before it is healthy and debt-free now makes the next replacement rebuild from local state, and resuming stays off until a reply beyond the recovery briefing or a healthy socket proves recovery. The briefing's own answer no longer resets the restart budget, since every replacement asks for one. While no socket transcribes the candidate, a candidate who kept talking was counted silent and nudged, so until a socket has transcribed them their microphone level also counts as speech for the idle timers. The log now says when a cold rebuild fails too, and the size of what each replacement was handed, which tells a failing checkpoint from input that fails on any socket. Close #198 --- docs/provider-cost-and-degradation.md | 16 +++- src/livekit.rs | 125 ++++++++++++++++++++++---- src/livekit/media.rs | 7 +- src/livekit/session.rs | 2 +- src/livekit/turn.rs | 19 ++++ tests/unit/livekit.rs | 114 +++++++++++++++++++++++ tests/unit/livekit/turn.rs | 30 +++++++ 7 files changed, 290 insertions(+), 23 deletions(-) diff --git a/docs/provider-cost-and-degradation.md b/docs/provider-cost-and-degradation.md index dffce2e3..70dd0268 100644 --- a/docs/provider-cost-and-degradation.md +++ b/docs/provider-cost-and-degradation.md @@ -15,7 +15,13 @@ moment, so the advisory is spent then rather than held for a turn boundary that Gemini will not send. The replacement resumes the same conversation when the server issued a handle; if the handle is unavailable or refused, the restart is logged as degraded and is grounded from the bounded local transcript tail, -editor, round and evidence state instead. +editor, round and evidence state instead. A resumed socket replaced +before reaching the healthy, debt-free age is rebuilt locally on the next +attempt, even if it answered once before failing. Resumption stays disabled +through that failure run until completed output or a healthy, debt-free socket +proves recovery; these cold opens consume the same restart budget. The reply to +a recovery briefing is not that proof: every replacement asks for one, so a run +of sockets that each answer their briefing and die still ends the interview. A candidate turn, required prompt (including the opening greeting), or tool continuation that produces no output for 45 seconds replaces the socket even @@ -31,7 +37,10 @@ cut off mid-generation is owed too, and the replacement is told not to repeat what was already said. Optional editor reviews accept silence; queued audio and a paused interview do not trigger this watchdog. A generation that stops producing output for the same interval also recovers. Periodic nudges wait -while a reply is owed rather than replacing its debt. +while a reply is owed rather than replacing its debt. Until a socket has +transcribed the candidate, the first one and every replacement alike, the idle +timers also count the candidate's microphone level as speech, so a socket that +is not hearing them does not take their talking for silence and nudge them. A close the interviewer asked for is not recovered either: it waits on the tool acknowledgement, and if that never comes, or starts and then stops, the @@ -48,7 +57,8 @@ respond to a silence that goes on. It is logged as a deliberate silence, so a report of the interviewer going quiet can be told apart from a stalled socket. `GEMINI_RESTART_LIMIT` bounds a failing endpoint rather than a long interview. -It allows 8 opens in a row. A completed turn with output clears the run, and so +It allows 8 opens in a row. A completed turn with output clears the run, unless +it answered a recovery briefing, and so does replacing a socket that lived past a minute and owed nothing. A socket replaced while it owed a reply never counts as healthy, however long it stayed connected and whether the watchdog, a `GoAway` or the server closed it, so diff --git a/src/livekit.rs b/src/livekit.rs index 44c028d5..bd516fba 100644 --- a/src/livekit.rs +++ b/src/livekit.rs @@ -254,6 +254,64 @@ fn take_restart_attempt(restarts: &mut usize, socket_age: Duration) -> bool { true } +/// A checkpoint can be accepted yet restore a conversation that fails again. +/// Once that happens, rebuild locally until output proves recovery worked. +/// +/// `resumed` describes the socket now open and only a connect writes it: a +/// resumed socket that answers once and then dies young is still a failed +/// resumption, so completed output lifts the suppression without forgetting +/// where the open socket came from. +/// +/// The answer to a recovery briefing proves nothing either way. Every +/// replacement asks for one, so counting it would let a run of sockets that +/// each answer their briefing and die reset the budget forever. It is the +/// first turn to end after the briefing rather than a reply tagged with the +/// briefing's cause, because a tool call in that answer retags the +/// continuation that completes it. +#[derive(Default)] +struct ResumeRecovery { + resumed: bool, + suppressed: bool, + briefing_reply_pending: bool, +} + +impl ResumeRecovery { + /// Judged on every replacement, whether or not a handle exists to use. + fn allow_resume(&mut self, age: Duration) -> bool { + if age >= HEALTHY_GEMINI_SOCKET { + self.suppressed = false; + } else if self.resumed { + self.suppressed = true; + } + !self.suppressed + } + + /// The replacement is up, and `briefed` says its briefing asked for a + /// reply. + fn connected(&mut self, resumed: bool, briefed: bool) { + self.resumed = resumed; + self.briefing_reply_pending = briefed; + } + + fn note_event( + &mut self, + restarts: &mut usize, + state: &RuntimeState, + activity: &RuntimeActivity, + event: &GeminiEvent, + ) { + // Whichever way the briefing's turn ends, silent or cut off by the + // candidate, it is over, and the next reply is not its answer. + if session::ends_turn(event) && std::mem::take(&mut self.briefing_reply_pending) { + return; + } + if completed_live_reply(state, activity, event) { + *restarts = 0; + self.suppressed = false; + } + } +} + /// The agent's output has settled: nothing generating, and nothing left in the /// LiveKit playout queue. The floor alone is not enough, because it is stamped /// once when a turn completes and the queue drains on its own afterwards. @@ -489,6 +547,9 @@ async fn replace_gemini_session( leave_room(room).await; return Ok(ControlFlow::Break(())); } + + // Reading the handle also records closing-key failures for cold opens and + // reporting, even when the handle or the restart budget cannot be used. let handle = context.gemini.recovery_handle(interview.keys); let age = replaced_socket_age(context.gemini.age(), reply_timeout, context.activity); if !take_restart_attempt(&mut loops.restarts, age) { @@ -500,6 +561,27 @@ async fn replace_gemini_session( return end_without_interviewer(room, context, interview, loops).await; } + // Judged before the handle is looked at: `Option::filter` skips its closure + // on `None`, which would leave suppression stale across a replacement that + // had no handle to offer. + let allow_resume = loops.resume_recovery.allow_resume(age); + let handle = handle.filter(|_| allow_resume); + + // The cold-rebuild line tells the two causes of a 1011 run apart: a resumed + // conversation that keeps failing stops failing once it is rebuilt, while + // input that fails on any socket fails the rebuild too. + if loops.resume_recovery.suppressed { + eprintln!( + "{}; rebuilding from local state room={}", + if loops.resume_recovery.resumed { + "Gemini resumed session failed before recovery" + } else { + "Gemini cold rebuild also failed before recovery" + }, + interview.boot.room_name + ); + } + // Said before the attempt, not after it: the whole point is to cover the // gap, and the gap starts here. Connect and setup are bounded at fifteen // seconds each, and for that long the candidate is talking to a socket that @@ -511,11 +593,12 @@ async fn replace_gemini_session( // one being replaced ahead of its `GoAway` is still live and this is the // orderly hang-up. Ignored either way for that reason. A socket that // resumed and was offered nothing new still holds the checkpoint it - // inherited, whose age this process does not know. - let checkpoint_age = match (context.gemini.checkpoint_age(), &handle) { - (Some(age), _) => format!("{}s", age.as_secs()), - (None, Some(_)) => "inherited".to_string(), - (None, None) => "none".to_string(), + // inherited, whose age this process does not know. A checkpoint that is not + // offered is not the one the replacement starts from, so it is not named. + let checkpoint_age = match (&handle, context.gemini.checkpoint_age()) { + (None, _) => "none".to_string(), + (Some(_), Some(age)) => format!("{}s", age.as_secs()), + (Some(_), None) => "inherited".to_string(), }; let _ = context.gemini.shutdown().await; session::drain_live_usage(room, context); @@ -584,8 +667,12 @@ async fn replace_gemini_session( // how old the checkpoint it resumed from is, what was owed, and what the // local record holds that the checkpoint may predate. eprintln!( - "replacement: at={} resumed={resumed} checkpoint_age={checkpoint_age} owed={owed} ({debt}) {} room={}", + "replacement: at={} resumed={resumed} checkpoint_age={checkpoint_age} owed={owed} ({debt}) owed_chars={} editor_chars={} {} room={}", log_clock(context.state), + owed_prompt + .as_deref() + .map_or(0, |prompt| prompt.chars().count()), + context.state.code.chars().count(), prompt_fields(context.state, context.activity), interview.boot.room_name ); @@ -604,6 +691,7 @@ async fn replace_gemini_session( owed_prompt.as_deref(), ) .await; + loops.resume_recovery.connected(resumed, spoke); if spoke { eprintln!( "{}", @@ -733,7 +821,7 @@ async fn brief_replacement( eprintln!("Gemini session resumed; the interview continues where it left off"); } else { eprintln!( - "Gemini session restart degraded; resumption was unavailable and the interviewer is rebuilding from local transcript, editor and interview state" + "Gemini session restart degraded; resumption was not used and the interviewer is rebuilding from local transcript, editor and interview state" ); } @@ -1204,6 +1292,7 @@ struct RoomLoop { /// How many sockets this interview has been through. restarts: usize, deferred_restart: DeferredRestart, + resume_recovery: ResumeRecovery, /// At most one idle-window review at a time, collected on the watch tick. /// A tick of latency on a note nobody is waiting for is not worth an arm @@ -1717,11 +1806,12 @@ async fn on_gemini_event( return Ok(ControlFlow::Continue(())); }; - // Completed output proves recovery worked; merely staying connected does - // not clear consecutive watchdog failures. - if completed_live_reply(context.state, context.activity, &event) { - loops.restarts = 0; - } + // Completed output proves recovery worked, unless it answered the recovery + // briefing; merely staying connected does not clear consecutive watchdog + // failures. + loops + .resume_recovery + .note_event(&mut loops.restarts, context.state, context.activity, &event); if let GeminiEvent::GoAway { time_left } = &event { eprintln!( @@ -1937,6 +2027,7 @@ pub async fn run_room( let mut loops = RoomLoop { restarts, deferred_restart: DeferredRestart::default(), + resume_recovery: ResumeRecovery::default(), interim_review: InterimReview::default(), presence: CandidatePresence::default(), }; @@ -2034,13 +2125,15 @@ pub async fn run_room( // need; propagating cost the interview. let ended = frame.is_none(); - // Only a checkpoint reads it, and only under a compression - // window; without one the per-frame level is not worth + // Read by a checkpoint under a compression window, and by + // the idle timers until the socket has transcribed the + // candidate; otherwise the per-frame level is not worth // computing. - if turn.state.context_compression.is_some() + if (turn.state.context_compression.is_some() + || !turn.activity.transcribed_on_socket) && frame.as_ref().is_some_and(media::frame_has_voice) { - turn.activity.candidate_voice_at = Some(Instant::now()); + turn.activity.note_candidate_voice(Instant::now()); } if turn.state.paused { discard_paused_audio(&mut media); diff --git a/src/livekit/media.rs b/src/livekit/media.rs index 4672975c..eb882339 100644 --- a/src/livekit/media.rs +++ b/src/livekit/media.rs @@ -459,9 +459,10 @@ pub(super) async fn publish_output_audio( /// Root-mean-square level, in PCM16 units, above which a frame counts as the /// candidate speaking: about -40 dBFS, the top of a quiet room and below soft -/// speech. Only used to hold a checkpoint back, so it errs toward calling -/// sound speech: a noisy room waits longer, which the watch tick retries, -/// where a soft speaker taken for silence would be interrupted. +/// speech. Used to hold a checkpoint back, and a silence nudge while the socket +/// has not transcribed the candidate, so it errs toward calling sound speech: a +/// noisy room waits longer, which the watch tick retries, where a soft speaker +/// taken for silence would be interrupted. const VOICE_RMS: f64 = 316.0; pub(super) fn frame_has_voice(frame: &AudioFrame<'_>) -> bool { diff --git a/src/livekit/session.rs b/src/livekit/session.rs index 4f34fae1..70063302 100644 --- a/src/livekit/session.rs +++ b/src/livekit/session.rs @@ -304,7 +304,7 @@ fn record_live_usage( } /// A turn ending, whichever way it ends. -fn ends_turn(event: &GeminiEvent) -> bool { +pub(super) fn ends_turn(event: &GeminiEvent) -> bool { matches!(event, GeminiEvent::TurnComplete | GeminiEvent::Interrupted) } diff --git a/src/livekit/turn.rs b/src/livekit/turn.rs index a2b7442a..93e28775 100644 --- a/src/livekit/turn.rs +++ b/src/livekit/turn.rs @@ -238,6 +238,10 @@ pub(super) struct RuntimeActivity { /// transcript is no guide here: it arrives after the speech it transcribes, /// and a checkpoint sent in that gap lands in the middle of an utterance. pub(super) candidate_voice_at: Option, + /// Whether the open socket has transcribed the candidate yet. Until it + /// has, it may not be hearing them at all, as a socket on its way to a + /// 1011 does not, so their voice on the track stamps `last_user_speech`. + pub(super) transcribed_on_socket: bool, } /// How long the microphone has to stay quiet before a checkpoint may go out: @@ -421,6 +425,19 @@ impl RuntimeActivity { } self.latest_prompt_tokens = None; self.context_refresh_pending = false; + self.transcribed_on_socket = false; + } + + /// The candidate's microphone carried more than room noise. Until the + /// socket has transcribed them, this is also the idle timers' evidence of + /// speech: a candidate talking to a socket that cannot hear them would + /// otherwise be counted silent and nudged. Only until then, because room + /// noise above `VOICE_RMS` would hold every nudge back for good. + pub(super) fn note_candidate_voice(&mut self, now: Instant) { + self.candidate_voice_at = Some(now); + if !self.transcribed_on_socket { + self.last_user_speech = now; + } } #[cfg(test)] @@ -481,6 +498,7 @@ impl RuntimeActivity { peak_prompt_tokens: 0, context_refresh_pending: false, candidate_voice_at: None, + transcribed_on_socket: false, // Seeded at `now` rather than in the past: the first minutes of an // interview are the greeting and the problem statement, and there @@ -799,6 +817,7 @@ impl RuntimeActivity { /// back, so this is the only place that can tell it from a new answer. pub(super) fn note_candidate_finished(&mut self, now: Instant, audio_playing: bool) { self.last_user_speech = now; + self.transcribed_on_socket = true; // More speech gets its own native reply; asking as well would answer // twice, or over the candidate. diff --git a/tests/unit/livekit.rs b/tests/unit/livekit.rs index 2f4459f9..6a160b8b 100644 --- a/tests/unit/livekit.rs +++ b/tests/unit/livekit.rs @@ -5158,3 +5158,117 @@ fn a_thinking_hold_is_not_an_interviewer_stall() { ReplyWatch::Recover ); } + +/// A completed reply the room loop counts as output, ended by `event`. +fn reply_event(recovery: &mut ResumeRecovery, restarts: &mut usize, event: &GeminiEvent) { + let mut activity = RuntimeActivity::new(Instant::now()); + activity.note_output(Instant::now()); + recovery.note_event(restarts, &RuntimeState::default(), &activity, event); +} + +fn complete_reply(recovery: &mut ResumeRecovery, restarts: &mut usize) { + reply_event(recovery, restarts, &GeminiEvent::TurnComplete); +} + +const YOUNG: Duration = Duration::from_secs(12); + +#[test] +fn failed_resumption_rebuilds_until_recovery() { + let mut recovery = ResumeRecovery::default(); + let mut restarts = 0; + assert!(recovery.allow_resume(Duration::ZERO)); + recovery.connected(true, false); + assert!(!recovery.allow_resume(YOUNG)); + recovery.connected(false, false); + assert!(!recovery.allow_resume(YOUNG)); + recovery.connected(false, false); + complete_reply(&mut recovery, &mut restarts); + assert!(recovery.allow_resume(YOUNG)); +} + +#[test] +fn healthy_socket_lifts_suppression() { + let mut recovery = ResumeRecovery::default(); + recovery.connected(true, false); + assert!(!recovery.allow_resume(YOUNG)); + recovery.connected(false, false); + assert!(recovery.allow_resume(HEALTHY_GEMINI_SOCKET)); + recovery.connected(true, false); + assert!(recovery.allow_resume(HEALTHY_GEMINI_SOCKET)); +} + +#[test] +fn resumed_socket_that_answers_then_dies_young_still_rebuilds() { + let mut recovery = ResumeRecovery::default(); + let mut restarts = 0; + recovery.connected(true, false); + complete_reply(&mut recovery, &mut restarts); + assert!(!recovery.allow_resume(YOUNG)); +} + +#[test] +fn refused_resumption_does_not_suppress_the_next_one() { + let mut recovery = ResumeRecovery::default(); + assert!(recovery.allow_resume(Duration::ZERO)); + recovery.connected(false, false); + assert!(recovery.allow_resume(YOUNG)); +} + +#[test] +fn sockets_that_only_answer_their_briefing_still_end_the_interview() { + let mut recovery = ResumeRecovery::default(); + let mut restarts = 0; + let mut attempts = 0; + while take_restart_attempt(&mut restarts, YOUNG) { + let resumed = recovery.allow_resume(YOUNG); + recovery.connected(resumed, true); + complete_reply(&mut recovery, &mut restarts); + attempts += 1; + assert!(attempts <= GEMINI_RESTART_LIMIT); + } + assert_eq!(attempts, GEMINI_RESTART_LIMIT); +} + +#[test] +fn only_output_beyond_the_briefing_proves_recovery() { + let mut recovery = ResumeRecovery::default(); + let mut restarts = GEMINI_RESTART_LIMIT; + recovery.connected(true, false); + assert!(!recovery.allow_resume(YOUNG)); + recovery.connected(false, true); + + // A tool call in the briefing's answer ends nothing; the continuation's + // completion is still that answer, whatever cause the tool response gave + // it. + reply_event( + &mut recovery, + &mut restarts, + &GeminiEvent::ToolCall(Vec::new()), + ); + complete_reply(&mut recovery, &mut restarts); + assert_eq!(restarts, GEMINI_RESTART_LIMIT); + assert!(recovery.suppressed); + + complete_reply(&mut recovery, &mut restarts); + assert_eq!(restarts, 0); + assert!(!recovery.suppressed); +} + +#[test] +fn a_silent_or_interrupted_briefing_does_not_swallow_the_next_reply() { + for ending in [GeminiEvent::TurnComplete, GeminiEvent::Interrupted] { + let mut recovery = ResumeRecovery::default(); + let mut restarts = GEMINI_RESTART_LIMIT; + recovery.connected(false, true); + recovery.note_event( + &mut restarts, + &RuntimeState::default(), + &RuntimeActivity::new(Instant::now()), + &ending, + ); + assert_eq!(restarts, GEMINI_RESTART_LIMIT); + + complete_reply(&mut recovery, &mut restarts); + assert_eq!(restarts, 0); + } +} diff --git a/tests/unit/livekit/turn.rs b/tests/unit/livekit/turn.rs index 7864f760..c5857135 100644 --- a/tests/unit/livekit/turn.rs +++ b/tests/unit/livekit/turn.rs @@ -1414,3 +1414,33 @@ fn audio_still_draining_spares_no_deadline_once_the_floor_is_the_candidates() { activity.note_candidate_finished(now, true); assert!(activity.reply_in_flight()); } + +#[test] +fn voice_holds_the_silence_nudge_until_the_socket_transcribes() { + let start = Instant::now(); + let now = start + Duration::from_secs(300); + let mut state = RuntimeState::default(); + let mut activity = RuntimeActivity::new(start); + activity.note_candidate_voice(now - Duration::from_secs(5)); + assert!(activity.watch_prompt(&mut state, now).is_none()); + + // Once the socket has transcribed the candidate, the transcript is the + // evidence and a voice level alone, room noise included, holds nothing. + let mut activity = RuntimeActivity::new(start); + activity.transcribed_on_socket = true; + activity.note_candidate_voice(now - Duration::from_secs(5)); + assert!(activity.watch_prompt(&mut state, now).is_some()); +} + +#[test] +fn a_replacement_socket_has_to_transcribe_the_candidate_again() { + let start = Instant::now(); + let mut activity = RuntimeActivity::new(start); + activity.note_candidate_finished(start, false); + activity.note_candidate_voice(start + Duration::from_secs(5)); + assert_eq!(activity.last_user_speech, start); + + activity.reset_context_observations(true); + activity.note_candidate_voice(start + Duration::from_secs(9)); + assert_eq!(activity.last_user_speech, start + Duration::from_secs(9)); +} From 7742a4704d6f625f53016fa38f2b196661cab958 Mon Sep 17 00:00:00 2001 From: Jim Huang Date: Fri, 2 Oct 2026 13:33:07 +0800 Subject: [PATCH 2/2] Judge resumption by how long the socket stayed up Resumption was suppressed by the debt-adjusted age the restart budget uses, so a resumed socket that ran for minutes and was replaced while owing a reply lost its handle, and the log blamed the checkpoint. The resume decision now uses how long the socket stayed up. A recovery briefing held through a pause or a thinking hold was answered after the unpause or the hold's end, and that answer reset the budget, so a held briefing now marks its reply the way an immediate one does. --- docs/provider-cost-and-degradation.md | 15 ++-- src/agent.rs | 5 ++ src/agent/turn_taking.rs | 8 +- src/livekit.rs | 26 ++++--- tests/unit/livekit.rs | 103 ++++++++++++++++++++------ 5 files changed, 115 insertions(+), 42 deletions(-) diff --git a/docs/provider-cost-and-degradation.md b/docs/provider-cost-and-degradation.md index 70dd0268..d69d79cd 100644 --- a/docs/provider-cost-and-degradation.md +++ b/docs/provider-cost-and-degradation.md @@ -16,12 +16,15 @@ Gemini will not send. The replacement resumes the same conversation when the server issued a handle; if the handle is unavailable or refused, the restart is logged as degraded and is grounded from the bounded local transcript tail, editor, round and evidence state instead. A resumed socket replaced -before reaching the healthy, debt-free age is rebuilt locally on the next -attempt, even if it answered once before failing. Resumption stays disabled -through that failure run until completed output or a healthy, debt-free socket -proves recovery; these cold opens consume the same restart budget. The reply to -a recovery briefing is not that proof: every replacement asks for one, so a run -of sockets that each answer their briefing and die still ends the interview. +within a minute of opening is rebuilt locally on the next attempt, even if it +answered once before failing; one that stayed up longer did not fail at its +checkpoint, whatever it owed when it was replaced. Resumption stays disabled +through that failure run until completed output or a socket that stays up past +a minute proves recovery; these cold opens consume the same restart budget. The +reply to a recovery briefing is not that proof, whether the briefing went out +at once or was held for an unpause or the end of a thinking hold: every +replacement asks for one, so a run of sockets that each answer their briefing +and die still ends the interview. A candidate turn, required prompt (including the opening greeting), or tool continuation that produces no output for 45 seconds replaces the socket even diff --git a/src/agent.rs b/src/agent.rs index 47c0b754..2ddd9552 100644 --- a/src/agent.rs +++ b/src/agent.rs @@ -753,6 +753,10 @@ pub struct RuntimeState { /// A reply the hold dropped is still in the model's history, and the /// release has to say so; see `thinking_resume`. pub thinking_unheard_reply: bool, + /// A recovery briefing, sent at once or held for an unpause or the end of + /// a hold, asked for a reply that has not ended yet. Its answer proves + /// nothing about the socket, so it does not reset the restart budget. + pub recovery_reply_pending: bool, pub framework_evidence: Vec, /// Deterministic, bounded facts derived from the live session. The ledger /// deliberately holds no editor text or runner diagnostics: those remain @@ -935,6 +939,7 @@ impl Default for RuntimeState { thinking_released_at: None, thinking_notice: None, thinking_unheard_reply: false, + recovery_reply_pending: false, framework_evidence: Vec::new(), evidence_ledger: EvidenceLedger::default(), code: String::new(), diff --git a/src/agent/turn_taking.rs b/src/agent/turn_taking.rs index 46303818..9a602c33 100644 --- a/src/agent/turn_taking.rs +++ b/src/agent/turn_taking.rs @@ -174,9 +174,15 @@ impl RuntimeState { /// The debt `thinking_debt` names has been delivered. pub(crate) fn clear_thinking_debt(&mut self) { - if std::mem::take(&mut self.needs_cold_brief) { + let cold_brief = std::mem::take(&mut self.needs_cold_brief); + if cold_brief { self.code_shown = self.code.clone(); } + + // A held briefing, or a reply owed across a pause or a replacement, is + // answered by the next turn. A pause can owe one with no replacement + // behind it; taking that answer as a briefing's only delays a reset. + self.recovery_reply_pending |= cold_brief || self.owed_reply_on_resume.is_some(); self.owed_reply_on_resume = None; self.thinking_unheard_reply = false; } diff --git a/src/livekit.rs b/src/livekit.rs index bd516fba..69e8e378 100644 --- a/src/livekit.rs +++ b/src/livekit.rs @@ -265,14 +265,17 @@ fn take_restart_attempt(restarts: &mut usize, socket_age: Duration) -> bool { /// The answer to a recovery briefing proves nothing either way. Every /// replacement asks for one, so counting it would let a run of sockets that /// each answer their briefing and die reset the budget forever. It is the -/// first turn to end after the briefing rather than a reply tagged with the -/// briefing's cause, because a tool call in that answer retags the -/// continuation that completes it. +/// first turn to end after the briefing (`recovery_reply_pending`) rather than +/// a reply tagged with the briefing's cause, because a tool call in that answer +/// retags the continuation that completes it. +/// +/// Judged by how long the socket stayed up, not by the debt-adjusted age the +/// restart budget uses: a resumed socket that carried minutes of interview and +/// was replaced owing a reply did not fail at its checkpoint. #[derive(Default)] struct ResumeRecovery { resumed: bool, suppressed: bool, - briefing_reply_pending: bool, } impl ResumeRecovery { @@ -286,23 +289,20 @@ impl ResumeRecovery { !self.suppressed } - /// The replacement is up, and `briefed` says its briefing asked for a - /// reply. - fn connected(&mut self, resumed: bool, briefed: bool) { + fn connected(&mut self, resumed: bool) { self.resumed = resumed; - self.briefing_reply_pending = briefed; } fn note_event( &mut self, restarts: &mut usize, - state: &RuntimeState, + state: &mut RuntimeState, activity: &RuntimeActivity, event: &GeminiEvent, ) { // Whichever way the briefing's turn ends, silent or cut off by the // candidate, it is over, and the next reply is not its answer. - if session::ends_turn(event) && std::mem::take(&mut self.briefing_reply_pending) { + if session::ends_turn(event) && std::mem::take(&mut state.recovery_reply_pending) { return; } if completed_live_reply(state, activity, event) { @@ -564,7 +564,7 @@ async fn replace_gemini_session( // Judged before the handle is looked at: `Option::filter` skips its closure // on `None`, which would leave suppression stale across a replacement that // had no handle to offer. - let allow_resume = loops.resume_recovery.allow_resume(age); + let allow_resume = loops.resume_recovery.allow_resume(context.gemini.age()); let handle = handle.filter(|_| allow_resume); // The cold-rebuild line tells the two causes of a 1011 run apart: a resumed @@ -657,6 +657,7 @@ async fn replace_gemini_session( .record_model_input(ModelInputKind::LiveSetup, &interview.boot.instructions); *context.gemini = session; + loops.resume_recovery.connected(resumed); context.activity.live_socket += 1; context.activity.reset_context_observations(resumed); @@ -691,7 +692,6 @@ async fn replace_gemini_session( owed_prompt.as_deref(), ) .await; - loops.resume_recovery.connected(resumed, spoke); if spoke { eprintln!( "{}", @@ -837,6 +837,7 @@ async fn brief_replacement( // naming the briefing there would nest it, transcript and all, // inside the next one. activity.mark_prompted(Instant::now(), owed_prompt, false); + state.recovery_reply_pending = true; true } Ok(false) => false, @@ -976,6 +977,7 @@ fn hold_owed_reply(state: &mut RuntimeState, owed_prompt: Option<&str>) { /// not create. fn clear_abandoned_socket_work(state: &mut RuntimeState, activity: &mut RuntimeActivity) { activity.discarding_output = false; + state.recovery_reply_pending = false; // A provisional request waits on its utterance's end, which the closed // socket will not send. A declared hold is the candidate's and survives. diff --git a/tests/unit/livekit.rs b/tests/unit/livekit.rs index 6a160b8b..6d195cec 100644 --- a/tests/unit/livekit.rs +++ b/tests/unit/livekit.rs @@ -5160,14 +5160,19 @@ fn a_thinking_hold_is_not_an_interviewer_stall() { } /// A completed reply the room loop counts as output, ended by `event`. -fn reply_event(recovery: &mut ResumeRecovery, restarts: &mut usize, event: &GeminiEvent) { +fn reply_event( + recovery: &mut ResumeRecovery, + state: &mut RuntimeState, + restarts: &mut usize, + event: &GeminiEvent, +) { let mut activity = RuntimeActivity::new(Instant::now()); activity.note_output(Instant::now()); - recovery.note_event(restarts, &RuntimeState::default(), &activity, event); + recovery.note_event(restarts, state, &activity, event); } -fn complete_reply(recovery: &mut ResumeRecovery, restarts: &mut usize) { - reply_event(recovery, restarts, &GeminiEvent::TurnComplete); +fn complete_reply(recovery: &mut ResumeRecovery, state: &mut RuntimeState, restarts: &mut usize) { + reply_event(recovery, state, restarts, &GeminiEvent::TurnComplete); } const YOUNG: Duration = Duration::from_secs(12); @@ -5175,25 +5180,26 @@ const YOUNG: Duration = Duration::from_secs(12); #[test] fn failed_resumption_rebuilds_until_recovery() { let mut recovery = ResumeRecovery::default(); + let mut state = RuntimeState::default(); let mut restarts = 0; assert!(recovery.allow_resume(Duration::ZERO)); - recovery.connected(true, false); + recovery.connected(true); assert!(!recovery.allow_resume(YOUNG)); - recovery.connected(false, false); + recovery.connected(false); assert!(!recovery.allow_resume(YOUNG)); - recovery.connected(false, false); - complete_reply(&mut recovery, &mut restarts); + recovery.connected(false); + complete_reply(&mut recovery, &mut state, &mut restarts); assert!(recovery.allow_resume(YOUNG)); } #[test] fn healthy_socket_lifts_suppression() { let mut recovery = ResumeRecovery::default(); - recovery.connected(true, false); + recovery.connected(true); assert!(!recovery.allow_resume(YOUNG)); - recovery.connected(false, false); + recovery.connected(false); assert!(recovery.allow_resume(HEALTHY_GEMINI_SOCKET)); - recovery.connected(true, false); + recovery.connected(true); assert!(recovery.allow_resume(HEALTHY_GEMINI_SOCKET)); } @@ -5201,8 +5207,8 @@ fn healthy_socket_lifts_suppression() { fn resumed_socket_that_answers_then_dies_young_still_rebuilds() { let mut recovery = ResumeRecovery::default(); let mut restarts = 0; - recovery.connected(true, false); - complete_reply(&mut recovery, &mut restarts); + recovery.connected(true); + complete_reply(&mut recovery, &mut RuntimeState::default(), &mut restarts); assert!(!recovery.allow_resume(YOUNG)); } @@ -5210,19 +5216,21 @@ fn resumed_socket_that_answers_then_dies_young_still_rebuilds() { fn refused_resumption_does_not_suppress_the_next_one() { let mut recovery = ResumeRecovery::default(); assert!(recovery.allow_resume(Duration::ZERO)); - recovery.connected(false, false); + recovery.connected(false); assert!(recovery.allow_resume(YOUNG)); } #[test] fn sockets_that_only_answer_their_briefing_still_end_the_interview() { let mut recovery = ResumeRecovery::default(); + let mut state = RuntimeState::default(); let mut restarts = 0; let mut attempts = 0; while take_restart_attempt(&mut restarts, YOUNG) { let resumed = recovery.allow_resume(YOUNG); - recovery.connected(resumed, true); - complete_reply(&mut recovery, &mut restarts); + recovery.connected(resumed); + state.recovery_reply_pending = true; + complete_reply(&mut recovery, &mut state, &mut restarts); attempts += 1; assert!(attempts <= GEMINI_RESTART_LIMIT); } @@ -5232,24 +5240,27 @@ fn sockets_that_only_answer_their_briefing_still_end_the_interview() { #[test] fn only_output_beyond_the_briefing_proves_recovery() { let mut recovery = ResumeRecovery::default(); + let mut state = RuntimeState::default(); let mut restarts = GEMINI_RESTART_LIMIT; - recovery.connected(true, false); + recovery.connected(true); assert!(!recovery.allow_resume(YOUNG)); - recovery.connected(false, true); + recovery.connected(false); + state.recovery_reply_pending = true; // A tool call in the briefing's answer ends nothing; the continuation's // completion is still that answer, whatever cause the tool response gave // it. reply_event( &mut recovery, + &mut state, &mut restarts, &GeminiEvent::ToolCall(Vec::new()), ); - complete_reply(&mut recovery, &mut restarts); + complete_reply(&mut recovery, &mut state, &mut restarts); assert_eq!(restarts, GEMINI_RESTART_LIMIT); assert!(recovery.suppressed); - complete_reply(&mut recovery, &mut restarts); + complete_reply(&mut recovery, &mut state, &mut restarts); assert_eq!(restarts, 0); assert!(!recovery.suppressed); } @@ -5258,17 +5269,63 @@ fn only_output_beyond_the_briefing_proves_recovery() { fn a_silent_or_interrupted_briefing_does_not_swallow_the_next_reply() { for ending in [GeminiEvent::TurnComplete, GeminiEvent::Interrupted] { let mut recovery = ResumeRecovery::default(); + let mut state = RuntimeState { + recovery_reply_pending: true, + ..RuntimeState::default() + }; let mut restarts = GEMINI_RESTART_LIMIT; - recovery.connected(false, true); recovery.note_event( &mut restarts, - &RuntimeState::default(), + &mut state, &RuntimeActivity::new(Instant::now()), &ending, ); assert_eq!(restarts, GEMINI_RESTART_LIMIT); - complete_reply(&mut recovery, &mut restarts); + complete_reply(&mut recovery, &mut state, &mut restarts); assert_eq!(restarts, 0); } } + +#[test] +fn a_briefing_held_for_an_unpause_or_a_hold_is_still_a_briefing() { + for held in [ + RuntimeState { + needs_cold_brief: true, + ..RuntimeState::default() + }, + RuntimeState { + owed_reply_on_resume: Some(None), + ..RuntimeState::default() + }, + ] { + let mut state = held; + let mut recovery = ResumeRecovery::default(); + let mut restarts = GEMINI_RESTART_LIMIT; + + // What the unpause and the end of a hold call once the held debt has + // gone out with the reply it asks for. + state.clear_thinking_debt(); + complete_reply(&mut recovery, &mut state, &mut restarts); + assert_eq!(restarts, GEMINI_RESTART_LIMIT); + complete_reply(&mut recovery, &mut state, &mut restarts); + assert_eq!(restarts, 0); + } + + let mut state = RuntimeState { + thinking_unheard_reply: true, + ..RuntimeState::default() + }; + state.clear_thinking_debt(); + assert!(!state.recovery_reply_pending); +} + +#[test] +fn a_replacement_drops_the_briefing_its_predecessor_owed() { + let mut state = RuntimeState { + recovery_reply_pending: true, + ..RuntimeState::default() + }; + clear_abandoned_socket_work(&mut state, &mut RuntimeActivity::new(Instant::now())); + assert!(!state.recovery_reply_pending); +}