Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 16 additions & 3 deletions docs/provider-cost-and-degradation.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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
Expand Down
5 changes: 5 additions & 0 deletions src/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<FrameworkEvidence>,
/// Deterministic, bounded facts derived from the live session. The ledger
/// deliberately holds no editor text or runner diagnostics: those remain
Expand Down Expand Up @@ -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(),
Expand Down
8 changes: 7 additions & 1 deletion src/agent/turn_taking.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down
127 changes: 111 additions & 16 deletions src/livekit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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) {
Expand All @@ -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
Expand All @@ -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);
Expand Down Expand Up @@ -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);

Expand All @@ -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
);
Expand Down Expand Up @@ -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"
);
}

Expand All @@ -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,
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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!(
Expand Down Expand Up @@ -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(),
};
Expand Down Expand Up @@ -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);
Expand Down
7 changes: 4 additions & 3 deletions src/livekit/media.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion src/livekit/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}

Expand Down
19 changes: 19 additions & 0 deletions src/livekit/turn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Instant>,
/// 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:
Expand Down Expand Up @@ -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)]
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down
Loading