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
25 changes: 17 additions & 8 deletions .cargo/mutants.toml
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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",
Expand Down
76 changes: 62 additions & 14 deletions docs/provider-cost-and-degradation.md
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -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

Expand All @@ -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,
Expand Down
5 changes: 5 additions & 0 deletions scripts/gen-wire-fixtures.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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() },
Expand Down Expand Up @@ -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),
Expand Down
53 changes: 41 additions & 12 deletions src/dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
///
Expand All @@ -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<Mutex<HashSet<String>>>,

// True while assessing; false while finalizing the report. Both hold
// capacity.
pub live: Arc<Mutex<HashMap<String, bool>>>,
/// 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,
};
Expand All @@ -48,24 +55,26 @@ 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(())
}
}

enum Reservation {
Existing,
New(Slot),
Full,
Finalizing,
}

impl LocalDispatcher {
Expand All @@ -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(),
Expand All @@ -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<Mutex<HashSet<String>>>,
pub live: Arc<Mutex<HashMap<String, bool>>>,
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
Expand Down
Loading