Files
hermes-agent/agent/turn_context.py
Siddharth Balyan afc3b7c6f3 feat(connectors): one backend-owned connection operation, with a setup card on Desktop, TUI and CLI (#111008)
* feat(connectors): the desktop connects apps through one backend-owned operation

Re-based onto main after #109517, #110368, #110574 and #110843 landed as squash
merges (e0ef0eb9c3, d105376b21, ee2f5629b8, 1ab32b212b): the branch's history
no longer shared a base with main, so this is the PR's exact delta against
+3762/-1470, identical to the branch tip 05d7f2d3d5.

The fifteen commits it carried, in order:

--- feat(connectors): the wire follows the connector contract (six-state status, required toolkit metadata, connectionId on mint)

Hermes types exactly what the contract page writes: /tmp/magic/CONTRACT-TOOLKIT-METADATA.resolved.md
(carriers: portal PR2 `sid/connection-api` for the enum and account rows, the portal contract branch
on top of it for the list metadata). No optional-for-compat fields, no fallback branch, no seven-state
word left in the tree. A gateway that does not speak this contract fails validation loudly.

Wire (tools/connectors/gateway/wire.py): `ConnectionStatus` is the six values, `pending` covering the
vendor's INITIALIZING and INITIATED; `ConnectorListItem` requires title, description, https iconUrl and
authKind, and carries activeConnectionId when the session binds an account; `ConnectorListResponse`
types a page whole with its total; `ConnectorAccount` and `ConnectorAccountsResponse` type the account
routes; a mint result names its connectionId on `initiated` and never on `failed`; CONNECTION_REQUIRED
carries connectUrl and connectionId together or not at all; `account` exists on the execute call and is
never sent while multi-account is off.

Client: `list_connectors(search, connected)` refuses a search under three characters before any
request and parses every page whole; `list_accounts` and `account_status` read the account routes,
404 `connection_not_found` is None and 429 raises `RateLimited(retry_after)`.

Watcher: the target keeps the account id the mint named and the toolkit's title and icon from one
list read before the card is emitted; `_status_for` is the one seam between the watcher and the
gateway's status source (the list walk today; the per-account route replaces that body only).
`pending` and a missing status move nothing; revoked and inactive are failed.

Desktop: `ConnectorRow` is the strict list item; `ConnectorRowSeed` is what a tool call's args can
say; `connection.request`/`update` targets carry `connection_id`, `title`, `icon_url`; the mark
ladder is glyph, vendor icon as a plain image, favicon, monogram.

Tests run red first: wire rejects the retired words and the missing metadata; client search
minimum, execute omits account, account_status 404/429; managed targets carry the metadata and the
account id, a pending row moves nothing; the logo ladder order; the store passes the fields through.

--- feat(connectors): onboarding connect-first rides the connection operation

The guided first build had its own connector surface: a 2 s connectors.list poll loop, an auto-open
of every minted link, a hidden "[setup] links opened" note that painted as a user bubble, CONNECT
FIRST rows in the transcript, and a "Start with N apps connected" composer pill that sent the go
signal early. Delete all of it. The build session's first manage_connections connect now shows the
same card as any chat, and the settled tool result is the go signal.

Deleted: store/first-build-connectors.ts, assistant-ui/first-build-connectors.tsx (nothing rendered
it any more), lib/first-build-start.ts, onboarding-chat/start.tsx, the first-build session markers in
handoff-receipt.ts, the connectionRows / latestConnectorPart helpers they were the last readers of,
and the five strings only they used.

The runbook (setup-profile.ts::connectFirstRunbook) now describes the operation's real contract: one
connect call with every picked slug, one card with a row per app, blocks until settled, a per-app
result of connected / skipped / not_connected. No wait, no "Start with", no "already active" branch
(the result never carried it), and no offer to re-mint: Try again and Continue live on the card, and a
Continue means the user moved on.

Tests (red first): the runbook names one connect call and its per-app result and never the retired
model actions; the no-account rule survives an empty pick.

--- feat(connectors): the registry is keyed by profile, every watch read is bounded by the deadline, connectionId is optional, connect and execute carry returnTo and op

Registry (P1-14): tools/connectors/live.py keys an operation by (profile home, session key). Two multiplexed
profiles can carry the same timestamp-based session key; the tool thread opens under its turn's profile
override and the RPC side names the session record's profile home, the default profile resolving to the
process home on both sides.

Watch reads (P1-8 residual): list_connectors takes a per-page timeout and the watcher bounds each read by
the operation's remaining deadline (floor one second), so a stalled gateway page cannot hold the operation
past its deadline.

Contract follow-up from the portal handoff (/tmp/magic/HANDOFF-HERMES-PR3-PORTAL-CONTRACT.md): connectionId
on a mint result and on CONNECTION_REQUIRED is plain Optional (a no-auth toolkit answers active with none; a
failed mint answers with neither link nor id); a target without one has nothing to watch. Connect and
execute requests carry returnTo (hermes-desktop | hermes-desktop-dev | portal) and op, so the vendor's done
page can send the browser back to the app; the call sites and the deep-link handler follow in the next
commit.

Tests, red first: two profiles share a session key without seeing each other; each watch read shrinks with
the deadline; the client honours a per-page timeout; the optional id; the two request fields. The fakes'
list_connectors accept the timeout keyword; two wrapper spies that did not were the cause of a 300 s hang
under the per-file harness.

--- feat(connectors): an MCP target runs on the shared connection operation

manage_connections ran MCP targets as a renderer errand: the desktop card ran
the catalog lookup, the install and the OAuth flow, then told the backend what
had happened, and every surface without a card got `unavailable` plus two
terminal commands. That left the outcome in the renderer's word, and left the
model unable to connect an MCP server anywhere but the desktop.

The backend now owns the work, the way it already owns a managed connector.
`prepare` starts the OAuth flow in process and points the browser at the
backend's own callback route; an install that still needs credentials waits
pending and publishes their names as `required_env`; `observe` reads the flow
or the install worker on every tick. The worker itself moved out of
tui_gateway/mcp_oauth_sessions.py into tools/connectors/mcp_oauth.py, and that
module now calls it, so the Capabilities tab keeps its RPC session table.

The card may only say approved, skipped or continue: a claim of any other
state moves nothing. With that, `renderer_flow` and the `unavailable` state
have no producer left and are gone from the contract. Off the desktop there is
no card, so the action runs at once and the result carries the authorization
URL for the user to open.

Try again reaches the same work: `_reissue` dispatches on `target.kind`, so an
MCP row re-runs its own install, enable or OAuth instead of minting a managed
gateway link.

(cherry picked from commit fa0438c4121807604e7c983ba42b901314a1305b)

--- feat(connectors): the MCP card projects the operation instead of running the flow

The MCP setup card used to own the work: it called the catalog, ran
installMcpCatalogEntry, polled the action, drove the OAuth window, and then told
the backend what had happened through connection.respond. That made the renderer
a second authority on target state, so an install that finished after the window
closed, or a card that never mounted, left the operation with a state nobody
could correct. PR3 moves that work to the backend, so the card has one job left:
show the operation and send the user's consent.

McpSetupPending now renders request.targets the way ConnectorOffer renders them:
one row per target, one verb by (action, state), Continue below the shell.
Install and Enable send {status: 'approved'} and nothing else; an authorize row
opens the link the backend minted, with no OAuth RPC of its own; a running
install holds its verb with the busy mark; a failed row sends
connectors.connect {reconnect: true} on the open operation, the same call the
managed card makes. The owner lookup and that call are now shared helpers on
connector-tool.tsx rather than a second copy here.

ConnectionTargetOutcome loses connected / initiated / failed: those were the
renderer reporting state, which it no longer may do. ConnectionActor loses
renderer_flow for the same reason. ConnectionOperationTarget gains required_env,
the credentials an MCP install is still waiting for; the card renders a field per
entry under the pending row and holds Install until every required one has text,
so the values travel with the approval instead of through a separate catalog call.

(cherry picked from commit 2e42ad19f63e41f274a87199d07eb99a7c995cbe)

--- fix(connectors): the toolkit metadata leaves the wire; the vendor logo is derived from the slug

Sid cancelled portal PR3 (the toolkit metadata on the list item) on 2026-09-14, so the fields the contract
commit typed as required come off the wire: no title, description, iconUrl, authKind or activeConnectionId
on the list item, no total on the page, no search or connected query, no title or icon_url on a target
snapshot. The list item is exactly what portal #1220 emits.

The desktop derives the vendor logo from the toolkit slug instead (`connectorIconUrl`,
https://logos.composio.dev/api/{slug}; the gateway slug is the vendor slug, checked for every lead-order
pick), rendered as a plain image because the host sends no CORS header. The title stays
`connectorTitle(slug)`. The mark ladder is unchanged: glyph, derived vendor icon, favicon, monogram.

Everything else in the contract stands: six-state status, optional connectionId on mint results and on
CONNECTION_REQUIRED, returnTo and op, the account routes. The contract commit's message still names the
metadata; this commit is the correction.

--- fix(connectors): an MCP target survives the card's Continue, and the model is not sent to an action that refuses it

The review of the MCP backend found eight defects; each one is a test first.

- The off-desktop note told the model to confirm with action 'status', which refuses MCP
  names outright. It now says the authorization finishes in the background and the tools
  arrive on the next turn.
- The card's existence is a property of the session. A missing callback no longer routes a
  desktop session down the off-desktop path, where the model would be handed a live link.
- A second failure of a backend attempt was published with actor 'user'. ``refresh`` takes
  the actor, so the frame says who produced the text.
- Continue on the RPC thread can settle the operation between any read of a target's state
  and the transition that follows it. A settled operation has a frozen result, so the lost
  move is dropped; an IllegalTransition no longer escapes into the tool result. The same
  rule covers a worker whose outcome arrives late and a Try again that arrives after Continue.
- Try again on an install carried no credentials, so the second install ran with an empty
  env. The approved map is kept on the runner (never on the target: target fields are
  serialised to the model) and reused.
- An approval that does not cover a required credential no longer reaches the worker, where
  ``install_entry``'s prompt would block on stdin forever; the row waits for the card's
  fields instead.
- ``enable`` writes ``mcp_servers`` under the scope and lock the dashboard's toggle route
  uses, so the two read-modify-write paths in one process cannot drop each other's write.
- Several authorize targets start their flows together and share one wait; one wait per
  target kept the card empty for minutes.
- The off-desktop operation is in no session's registry, so it now emits no connection.update.

Try again is also refused for a target the transition table cannot move back to 'initiated'
(an MCP target has no move out of 'expired') and for an operation that settled during the call.

The contract test that froze the Actor enum is replaced by two behaviour tests: no actor but
the backend watcher can connect an MCP target, and the card's word never moves one.

(cherry picked from commit 081c16412b48667758a77a7e7c29eaf2be6a93dd)

--- fix(connectors): the watcher reads one account route per target, and every frame carries its seq

The watch loop walked the whole toolkit list once per tick to learn whether one
target had connected. That read costs a vendor call per page, cannot tell one
account from another, and forced `awaiting_new_attempt`: after a forced
reconnect the list still reported the OLD account `active`, so the row had to be
disbelieved until it read as something else once. The gateway now serves
`GET /v1/connectors/accounts/{connectionId}`, so each pending target reads its
own account: the id the mint named, one read per target per tick at 1 Hz (the
route's bucket is 180/min per principal), each bounded by the operation's
remaining deadline. A 429 parks that one target until its Retry-After passes and
leaves the others reading. A 404 is "not yet" until the deadline. A target the
mint gave no account for has nothing to read, so it is not read. A forced
reconnect watches the new account, which is why `awaiting_new_attempt` and its
two tests are gone.

`mint` and `run_remote` now name where the browser should come back to
(`returnTo`, plus the operation id on a mint), so the vendor's done page can
hand the user back to the desktop app that asked instead of stranding them on a
web page. Only the desktop registers that URL scheme, so no other surface sends
either field.

`connection.update` frames were built by re-reading the operation after the lock
was released, so a second writer could give an older frame a newer state and the
renderer could not tell which frame was last. Every write now advances a
monotonic `seq` and takes its snapshot under the same lock, and the emitter sends
that snapshot; a renderer that keeps the highest seq per operation can drop a
frame that arrives out of order. `connectors.operation.wake` lets the desktop's
deep-link return ask for a read now instead of at the next tick; it only shortens
the wait and trusts nothing else in the link.

(cherry picked from commit f5423cf72302b92b26c043495184ad6426fa2b36)

--- fix(connectors): the desktop card follows the operation, and says so out loud

The MCP card was a second implementation of the connector card with the
review's defects: it painted for any request on the session, kept its
controls after the operation settled, offered an approve verb while the
backend was still minting an authorize link, and re-enabled Install when
the RPC returned rather than when the state frame moved the row. A second
click in that window sent the consent twice.

The renderer now reads the operation's `seq`: an update or a status frame
whose sequence is not greater than the one already applied changes
nothing, and a resume snapshot neither revives a settled card nor puts
back a row a newer frame has moved. Without it the transport's ordering
decided what the user saw.

`hermes://connections/done?op=…` brings the user back from the browser to
the session that opened the operation and wakes its watcher, so the row
moves at once instead of at the watcher's next tick. Only the op id is
used; the link's status moves no row.

Each row's mark and cue are one polite live region, so a row that flips is
heard and not only seen, focus follows the row the backend moved while the
card holds it, and the waiting mark stops spinning under reduced motion.

(cherry picked from commit 7c30d252556d7d496ddb354b50bcecc41bceef27)

--- fix(connectors): the renderer types seq as the wire carries it

The watcher commit made `seq` a required field on every operation frame. The
desktop store still declared its own optional `seq` so it would compile against a
backend that predates the field; that backend no longer exists on this branch, so
the hedge is dead code and the fixtures were short one field. The store now reads
`seq` from the shared types, and the fixtures count the way the backend does.

The repeat-frame test asserted the whole request keeps its reference. With a real
rising `seq` the request must change; the invariant the test guards is that the
target row keeps its identity so open credential inputs do not remount.

--- fix(connectors): a URL that arrives after the shared wait still lands on its row

The shared prepare wait failed every row still pending when it ran out, while
that row's own thread was still waiting on the provider. When the URL arrived a
moment later the thread's move raised inside the daemon thread, the link was
lost, and Try again started a second flow. The wait now bounds only how long
prepare blocks; a row still pending afterwards is left to its own thread, which
is the only writer of that row and ends with the URL or the flow's own failure.

The comment on the approved credentials said every target field is serialised to
the model; it is not (the snapshot names its keys). The reason they live on the
runner is that they are secrets and the runner's life is exactly the operation's.

--- fix(connectors): the review findings the operation must survive before the fold

The MCP prepare threads and the install worker started with an empty context,
so a named-profile turn's home override never reached them: the flow resolved
the process home's `mcp_servers` and stored the token there. Each thread now
runs in a copy of the calling thread's context. `connection.respond` had the
same gap on the RPC thread: an approval ran the enable, which writes
config.yaml, with no profile bound, so the flag landed in the launch home. The
handler binds the session record's profile the way `_connector_rpc` does;
`config_write_scope(None)` keeps that override, so the enable needs no change.

The watcher raised out of the tool when a per-target Skip landed while that
target's read was in flight: the operation stayed open with `live` closed, and
every later answer got 4004. A read for a row that is no longer live is dropped
at debug; only a refusal on a live row is still a contract violation. The same
skip from the card raced a row the backend had just connected and aborted the
rest of the answer; a skip for a resolved row is ignored and every entry, then
the settle check, still runs.

A read was bounded by the whole remaining deadline, so a hung gateway held the
first read for 300 s and Continue could not return the tool; a read now waits
ten seconds at most. A 429 parked only the target that read, but the budget is
the principal's, so the next target's read in the same tick spent it again:
every live target waits out the one Retry-After.

The MCP surface rule read the platform alone, so a desktop call without the
callback (registry dispatch from execute_code) opened an operation nobody
rendered and blocked for the deadline; it now uses the managed rule, surface
and callback. A write after settlement advanced `seq` while emitting the
frozen frame, so the resume snapshot named a seq no frame carried; the counter
stops at the settle frame. A mint that reports `initiated` with no account id
logs that the watcher cannot read the row.

(cherry picked from commit 587b228016bf8ee2b022e7fa58c2d9f6057809e9)

--- fix(connectors): the desktop card holds a verb until the backend answers, and never takes the keyboard from a credential field

The review of the desktop card found five defects; each one is a test first.

- A resume snapshot was refused whenever the cache held a settled operation, whichever
  operation it was, so a session that opened a second operation after settling the first
  never got its card back from a resume. And the refused snapshot handed the caller the
  settled cache as "the request", so the session was flagged as needing input behind a
  summary with no controls. Only the same operation can refuse the snapshot now, and a
  refused one is no pending card.
- The focus handoff picked the row's first button, which after pending -> initiated is the
  disabled working verb; the focus call was a no-op and the keyboard landed on the document
  body. It picks the first control that can take focus, else the row. It also moved focus out
  of a credential field the user was typing in whenever another row moved; it leaves an
  editable alone. The "focus Continue once every row resolved" branch was dead (Continue
  unmounts the moment nothing is unresolved), so it and its ref plumbing are gone.
- The done link navigated to a settled operation's session and rejected when the wake RPC
  did (4004 once the operation left the live registry). A settled request is ignored, and a
  refused wake is nothing: the wake only shortens the wait, the watcher still ticks.
- Install spun forever when the store refused to send the consent (the operation gone or
  settled under the card): `respondToConnectionRequest` resolves false in that case and the
  verb was only released in `catch`.
- After a partial approval (a required credential missing) the backend answers with a
  same-state frame whose detail names what is missing; nothing released the verb because it
  was held until the row's state moved. The row now remembers the seq the click saw and holds
  the verb only until a frame past it arrives, which is the backend's word on the click
  whether or not the row moved.

(cherry picked from commit 4c534d6e2af979778d9a2423cf97d717edeaa1a9)

--- fix(connectors): a skip that loses the race to the watcher is ignored, not raised

The skip guard read the row's state and then moved it; the watcher can connect the
row between the two, and the refused move aborted the rest of the card's answer.
The refusal itself is now the witness: a move refused for a row that is resolved,
or on a settled operation, is the same nothing-to-do as a row resolved earlier.

Two recording fakes in the managed tests kept their lists on the class; they now
start per instance so a lifted fake cannot share reads between tests.

--- fix(connectors): a resume that lands behind a newer live frame still reports the pending card

The refused-snapshot branch answered "no pending card" for both reasons it can
refuse: the operation settled, or a newer frame already moved a row. Only the
first is no card. For the second the live card is still open and blocking the
turn, so the caller must keep the session flagged as waiting on it.

* test(connectors): defer new connection coverage until implementation settles

Remove PR-added test cases and their unused helpers while retaining
existing tests adapted to the changed connection contract. The three
PR-only renderer test files are removed for now.

Focused behavioral coverage will be added as the final implementation
step before verification. Existing main coverage is not being removed
wholesale, and this does not declare the feature merge-ready.

* fix(connectors): commit MCP authorization at initialize, save setup values after success

The OAuth probe treated one exception as one outcome: any failure after the browser
step restored the token snapshot and manager entry, so a server that accepted the
token but failed tools/list discarded a completed consent. Now the probe reports
whether initialize succeeded (details["initialized"], read from the claimed
MCPServerTask). Failure before that point rolls back as before. Failure after it
saves the server config, keeps the tokens, and reports tools unavailable through
flow.discovery_error; the card can retry discovery without repeating consent.

Catalog install wrote the submitted values to .env before install_entry and the
probe ran. The values now live in the secret scope for the duration of the install
(get_env_value reads through get_secret, so install_entry finds them without a
prompt), and .env is written only after the probe returns tools. A failed probe
removes the server block and writes nothing. _probe_tool_names returns None on a
failed probe instead of an empty list, so failure and a valid empty listing are
distinct.

Failure text is redacted before it reaches target detail: every value the card
submitted for that target is replaced by exact match, then the pattern redactor
runs.

required_env now carries the manifest's secret and default flags; Target carries the
manifest's post_install text as instructions. The wire contract gains secret,
default, instructions and discovery_error; generated TS and OpenRPC regenerated.

* fix(connectors): one OAuth callback receiver picker for the connection card

The card's authorize target built its redirect from the dashboard web server and
raised when none was bound in the process, so a standalone hermes --tui session
could never authorize an MCP server. The receiver is now chosen in one place
(choose_callback_receiver): a pinned pre-registered client keeps the SDK's own
listener on the registered port; a client-advertised loopback URI is used as-is and
its callback arrives through the mcp.servers.oauth.callback relay; otherwise the
backend binds a one-shot loopback listener and feeds it into the flow. The dashboard
route stays with the dashboard web page, which cannot bind a port.

tui_gateway/mcp_oauth_sessions.py had a second copy of the loopback listener and a
_worker that referenced _probe_with_rollback, set_hermes_home_override,
reset_hermes_home_override and Path without importing them, so every RPC-started
flow raised NameError. Both are deleted; start_flow spawns run_worker directly and
uses the same receiver picker. Flow registration is shared (register_flow /
finish_flow) so a card-started flow with a client URI is reachable by the relay.

Under an SSH session with no client listener the attempt's detail carries the
existing paste-the-redirect instructions. Nothing on the card path opens a browser.

* feat(desktop): setup-form modal for MCP connection cards

The connection card rendered an MCP server's setup fields inline: every field as a
password input, no default value, no instructions. A URL such as the n8n MCP
server URL was typed blind, and the manifest's setup text never reached the user.

Two components carry the form now. SetupFieldList renders the ordered fields the
backend declares (a plain field as text prefilled with the manifest default, a
secret field masked and empty). SetupFormDialog composes it with a one-line title,
the manifest instructions, an inline error for a failed attempt, and Cancel /
Connect. The row's Install action opens the dialog when the target has fields;
Connect sends {status: approved, env}; Cancel sends {status: skipped}. A failed
attempt keeps the dialog open with the draft intact. The draft lives in the dialog
component only; nothing reaches the store or the resume snapshot.

Once the backend publishes the authorization URL the dialog shows it as text with
an Open in browser button. Nothing opens a browser on a state change: the Try
again path on both cards used to open the re-minted link at once; it now waits for
the row's update frame and the user's click.

The store types gain the wire's secret, default, instructions and discoveryError
fields and normalise them; a connected target with discoveryError renders as
authorized with tools unavailable.

* feat(tui): connection card in the Ink TUI

The Ink TUI had no client for the connection operation: connection.request,
connection.update, connection.respond and the pending_connection resume snapshot
were unhandled, so a manage_connections call in a TUI session could only print a
link through the model.

connectionOperationStore.ts holds the backend snapshot: a request opens only for a
new operation id, an update applies only to the live operation with a higher seq,
a settled update freezes the id so a late request frame cannot reopen the card.
The gateway event handler feeds it; session resume hydrates it from
pending_connection.

connectionSetupOverlay.tsx renders one callout in the prompt zone: a one-line
title, the manifest instructions, every field (plain rows prefilled with the
default, secret rows masked), then a Connect / Cancel selector. Connect sends the
draft through connection.respond; Cancel skips the active target. When the backend
publishes the authorization URL the callout shows it as text under "Press Enter to
open in browser"; Enter is the only thing that opens it. A failed attempt unlocks
the fields with the draft kept. An accepted secret renders as "Set" and is never
echoed. A target authorized without tools shows that state and Continue.

The overlay joins the existing input-owner set in overlayStore so typing and other
prompts are blocked while it is up.

* feat(cli): connection panel in the classic CLI

The classic hermes CLI passed no connection_callback, so a manage_connections
call could only print an authorization link through the model and could not take
a setup value at all.

The CLI now renders the connection operation as a prompt_toolkit panel, on the
same queue mechanism as the clarify panel: the agent thread's callback opens the
panel, blocks until the first decision, then returns so the operation's watch loop
runs; every later state reaches the panel through ConnectionOperation.on_change
(installed only when no gateway hook is set, restored on close).

Panel: one-line title, the manifest instructions, one row per field (plain rows
prefilled with the default, secret rows masked in the input buffer and rendered as
"Set" once accepted), then Connect / Cancel. Connect sends the draft through
apply_answer; Esc skips the active target; Ctrl-C sets the tool-thread interrupt so
the operation settles as interrupted. Once the backend publishes the authorization
URL the panel shows it with the target detail and "Press Enter to open in browser";
Enter is the only thing that opens it. A failed attempt unlocks the fields with the
draft kept; an authorized target without tools offers Retry discovery / Continue.

Up/Down move between rows, Left/Right toggle the action, and the panel joins the
blocking-overlay guards so chat input, history and voice stay out while it is up.
The single-query (headless) mode passes no callback, as it does for clarify.

* feat(connectors): register a connected MCP server and report its tools in the result

After a successful install or authorization the target carried the probe's tool
names and the model was told the tools "become available on your next turn". MCP
tool schemas are deferrable by construction, so nothing about them lives in the
sent tool array; a server registered in the scoped registry is callable through
tool_describe/tool_call in the same turn. The operation now registers the server
(register_mcp_servers under the owner's home scope) once authorization is
committed, records the registered names on the target, and the settled result
carries a tools_listing block in the deferred-catalog format plus a note that the
tools are callable now. A registration failure keeps the target connected with
tools: [] and a sanitized discovery_error. agent.tools and the system prompt are
untouched; the between-turns refresh updates the catalog block as before.

The card gate no longer asks for the desktop platform. Every surface that renders
the card attaches a connection callback (Desktop, the Ink TUI, the classic CLI);
registry dispatch and messaging sessions attach none and keep the link result.

Pre-commit rollback in probe_with_rollback used restore(only_if_absent=True),
which skips the rollback when a token file exists. On a first-ever authorization
the only file is the one this attempt wrote, so a token the resource rejected was
kept. Live E2E (controlled provider answering 403 to the issued token) showed the
row fail and the token survive; the rollback now restores the snapshot outright,
and the same run shows the token file removed.

* fix(connectors): install an OAuth catalog entry through the card's own flow

The card's install ran install_entry and then a plain probe. For an OAuth entry
that probe has no card flow around it: in the desktop backend it failed at once
("non-interactive environment and no cached tokens"), and in the classic CLI it
saw a TTY, opened a browser by itself and drew install_entry's curses tool
checklist over the panel. Since a probe failure is now an error, the failed
install was rolled back with _remove_mcp_server, and authorize refuses a server
that is not configured. A clean home had no path to a connected OAuth entry, which
is 55 of the 65 catalog entries. Found by the live three-surface run.

An install now builds the entry's configuration in memory (card_install_config:
no prompts, no probe, no checklist; a prior tool selection or the manifest's
curated filter applies). An OAuth entry starts the same flow authorize uses with
that configuration and the setup values in the attempt's secret scope, so the row
reaches the URL step, and the configuration and the setup values are saved
together when initialize accepts the token. Every other entry is probed in memory
and saved after the server answers. A failure writes nothing, so a failed
reinstall keeps the previous configuration, and the failed row asks for its fields
again so the card can reopen the form over the draft it kept.

* fix(agent): make a server connected in this turn callable in this turn

tool_describe and tool_call resolve names inside the agent's toolset selection,
which is fixed when the agent is built. A server that manage_connections had just
registered was therefore "not found" for the rest of the turn whose result calls
its tools available, and stayed out of the next turn's catalog too. The live
Desktop run showed it: the result listed mcp__fx_oauth_fields__echo and the
tool_describe that followed answered not_found.

The executor now adds the MCP servers the call connected to the selection. Only
the selection changes; agent.tools does not, so the sent tool schema bytes stay
the same. A selection of None (every toolset) and the no_mcp sentinel are left
alone.

* fix(cli): reopen the connection form when a required field is still empty

When the backend refuses an answer because a required field is missing it keeps
the row pending and names the fields. The panel mapped that frame to its waiting
phase, which draws the detail line and nothing else, so the user saw "waiting for
FX_API_KEY" with no fields and no buttons until the 300 s deadline. The panel now
returns to the form on the first missing field, over the draft it kept.

* test(connectors): the no-card path is the one with no callback attached

The card gate is "a connection callback is attached" since the classic CLI got
its own panel. This test still attached one under platform "cli" and expected no
card, so it opened a real operation, waited out the 300 s deadline and failed.
It now drives the path that has no card: no callback.

* fix(connectors): order a cancel against the commit and stop deleting a working grant

Found by the live runs on Desktop, the Ink TUI and the classic CLI.

A user's skip did not stop the OAuth attempt. With the worker parked in the token
request, the row settled "skipped" and the token was written 33 s later. A skip
now cancels the attempt, and one lock orders that cancel against the commit: the
attempt is either canceled with the earlier tokens restored, or committed and
kept. When the commit won, the skipped row says the authorization was kept. An
interrupted turn cancels its attempts the same way.

Every attempt began by deleting the saved tokens, so retrying discovery for an
authorized server demanded consent again, and a cancel in between left no grant
at all. The card's flow now connects with the saved tokens first, with no browser
step; only when they do not work does it replace them. The RPC session surface
keeps the old behavior because its caller waits for an authorization URL.

A second Connect from the form a failed row reopened was dropped without a frame,
which left the Desktop dialog with Connect loading and Cancel disabled until the
deadline. An approval on a failed or expired row is now a retry with the new
values.

The reported tool names came from the registration call, which returns nothing
for a server the process already holds. A retry after a failed listing therefore
said "no tools" while the server had them. The names are read from the registry,
and a parked server is woken first.

An authorization that commits after its card closed by deadline was never
registered, so the next turn still could not reach it. The runner keeps such
attempts, and the between-turns refresh adopts the ones that were approved.

An error with no message reached the user as a class name ("CancelledError").
The tool description still said a server's tools arrive on the next turn.

* fix(desktop): let an OAuth install open its link, and keep the form's draft

An install of an OAuth catalog entry with no setup fields reached the URL step
and gave the user nothing to click: an initiated install was always drawn as a
disabled spinner, and the only other place the link is shown is the setup dialog,
which opens for entries with fields. That is most of the catalog. An initiated row
that carries a link now offers Open, whatever the action.

The setup dialog reset its draft whenever the field list changed identity. The
backend sends a fresh list with every frame and an empty one while an attempt
runs, so a failed Connect erased what the user had typed. Fields now only fill in
what the draft lacks, and closing the dialog drops the draft.

The settled summary dropped the "tools unavailable" fact, and the live row offered
a Try again for that state which the gateway refuses for a connected target. The
summary keeps the fact and the row offers no dead control.

* fix(tui): keep the typed draft when a Connect fails

The overlay reset its draft whenever required_env changed identity, and every
backend snapshot delivers a freshly parsed array. A failed Connect therefore came
back as a form with the default region and an empty secret. The draft now resets
per target only, and a snapshot's fields fill in what the draft lacks.

* fix(cli): show the URL step and start an install that has no fields

The panel applied the user's answer and then set its phase to "waiting". The
backend applies the answer on the same thread and its change hook had already set
the next phase, so the URL step of an OAuth install and the reopened form for a
missing field were both overwritten, and the card sat on "Waiting…" until the
deadline. The waiting phase is now set before the answer is applied.

A pending install or enable with no fields opened in the waiting phase, which
sends no approval, so the flow never started. It now opens on Connect/Cancel. A
target that already carries its link opens on the URL step.

Connect on a failed row re-ran the attempt without the values now in the draft; it
sends them. The authorized-without-tools phase offered a "Retry discovery" that a
connected target cannot run inside the same operation; it offers Continue.

* fix(connectors): a newer OAuth attempt replaces the older one for the same server

A card that closed by deadline leaves its worker waiting on the browser for up to
300 s, and a tampered callback leaves one waiting too. A retry or a new operation
for the same server then ran beside it: both wrote the same token files, and the
older one's rollback could write over the newer attempt's grant.

The newest attempt per home and server is recorded. Starting one cancels the older
attempt, takes over its pre-attempt snapshot so a later failure still restores the
state from before either, and the older attempt's rollback leaves the files alone.

* fix(desktop): label an install's link control "Open in browser"

The control that hands an install's authorization link to the browser reused the
action's verb, so the row showed "Install" before the click and "Install" again at
the link step, told apart only by the cue. It now reads "Open in browser", the
label the setup dialog already uses for the same act. Authorize keeps its verb.

* docs(mcp): describe the setup card on the desktop, the terminal UI and the CLI

The MCP guide said the chat install exists only in the desktop app and that the
CLI relays commands. All three surfaces now show the same card: fields, Connect or
Cancel, an authorization link the user opens, one save when the server has
accepted the token, and tools the agent can call in the same turn. The tools
reference gains the result fields (tools, tools_listing, discovery_error) and the
no-card behavior of an OAuth install.

* fix(cli): keep the layout hook callable without a connection widget

The CLI panel commit added connection_widget to _build_tui_layout_children as a
required keyword. That method is the documented override point for wrappers, and
five existing tests (extension hooks, prompt stash, subagent dock) call it without
the new argument, so CI failed with a TypeError. The argument is now optional, like
the other widgets added after the hook was published; a missing widget is left out
of the layout.

The settled tool result also carried the target's setup instructions. Those are
the card's text for the user, and a catalog entry's notes can predate this flow
("restart your session so the tools are loaded"), which contradicts a result that
says the tools are callable now. The model-facing result drops them; the card
payloads keep them. This restores test_connector_local_batches.

* fix(connectors): wait for an in-progress registration before reporting a server's tools

Found with the real Vercel MCP server. Saving the configuration wakes the config
watcher, which starts its own connect for the new server. The operation's
registration call then skips the server as "already connecting" and returns at
once, so the card settled "connected" with no tools, and the 214 tools were
registered three seconds later. With no names in the result the model searched,
found the hosted connector of the same vendor and asked the user to connect that
instead.

The registered names are now read once the registration has finished: while
another task is connecting the server the read waits (30 s at most), a parked
server is woken once, and a server that finished registering with no tools is
still a valid empty list.
2026-09-21 10:04:32 +05:30

1247 lines
62 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Per-turn setup for ``run_conversation`` (the turn prologue).
``build_turn_context`` runs the once-per-turn setup (stdio guard, sanitization, prompt
restore-or-build, session row, idle/preflight compaction via ``turn_context_compaction``,
pre_llm_call hook, prefetch, persistence), mutating ``agent`` as the loop expects, and
returns a ``TurnContext`` with only the locals the loop reads back.
``build_api_messages`` builds the wire copy for one API call."""
from __future__ import annotations
import logging
import sys
import threading
import time
import uuid
from contextlib import suppress
from dataclasses import dataclass
from typing import Any, Dict, List, Mapping, Optional, Tuple
from agent.conversation_compression import recover_rotated_compression_session
from agent.iteration_budget import IterationBudget
from agent.memory_manager import build_memory_context_block
from agent.memory_provider import is_trivial_prompt
from agent.message_metadata import append_message, stamp_message_timestamp
from agent.model_metadata import estimate_messages_tokens_rough, estimate_request_tokens_rough
from agent.image_token_cost import bind_image_token_cost
from agent.usage_anchor import anchored_context_tokens, restore_usage_anchor
from agent.turn_author import parse_turn_author
logger = logging.getLogger(__name__)
def _str_attr(agent: Any, name: str) -> str:
"""``getattr(agent, name, "") or ""`` — route facts read off partial agents/doubles."""
return getattr(agent, name, "") or ""
def _preflight_request_tokens(
agent: Any, messages: List[Dict[str, Any]], system_prompt: str
) -> int:
"""Token estimate for automatic preflight compression: a valid provider usage anchor,
else the checkpoint-pruned native wire payload, else the generic estimator."""
anchored = anchored_context_tokens(messages, getattr(agent, "_usage_anchor", None))
agent._request_pressure_anchored = anchored is not None
if anchored is not None:
return anchored
tools = getattr(agent, "tools", None) or None
try:
from agent.codex_responses_adapter import estimate_native_responses_preflight_tokens
native = estimate_native_responses_preflight_tokens(
agent, messages, system_prompt=system_prompt or "", tools=tools
)
if isinstance(native, int) and not isinstance(native, bool) and native >= 0:
return native
except Exception:
logger.debug(
"native Responses preflight estimate unavailable; "
"using generic transcript estimate",
exc_info=True,
)
return estimate_request_tokens_rough(
messages, system_prompt=system_prompt or "", tools=tools,
charge_stale_thinking=_agent_stale_thinking_on_wire(agent),
)
def _agent_stale_thinking_on_wire(agent: Any) -> bool:
"""Whether the active route replays stale thinking text; ``True`` (conservative full
charge) when route facts are unavailable."""
try:
from agent.message_sanitization import stale_thinking_reaches_wire
return stale_thinking_reaches_wire(
*(_str_attr(agent, k) for k in ("api_mode", "provider", "model", "base_url"))
)
except Exception:
return True
def compose_user_api_content(
content: Any, ext_prefetch_cache: str, plugin_user_context: str
) -> Optional[str]:
"""Compose the API-bound content of the current turn's user message.
Single source for the ``api_content`` sidecar and the wire bytes so they never drift
(what turn N sends is what turn N+1 replays). ``None`` when nothing is injected."""
if not isinstance(content, str):
return None
fenced = build_memory_context_block(ext_prefetch_cache) if ext_prefetch_cache else ""
injections = [part for part in (fenced, plugin_user_context) if part]
if not injections:
return None
return content + "\n\n" + "\n\n".join(injections)
def substitute_api_content(api_msg: Dict[str, Any]) -> Optional[str]:
"""Pop the ``api_content`` sidecar and substitute it into ``content`` (keeps the
prompt-cache prefix byte-stable). Returns the popped sidecar, or ``None``."""
sidecar = api_msg.pop("api_content", None)
if isinstance(sidecar, str) and sidecar and api_msg.get("role") in ("user", "assistant"):
api_msg["content"] = sidecar
return sidecar
def drop_stale_api_content(msg: Dict[str, Any]) -> None:
"""Drop the ``api_content`` sidecar from a message whose content was rewritten
(replaying it would resend what the rewrite removed; cost is one cache miss)."""
msg.pop("api_content", None)
def extract_api_content_sidecar(msg: Mapping[str, Any]) -> Optional[str]:
"""Extract the ``api_content`` sidecar; ``None`` when absent/non-string."""
v = msg.get("api_content")
return v if isinstance(v, str) else None
def _pop_turn_note(agent: Any, attr: str) -> str:
"""One-shot per-turn note: read and clear, so the system prompt stays byte-stable and a
cached agent never replays a stale note."""
note = getattr(agent, attr, "") or ""
if hasattr(agent, attr):
with suppress(Exception):
setattr(agent, attr, "")
return note if isinstance(note, str) else ""
def consume_gateway_turn_context_notes(agent: Any) -> str:
"""Pop the gateway's per-turn must-deliver notes."""
return _pop_turn_note(agent, "_gateway_turn_context_notes")
def consume_surface_switch_note(agent: Any) -> str:
"""Pop the surface-switch note staged by the system-prompt restore (#104414); rides the same
user-message channel as the gateway notes, behind the cached prefix."""
return _pop_turn_note(agent, "_surface_switch_note")
def append_notes_to_multimodal_content(content: Any, notes: str) -> bool:
"""Append must-deliver notes as a durable text part on a multimodal (list) user
message (the sidecar path returns ``None`` for non-string content)."""
if not notes or not isinstance(content, list):
return False
with suppress(Exception):
content.append({"type": "text", "text": notes})
return True
return False
# Surfaces whose sessions must not be auto-titled: cron names its own session and
# its opener is a delivery hint; subagent sessions are hidden from every picker.
_UNTITLED_PLATFORMS = frozenset({"cron", "subagent"})
def _maybe_title_session_at_turn_start(agent: Any, messages: List[Any]) -> None:
"""Kick off auto-titling for the session's first user message; never fatal."""
session_db = getattr(agent, "_session_db", None)
session_id = getattr(agent, "session_id", None)
if not session_db or not session_id:
return
if str(getattr(agent, "platform", "") or "").lower() in _UNTITLED_PLATFORMS:
return
try:
from agent.message_content import flatten_message_text
from agent.title_generator import maybe_auto_title
# Turn's user message as text; image-only turns yield "" and are skipped.
user_text = ""
title_preview = None
for msg in reversed(messages or []):
if isinstance(msg, dict) and msg.get("role") == "user":
user_text = flatten_message_text(msg.get("content")).strip()
metadata = msg.get("display_metadata")
if isinstance(metadata, dict) and isinstance(metadata.get("title_preview"), str):
title_preview = metadata["title_preview"]
break
if not user_text:
return
# The session row is created lazily; force it now or the title write matches
# zero rows.
if not getattr(agent, "_session_db_created", False):
ensure = getattr(agent, "_ensure_db_session", None)
if callable(ensure):
ensure()
if not getattr(agent, "_session_db_created", False):
return
# Snapshot runtime identity so the background titler can skip if the user
# switches models before it fires.
# ``session_id`` rides along so the background titler's OpenCode request carries the
# same ``x-opencode-session`` affinity as the turn it belongs to (#112717).
main_runtime = {
k: getattr(agent, k, None)
for k in ("model", "provider", "base_url", "api_key", "api_mode", "session_id")
}
# See #19027.
upgrade = maybe_auto_title(
session_db,
session_id,
user_text,
conversation_history=messages,
failure_callback=(
getattr(agent, "_title_failure_callback", None)
or getattr(agent, "_emit_auxiliary_failure", None)
),
main_runtime=main_runtime,
title_callback=getattr(agent, "_on_session_title", None),
runtime_validator=lambda: (
getattr(agent, "model", None) == main_runtime["model"]
and getattr(agent, "provider", None) == main_runtime["provider"]
),
title_preview=title_preview,
)
# Unstarted = the title call would share a self-hosted endpoint with this turn's request
# (#117296); ``finalize_turn`` starts it once the model has answered.
if upgrade is not None and upgrade.ident is None:
agent._deferred_title_upgrade = upgrade
except Exception:
logger.debug("Turn-start auto-title dispatch failed", exc_info=True)
def start_deferred_title_upgrade(agent: Any) -> None:
"""Fire the title upgrade ``_maybe_title_session_at_turn_start`` held back; no-op when none."""
upgrade = getattr(agent, "_deferred_title_upgrade", None)
if upgrade is None:
return
agent._deferred_title_upgrade = None
from agent.title_generator import start_title_upgrade
start_title_upgrade(upgrade)
def reanchor_current_turn_user_idx(messages: List[Any], user_message: Any) -> int:
"""Locate this turn's user message after compaction rebuilt ``messages``.
Prefers the LAST user message whose content exactly matches this turn's text, else
the last user-originated turn; compaction handoffs are never the fallback.
Returns -1 when there is no user-originated message.
Compression replaces list entries with fresh copies (and may append a todo-snapshot user message or a
restored user turn AFTER the surviving copy of the current turn's message), so a pre-compression index
is meaningless. Prefer the LAST user message whose content exactly matches this turn's text — the
surviving copy in the common case — so the injection stamp and the #48677 persist override can't land on
a todo-snapshot or historical row. Fall back to the last *user-originated* turn when no exact match
survives (merge-summary-into-tail rewrites the content but the trackers still need a live anchor).
Compaction handoffs must never become the fallback anchor (#80622) — they are reference-only
scaffolding, not the active ask.
"""
from agent.context_compressor import user_originated_turn_view
fallback = -1
for i in range(len(messages) - 1, -1, -1):
msg = messages[i]
if not (isinstance(msg, dict) and msg.get("role") == "user"):
continue
# Typed synthetic current events keep their persistence anchor when raw
# content is unchanged; not eligible for the human-only fallback below.
if msg.get("content") == user_message:
return i
live_view = user_originated_turn_view(msg)
if live_view is None:
continue
if live_view.get("content") == user_message:
return i
# Prefer a real human turn over a synthetic handoff / continuation marker
# when the exact content was rewritten by merge-into-tail.
if fallback < 0:
fallback = i
return fallback
def export_current_turn_boundary(agent: Any, result: Any, user_message: Any) -> Any:
"""Stamp ``{turn_id, current_turn_user_idx}`` on a result envelope, proven against the
exact ``result["messages"]`` projection it travels with.
Hosts that settle their own transcript by index (hermes-webui) must never guess which
row is the current user turn after this loop rewrote history (alternation repair,
compaction, post-turn micro-compaction): a guessed index or a text match can relabel an
identical historical prompt and claim its old answer as this turn's. So the producer
exports the coordinate, computed on the final list, only when the addressed row is this
turn's user message verbatim. Otherwise the keys are omitted and hosts fail closed.
A preflight-timeout envelope carries the prior history without this turn's row (#7100), so a
repeated prompt would resolve to its historical copy: nothing is exported there.
"""
if not isinstance(result, dict) or result.get("turn_exit_reason") == "context_compression_timeout":
return result
messages = result.get("messages")
turn_id = str(getattr(agent, "_current_turn_id", "") or "")
if not isinstance(messages, list) or not turn_id or user_message is None:
return result
idx = reanchor_current_turn_user_idx(messages, user_message)
if idx < 0 or idx >= len(messages):
return result
row = messages[idx]
if not (isinstance(row, dict) and row.get("role") == "user"):
return result
from agent.context_compressor import user_originated_turn_view
live_view = user_originated_turn_view(row)
if row.get("content") != user_message and not (
isinstance(live_view, dict) and live_view.get("content") == user_message
):
return result # rewritten (merge-into-tail) row: not a proven boundary
result["turn_id"] = turn_id
result["current_turn_user_idx"] = idx
return result
def compression_made_progress(
orig_len: int, new_len: int, orig_tokens: int, new_tokens: int
) -> bool:
"""``True`` if a compression pass materially reduced the request: fewer rows, or a
>5% token cut with the same rows (same floor as the overflow-handler retry).
Compression can succeed by summarising message contents — reducing the estimated request token count —
without reducing the message row count. Treating row count as the sole progress signal false-positives
on size-only wins and surfaces a misleading "Cannot compress further" failure even when post-compression
tokens are well below the model context window. See issue #39548 for an observed case: 220 → 220
messages, ~288k → ~183k tokens on a 1M-context model still triggered auto-reset.
The token reduction must be *material* (>5%) to count as progress — the same floor the overflow-handler
retry path uses (conversation_loop.py, 39550) — so a sub-5% wobble doesn't keep the multi-pass loop
spinning. See #39550.
"""
return new_len < orig_len or (orig_tokens > 0 and new_tokens < orig_tokens * 0.95)
class PreflightCompressionTimedOut(RuntimeError):
"""Raised when an oversized turn cannot safely finish preflight."""
def _fail_closed_after_preflight_timeout(agent, request_tokens: int) -> None:
"""Stop an oversized turn instead of sending its unchanged provider payload.
Only a request the model cannot accept (above its context window, or of unknown fit) is stopped: a
request that merely sits above the compression threshold is sent unchanged, exactly as the
cooldown-blocked path sends it every turn — otherwise a slow summariser turns a session that still
fits its window into a turn that can never run (#113646, #114594)."""
from agent.conversation_compression import context_compression_timed_out, request_exceeds_model_window
if not context_compression_timed_out(agent):
return
if request_exceeds_model_window(agent, request_tokens) is False:
logger.warning(
"Preflight compression timed out but the request (~%s tokens) fits the model window (%s); "
"sending it uncompressed this turn",
f"{request_tokens:,}", f"{agent.context_compressor.context_length:,}",
)
return
raise PreflightCompressionTimedOut(
"Context compression timed out before it could commit while the request "
f"was still approximately {request_tokens:,} tokens. The provider call "
"was not sent. Run /compress and wait for it to finish, then retry."
)
def _fail_closed_on_insufficient_progress(agent, request_tokens: int) -> None:
"""Stop an over-window turn the moment preflight proves it cannot shrink the session, with
"start a new session" guidance, instead of sending a request the model cannot accept.
``_fail_closed_after_preflight_timeout`` only stops a turn whose compression wait timed out. A
pass that ran and reclaimed nothing (or under 5%) on a request still above the model window used
to fall through to the provider call: the provider rejected it, the overflow handler forced
another compression pass, and each pass re-waited its budget while the UI sat blocked (#116472:
~356k tokens on a 131k window). Only a ``True`` verdict fails closed — an unknown window or a
fitting request keeps the send-as-is behaviour — and a pass skipped by the summary-failure
cooldown is a defer, not proof of incompressibility, so it keeps its typed cooldown result.
"""
from agent.conversation_compression import compression_blocked_transiently, request_exceeds_model_window
if request_exceeds_model_window(agent, request_tokens) is not True:
return
if compression_blocked_transiently(agent):
return
window = agent.context_compressor.context_length
raise PreflightCompressionTimedOut(
"Context compression could not bring this session under the model's context window "
f"(~{request_tokens:,} tokens vs {window:,}). The provider call was not "
"sent. Start a new session with /new; this session is too large to compress further."
)
def _review_fork_first_request_pending(agent: Any) -> bool:
"""Whether a detached review fork has yet to send its first provider request: it
replays the parent's FULL snapshot as a warm cache read, so compaction must wait
for that first response. Dormant without the attribute."""
return bool(
getattr(agent, "_review_defer_compaction_before_first_response", False)
and not getattr(agent, "_turn_received_provider_response", False)
)
def _compression_warrants_another_preflight_pass(
orig_tokens: int, new_tokens: int, threshold_tokens: int
) -> bool:
"""Another immediate summary only if still over threshold AND the previous pass cut
tokens by >5%."""
return new_tokens >= threshold_tokens and orig_tokens > 0 and new_tokens < orig_tokens * 0.95
def _should_run_preflight_estimate(
messages: List[Dict[str, Any]], protect_first_n: int, protect_last_n: int, threshold_tokens: int
) -> bool:
"""Cheap gate for the (expensive) full preflight estimate: message count exceeds the
protected ranges OR a rough char-based estimate crosses the threshold (few-but-huge
case). The estimator undercounts by design (omits system/tools) so one large base64
image is not mistaken for ~250K tokens."""
return (
len(messages) > protect_first_n + protect_last_n + 1
or estimate_messages_tokens_rough(messages) >= threshold_tokens
)
def _should_idle_compact(
*, enabled: bool, idle_after_seconds: int, idle_gap_seconds: float, tokens: int,
floor_tokens: int, cooldown_active: bool, last_compaction_tokens: int = 0,
) -> bool:
"""Pure predicate: idle compaction fires after a wall-clock gap of
``idle_after_seconds`` (opt-in, <= 0 disables), independent of ``threshold_tokens``;
never at/below ``floor_tokens`` or during a compression-failure cooldown.
``floor_tokens`` (``threshold_tokens × summary_target_ratio``) is a theoretical target a
real pass routinely misses (system prompt, tool schemas and protected head/tail are
incompressible), so a session compacted to above it would re-summarise on every idle
resume without growing. ``last_compaction_tokens`` — what the previous pass actually
produced (``ContextCompressor.last_compression_rough_tokens``, same rough shape as
``tokens``) — raises the floor to ``last + floor_tokens`` so the transcript must gain a
floor's worth of NEW content first. ``0`` (nothing compacted yet / counter reset) keeps
the original semantics exactly.
A session that compacted to well above that target therefore stays above it forever, so every later idle
resume re-runs a full summarisation over a transcript that has not grown — minutes of silently blocked
prompt on a slow route, reclaiming nothing (#97239).
"""
if not enabled or idle_after_seconds <= 0 or idle_gap_seconds < idle_after_seconds or cooldown_active:
return False
effective_floor = floor_tokens
if last_compaction_tokens > 0:
effective_floor = max(effective_floor, last_compaction_tokens + floor_tokens)
return tokens > effective_floor
@dataclass
class TurnContext:
"""Values produced by the turn prologue and consumed by the turn loop."""
user_message: str # sanitized inbound message (surrogates stripped)
original_user_message: Any # clean text for transcripts / memory queries (no nudges)
messages: List[Dict[str, Any]] # working list for this turn (loop appends to it)
conversation_history: Optional[List[Dict[str, Any]]] # None after rotation
active_system_prompt: Optional[str] # may be rebuilt by compression
effective_task_id: str
turn_id: str
current_turn_user_idx: int # index of the current user turn within ``messages``
should_review_memory: bool = False # post-turn memory review should fire
plugin_user_context: str = "" # ``pre_llm_call`` context (appended to user message)
ext_prefetch_cache: str = "" # external-memory prefetch, reused across iterations
preflight_compression_blocked: bool = False # immediate retry proved ineffective
def _persist_under_lock(agent: Any, fn, failure_msg: str, pending_cli_message: Any) -> None:
"""Run ``fn`` under the session persist lock (when the agent has one), log-and-swallow
failures, then drop staged CLI input — unless it is an unmarked handoff kept for a
close retry (once ``_db_persisted`` the close path must not treat it as pre-worker
UI input). Eager clearing keeps a preflight crash from leaking stale input."""
try:
lock = getattr(agent, "_session_persist_lock", None)
if lock is None:
fn()
else:
with lock:
fn()
except Exception:
logger.warning(failure_msg, agent.session_id or "none", exc_info=True)
finally:
if not isinstance(pending_cli_message, dict) or pending_cli_message.get("_db_persisted"):
agent._pending_cli_user_message = None
def _publish_runtime_main(agent: Any) -> None:
"""Tell auxiliary_client the live main provider/model for this turn (after primary
restoration settled the runtime). Never raises: failure loses only the scope."""
with suppress(Exception):
from agent.auxiliary_client import set_runtime_main
from agent.prompt_cache_scope import resolve_prompt_cache_scope_safe
# Rotation-stable prompt-cache scope (lineage root), memoized per segment; a new
# session uses the physical id until build_api_kwargs re-resolves.
# Memoized per segment on the agent, so this is a DB walk at most once per segment — except a
# brand-new session whose row lands later in turn setup (_ensure_db_session); that first turn falls
# back to the physical id here and the first build_api_kwargs re-resolves. Stays valid through a
# mid-turn compression rotation because the lineage root is by definition rotation-invariant
# (#79017). Resolved with the never-raising variant OUTSIDE the argument list, so a resolution
# failure can only lose the scope — never the whole runtime binding.
_cache_scope = resolve_prompt_cache_scope_safe(agent) or ""
set_runtime_main(
_str_attr(agent, "provider"), _str_attr(agent, "model"),
**{k: _str_attr(agent, k) for k in (
"requested_provider", "base_url", "api_key", "api_mode", "auth_mode", "session_id"
)},
cache_scope=_cache_scope,
)
def _refresh_mcp_tools_between_turns(agent: Any) -> None:
"""Late-connecting MCP servers land in THIS turn's snapshot, before the first API
call assembles ``tools=``. ``preserve_prefix`` keeps the tool array append-only so a
flapping ``check_fn`` can't fork the cache."""
try:
# An authorization that committed after its connection card closed: same import-cost gate,
# the module is loaded only in a process that ran a connection operation.
if "tools.connectors.mcp" in sys.modules:
from tools.connectors.mcp import adopt_late_connections
adopt_late_connections(agent)
# Import-cost gate: MCP tools are only registered by code that already imported
# ``tools.mcp_tool`` (~0.4s); not in sys.modules => nothing to do.
if not getattr(agent, "_skip_mcp_refresh", False) and "tools.mcp_tool" in sys.modules:
from tools.mcp_tool_discovery import has_registered_mcp_tools
from tools.mcp_tool_agent import refresh_agent_mcp_tools
if has_registered_mcp_tools():
refresh_agent_mcp_tools(agent, quiet_mode=True, preserve_prefix=True)
except Exception:
logger.debug("between-turns MCP tool refresh skipped", exc_info=True)
def _bind_turn_identity(
agent: Any, task_id: Optional[str], stream_callback, persist_user_message: Any,
persist_user_timestamp: Optional[float], persist_user_platform_id: Optional[str],
) -> Tuple[str, str]:
"""Stage callback/persist overrides on the agent and bind this turn's task and turn
ids. Returns ``(effective_task_id, turn_id)``."""
agent._stream_callback = stream_callback # picked up by _interruptible_api_call
agent._persist_user_message_idx = None
agent._persist_user_message_override = persist_user_message
agent._persist_user_message_timestamp = persist_user_timestamp
agent._persist_user_message_platform_id = persist_user_platform_id
# Unique task_id when not provided isolates VMs between tasks.
effective_task_id = task_id or str(uuid.uuid4())
agent._current_task_id = effective_task_id
agent._process_owner_task_ids = {*getattr(agent, "_process_owner_task_ids", ()), effective_task_id}
turn_id = str(getattr(agent, "_relay_pending_turn_id", "") or "") or (
f"{agent.session_id or 'session'}:{effective_task_id}:{uuid.uuid4().hex[:8]}"
)
agent._relay_pending_turn_id = None
agent._current_turn_id = turn_id
agent._current_api_request_id = ""
# Tripwire: warn when this turn starts before the previous turn-end persist
# (concurrent turns interleave transcript writes). Cleared in _persist_session.
from agent.agent_runtime_helpers import note_turn_start
note_turn_start(agent, turn_id)
return effective_task_id, turn_id
# Per-turn agent state reset at turn start (retry counters, guardrail halt, file-mutation
# verifier). ``_turns_since_memory`` / ``_iters_since_skill`` are deliberately NOT reset.
_PER_TURN_RESET_STATE: Tuple[Tuple[str, Any], ...] = (
("_invalid_tool_retries", 0), ("_invalid_json_retries", 0), ("_empty_content_retries", 0),
("_incomplete_scratchpad_retries", 0), ("_codex_incomplete_retries", 0),
# Consecutive Codex reasoning-only (no answer, no tool call) responses, kept apart from
# the aggregate incomplete count so a visible partial resets it (#67321).
("_codex_reasoning_only_streak", 0),
("_thinking_prefill_retries", 0), ("_post_tool_empty_retried", False),
("_last_content_with_tools", None), ("_last_content_tools_all_housekeeping", False),
("_mute_post_response", False), ("_unicode_sanitization_passes", 0),
("_tool_guardrail_halt_decision", None), ("_vision_supported", True),
("_iteration_budget_warning_injected", False),
("_run_budget_wrapup_injected", False), ("_verification_stop_nudges", 0),
("_pre_verify_nudges", 0),
)
def _reset_per_turn_agent_state(agent: Any) -> None:
"""Reset retry counters, guardrails, iteration and run budgets at turn start."""
for name, value in _PER_TURN_RESET_STATE:
setattr(agent, name, value)
agent._turn_failed_file_mutations = {}
agent._turn_file_mutation_paths = set()
agent._tool_guardrails.reset_for_turn()
_reset_consol = getattr(agent._memory_store, "reset_consolidation_failures", None)
if callable(_reset_consol):
_reset_consol()
# Expiry clock for build_api_messages: admission time (not the input's platform-event
# stamp, which can predate admission by minutes), frozen so every request this turn
# sends identical bytes. Distinct from note_turn_start's _inflight_turn_started, a
# tripwire slot cleared at persist.
agent._current_turn_timestamp = time.time()
# Pre-turn connection health check: clean up dead TCP connections.
if agent.api_mode != "anthropic_messages":
with suppress(Exception):
if agent._cleanup_dead_connections():
agent._emit_diagnostic_status(
"🔌 Detected stale connections from a previous provider "
"issue — cleaned up automatically. Proceeding with fresh "
"connection."
)
# Replay compression warning through status_callback for gateway platforms.
if agent._compression_warning:
agent._replay_compression_warning()
agent._compression_warning = None # send once
agent.iteration_budget = IterationBudget(agent.max_iterations)
# Wall-clock run budget: stamped only when configured (one wrap-up notice per run).
agent._run_budget_started_at = (
time.time() if getattr(agent, "run_budget_seconds", None) else None
)
# Reset the streaming context / think scrubbers at the top of each turn.
for name in ("_stream_context_scrubber", "_stream_think_scrubber"):
scrubber = getattr(agent, name, None)
if scrubber is not None:
scrubber.reset()
def _stage_turn_user_message(
agent: Any, user_message: Any, persist_user_message: Any,
persist_user_timestamp: Optional[float], persist_user_platform_id: Optional[str],
persist_user_display_kind: Optional[str],
persist_user_display_metadata: Optional[Dict[str, Any]],
) -> Tuple[Dict[str, Any], Any]:
"""Build this turn's user dict, reusing CLI-staged input only when its clean text
matches this turn (a stale handoff must not replace later input; voice turns
compare the clean override). Returns ``(user_msg, pending_cli_message)``."""
pending_cli_message = getattr(agent, "_pending_cli_user_message", None)
expected_persist_content = (
persist_user_message if persist_user_message is not None else user_message
)
if (
isinstance(pending_cli_message, dict)
and pending_cli_message.get("content") == expected_persist_content
):
user_msg = pending_cli_message
# CLI-staged value is the clean text; restore the API-facing variant (e.g. voice
# prefix) on the same dict, keeping any close-path durable marker.
user_msg["content"] = user_message
else:
user_msg = {"role": "user", "content": user_message}
if isinstance(pending_cli_message, dict):
agent._pending_cli_user_message = None
# CLI input is stamped when staged; gateway input may carry the platform event
# time. Preserve either value and cover any legacy unstamped handoff.
stamp_message_timestamp(user_msg, timestamp=persist_user_timestamp)
# Synthesized turns stamp their transcript type so the crash persist writes a typed
# row; the model still receives role/content unchanged (api_messages strips both).
if persist_user_display_kind:
user_msg["display_kind"] = persist_user_display_kind
if persist_user_display_metadata:
user_msg["display_metadata"] = persist_user_display_metadata
# The platform message id survives the turn-start flush; restart drain-window
# recovery dedups via ``has_platform_message_id`` against this row.
if persist_user_platform_id is not None:
user_msg["platform_message_id"] = persist_user_platform_id
return user_msg, pending_cli_message
def _hydrate_from_history(agent: Any, conversation_history: Optional[List[Any]]) -> None:
"""Hydrate process-local state from persisted history on the first resumed turn."""
if not conversation_history:
return
if not agent._todo_store.has_items():
agent._hydrate_todo_store(conversation_history)
# A live native checkpoint arms this latch while its response is captured. A
# restarted agent must recover the same one-response deferral before turn-start
# compression can rewrite the restored opaque checkpoint. Reuse the adapter's
# exact route/issuer/replay filtering and tolerate plugin compressors without the
# optional hook.
if agent._user_turn_count == 0:
# A fresh process has no in-memory anchor; the persisted one is honored only while the
# restored transcript still carries the priced prefix (see agent/usage_anchor.py).
restore_usage_anchor(agent, conversation_history)
note_checkpoint = getattr(
getattr(agent, "context_compressor", None),
"note_native_compaction_checkpoint",
None,
)
if callable(note_checkpoint):
try:
from agent.codex_responses_adapter import (
has_replayable_native_compaction_checkpoint,
)
if has_replayable_native_compaction_checkpoint(
agent, conversation_history
):
note_checkpoint()
except Exception:
logger.debug(
"restored native checkpoint hydration skipped", exc_info=True
)
# Hydrate per-session nudge counters from persisted history.
prior_user_turns = sum(1 for m in conversation_history if m.get("role") == "user")
if prior_user_turns > 0:
agent._user_turn_count = prior_user_turns
if agent._memory_nudge_interval > 0 and agent._turns_since_memory == 0:
agent._turns_since_memory = prior_user_turns % agent._memory_nudge_interval
def _tick_memory_nudge(agent: Any) -> bool:
"""Advance the turn-based memory nudge counter; ``True`` when the review should fire."""
if (agent._memory_nudge_interval > 0
and "memory" in agent.valid_tool_names
and agent._memory_store):
agent._turns_since_memory += 1
if agent._turns_since_memory >= agent._memory_nudge_interval:
agent._turns_since_memory = 0
return True
return False
def _emit_reaction(agent: Any, original_user_message: Any) -> None:
"""Cosmetic side-signal: detect an affection reaction so the host can play hearts.
Token-free, never touches the conversation, never fatal."""
reaction_callback = getattr(agent, "reaction_callback", None)
if reaction_callback is None:
return
with suppress(Exception):
from agent.reactions import detect_reaction
kind = detect_reaction(original_user_message)
if kind:
reaction_callback(kind)
def _ensure_session_row(agent: Any, pending_cli_message: Any) -> None:
"""Create the DB row now (system prompt populated => non-NULL) and BEFORE preflight
compression: compaction/rotation INSERTs reference this row under PRAGMA
foreign_keys=ON. Idempotent; the user-turn crash persist runs later."""
_persist_under_lock(
agent, agent._ensure_db_session,
"Turn-start session row creation failed for session=%s", pending_cli_message,
)
def _collect_pre_llm_call_context(
agent: Any, *, effective_task_id: str, turn_id: str, original_user_message: Any,
messages: List[Any], conversation_history: Optional[List[Any]],
) -> str:
"""Run ``pre_llm_call`` plugins; their context is injected into the user message
(never the system prompt). Oversized per-hook context is spilled to disk so a
runaway plugin can't inflate every subsequent turn's prompt."""
if getattr(agent, "_persist_disabled", False):
return ""
try:
from hermes_cli.lifecycle import invoke_hook as _invoke_hook
_pre_results = _invoke_hook(
"pre_llm_call",
session_id=agent.session_id,
task_id=effective_task_id,
turn_id=turn_id,
user_message=original_user_message,
conversation_history=list(messages),
is_first_turn=(not bool(conversation_history)),
model=agent.model,
platform=getattr(agent, "platform", None) or "",
parent_session_id=getattr(agent, "_parent_session_id", None) or "",
sender_id=getattr(agent, "_user_id", None) or "",
)
try:
# Spill oversized per-hook context to disk so a runaway plugin can't inflate every subsequent
# turn's prompt. Ported from openai/codex PR #21069 ("Spill large hook outputs from context").
from tools.hook_output_spill import (
get_spill_config as _spill_cfg, spill_if_oversized as _spill_if_oversized
)
_spill_config_cached = _spill_cfg()
except Exception:
_spill_if_oversized = None # type: ignore[assignment]
_spill_config_cached = None
_ctx_parts: list[str] = []
for r in _pre_results:
if isinstance(r, dict) and r.get("context"):
_piece = str(r["context"])
elif isinstance(r, str) and r.strip():
_piece = r
else:
continue
if _spill_if_oversized is not None:
try:
_piece = _spill_if_oversized(
_piece, session_id=agent.session_id, source="plugin hook",
config=_spill_config_cached,
)
except Exception as _spill_exc:
logger.warning("hook context spill failed: %s", _spill_exc)
_ctx_parts.append(_piece)
return "\n\n".join(_ctx_parts)
except Exception as exc:
logger.warning("pre_llm_call hook failed: %s", exc)
return ""
def _merge_gateway_notes(
agent: Any, messages: List[Any], current_turn_user_idx: int, plugin_user_context: str
) -> str:
"""Must-deliver per-turn notes ride the user-message injection channel (one-shot) so the
ephemeral system prompt stays byte-stable: the gateway's staged notes, then the
surface-switch correction. Multimodal (list) content can't take the string sidecar —
append a durable text part instead."""
_turn_notes = "\n\n".join(
part for part in (consume_gateway_turn_context_notes(agent),
consume_surface_switch_note(agent)) if part
)
if not _turn_notes:
return plugin_user_context
_gw_turn_content = (
messages[current_turn_user_idx].get("content")
if 0 <= current_turn_user_idx < len(messages)
and isinstance(messages[current_turn_user_idx], dict)
else None
)
if isinstance(_gw_turn_content, list):
append_notes_to_multimodal_content(_gw_turn_content, _turn_notes)
return plugin_user_context
return (
plugin_user_context + "\n\n" + _turn_notes if plugin_user_context else _turn_notes
)
def _bind_interrupt_scope(agent: Any, ra) -> None:
"""Record the execution thread so interrupt()/clear_interrupt() scope the tool-level
signal to THIS agent's thread; clear stale state, preserving a pending interrupt."""
agent._execution_thread_id = threading.current_thread().ident
ra()._set_interrupt(False, agent._execution_thread_id)
if agent._interrupt_requested:
ra()._set_interrupt(
True, agent._execution_thread_id, reason=getattr(agent, "_tool_interrupt_reason", None)
)
else:
agent._interrupt_message = None
agent._tool_interrupt_reason = None
agent._interrupt_thread_signal_pending = False
def _memory_turn_start_and_prefetch(
agent: Any, original_user_message: Any, turn_author: Optional[Dict[str, Any]] = None,
) -> str:
"""Notify memory providers of the new turn, then prefetch external memory once
before the tool loop (skipped on trivial prompts with no semantic signal).
Returns the prefetch text (``""`` when nothing was injected)."""
if not agent._memory_manager:
return ""
_query = original_user_message if isinstance(original_user_message, str) else ""
# The author rides along so a provider can attribute THIS turn, not whoever opened the session.
_author = turn_author if isinstance(turn_author, dict) else {}
with suppress(Exception):
agent._memory_manager.on_turn_start(
agent._user_turn_count, _query,
author_id=_author.get("id") or None, author_name=_author.get("name") or None,
author_is_bot=bool(_author.get("is_bot")),
)
ext_prefetch_cache = ""
with suppress(Exception):
if not is_trivial_prompt(_query):
ext_prefetch_cache = agent._memory_manager.prefetch_all(_query, session_id=agent.session_id) or ""
# Deterministic recall indicator via _emit_status so the model can't silently
# drop injected memory.
if ext_prefetch_cache:
with suppress(Exception):
_recall_indicator = agent._memory_manager.describe_recall()
if _recall_indicator:
agent._emit_status(_recall_indicator)
return ext_prefetch_cache
def _stamp_api_content_sidecar(
agent: Any, messages: List[Any], current_turn_user_idx: int, ext_prefetch_cache: str,
plugin_user_context: str, *, preflight_compressed: bool,
) -> None:
"""api_content sidecar — persist what you send: injected context lives only in the
API copy, so stamp the exact sent bytes on the live dict for replay."""
_turn_user_msg = messages[current_turn_user_idx]
live_content = _turn_user_msg.get("content")
from agent.session_persistence import _persist_lock, durable_user_row_content
# Match the row the flush wrote (persist override = clean transcript), not the live bytes.
durable_content, _api_content = durable_user_row_content(
agent, _turn_user_msg, live_content,
compose_user_api_content(live_content or "", ext_prefetch_cache, plugin_user_context),
)
if _api_content is None or _api_content == durable_content:
return
_turn_user_msg["api_content"] = _api_content
# When another writer materialized this turn's user row BEFORE the sidecar existed — in-place
# preflight compaction, or a close/early flush that raced the prologue (#102194) — the crash
# persist marker-skips the message and the stamp never reaches the DB, so the next turn replays
# clean content and the request prefix diverges here. Both writers stamp ``_row_id`` on the live
# dict, which is at once the proof a row exists and the address to update.
#
# Never widen this to an unconditional positional backfill — see set_latest_user_api_content.
#
# ``_row_id`` is read under ``_session_persist_lock``: a close flush holds it while it commits
# the row and only then writes ``_row_id`` back (``sync_flushed_message_markers``). Read outside
# it, the stamp can land in between, see no id, return — and the flush then marks the message
# persisted with ``api_content = NULL``, leaving no writer to correct the row.
with _persist_lock(agent):
_row_id = _turn_user_msg.get("_row_id")
_in_place_compacted = preflight_compressed and bool(getattr(agent, "_last_compaction_in_place", False))
_db = getattr(agent, "_session_db", None)
if _db is None or not (isinstance(_row_id, int) or _in_place_compacted):
return
try:
if isinstance(_row_id, int):
_db.set_message_api_content(agent.session_id, _row_id, durable_content, _api_content)
else:
# Compacted copies carry no row id; positional is safe only because
# archive_and_compact just made this message the newest active user row.
_db.set_latest_user_api_content(agent.session_id, durable_content, _api_content)
except Exception:
logger.warning("api_content backfill failed for session=%s", agent.session_id or "none", exc_info=True)
def _persist_turn_start(
agent: Any, messages: List[Any], conversation_history: Optional[List[Any]],
pending_cli_message: Any,
) -> None:
"""Crash-resilience: persist the inbound user turn once, with final api_content,
before the first LLM call. Same critical section as CLI close persistence; retries
the row create if the pre-compression attempt failed transiently."""
def _ensure_and_persist() -> None:
agent._ensure_db_session()
agent._persist_session(messages, conversation_history)
_persist_under_lock(
agent, _ensure_and_persist,
"Early turn-start session persistence failed for session=%s", pending_cli_message,
)
def build_turn_context(
agent, user_message: Any, system_message: Optional[str],
conversation_history: Optional[List[Dict[str, Any]]], task_id: Optional[str], stream_callback,
persist_user_message: Optional[Any], persist_user_timestamp: Optional[float]=None,
persist_user_platform_id: Optional[str]=None, *, persist_user_display_kind: Optional[str]=None,
persist_user_display_metadata: Optional[Dict[str, Any]]=None, turn_author: Optional[Dict[str, Any]]=None,
restore_or_build_system_prompt,
install_safe_stdio, sanitize_surrogates, summarize_user_message_for_log, set_session_context,
set_current_write_origin, ra, moa_active: bool=False,
) -> TurnContext:
"""Run the once-per-turn setup and return the loop's input context.
Helpers are passed in to avoid an import cycle with ``agent.conversation_loop``.
Order matters: the DB session row is created only AFTER the system prompt is built
(else it persists system_prompt=NULL and costs a cache miss) and BEFORE preflight
compression."""
from agent.turn_context_compaction import run_turn_start_compaction
# Guard stdio against OSError from broken pipes (systemd/headless/daemon).
install_safe_stdio()
# Reset first: a cached gateway agent must never carry the previous turn's bot author into a human turn.
turn_author = parse_turn_author(turn_author)
agent._turn_author = turn_author
# Recover a rotated session before binding log/turn ids or copying client history so
# everything in this turn belongs to the canonical child.
recovered_history = recover_rotated_compression_session(agent)
if recovered_history is not None:
conversation_history = recovered_history
# Tag log records on this thread with the session ID for ``hermes logs``; bind the
# skill write-origin ContextVar; restore the primary runtime after a fallback turn.
# NOTE: the DB session row is created later, AFTER the system prompt is restored/built (see
# _ensure_db_session() below the system-prompt block). Creating it here — before _cached_system_prompt
# is populated — inserts a row with system_prompt=NULL on a fresh API/gateway agent that carries
# client-managed history, which then trips the "stored system prompt is null; rebuilding from scratch"
# warning and a needless first-turn prefix cache miss. (Issue #45499.)
set_session_context(agent.session_id)
set_current_write_origin(getattr(agent, "_memory_write_origin", "assistant_tool"))
from tools.skill_provenance import set_review_attended
set_review_attended(getattr(agent, "_review_attended", False))
agent._restore_primary_runtime()
_publish_runtime_main(agent)
_refresh_mcp_tools_between_turns(agent)
if isinstance(user_message, str):
user_message = sanitize_surrogates(user_message)
if isinstance(persist_user_message, str):
persist_user_message = sanitize_surrogates(persist_user_message)
effective_task_id, turn_id = _bind_turn_identity(
agent, task_id, stream_callback, persist_user_message,
persist_user_timestamp, persist_user_platform_id,
)
_reset_per_turn_agent_state(agent)
_preview_text = summarize_user_message_for_log(user_message)
_msg_preview = _preview_text[:80] + ("..." if len(_preview_text) > 80 else "")
logger.info(
"conversation turn: session=%s model=%s provider=%s platform=%s history=%d msg=%r",
agent.session_id or "none", agent.model, agent.provider or "unknown",
agent.platform or "unknown", len(conversation_history or []),
_msg_preview.replace("\n", " "),
)
# Copy so the caller's list is never mutated.
messages = list(conversation_history) if conversation_history else []
user_msg, pending_cli_message = _stage_turn_user_message(
agent, user_message, persist_user_message, persist_user_timestamp,
persist_user_platform_id, persist_user_display_kind, persist_user_display_metadata,
)
_hydrate_from_history(agent, conversation_history)
# Every estimator this turn prices images at the cost learned from this model's real usage.
bind_image_token_cost(agent)
# Append the user message now that close persistence is safe.
append_message(messages, user_msg)
current_turn_user_idx = len(messages) - 1
agent._persist_user_message_idx = current_turn_user_idx
agent._user_turn_count += 1
# Copilot x-initiator: the first API call of this user turn is user-initiated;
# tool-loop follow-ups revert to "agent".
agent._is_user_initiated_turn = True
# Preserve the original user message (no nudge injection).
original_user_message = persist_user_message if persist_user_message is not None else user_message
should_review_memory = _tick_memory_nudge(agent)
_emit_reaction(agent, original_user_message)
if not agent.quiet_mode:
agent._safe_print(
f"💬 Starting conversation: '{_preview_text[:60]}"
f"{'...' if len(_preview_text) > 60 else ''}'"
)
# System prompt is cached per session for prefix caching.
if agent._cached_system_prompt is None:
restore_or_build_system_prompt(agent, system_message, conversation_history)
active_system_prompt = agent._cached_system_prompt
# Bot Mode DM tool — injected ONLY into a bot's canonical "Bot Chat" session (same
# gate as the protocol section); gate is session-stable, so cache-safe.
try:
from tools.bot_mode_dm import ensure_message_agent_tool
ensure_message_agent_tool(agent)
except Exception:
logger.debug("message_agent injection skipped", exc_info=True)
_ensure_session_row(agent, pending_cli_message)
# A turn interrupted before admission could not write its accepted input because
# it did not own the session lease. Persist that carried-forward row now, before
# compaction can rewrite or drop it.
from agent.session_persistence import _PERSIST_AFTER_ADMISSION_INTERRUPT
if conversation_history and any(
isinstance(msg, dict) and msg.get(_PERSIST_AFTER_ADMISSION_INTERRUPT)
for msg in conversation_history
):
agent._flush_messages_to_session_db(conversation_history, conversation_history)
compaction = run_turn_start_compaction(
agent, messages=messages, system_message=system_message,
active_system_prompt=active_system_prompt, conversation_history=conversation_history,
current_turn_user_idx=current_turn_user_idx, user_message=user_message,
effective_task_id=effective_task_id,
)
messages = compaction.messages
active_system_prompt = compaction.active_system_prompt
conversation_history = compaction.conversation_history
current_turn_user_idx = compaction.current_turn_user_idx
plugin_user_context = _collect_pre_llm_call_context(
agent, effective_task_id=effective_task_id, turn_id=turn_id,
original_user_message=original_user_message, messages=messages,
conversation_history=conversation_history,
)
plugin_user_context = _merge_gateway_notes(
agent, messages, current_turn_user_idx, plugin_user_context
)
_bind_interrupt_scope(agent, ra)
ext_prefetch_cache = _memory_turn_start_and_prefetch(agent, original_user_message, turn_author)
# Sidecar skipped for codex_app_server/MoA.
if (
not moa_active
and getattr(agent, "api_mode", None) != "codex_app_server"
and 0 <= current_turn_user_idx < len(messages)
and messages[current_turn_user_idx].get("role") == "user"
):
_stamp_api_content_sidecar(
agent, messages, current_turn_user_idx, ext_prefetch_cache,
plugin_user_context, preflight_compressed=compaction.compressed,
)
_persist_turn_start(agent, messages, conversation_history, pending_cli_message)
# Title the session now: the row exists and titling depends only on the user's ask,
# so it runs concurrently with the turn. Daemon thread, no-op once titled.
_maybe_title_session_at_turn_start(agent, messages)
return TurnContext(
user_message=user_message, original_user_message=original_user_message, messages=messages,
conversation_history=conversation_history, active_system_prompt=active_system_prompt,
effective_task_id=effective_task_id, turn_id=turn_id,
current_turn_user_idx=current_turn_user_idx, should_review_memory=should_review_memory,
plugin_user_context=plugin_user_context, ext_prefetch_cache=ext_prefetch_cache,
preflight_compression_blocked=compaction.blocked,
)
def _sanitize_model_for(agent: Any, moa_config: Any) -> Any:
"""Model name for strict-API tool-call sanitization. In MoA mode ``agent.model`` is
the virtual preset name; use the resolved aggregator so Gemini keeps
thought_signature (extra_content)."""
_sanitize_model = agent.model
if agent.provider == "moa":
if moa_config:
_agg = moa_config.get("aggregator") or {}
if _agg.get("model"):
_sanitize_model = _agg["model"]
if _sanitize_model == agent.model:
# Virtual-provider mode: no moa_config is threaded through; ask the facade
# for the aggregator slot from the previous create().
_agg_slot = getattr(getattr(agent, "client", None), "last_aggregator_slot", None)
if _agg_slot and _agg_slot.get("model"):
_sanitize_model = _agg_slot["model"]
return _sanitize_model
def build_api_messages(
agent: Any, messages: List[Dict[str, Any]], *, current_turn_user_idx: Any,
ext_prefetch_cache: Any, plugin_user_context: Any, moa_config: Any, active_system_prompt: Any,
) -> Tuple[List[Dict[str, Any]], str]:
"""Build the wire copy of ``messages`` for one API call plus the effective system
message. Returns ``(api_messages, effective_system)``.
Prompt-cache invariant: historical user/assistant rows replay their ``api_content``
sidecar (the exact bytes sent live) so the prefix stays byte-stable; the current
user turn reuses the prologue's stamp (or composes live when a caller bypassed the
prologue). Ephemeral context (prefetch, ``pre_llm_call`` hooks,
``ephemeral_system_prompt``) is added at API time only — ``messages`` stays untouched
beyond the sidecar stamp, and the system prompt is built ONCE per session and
replayed verbatim."""
from agent.agent_runtime_helpers import fill_empty_non_final_wire_payload
from agent.conversation_loop import _clone_message_for_send
from agent.replay_cleanup import canonicalize_replay_history
has_current = isinstance(current_turn_user_idx, int) and 0 <= current_turn_user_idx < len(messages)
current_turn_message = messages[current_turn_user_idx] if has_current else None
# Replay consumers canonicalize the persisted prefix on read; the request copy must
# carry the same bytes or a resume diverges mid-prefix. Only the rows BEFORE this
# turn's user message are the replayed prefix — rows this turn appended (its tool
# calls/results) are live and must never be rewritten between iterations. The
# expiry clock is the turn's admission time, frozen in _reset_per_turn_agent_state.
# Without an anchor (compaction found no surviving user row) there is no provable
# persisted prefix, so nothing is canonicalized. The clock is stamped once per turn in
# _reset_per_turn_agent_state; a caller that skipped the prologue fails loudly here
# rather than silently un-freezing it.
turn_now = agent._current_turn_timestamp
split = current_turn_user_idx if has_current else 0
canonical_messages = canonicalize_replay_history(messages[:split], now=turn_now) + messages[split:]
api_messages = []
for idx, msg in enumerate(canonical_messages):
# Structural clone, NOT msg.copy(): in-place transforms below must not reach
# persisted history via nested containers; see _clone_message_for_send.
api_msg = _clone_message_for_send(msg)
# api_content is bookkeeping (exact bytes sent), never a provider field — pop
# it from EVERY outgoing copy. display_* is display-only timeline metadata
# (strict OpenAI backends reject unknown keys); _row_id is the durable row id
# from _rows_to_conversation and only chat-completions strips underscore keys.
_api_content = api_msg.pop("api_content", None)
for key in ("display_kind", "display_metadata", "_row_id"):
api_msg.pop(key, None)
# Inject ephemeral context (memory prefetch + pre_llm_call user hooks)
# at API time only; `messages` is untouched beyond the api_content stamp.
if msg is current_turn_message and msg.get("role") == "user":
if isinstance(_api_content, str) and _api_content:
# Reuse the prologue's stamp so sidecar and wire cannot drift
# and every pass this turn sends identical bytes.
api_msg["content"] = _api_content
else:
# Callers that bypass the prologue stamping: compose live.
_composed = compose_user_api_content(
api_msg.get("content", ""), ext_prefetch_cache, plugin_user_context
)
if _composed is not None:
api_msg["content"] = _composed
elif (
isinstance(_api_content, str) and _api_content
and msg.get("role") in ("user", "assistant")
):
# Historical row: replay the exact bytes sent live so the prompt-cache
# prefix stays byte-stable. User rows carry the injection sidecar; user
# and assistant rows may carry a sanitize-divergence sidecar.
api_msg["content"] = _api_content
# Pass reasoning back to the API for ALL assistant messages so multi-turn
# reasoning context is preserved.
agent._copy_reasoning_content_for_api(msg, api_msg)
# 'reasoning' is trajectory-only (copied to 'reasoning_content' above);
# finish_reason is rejected by strict APIs (e.g. Mistral).
api_msg.pop("reasoning", None)
api_msg.pop("finish_reason", None)
# Fill empty non-final user/assistant wire copies so the pre-call sanitizer
# stops re-healing and flooding errors.log; durable history is untouched.
# After the reasoning copy so thinking-only turns keep payload.
fill_empty_non_final_wire_payload(api_msg, is_final=(idx == len(canonical_messages) - 1))
# _thinking_prefill survives intentionally: the drop pass below needs it.
# Strip length-continuation marks; some transports keep underscore keys.
api_msg.pop("_length_continuation_fragment", None)
api_msg.pop("_length_continuation_nudge", None)
# Strip Codex Responses fields (call_id, response_item_id): strict providers
# reject unknown fields. New dicts keep the internal list intact for Codex.
if agent._should_sanitize_tool_calls():
agent._sanitize_tool_calls_for_strict_api(
api_msg, model=_sanitize_model_for(agent, moa_config)
)
# 'reasoning_details' is kept here; the chat-completions transport drops it on the
# wire for every route that does not replay it (OpenRouter/Nous do).
api_messages.append(api_msg)
# Final system message = cached prompt + ephemeral additions (API-time only).
# Plugin/recall context goes into the user message, never the system prompt: the
# prompt is built ONCE per session and replayed verbatim (stable cache prefix).
effective_system = active_system_prompt or ""
if agent.ephemeral_system_prompt:
effective_system = (effective_system + "\n\n" + agent.ephemeral_system_prompt).strip()
if effective_system:
api_messages = [{"role": "system", "content": effective_system}] + api_messages
return api_messages, effective_system
# ---- BEGIN PLUGIN-COMPAT (revert-scheduled; see COMPAT_MANIFEST.md) ----
# Names external plugins imported from this module before the Sep 2026 decomposition.
# Internal code MUST NOT use these (scripts/check_compat_pointers.py fails CI if it does).
# The whole block is removed by reverting the commit that added it.
_PLUGIN_COMPAT_LAZY = {
'IDLE_COMPACTION_STATUS_TEMPLATE': ('agent.conversation_compression', 'IDLE_COMPACTION_STATUS_TEMPLATE'),
'PREFLIGHT_COMPRESSION_STATUS_TEMPLATE': ('agent.conversation_compression', 'PREFLIGHT_COMPRESSION_STATUS_TEMPLATE'),
'automatic_compaction_status_message': ('agent.context_engine', 'automatic_compaction_status_message'),
'compression_skipped_due_to_lock': ('agent.conversation_compression', 'compression_skipped_due_to_lock'),
'conversation_history_after_compression': ('agent.conversation_compression', 'conversation_history_after_compression'),
}
def __getattr__(name): # PEP 562 — lazy so no import cycles
target = _PLUGIN_COMPAT_LAZY.get(name)
if target is None:
raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
import importlib
from hermes_cli.plugin_compat import warn_once
warn_once(__name__, name, *target)
return getattr(importlib.import_module(target[0]), target[1])
# ---- END PLUGIN-COMPAT ----