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) backingnew_session/fork/switch_session. Optional —Noneis 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 withruntime=Nonegets a clearRuntimeError(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 (seeDEFAULT_OUTPUT_QUEUE_EVENT_BOUNDmodule 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 "_runningis still True, or the queue still has items", and only the reader ever flips_runningoff — so this reduces to "wait for the reader" exactly as before, avoiding the deadlock a blindgather()would cause (the writer blocked onqueue.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 (acompactwhose result was durably written and never announced) and for why thefinally'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
CancelledErrorraised at theasyncio.wait(...)call itself (as opposed to one carried by the reader/writer tasks it's waiting on) meansrun()'s own task was cancelled from the outside — e.g. a supervisingTaskGrouporgather()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 ofrun()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.