diff --git a/docs/provider-cost-and-degradation.md b/docs/provider-cost-and-degradation.md index dffce2e3..d69d79cd 100644 --- a/docs/provider-cost-and-degradation.md +++ b/docs/provider-cost-and-degradation.md @@ -15,7 +15,16 @@ 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 +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 @@ -31,7 +40,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 +60,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/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 44c028d5..69e8e378 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 (`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, +} + +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 + } + + fn connected(&mut self, resumed: bool) { + self.resumed = resumed; + } + + fn note_event( + &mut self, + restarts: &mut usize, + 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 state.recovery_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(context.gemini.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); @@ -574,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); @@ -584,8 +668,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 ); @@ -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" ); } @@ -749,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, @@ -888,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. @@ -1204,6 +1294,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 +1808,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 +2029,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 +2127,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..6d195cec 100644 --- a/tests/unit/livekit.rs +++ b/tests/unit/livekit.rs @@ -5158,3 +5158,174 @@ 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, + state: &mut RuntimeState, + restarts: &mut usize, + event: &GeminiEvent, +) { + let mut activity = RuntimeActivity::new(Instant::now()); + activity.note_output(Instant::now()); + recovery.note_event(restarts, state, &activity, event); +} + +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); + +#[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); + assert!(!recovery.allow_resume(YOUNG)); + recovery.connected(false); + assert!(!recovery.allow_resume(YOUNG)); + 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); + assert!(!recovery.allow_resume(YOUNG)); + recovery.connected(false); + assert!(recovery.allow_resume(HEALTHY_GEMINI_SOCKET)); + recovery.connected(true); + 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); + complete_reply(&mut recovery, &mut RuntimeState::default(), &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); + 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); + state.recovery_reply_pending = true; + complete_reply(&mut recovery, &mut state, &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 state = RuntimeState::default(); + let mut restarts = GEMINI_RESTART_LIMIT; + recovery.connected(true); + assert!(!recovery.allow_resume(YOUNG)); + 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 state, &mut restarts); + assert_eq!(restarts, GEMINI_RESTART_LIMIT); + assert!(recovery.suppressed); + + complete_reply(&mut recovery, &mut state, &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 state = RuntimeState { + recovery_reply_pending: true, + ..RuntimeState::default() + }; + let mut restarts = GEMINI_RESTART_LIMIT; + recovery.note_event( + &mut restarts, + &mut state, + &RuntimeActivity::new(Instant::now()), + &ending, + ); + assert_eq!(restarts, GEMINI_RESTART_LIMIT); + + 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); +} 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)); +}