Skip to content

The RPC protocol

RPCEvent

class RPCEvent(jsonrpc: Literal['2.0'] = '2.0', method: Literal['event'] = 'event', params: dict[str, Any] = dict())

tau_agent_core.rpc.dialect.RPCEvent

A JSON-RPC 2.0 event notification (fire-and-forget).

Events use method="event" and carry an AgentEvent as params. They have no request ID and expect no response.

Constructor parameters

  • jsonrpc: Literal['2.0'] = '2.0' — JSON-RPC protocol version (always "2.0").
  • method: Literal['event'] = 'event' — Always "event" for notifications.
  • params: dict[str, Any] = dict() — The event payload (AgentEvent serialized as dict).

from_json_line

from_json_line(line: str) -> RPCEvent

tau_agent_core.rpc.dialect.RPCEvent.from_json_line

Deserialize from a LF-delimited JSON line.

Parameters

  • line: str — A single LF-delimited JSON string.

Returns

An RPCEvent instance.

to_json_line

to_json_line() -> str

tau_agent_core.rpc.dialect.RPCEvent.to_json_line

Serialize to a single LF-delimited JSON line.

Returns

JSON string suitable for LF-delimited framing.

RPCHandler

class RPCHandler(session: 'AgentSession', *, runtime: 'AgentSessionRuntime | None' = None, output_queue_event_bound: int = DEFAULT_OUTPUT_QUEUE_EVENT_BOUND)

tau_agent_core.rpc.handler.RPCHandler

JSON-RPC 2.0 server over stdin/stdout.

Implements the RPC protocol for τ-agent-core. External tools, custom UIs, and CI/CD pipelines can connect via this handler.

Dispatch (block [3]) is table-driven — see tau_agent_core.rpc.commands .COMMAND_TABLE for the verbs and their tiers/schemas/notes. This class owns lifecycle, framing glue, and the JSON-RPC envelope only; it no longer hardcodes verb names or behaviour.

Reference: docs/REMOTE-CONTROL.md §3/§4/§6 Reference: docs/PHASE-6-SUBPHASE-1.md Reference: SUBPHASE-0.0.md AgentSession interface

Constructor parameters

  • session: 'AgentSession' — The AgentSession to manage.
  • runtime: 'AgentSessionRuntime | None' = None — The session-lifecycle layer (phase 3, H1) backing new_session/fork/switch_session. Optional — None is a legitimate construction for a handler that never needs those three verbs (most of this module's own test suite), and every OTHER verb works identically without one. A host that calls one of the three verbs against a handler built with runtime=None gets a clear RuntimeError (commands._require_runtime), never a silent no-op — production wiring (tau_coding_agent .rpc_mode.run_rpc) always supplies one.
  • output_queue_event_bound: int = DEFAULT_OUTPUT_QUEUE_EVENT_BOUND — T3's bound (see DEFAULT_OUTPUT_QUEUE_EVENT_BOUND module constant for the default and its rationale). A constructor param rather than a bare module constant so a test can drive the bound with a handful of items instead of manufacturing hundreds.

abort_compaction

abort_compaction() -> str | None

tau_agent_core.rpc.handler.RPCHandler.abort_compaction

Deliver abort's signal to the in-flight compaction, if any.

Returns the compaction_id the signal was delivered to, or None when no compaction was in flight. commands._handle_abort puts that value on its response so a host knows to expect a compaction_end with cancelled: true for it.

A SIGNAL, exactly like the AgentSession.abort() beside it: this returns before the compaction has unwound, and whether it actually stopped is reported by that compaction_end, never here (see commands.ABORT_RESULT_SCHEMA's notes on why this verb reports no cursor either — same reason, same shape).

The aborter is cleared here as well as in commands._handle_compact's finally, so a second abort against the same compaction reports None rather than claiming a second delivery: abort stays idempotent, and its answer stays true.

bind_compaction_aborter

bind_compaction_aborter(cancel: Callable[[], None]) -> None

tau_agent_core.rpc.handler.RPCHandler.bind_compaction_aborter

Make the in-flight compaction reachable by abort (finding 5).

Called by commands._handle_compact from inside _acknowledge, i.e. at the instant the acknowledgement carrying compaction_id is enqueued. Not earlier, and the timing is load-bearing rather than incidental: before that instant the compaction is still waiting on D-1's turn_lock, the dispatch coroutine is still awaiting its own acknowledgement future, and cancelling the task there would resolve that future by CANCELLING it — wedging the compact request instead of answering it. It is also the instant the host first learns the compaction exists, so there is nothing it could have aborted before.

Parameters

  • cancel: Callable[[], None] — (no description)

compaction_in_flight

tau_agent_core.rpc.handler.RPCHandler.compaction_in_flight: str | None

No description. This object is marked but undocumented.

exit_code

tau_agent_core.rpc.handler.RPCHandler.exit_code: int | None

No description. This object is marked but undocumented.

output_is_deliverable

tau_agent_core.rpc.handler.RPCHandler.output_is_deliverable: bool

Whether an item put on _output_queue from right now can still reach the host (finding 3, Tier B review).

False means the answer is a definite no: run() started a writer for this handler and that writer has since finished — drained and exited on a clean shutdown, cancelled on SIGTERM (P1) or by the post-EOF flush deadline, or dead of a broken pipe (T6). Nothing else will ever dequeue _output_queue again, so a put_nowait past that point is not a slow delivery, it is a silent drop. A caller that has an outcome to report — commands._handle_compact's _complete is the one today — checks this and says so on stderr (T4) instead, exactly as the cancellation arm beside it already does: an outcome nobody can read is not an outcome, and Fail Early forbids pretending otherwise.

True when _stdout_task is None, and that is not a hedge: it means no writer was ever started for this handler at all — every white-box test in this package's own suite, and any embedded caller driving RPCHandler without run() — in which case _output_queue belongs to whoever constructed the handler and enqueueing is the only correct thing to do. "Not yet started" and "already gone" are genuinely different answers, and this property must not conflate them; conflating them would silence every in-process completion instead of the one that has nowhere to go.

Deliberately NOT the writer's own loop condition (_running or the queue is non-empty): between _running going False and the writer task actually completing, the writer is still draining, and an item enqueued in that gap is still written. This property is about the state after that, which Task.done() is the exact witness for.

prepare_outbound

prepare_outbound(item: dict[str, Any]) -> None

tau_agent_core.rpc.handler.RPCHandler.prepare_outbound

Last-moment fixups on a queued item, applied just before it is framed.

Called by the writer (transport._write_stdout) at DEQUEUE time. The transport calls it without knowing what it does — block [1] frames and orders bytes, it does not know an agent_end from a turn_start (§7.3 X1's separation, which a socket transport would otherwise have to carry). Everything wire-shaped lives on this side of that line.

Today it does exactly one thing: E5's cursor (see below).

Parameters

  • item: dict[str, Any] — (no description)

release_compaction

release_compaction() -> None

tau_agent_core.rpc.handler.RPCHandler.release_compaction

Free the single-flight slot and drop abort's handle on it.

The one place both halves of the in-flight compaction's state are cleared, called from commands._handle_compact's finally so the slot cannot be freed while abort still holds a live handle on a task that has already finished.

run

run() -> None

tau_agent_core.rpc.handler.RPCHandler.run

Run the RPC server until stdin closes or a shutdown signal fires.

Reads JSON-RPC requests from stdin and writes responses/events to stdout (see module docstring for framing). Claims stdout exclusively for the duration of the call (T2) and installs SIGTERM/SIGHUP handlers (P1) so a host can request shutdown out-of-band.

Setup (stdout takeover, signal registration, task creation) lives inside the outer try specifically so the finally still runs — stdout gets released and _running/_stopped_event still get reset — even if a setup step itself raises (e.g. signal registration failing on a platform where it was expected to work, see _register_signal_handlers). Without that, a mid-setup exception would wedge the handler: stdout hijacked forever, and every later run() call raising "already running".

Shutdown ordering is deliberate, and easy to get backwards:

  • We race the reader and writer with asyncio.wait(..., return_when=FIRST_COMPLETED) instead of gathering them unconditionally. Under normal operation the writer never finishes before the reader — its loop condition is "_running is still True, or the queue still has items", and only the reader ever flips _running off — so this reduces to "wait for the reader" exactly as before, avoiding the deadlock a blind gather() would cause (the writer blocked on queue.get() forever after EOF, because nothing else ever ends its loop). What the race adds: if the WRITER finishes first, that can only mean it raised (T6: a broken pipe — the peer is gone). In that case there is nothing left to read for, so we cancel the reader too instead of continuing to parse, execute, and queue requests into a sink nobody drains, then re-raise the write failure.
  • Once the reader is done normally, the background tasks (C3's post-acknowledgement turns and compactions) are reaped BEFORE the writer is drained, not after — finding 3, Tier B review. Reaping them last meant anything they completed with during _cancel_background_tasks's no-cancel grace period was enqueued onto a queue whose writer had already exited: silently lost, on a run that reported rc 0 and an empty stderr. See the call site for the measured case (a compact whose result was durably written and never announced) and for why the finally's reap stays.
  • Then we let the writer drain whatever is queued (P4: a clean shutdown flushes pending output) — unless the shutdown reason was SIGTERM, in which case we cancel the writer outright instead: SIGTERM means the host is already impatient and does not want to wait on us (P1). That drain is bounded by a no-progress deadline (_flush_stdout_with_deadline), because EOF says only that the host closed the end it writes to, not that it is still reading — and an unbounded wait on a peer that stopped reading is a τ process no clean-shutdown trigger can end.
  • A CancelledError raised at the asyncio.wait(...) call itself (as opposed to one carried by the reader/writer tasks it's waiting on) means run()'s own task was cancelled from the outside — e.g. a supervising TaskGroup or gather() tearing down — not _on_signal/stop() cancelling the reader internally. Those two must not be conflated: swallowing this one the same way as the internal case would report a normal, clean completion for a task that was actually cancelled, the standard asyncio footgun. So it stops both tasks and re-raises instead.
  • Cleanup (removing signal handlers, releasing the stdout claim) is unconditional — it runs whether the reader/writer finished cleanly, were cancelled, or raised — via the outer finally. A write failure (T6: broken pipe) is deliberately NOT caught here; it propagates out of run() after cleanup still runs.

stop

stop() -> None

tau_agent_core.rpc.handler.RPCHandler.stop

Request a clean (non-signal) shutdown, and wait for it to finish.

Cancels the stdin reader; run()'s own shutdown path (drain the writer, remove signal handlers, release stdout) takes it from there — identical to what happens on stdin EOF. A no-op if run() was never started.

track_background_task

track_background_task(task: 'asyncio.Task[Any]') -> None

tau_agent_core.rpc.handler.RPCHandler.track_background_task

Hold a strong reference to a background task until it finishes.

Two callers today, both C3's dual completion: commands._submit_and_acknowledge (the post-admission turn a submit/prompt response has already returned for) and commands._handle_compact (the summarization call a compact acknowledgement has already returned for — Blocker 1, Tier B review: that call is bounded only by the provider, so running it inline on the dispatch path stopped transport._read_stdin from PARSING the next line, abort included, for its whole duration).

Without this, nothing holds either task alive once the handler function's local variable goes out of scope, and asyncio is explicit that it may then be garbage-collected before it completes. run()'s teardown reaps whatever is still tracked here (_cancel_background_tasks), so a task registered through this method can never outlive the process — which is the other half of why a background verb must register here rather than fire off a bare create_task.

Parameters

  • task: 'asyncio.Task[Any]' — (no description)

RPCRequest

class RPCRequest(jsonrpc: Literal['2.0'] = '2.0', id: int | None = None, method: str = '', params: dict[str, Any] | None = None)

tau_agent_core.rpc.dialect.RPCRequest

A JSON-RPC 2.0 request message.

Constructor parameters

  • jsonrpc: Literal['2.0'] = '2.0' — JSON-RPC protocol version (always "2.0").
  • id: int | None = None — Request ID (int) for matching responses. None for notifications.
  • method: str = '' — RPC method name ("send_prompt", "send_tool_result", "get_commands", etc.).
  • params: dict[str, Any] | None = None — Method-specific parameters, or None.

from_json_line

from_json_line(line: str) -> RPCRequest

tau_agent_core.rpc.dialect.RPCRequest.from_json_line

Deserialize from a LF-delimited JSON line.

Parameters

  • line: str — A single LF-delimited JSON string.

Returns

An RPCRequest instance.

to_json_line

to_json_line() -> str

tau_agent_core.rpc.dialect.RPCRequest.to_json_line

Serialize to a single LF-delimited JSON line.

Returns

JSON string suitable for LF-delimited framing.

RPCResponse

class RPCResponse(jsonrpc: Literal['2.0'] = '2.0', id: int | None = None, result: dict[str, Any] | None = None, error: dict[str, Any] | None = None)

tau_agent_core.rpc.dialect.RPCResponse

A JSON-RPC 2.0 response message.

Either result or error must be set (never both). For error responses, result is None and error is an error dict. For success responses, error is None and result contains the result.

Constructor parameters

  • jsonrpc: Literal['2.0'] = '2.0' — JSON-RPC protocol version (always "2.0").
  • id: int | None = None — Request ID matching the original request. None for notifications.
  • result: dict[str, Any] | None = None — The response result, or None on error.
  • error: dict[str, Any] | None = None — The error dict on failure, or None on success.

from_json_line

from_json_line(line: str) -> RPCResponse

tau_agent_core.rpc.dialect.RPCResponse.from_json_line

Deserialize from a LF-delimited JSON line.

Parameters

  • line: str — A single LF-delimited JSON string.

Returns

An RPCResponse instance.

is_error

is_error() -> bool

tau_agent_core.rpc.dialect.RPCResponse.is_error

Check if this response represents an error.

to_json_line

to_json_line() -> str

tau_agent_core.rpc.dialect.RPCResponse.to_json_line

Serialize to a single LF-delimited JSON line.

Returns

JSON string suitable for LF-delimited framing.