Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Concurrency Locking

The agent runs one session as an event loop. SessionLock (crates/core/src/lock.rs) keeps LLM requests, tool execution and tool-output summarization from interleaving, and queues user input that arrives while any of them is in flight. Every hold on the lock is bounded: one that goes too long without evidence of progress is reclaimed by a sweep the heartbeat runs.

What holds the lock

SessionLock tracks three kinds of hold:

  • llm: Option<Lease> — an LLM request is in flight.
  • summarizing: Option<Lease> — a tool output is at the summary model.
  • active_tools: HashMap<String, ActiveToolState> — tool executions keyed by their unique call_id, each with its own started_at and timeout.

is_locked() is true while any of them is set. at_tool_boundary() is the inverse — no request, no summarization, no registered tool — and is the only point at which queued input is replayed.

Leases and epochs

A Lease is a hold with a deadline: started_at, a timeout, and an epoch naming this particular hold. The budgets live in lock::lease: LLM (300s) covers one round trip with no intermediate signal, SUMMARIZING (180s) one cheap-model summarization of a single tool output.

The deadline means “time since the work last moved”, not “time since dispatch”. touch_llm() resets started_at on evidence the turn advanced — a stream poll whose transcript grew, a compaction reply that starts its second round trip — so a long but healthy streamed turn is never reaped, while a relay that answers every poll with an unchanged buffer does not push the deadline out.

try_acquire_llm(timeout) takes the LLM lease and returns its epoch, or None when a request is already out. acquire_or_current_llm(timeout) returns the live epoch when the lease is already held and takes it otherwise; compaction uses it because it is reached both from inside a turn and from a free lock. set_summarizing(timeout) takes the summarization lease and returns its epoch. force_release_llm() and clear_summarizing() release them.

Epochs come from one monotonic counter shared by both leases, so an epoch names exactly one hold for the life of the session. A dispatch stamps its epoch into the request’s Ctx (dispatch::stamp_epoch, key EPOCH_KEY) and the host echoes it back with the reply. Before any reply reaches a handler, AgentHarness reads it (dispatch::ctx_epoch) and asks epoch_is_current(epoch); a reply whose lease has been reaped or re-taken belongs to a turn that has moved on and is dropped. A reply carrying no epoch passes through — failing closed there would turn a host that drops the key into a starved agent.

Tool tracking

register_tool(call_id, tool_name, timeout) records an ActiveToolState; complete_tool(call_id) removes and returns it. Results are matched by call_id, never by tool name, so concurrent calls to the same tool are tracked independently. set_exec_kind(call_id, kind) attaches the execution-policy classification made at dispatch, so a continuation of the same call is held to the same policy. release_all_tools() drops every registration at once and hands the states back: an interrupt settles the whole set, and the caller still owes each one an abandoned-process record.

The sweep

sweep_expired() reclaims every hold whose deadline has passed — the LLM lease, the summarization lease and timed-out tools alike — and returns them as ExpiredLock values (Llm { waited, epoch }, Summarizing { waited, epoch }, Tool(ActiveToolState)). The leases come first so their handlers settle the turn before a tool result handled afterwards can start the next request.

AgentHarness::sweep_session_lock (crates/core/src/heartbeat.rs) runs it at the head of every timer fire, before any branch can return, and settles each entry:

  • an expired LLM lease drops any stranded relay stream, resets an active status to WaitingForInput and drains queued input; a late reply is then dropped by the epoch guard;
  • an expired summarization routes a synthetic 504 through handle_summarization_response, which truncates the output it still holds, appends the tool result the model is waiting on and continues the turn;
  • a timed-out tool gets a synthetic tool result saying the harness stopped waiting, naming the PID that was not cancelled — Host::exec is fire-and-forget, so a timeout can only stop waiting — and the process is tracked as abandoned.

hold_summary() renders what holds the lock and how long each hold has left; submit logs it whenever input is queued, so a quiet agent says what it is waiting on.

Input queueing and coalescing

submit (crates/core/src/intent.rs) calls queue_user_input(content) when input_can_run_now() is false. boundary_reached() returns None unless at_tool_boundary() holds, and otherwise drains the queue through drain_ready_inputs():

  • a single queued input is returned verbatim;
  • several are coalesced into one message under a [Multiple inputs received while processing:] header, joined by ---.

execute_next_tool_call calls boundary_reached() after a tool result is recorded, gated on the same input_can_run_now() predicate submit queues on, so a rate-limit backoff — which holds no lease — cannot let queued input through. The heartbeat’s tick also drains the queue (check_queued_inputs, behind the same predicate), so input stranded by a lock-free window is picked up without another turn event.

Tests

lock.rs covers LLM mutual exclusion, boundary detection, coalescing, tool timeout removal, the combined locked-state predicate, the sweep reclaiming all three kinds of hold in lease-first order, touch_llm keeping a live stream from being reaped, and an epoch ceasing to be current once its lease is gone.