This document explains how agentic-api is put together: the crate boundaries, the
request lifecycle, and where to make a change for common contribution tasks. It
assumes you've read the crate overview in AGENTS.md and complements it —
AGENTS.md covers tooling and conventions, this document covers the mental model.
Background on why the system is shaped this way lives in the ADRs
(ADR-01, ADR-02,
ADR-03) and the design docs under
docs/design/. Those documents record proposals and drift over time; this document
describes the code as it exists today and will be kept current as the code moves.
Where a design doc's "as-built" notes and the code agree, this document just states
the outcome.
agentic-api/
crates/
agentic-server-core/ # "agentic_core" — pure Rust orchestration library
agentic-server/ # axum HTTP/WS gateway + the `agentic` CLI launcher
agentic-llm-d/ # split-execution state backend for the llm-d coordinator
agentic-praxis/ # placeholder: future Praxis gateway adapter
agentic-server-core(library crate nameagentic_core) is where all domain logic lives: request/response types, SSE parsing, the agentic loop, tool framework, and the storage layer. It has no HTTP framework dependency.agentic-serveris a thin transport layer: an axum binary that parses HTTP/WS, calls intoagentic_core, and streams the result back. It also happens to host a second, unrelated binary — a CLI launcher (agentic) that spawns the gateway and a coding harness (Codex/Claude Code) as subprocesses for local use.agentic-llm-dis a separate axum backend for the llm-d coordinator, which runs inference itself. Its router exposes/v1alpha/responses/hydrateand/v1alpha/responses/persist, plus health/readiness probes. It usesagentic_corefor state services and does not proxy or call a model.agentic-praxisis currently a placeholder. Per ADR-03, the intent is for it to wrap eachagentic-server-corepublic function as anHttpFilterso Praxis can compose the agentic loop declaratively instead of going throughagentic-server's axum router. Nothing is implemented there yet.
The dependency direction is one-way: agentic-server and agentic-llm-d depend on
agentic-server-core, never the reverse. agentic-server-core has no axum dependency;
client-facing HTTP/WS transport belongs to the adapter crates, while upstream HTTP/SSE
I/O lives in core inference transport.
Client ──HTTP/WS──▶ agentic-server (handler/*)
│
▼
agentic_core::executor::ExecuteRequest::run()
│
┌─────────────┼──────────────────────────┐
▼ ▼ ▼
rehydrate() upstream call(s) + tool loop persist()
(storage read) (vLLM, gateway tool execution) (storage write)
│
▼
SSE / JSON back to client
Persistence uses sqlx against a driver-agnostic Any pool backed by SQLite or
Postgres (storage::pool). The upstream inference call targets vLLM's own stateless
Responses API — this project owns the state, vLLM owns tokenization and generation
(see ADR-01 §1.1).
One Responses turn may contain several inference rounds, but it is exposed and persisted as one response:
rehydrate history
│
▼
create AgentPipeline + prepare tool-search state
│
▼
build EngineOrchestration (ToolRegistry + response budget)
│
▼
┌─▶ optional compaction ─▶ one upstream inference round
│ │
│ ▼
│ AgentPipeline ingests JSON or SSE
│ and relays public stream events
│ │
│ ▼
│ resolve calls by ownership
│ │ │
│ │ └─ client-owned calls stay unresolved
│ ▼
│ GatewayScheduler plans and executes calls
│ with bounded fan-out and ordered results/events
│ │
│ ▼
│ classify_round
│ │ │ │
│ │ │ └─ client/incomplete/done: finalize
│ │ └─ append calls/results to continuation input
└──────────────┘ (`Continue`, at most 10 rounds)
│
▼
finalize one public response + persist one turn
A mixed round may contain both ownership classes. Gateway-owned calls still execute and are recorded, while the response returns the unresolved client-owned calls for the client to resolve. Streaming uses the same round loop and projects it through one continuous SSE lifecycle.
The crate produces a library plus two independent binaries. src/lib.rs exports
agentic_cli, agentic_harness, agentic_output, agentic_process, app, auth,
handler, and model_capabilities. On top of that:
| Binary | Entry point | Uses |
|---|---|---|
agentic-server (the gateway) |
src/main.rs |
app, auth, handler, model_capabilities, plus binary-private server.rs and config_file.rs |
agentic (the CLI launcher) |
src/bin/agentic.rs |
agentic_cli, agentic_harness, agentic_output, agentic_process, model_capabilities |
These are two unrelated concerns bundled in one crate. If you're working on request
handling, ignore agentic_cli*/agentic_harness.rs/agentic_output.rs/
agentic_process.rs entirely — they're the launcher that spawns the gateway binary and
a coding harness (Codex or Claude Code) as subprocesses for local, single-command use
(agentic serve <model>), and never touch the request path.
model_capabilities.rs is the one deliberate exception: it is the shared contract both
sides speak. The gateway resolves each model's InputModalities there and serves them in
the Codex catalog; the launcher parses that same catalog back through
CodexCatalogCapabilities before writing an isolated Codex home. Keeping one definition is
what stops the HTTP catalog and a launcher catalog from disagreeing about image support.
app.rs(library) builds the router:AppState(the per-request-shared state:exec_ctx: Arc<ExecutionContext>, proxy state, readiness/websocket trackers, config) andbuild_router_with_auth(state, server_config, authenticator), which wires every route and optionally layers OIDC auth (auth::require_oidc) onto the protected ones.server.rs(binary-private) owns process lifecycle:build_state(constructsExecutionContext::from_config, i.e. where the DB pool actually gets created),serve_gateway/serve_gateway_until_signal(bind, serve, graceful shutdown with a bounded drain), andrun/run_with_llm(standalone mode, optionally spawning a vLLM subprocess).main.rs(binary-private) is theclapCLI front end: parses config from flags/env/config.toml, then callsserver::runorserver::run_with_llm.
| Route | Handler | File |
|---|---|---|
POST /v1/responses |
responses |
handler/http/responses.rs |
POST /v1/responses/compact |
compact_response |
handler/http/responses.rs |
POST /v1/conversations |
create_conversation |
handler/http/conversations.rs |
GET/POST/DELETE /v1/conversations/{id} |
Conversation CRUD | handler/http/conversations.rs |
POST/GET /v1/conversations/{id}/items |
Batch append and ordered listing | handler/http/conversation_items.rs |
GET/DELETE /v1/conversations/{id}/items/{item_id} |
Item retrieval and removal | handler/http/conversation_items.rs |
POST /v1/messages |
messages |
handler/http/messages.rs |
POST /v1/messages/count_tokens |
count_tokens |
handler/http/messages.rs |
GET /v1/models |
models |
handler/http/models.rs |
GET /health |
health |
handler/http/models.rs |
GET /ready |
ready |
handler/http/models.rs |
Handlers make a request-scoped decision between two paths:
- Stateful/executor path — used when the request needs state (
store: true,previous_response_id,conversation_id, compaction, or a gateway-owned tool). Builds anExecuteRequest(Responses) or callsrun_messages_loop/run_messages_stream(Messages) againststate.exec_ctx. - Pass-through path — everything else is forwarded to vLLM unchanged via
agentic_core::proxy, with no state, no persistence.
The HTTP Responses handler first parses RequestPayload<RawValue> so routing fields
remain typed while provider-specific text formats stay opaque on the pass-through
path. Executor routes convert that raw field to ResponseTextConfig before
building ExecuteRequest, preserving strict validation for in-process execution.
GET /v1/responses upgrades to a WebSocket. Structurally this is not a one-shot
handler like the HTTP routes — responses_ws_loop is a long-lived session loop that
reads response.create messages off the socket and drives the same
ExecuteRequest::run() executor call the HTTP handler uses. Requests with distinct
stream_id values run concurrently, while requests in the same lane remain FIFO;
requests without a stream_id share a default FIFO lane. The session admits at most
64 active or queued requests and 12 MiB of aggregate request data. WebSocket sessions
force stream: true and honor the requested store value. Because axum's built-in graceful shutdown
doesn't wait for upgraded connections, AppState carries a separate
WebSocketTracker so shutdown can drain in-flight sessions.
Executor streams propagate downstream backpressure through a bounded event channel.
[responses] configures separate ceilings for upstream JSON bodies, upstream SSE
lines, retained response data, and client stream events. Wire ceilings default to
16 MiB each. Each request shares an 8 MiB retained-data budget across MCP discovery,
upstream rounds, and normalized gateway tool output. Ingestion charges logical retained
items rather than repeated SSE envelopes, so changing delta size does not change the
budget for the same output. These are logical limits, not total process memory bounds. MCP discovery participates in the same 16-permit
materialization window as gateway calls and is capped at 64 server declarations and
128 discovered tools per request.
The WebSocket transport queues only serialized, size-checked events, capped at the
configured max_stream_event_bytes including routing metadata, in its 64-entry outbound
queue. handler/websocket/responses/event.rs owns that envelope and its limit. The
executor receives the effective event ceiling after routing overhead and checks its
terminal event before persistence or session checkpoint publication. Local completion
validates both lifecycle events before persistence. A slow connection can retain up to
64 times the configured event ceiling in its outbound queue, in addition to executor
queues, parsed upstream payloads, and active response state. Authentication is
rechecked at request dispatch so queued work cannot start after identity expiry.
Errors are modeled by a dedicated WsError enum (handler/websocket/error.rs) rather
than reusing the HTTP JSON-error path, since some failure modes (a dead socket) must
not attempt to write a response.
Core callers can use ExecuteRequest::with_session or rehydrate_in_session to
retain response state without durable storage. The WebSocket multiplexer owns one
ResponseSessionGroup per connection and keeps one session per lane, including idle
lanes. No-session HTTP and split execution keep their existing behavior.
The connection retains at most 128 lanes (including the default lane), 32,768 items and 16 MiB per checkpoint, and 32 MiB of aggregate serialized checkpoints. These are retention ceilings, not measured process-memory bounds. New lanes receive an immediate 429 once the lifetime lane limit is reached. Request-count and request-byte overloads also return immediate 429 responses without mutating retained state, so rejected work cannot invalidate an earlier accepted queued continuation. Other validation errors execute in lane order and evict only a matching referenced parent when routing is valid. Parent lookup happens when execution begins; accepted queues do not reserve parent snapshots. Existing fork and execution-failure eviction rules still apply. Disconnect aborts and joins request tasks, waits for active leases to release pinned state, and drops the entire connection group before the close handshake.
A ResponseSession owns one latest canonical checkpoint and one execution slot.
The executor pins a parent before inference and publishes completed or incomplete
state before exposing terminal completion. Failed continuations discard only a
referenced checkpoint owned by that session. Dropping the owner closes the session
and rejects late publication. Callers still cancel and join active work explicitly;
wait_until_idle waits for the execution lease to end, but does not cancel it.
ResponseSessionGroup allows independent serial members to find and pin each
other's latest checkpoints. Failed forks cannot evict the source member's state.
The group bounds lifetime member count, each checkpoint's items and serialized
bytes, and aggregate retained bytes. Shared parent references count once; replaced
parents still pinned by active work and prepared checkpoints awaiting persistence
remain charged until their last reference is released. Reservation happens before
durable writes, and failure or cancellation returns unused capacity. Replacement
requires room for both old and new snapshots until publication. These retention
budgets are not a bound on temporary allocations, execution copies or process
memory. Scheduling, FIFO queues and active-work limits remain caller concerns.
Response-scoped session history records reasoning, messages and calls in inference-round order, then built-in call outputs. Public output is accumulated separately and is not appended to retained history twice. Compaction replaces the canonical window while preserving MCP discovery records needed for orchestration. Durable restoration and replayed compaction input select the effective compacted window before validating current calls; obsolete stored rows are not deleted or charged to that retained window.
For response-scoped session execution, store: false has no durable writes or
database fallback. A stored child of a transient parent persists the complete
canonical window without creating a row or dangling database reference for that
parent. Explicit conversation_id requests keep the existing durable Conversations
policy and append output only through the conversation handler. Their session lease
still serializes execution without recording a second copy. Core commit supports
prewarming without inference; WebSocket generate: false uses that same session
rehydration and commit path.
Transport helpers shared by the HTTP and WS handlers: body reading with a shared size
cap, generic JSON parsing, bearer-token extraction, SSE response wrapping, and
rendering an ExecutorError as a JSON error body.
OIDC bearer-token authentication: discovery, JWKS fetch/cache/refresh, and the
require_oidc axum middleware layered onto protected routes in app.rs. The
WebSocket handler also reads the resulting AuthenticatedPrincipal extension directly,
to detect token expiry mid-session.
Nothing in agentic-server's request-handling code (handler/*, app.rs,
server.rs) imports agentic_core::storage directly. All persistence goes through
AppState.exec_ctx: Arc<ExecutionContext> — e.g. ExecuteRequest::run(),
create_conversation(), persist_turn(), rehydrate_conversation(),
ExecutionContext::storage_ready(). ExecutionContext::from_config is the only place
that constructs the storage handlers, and it does so precisely so callers don't need
to depend on the storage layer:
let conv_handler = ConversationHandler::new(ConversationStore::new(pool.clone()));
let resp_handler = ResponseHandler::new(ResponseStore::new(pool.clone()));Two narrow, deliberate exceptions: the agentic validate CLI subcommand
(src/bin/agentic.rs) calls storage::create_pool_with_schema directly as a
connectivity pre-flight check, outside the request path; and integration tests /
benches under crates/agentic-server/{tests,benches} import ConversationStore/
ResponseStore directly for fixture setup and assertions. Production request-handling
code should never do either.
Per AGENTS.md, the internal dependency direction is: types/ owns
wire/domain data → events/ parses upstream events → tool/ owns tool discovery,
routing, and execution → executor/ orchestrates across inference, tools, and
storage → storage/ owns persistence. Handlers call executor APIs; the executor
coordinates events, tool, and storage; those share contracts through types.
In src/ code, reuse utils::common for JSON serialization/deserialization and
fallback behavior. Do not call serde_json directly when an existing strict,
optional, or defaulting helper expresses the required policy; add a focused helper
there when the policy is reused. Direct serde_json use is fine in tests, fixtures,
and cassette tooling. Keep Serde wire-format attributes on the owning type.
Before adding a conversion helper, wrapper API, or state abstraction, search the owning types and the standard Rust traits for an existing pattern. Use the standard conversion traits and consult the Rust Cookbook for established recipes for common operations. Prefer code whose ownership, failure mode, and lifecycle are visible in its types.
- Implement
From<T> for Ufor an infallible, obvious, and value-preserving consuming conversion. This providesInto<U>automatically; implementFromrather than implementingIntodirectly. - Implement
TryFrom<T> for Uwhen validation or conversion can fail. This providesTryInto<U>automatically and keeps the error type with the destination contract. - Use
AsRef, borrowed accessors, or typed views for cheap borrowed conversions that should not consume or clone the value. - Use iterator adapters such as
map,filter_map,flat_map, andcollectfor collection conversion. Keep the per-item conversion on the item type instead of duplicating a variant match in each caller. - Use newtypes to establish validated identities, indexes, names, and framed values at
construction.
OutputIndex,NonEmptyToolName, andSseLineare examples in this codebase. - Use enums and exhaustive matches for closed state machines,
Option/Resultcombinators for explicit absence and failure, and RAII guards for cleanup or cancellation tied to ownership.
Search types/io/ before introducing a helper named convert_*, into_*, or
append_*. Existing examples include TryFrom<&EventPayload> for typed output items,
From implementations between wire/domain types, and
OutputItem::to_input_item for the deliberately optional continuation conversion.
An inherent conversion method is appropriate when the operation has domain policy
that the standard trait cannot express clearly, as with output variants that are
intentionally omitted from continuation input.
This module's job is JSON ⇄ Rust type conversion and shape validation for the
Responses and Messages APIs. It is not where tool execution, state transitions, or DB
access happen — those live in tool/, executor/, and storage/ respectively.
types/request_response.rs—RequestPayloadis the deserialized incoming request. Itsto_upstream_request(&self, stream: bool) -> Result<UpstreamRequest<'_>, ToolError>is the seam between the OpenAI-shaped request and vLLM's contract. It: flattens Codex namespace tool members to model-visible names, validates every declared tool (ResponsesTool::validate()), and normalizes each supported model-visible tool toUpstreamTool::Function(ResponsesTool::to_function_tools()). File search, code interpreter, and unknown typed declarations currently normalize to no upstream tool; every declaration that does reach vLLM istype: "function", because that's the only tool type it speaks. The conversion also resolves/validatestool_choiceand appliesResponsesInput::model_input(). It's called fromexecutor/upstream.rs'sfetch_blocking_payloadandfetch_stream_payload— the two functions that actually build the outbound request to vLLM.types/io/—input.rs(inbound message/tool-call/tool-result shapes,ResponsesInput),output.rs(outbound output items: messages, function calls, web search/MCP calls, reasoning — plus theApplyDonetrait described below),tools.rs(the normalizedFunctionToolandToolChoice, distinct from tool declarations),usage.rs(token accounting structs).ResponsesInput::model_input()is the final model-visibility boundary used byRequestPayload::to_upstream_request: it removes orchestration-onlyMcpListToolsandCompactionTriggerinput items. A persistedCompactionitem is different: the latest checkpoint supersedes earlier model context and is converted into an assistantoutput_textsummary, while canonical retained user messages and items after the checkpoint remain. This keeps rich continuation state available to orchestration without sending unsupported public item types to vLLM.types/tools/params.rs— the tool declaration shapes a client sends:ResponsesTool(tagged enum:Function,ToolSearch,Mcp,WebSearch,FileSearch,CodeInterpreter,Namespace,Custom,Unknown) and each variant's param struct. This is a good concrete example of the module boundary:ResponsesToolis defined here as a pure shape, but its behavior —validate()andto_function_tools()— is implemented as animpl ResponsesToolblock physically living intool/normalize.rs, which delegates to per-type handlers. Types own the shape; tool owns what it means.types/messages/— a separate, parallel type layer for the Anthropic Messages API (MessagesRequest,ContentBlock, etc.).tool_seam.rsis the pure, I/O-free adapter that converts Anthropic tool blocks into the same internalResponsesTool/FunctionToolCallvocabulary the Responses-sideToolRegistryalready understands, so both APIs share one tool-routing mechanism without the Messages loop depending onRequestPayload/ResponsePayload.types/event.rs— small status enums (ResponseStatus,MessageStatus).
types/io/output.rs::OutputItem::to_input_item is the single semantic conversion
from a response output item to the input item seen by a later inference round. Code
that continues a response must use this method; it must not serialize an OutputItem
as input or introduce another match over OutputItem variants.
current round Vec<OutputItem>
│
▼
filter_map(OutputItem::to_input_item) ← the conversion policy
│
▼
ResponsesInput ← append only
│
▼
ResponsesInput::model_input() ← final upstream-visibility filter
There are two adapters around that one policy:
executor/gateway.rs::append_output_items_to_inputcallsOutputItem::to_input_itemand appends the result to the current request before the next round. Its handling ofResponsesInput::TextversusResponsesInput::Itemsis container mutation, not a second conversion policy.storage/types/item.rs::InOutItem::into_input_itemsuses the same method when rebuilding continuation input from persisted history. Already storedInputItems pass through unchanged.
Messages, reasoning, function calls, custom calls, tool-search calls, compaction
items, and MCP list metadata have explicit continuation representations. Public
web_search_call and mcp_call output items return None: gateway execution records
the canonical model-facing function call and its function_call_output separately,
so converting the public projection would duplicate the call or lose result details.
Gateway tool results are already InputItems and are appended through
append_tool_outputs; they do not need an output-to-input conversion.
This module normalizes raw upstream SSE lines into typed frames, decoupled from the executor so the accumulator doesn't do inline JSON parsing.
types.rs—SSEEventType(the wire event'stype, covering both OpenAI's and vLLM's naming, e.g.response.donevsresponse.completed),EventPayload(the typed, extracted payload — falls back toRaw(Value)for events not deeply parsed yet),WireEvent(the raw pass-through shape, used for re-serialization),EventFrame { event_type, payload, wire }(the normalized output),SSEItemType(output-item kind: reasoning, function call, MCP call, etc.).sse.rs—SseLineandClassifiedSseLine.SseLine::parseperforms only field-level SSE classification (Data,Done, orIgnore) and accepts the optional space afterdata:. Its redactedDebugimplementation reports only payload size.normalize.rs—normalize_sse_data_checkedparses one classified data payload into anEventFramewhile preserving invalid-output-index errors for the ingestion policy.extract_payloaddispatches to small per-eventextract_*helpers. The publicnormalize_sse_lineremains a convenience adapter for callers that do not need policy-aware errors.
To add support for a new SSE event, the touch points are, in order:
events/types.rs— add theSSEEventTypevariant, its wire-string mapping both directions, and (if it carries structured data) anEventPayloadvariant.events/normalize.rs— extendextract_payloadand add anextract_*helper if the payload needs real parsing (otherwise it can fall through toRaw).- Extend
executor/accumulator/'s typed transition matches. Streaming callers enter throughRoundIngestion::push; do not create a caller-specific validator or folding path. - If the event represents a client-executed function shape, extend the corresponding
translator under
executor/translate/. If it is gateway-synthesized, construct the typedEventFrameinexecutor/gateway.rs;pipeline/delivery.rsowns client relay.
This is the layer agentic-server talks to. It owns the request lifecycle: rehydrate,
call inference, run the tool loop, persist. agentic-server never reaches past it.
request.rs—RequestContext(per-turn state: original + enriched request, response/conversation IDs) andExecutionContext(long-lived deps: storage handlers, HTTP client, gateway tool executors, LLM base URL).ExecutionContextis whatAppStateholds; it exposesconv_handler/resp_handler(themodes/handlers below), never the raw stores.rehydrate.rs—rehydrate_conversation()loads prior history from either the conversation store or the response store depending on which ID the request carries, and builds the enrichedRequestContext. Rehydration retains internalInputItem::McpListToolsrecords soToolRegistrycan suppress repeated MCP discovery lifecycle output. After stored history and the new request input are combined,pending_calls.rsvalidates the complete continuation's function/custom call sequence. Every call and call output must have a non-emptycall_id; call IDs must be unique across the sequence; and each output must resolve exactly one currently pending call of the same item kind. An output without a pending call, a second output for an already resolved call, or a function/custom kind mismatch is an invalid request rather than evidence that the call was resolved. Valid unresolved calls remain ordered by their original emission, and the first unresolved client-executed call produces the existing missing-output error. Gateway-executed built-in tool calls are resolved and recorded within their originating round, so they do not remain pending at this boundary.upstream.rs— the narrow adapter between inference transport and the pipeline. It buildsUpstreamRequests, snapshots registry classification facts into an ownedTranslationContext, passes the shared retained-data budget into ingestion, and passes each live JSON or SSE body toAgentPipeline. It does not charge SSE wire bytes.inference.rs—call_inference(): the raw HTTP/SSE transport to vLLM. No parsing beyond splittingdata: ...lines and stopping at[DONE].pipeline.rs,pipeline/—AgentPipeline, the request-owned entry point for both response body formats.RoundIngestionowns the synchronous per-round semantic core;StreamDeliveryowns awaited, ordered client delivery across rounds.engine.rs— the top-level orchestrator:ExecuteRequest/execute(),create_conversation(), andEngineOrchestration, which owns the request-scopedToolRegistry, response budget, and mutableAgentPipelinewhile running the multi-round loop. Its localclassify_round/LoopDecisiondecides whether to loop again, finish, hand back to the client, or return an incomplete response (capped atMAX_GATEWAY_TOOL_ROUNDS = 10). It accumulates output and token usage across rounds, changes continuationtool_choicetoauto, and persists gateway calls plus their outputs as model-facingInputItems. Also home torun_compaction_trigger,run_blocking, andrun_stream(spawns the loop, forwards events as SSE, persists before yielding the terminal event).engine/streaming.rsowns the streaming task's cancellation, failure delivery, and terminal validation before persistence.persist.rs—persist_response/persist_turn, which route toConversationHandlerorResponseHandlerinmodes/depending on whether the turn is conversation-scoped or response-scoped.compaction.rs—compact_response()(the explicit/v1/responses/compactpath) andmaybe_compact_context()(automatic, threshold-triggered, called from the round loop before each inference call).modes/conversation.rs,modes/response.rs—ConversationHandlerandResponseHandler. Thin, 1:1 wrappers aroundstorage::ConversationStore/storage::ResponseStorethat translateRequestContextinto store calls andStorageErrorintoExecutorError. This is the sanctioned boundary between the executor and the storage stores — nothing above this layer touchesstorage::conversation/storage::responsedirectly. Today they only cover what the pipeline needs (get,get_or_create,create,rehydrate[_snapshot],execute_turn,validate_exists); any new CRUD operation beyond persist/rehydrate belongs here, added as a new method that delegates to the corresponding store.error.rs—ExecutorError, with the mapping methods (http_status(),error_type(),into_response_body(), ...) handlers use to render errors.
RFC #241 and issue #243 established one path for JSON and live SSE response bodies:
inference.rs
├─complete JSON───────────────────────┐
└─framed SSE lines─▶ SseLine::parse───┤
▼
AgentPipeline
│
▼
RoundIngestion
│
▼
ResponseAccumulator
(typed response state)
│ validated frame + typed call
▼
TranslationDispatcher
(per-tool translators)
│ translated SSE frames
▼
StreamDelivery
│
▼
GatewayStreamAccumulator
│
▼
client
EngineOrchestration surrounds the pipeline:
inference rounds → gateway execution → terminal policy → persistence
| Stage | Responsibility |
|---|---|
Inference transport (inference.rs) |
HTTP I/O, byte chunks, SSE framing, timeouts, and [DONE] detection. |
Event parsing (events/) |
Classify one SSE line and normalize one data payload into an EventFrame. |
Request pipeline (pipeline.rs) |
Hold one request context, tool-search state, cross-round delivery state, and the JSON/SSE body entry points. |
Round ingestion (pipeline/ingest.rs) |
Process one body with one ResponseAccumulator, one TranslationDispatcher, and final response normalization. |
Stream delivery (pipeline/delivery.rs) |
Provide awaited sender delivery, gateway-event deferral and release, response IDs, and cross-round stream accumulation. |
Orchestration (engine.rs) |
Own the request-scoped registry and response budget, round decisions, gateway execution, terminal policy, and persistence. |
AgentPipeline lives for the complete public response. EngineOrchestration creates
one registry and one response budget around it, then asks upstream.rs to run each
inference body through run_with_json_body or live run_with_stream_body. A new
RoundIngestion is created for every body and consumed by finalization, while
StreamDelivery and GatewayStreamAccumulator survive across inference rounds.
The live runner polls one framed line, performs synchronous ingestion and translation, then awaits delivery before polling the next line. This propagates bounded sender backpressure to the upstream body. Ingestion remains inline; moving it to a worker is the benchmark decision tracked by #245.
The boundary contract is one owner and one path per concern:
- Strict and lenient validation use the same typed accumulator transitions; the policy selects validation/disposition rules and end-of-stream behavior.
- Output-item lifecycle state is scoped to one inference round and keyed by validated
output_index. Item ID and kind must match on later events; active and completed slots remain distinct so index reuse and duplicate completion are rejected. - Translation consumes validated
EventFrames plus the accumulator's typed function call view and produces the public wire lifecycle for client-owned tools. - Delivery consumes translated frames and provides ordered, awaited client emission across inference rounds.
ResponseAccumulator owns validation, response lifecycle, typed output slots, delta
folding, terminal error/incomplete state, usage, and final ResponsePayload assembly.
slot.rs owns active/completed lifecycle and identity bindings, identity.rs extracts
semantic identities, and json.rs contains strict JSON response-shape validation.
active.rs dispatches exhaustively to per-kind state; active_text.rs owns message
parts and reasoning text/summary accounting. Both JSON and SSE ultimately use the same finalization
state.
Output items are constructed through their TryFrom<&EventPayload> implementations in
types/io/output.rs. Active slots fold deltas in place and use the type's ApplyDone
implementation when its completion event arrives. For an already parsed output-item
completion, the private MergeDone implementations in accumulator/completion.rs
merge the concrete item with its retained fields and buffers. slot.rs dispatches
to those implementations after validating identity and lifecycle. Finalization
promotes each completed typed item once and preserves validated output-index order.
The pipeline's streaming entry is process_line(ClassifiedSseLine). It normalizes and
validates a data line, applies the event to the slot keyed by its validated output
index, and returns at most one validated EventFrame for translation. For function
events, accumulated_function_call(output_index) exposes a borrowed typed call and its
folded arguments. finish applies SSE end-of-stream policy; finalize preserves the
status loaded from a complete JSON body.
When adding an output-item kind, extend its typed construction and completion logic in
types/io/output.rs, the variants and exhaustive dispatch in accumulator/active.rs,
and completed-item measurement in executor/response_budget.rs.
Retained-byte accounting stays in synchronous ingestion. response_budget.rs defines
one comprehensive RetainedSize measurement and RetainedAccount for charging growth
and reconciling completed items. Every unbounded collection entry has a structural
charge, including empty JSON values and web-search queries; unrestricted string fields
such as role, content type, and reasoning status count by length. Bounded enums
need no variable charge. Delta text and new part containers are charged before growth.
Completion uses the existing ApplyDone/MergeDone policy and measures its effect;
reasoning text/summary completion measures only the inserted part and its corresponding
streamed counter; shell command completion measures only its command and current buffer.
Both avoid rescanning preceding parts or static nested metadata. Full measurement is reserved
for item completion and final reconciliation. slot.rs charges supplied identities before
index insertion or late binding; IDs shared by a typed item and its indexes count once
logically, including pending kinds with no typed snapshot yet. accumulator/details.rs
charges retained terminal errors and incomplete reasons under the same JSON/SSE measurement.
Charges are cumulative without refunds.
A rejected completion may temporarily allocate its single bounded wire payload before
ingestion drops the failed round.
TranslationDispatcher is a synchronous inline dispatcher. It owns per-call
translation state, classifies each validated function call from an owned
TranslationContext, and routes client-executed tools to their specific
ToolTranslator implementation:
FunctionHandler→FunctionTranslatorCustomHandler→CustomTranslatorShellHandler→ShellTranslatorCodexNamespaceHandler→CodexNamespaceTranslatorToolSearchHandler→ToolSearchTranslator
This is the reverse side of upstream canonicalization. The model emits the canonical
function_call shape. For a client-owned tool, its associated translator restores the
public Responses wire item and its SSE lifecycle (output_item.added, type-specific
delta/done events, and output_item.done). The ordinary function translator preserves
the function-call shape; custom, namespace, and tool-search translators restore their
specific public forms.
The private HasTranslator association lives in the executor so tool handlers do not
depend on SSE or executor types. TranslationContext is an owned snapshot of the
registry facts and request metadata needed for classification and restoration. It
provides each round's tool classification and owns final restoration of public tool
declarations, custom/namespace shapes, tool choice, and tool-search metadata.
Gateway-executed types (Mcp, WebSearch, FileSearch, and CodeInterpreter) are
classified as gateway calls by the dispatcher, which suppresses their canonical
upstream function-call lifecycle. gateway.rs owns their execution result mapping and
synthesizes their public item lifecycle because the gateway knows when execution
starts, completes, or fails. GatewayStreamAccumulator then assigns cross-round
sequence numbers, rebases output indexes, and deduplicates response start events.
Events that arrive before a function name is known are buffered with a 256 KiB total
byte limit and replayed when the call resolves.
Shell functions are restored to shell_call and command SSE events for client
execution by default. When an application registers a shell executor, the owned
translation context carries the resolved gateway ownership so the dispatcher
suppresses the canonical function lifecycle. The gateway event plan then emits
the shell call’s public added/done lifecycle.
StreamDelivery, owned by AgentPipeline, withholds terminal upstream lifecycle
events for the engine, defers frames at and after the first hidden gateway-call
index, and later releases them in output order around synthesized gateway events.
Deferred frames are bounded by both 1024 entries and the existing 256 KiB
serialized-wire-data limit. This is separate from response assembly, retained
session state, and total process memory.
pipeline/delivery.rs owns the shared emission path for both translated upstream
events and gateway-synthesized events. Only the upstream adapter restores response
IDs and applies the round's output offset; gateway events already use public IDs
and absolute indexes. Sequence and response-start deduplication state is committed
after the bounded sender accepts the event. A failed serialization, closed receiver,
or cancelled send cannot consume that state. Enqueueing is not client receipt or
playback acknowledgement. Failed or cancelled frames are discarded by the caller;
this does not make an interrupted pipeline resumable or support retrying an already
rebased frame.
While a gateway-call defer window is active, index-less frames remain deferred
and are released after indexed frames. Flushes keep unsent frames and their byte
accounting in StreamDelivery until the bounded sender accepts each event, so a
cancelled or disconnected flush cannot silently discard the remainder.
When an upstream round fails, the engine releases its deferred public events through
the same delivery path before the terminal event, without executing gateway tools.
Its GatewayStreamAccumulator carries only cross-round presentation state: monotonic
sequence_numbers, public output_index rebasing, and deduplication of response start
events. It holds no OutputItem assembly state. Sending is awaited before ingestion
continues, so sender closure or delivery failure propagates through the live pipeline.
The round-by-round loop belongs to engine.rs::EngineOrchestration. gateway.rs
supplies the planning, execution, projection, and continuation helpers it calls each
round:
GatewayScheduler::plancreates one slot per gateway-owned function call. Each slot owns the original item index, public output index, typedGatewayBinding, and lifecycle projection; a missing executor is represented by an explicit slot rather than omitted from a parallel vector.GatewayScheduler::execute_with_budgetthen returns one orderedGatewayCallResultper slot.futures::future::join_allpolls all planned calls, while atokio::sync::Semaphorecreated by eachexecute_with_budgetinvocation limits active tool executions in that round usingtools.max_concurrent_gateway_calls(default5, configurable throughAGENTIC_MAX_CONCURRENT_GATEWAY_CALLS). This nonzero setting is carried by the owningExecutionContextinto each scheduler; the permits are local to the round, not a process-wide or cross-request limit. Completion may occur out of order, butjoin_allpreserves model call order in the collected results. The permit limit bounds execution, not the number of planned call futures waiting for permits.- A normalized
web_searchfunction call may batch at most five queries. The JSON Schema advertises the ceiling and the handler enforces it again because normalized web search currently uses non-strict arguments. Provider searches acquire a shared handler semaphore initialized fromtools.max_concurrent_gateway_calls, preventing batched calls from multiplying the configured outbound concurrency. Results remain collected in query order for the publicweb_search_call.action.queriesprojection. - Each bound call first acquires its optional same-tool exclusion permit, then a
round execution permit, then a materialization permit shared by cloned scheduler
policies (also used by MCP discovery, with a limit of 16). Waiting for same-tool
exclusion does not consume a round permit; waiting for materialization does. The
independent 60-second timeout wraps
GatewayBinding::executeonly after all permits are acquired. These permit waits do not count toward it, so it is not a deadline for the entire round or total call latency. Timeout, execution, and tool-config failures become failed tool outputs that can be fed back to the model instead of failing the whole response. A tool registered as gateway-owned without an implementation (currently file search/code interpreter) likewise produces an error tool result. - Parallel safety is a per-handler contract.
GatewayExecutor::supports_parallel_executiondefaults tofalse; registration turns that into aGatewayBinding::self_exclusionsemaphore. The semaphore serializes only simultaneous calls to the same model-visible tool name. It never blocks different tools from running concurrently. MCP and web search opt into same-tool parallel execution. - Each scheduler slot retains its
GatewayEventPlan;emit_gateway_start_eventsandemit_gateway_completed_eventssynthesize the OpenAI lifecycle for gateway-executed web search/MCP calls from those same slots. The ordinary path emits all planned start events, executes the round concurrently, then emits ordered completed/failed events. - Streaming may receive client-visible output interleaved with gateway calls. In that
case
engine.rs::execute_and_emit_ordered_output_callstemporarily groups deferred upstream frames byoutput_index, executes the sameGatewaySchedulerconcurrently, and then interleaves synthetic gateway lifecycle events with released upstream frames in original output order. Concurrency and wire ordering are therefore separate concerns. public_output_itemsis the public projection: custom function calls becomecustom_tool_call; gateway-owned internal function calls become their handler'sweb_search_call/mcp_calloutput; client-owned function calls remain function calls. The original gateway function calls andfunction_call_outputresults are retained separately for continuation persistence.
The round decision remains in engine.rs, after gateway execution:
| Decision | Condition and state transition |
|---|---|
RequiresClientAction |
At least one client-owned call exists. Any gateway calls from the mixed round have already executed; their internal calls/results are recorded before returning. |
Done |
No gateway result and no client-owned call remains. Finalize accumulated output and usage. |
Continue |
Gateway calls ran and round budget remains. Append the upstream output plus gateway results, set tool_choice: auto, and infer again. |
Incomplete |
Gateway calls ran on the tenth round. Record the final calls/results and return status: incomplete instead of leaving a dangling call. |
parallel_tool_calls is an upstream model-generation preference, not a gateway
scheduler switch. It is forwarded to vLLM for all supported declaration mixtures and
defaults to false when omitted. Whatever calls the model emits are executed under
the per-round execution permit limit and each handler's same-tool safety policy.
A parallel, independent implementation of the same shape of loop for the Anthropic
Messages API. messages_stream.rs's own header comment describes it as "structurally
the Anthropic-native analogue of GatewayStreamAccumulator, kept deliberately parallel
for a future consolidation" — it never touches RequestPayload/ResponsePayload/
AgentPipeline/ResponseAccumulator/GatewayStreamAccumulator/TranslationDispatcher, operating
directly on Anthropic-shaped JSON. The two loops share only the protocol-neutral
pieces: ToolRegistry::dispatch and types::messages::tool_seam. The round/timeout
constants (MAX_GATEWAY_TOOL_ROUNDS, GATEWAY_TOOL_TIMEOUT) are duplicated and
manually kept in sync with the Responses-side ones rather than shared — a known seam,
not an oversight, per the future-consolidation note.
Both loops take a MessagesRequestContext (messages_context.rs), the per-request
type that replaced a bare serde_json::Value at that boundary. It holds two views of
one request: typed fields for reading tools/stream/model, the current
MessagesToolChoice, and the raw JSON body that is actually forwarded upstream. The raw body is deliberately not
re-serialized from the typed view — ContentBlock catches unmodeled block types in
#[serde(other)] Unknown and models only the fields the gateway reads, so a typed
round-trip would drop cache_control and is_error and collapse image/
redacted_thinking into {"type":"unknown"}. The context owns every mutation the
loops make to that body (force_stream, append_round) and the native web-search
budget, so the two views cannot drift apart uncontrolled; messages and system are
reachable only through the raw body, never the typed view.
append_round also changes a fulfilled forced tool_choice (any, or tool
matching a returned gateway call result) to auto. The first round retains the
client's selector; later rounds can answer from the tool result or choose another
tool. Parallel-use settings and extension fields remain intact. Rounds containing
client-executed function tools return before this mutation, and Messages does not
persist this state.
MessagesToolChoice models auto, any, tool with a non-empty name, and
none as an exhaustive enum. It retains parallel-use settings and flattened
extension fields. The context performs fulfillment checks and transitions using
that enum, then serializes the changed selector into the raw body. Malformed
gateway-tool requests return HTTP 400 before inference; a parse failure cannot
bypass validation by falling back to the proxy. Requests without gateway tools
retain the transparent proxy contract and upstream validation.
vLLM can label a completed, explicitly named tool call end_turn. The shared
request context accepts that stop only when the selected gateway tool appears in
the round. Streaming additionally requires message_stop; client-executed
function tools and truncated rounds remain terminal. A completed client-executed
tool_use is surfaced with public stop_reason: tool_use, correcting vLLM's
end_turn in JSON and the final SSE message_delta. Mixed rounds hide gateway
calls and await the client's output. Token limits, other stop reasons and streams
without a completed round keep their original terminal semantics.
Hide-the-call covers every terminal round, not only a mixed one. A round can end while a
gateway tool_use is present whenever the stop reason is not a tool-call stop — a
max_tokens truncation mid-call, or an end_turn the context does not accept — and that
call is never executed. messages_loop.rs's deliver strips gateway-owned tool_use
blocks from the returned message, matching the streaming accumulator, which suppresses them
on every round. The assistant turn fed back to the model still carries the call.
Each round's usage is folded into messages_usage.rs's MessagesUsageTotals, so when
hidden gateway rounds ran, the returned message (JSON) and the terminal message_delta (SSE)
report the saturating sum of the Anthropic token counters across every round rather than the
final round alone. A streamed round counts its message_start.usage overlaid by its
message_delta.usage. Counters no round reported stay absent, the remaining usage fields
pass through from the final round, a single-round turn is returned unchanged, and a final
round that omits usage still reports the hidden rounds' counters.
-
pool.rs—DbPool = sqlx::Pool<sqlx::Any>, driver-agnostic across SQLite and Postgres.create_pool/create_pool_with_schemaand friends build and tune it (WAL mode + busy-timeout retry on SQLite, statement/lock timeouts on Postgres). -
backend.rs—DatabaseBackend(Postgres/Sqlite/Other) detection from a connection URL, plus URL redaction for safe logging. -
schema.rs— migrations and readiness (PoolWithSchema::ensure_schema_ready), including a path for a supervisor-managed schema that skips running migrations itself and just verifies compatibility. Readiness probes live inschema/readiness.rs. -
conversation/api.rs— tenant-scoped conversation and item CRUD. Item ownership is checked through the conversation so Responses-persisted items are visible too. Item writes sharelock_in_txwith turn persistence; a monotonic conversation revision makes deletions invalidate active snapshots. Removal detaches items rather than destroying history referenced by stored responses. -
models/— rawsqlx::FromRowrow structs per table (Conversation,Item,Response) plus their raw, transaction-aware SQL functions (create_in_tx,get,lock_in_tx, ...). This is the literal DB row shape: JSON columns are still strings here. -
types/— the conversion layer from those raw rows into business types, viaFrom/TryFromimpls:ConversationData/ConversationSnapshot,ResponseData/ResponseMetadata(parses the JSON metadata column into a typed struct),InOutItem(parses anItem.dataJSON blob back into a typedInputItemorOutputItem), andStorageError.InOutItem::into_input_itemsturns a full history into theVec<InputItem>used for continuation processing: storedInputItems pass through, while storedOutputItems go throughOutputItem::to_input_item(). Messages, reasoning, function/custom calls, compaction checkpoints, and MCP list-tools records are retained. Publicweb_search_callandmcp_calloutputs are deliberately omitted because their model-facing function calls and results are already persisted as input items; reconstructing them here would duplicate and lose information from that canonical pair.This conversion is not the model visibility boundary. The resulting enriched history still contains
InputItem::Compaction,InputItem::CompactionTrigger, andInputItem::McpListToolsfor executor/registry decisions. Immediately before an upstream request,RequestPayload::to_upstream_requestcallsResponsesInput::model_input(): the latest compaction checkpoint is converted to an assistant summary and supersedes older context, while compaction triggers and MCP list-tools records are removed. In particular, MCP list-tools remains available long enough for the registry to remember which server labels have already been listed, but it is never serialized to vLLM. -
conversation.rs,response.rs—ConversationStoreandResponseStore: the CRUD-with-transactions layer (create,get,get_or_create,rehydrate[_snapshot],persist/persist_if_version— each transactional, viapool.begin()/tx.commit()). These are not to be called outsideexecutor/andstorage/themselves. The only sanctioned callers areexecutor/modes/conversation.rsandexecutor/modes/response.rs, described above. (Integration tests and benches import them directly for fixtures — that's expected and fine; production code paths should not.)
Wire shapes for tool declarations live in types::tools (see above); this module owns
the behavioral layer — routing, handler traits, normalization, and execution.
Every supported tool, whether client-owned or gateway-owned, implements ToolHandler.
That trait is the authority for validating a public declaration and normalizing it to
the fixed FunctionTool format understood by the upstream model. The outbound path is:
public ResponsesTool declaration
│
▼
RequestPayload::to_upstream_request
│
├─▶ ResponsesTool::validate ──────────▶ ToolHandler::validate
└─▶ ResponsesTool::to_function_tools ─▶ ToolHandler::normalize
│
▼
canonical UpstreamTool::Function
RequestPayload::to_upstream_request is the only request-level seam that prepares
tools for vLLM. New callers must use it rather than rebuilding function schemas or
normalizing declarations in the executor. Declared placeholders that are not yet
supported, currently file search and code interpreter, produce no upstream function
declaration until they have a complete handler and execution path.
| Component | Responsibility |
|---|---|
ToolHandler |
Validate one supported public tool declaration and normalize it into one or more canonical model-visible FunctionTools. |
ToolRegistry |
Provide the request's read-only name lookup for all available tools, including tool type, client/gateway ownership, server label, and any typed gateway binding. |
Client ToolTranslator |
Convert a validated canonical function call back to that client-owned tool's public output shape and SSE lifecycle. |
GatewayExecutor and gateway.rs |
Execute gateway-owned calls and map their start, completion, failure, result, and public output lifecycle. |
GatewayStreamAccumulator |
Project gateway and upstream lifecycle frames into one continuous, correctly indexed and sequenced client stream. |
-
normalize.rs— theimpl ResponsesToolblock withvalidate()andto_function_tools(). These are the declaration-level validation and normalization entry points used byRequestPayload::to_upstream_request. Each supported variant's policy belongs to its correspondingToolHandler:FunctionHandler,ToolSearchHandler,McpHandler,WebSearchHandler,CodexNamespaceHandler, orCustomHandler. Web search's fixed canonical builder is shared withWebSearchHandler::normalize; it remains one schema even though it has no per-declaration normalization state. The method name is plural because namespace and MCP declarations may expand to several model-visible function tools.FileSearch/CodeInterpreterremain unsupported placeholders and normalize to nothing. -
handler.rs— the two traits every tool type reasons about:pub trait ToolHandler: Send + Sync { type ToolParams: Send + Sync; fn tool_type(&self) -> ToolType; fn validate(&self, params: &Self::ToolParams) -> Result<(), ToolError>; fn normalize(&self, params: &Self::ToolParams) -> Vec<FunctionTool>; } pub trait GatewayExecutor: ToolHandler + 'static { type ExecutionParams: Clone + Send + Sync + 'static; fn execute( &self, call_id: &str, tool_name: &str, arguments: &str, params: &Self::ExecutionParams, ) -> Pin<Box<dyn Future<Output = Result<ToolOutput, ToolError>> + Send + '_>>; fn supports_parallel_execution(&self) -> bool; fn plan_gateway_events( &self, call: &FunctionToolCall, params: &Self::ExecutionParams, ) -> GatewayToolEventPlan; fn public_output( &self, call: &FunctionToolCall, output: &ToolOutput, status: GatewayCallStatus, params: &Self::ExecutionParams, ) -> Option<OutputItem>; }
GatewayExecutorrequiresToolHandler: every executable gateway handler supports typed validation and normalization, but not everyToolHandleris gateway-executable. This inheritance is the ownership invariant: supported client and gateway tools share the same declaration contract before execution ownership matters.ToolParamsdescribes the public declaration;ExecutionParamsdescribes one model-visible executable entry. They intentionally differ for MCP: anMcpToolParamdeclares a server, while anMcpDiscoveredToolParamidentifies one tool returned by that server. Gateway-owned registry types may also lack an executor entirely. The trait owns three runtime hooks:supports_parallel_execution()controls same-tool self-exclusion,plan_gateway_events()creates the typed public lifecycle projection, andpublic_output()shapes the completed/failed client-visible item.- Client-owned tools implement
ToolHandlerand have a translator association: seefunction.rs(FunctionHandler),custom.rs(CustomHandler),codex.rs(CodexNamespaceHandler), andtool_search.rs(ToolSearchHandler). Their calls are returned for the client to resolve; the gateway does not execute them. - Gateway-owned / built-in tools implement both traits: see
web_search/mod.rs(WebSearchHandler, backed by the configuredWebSearchProviderinweb_search/you.rsorweb_search/brave.rs) andmcp/handler.rs(McpHandler, backed bymcp/client.rs's MCP protocol client andmcp/pool.rs's connection pool). They have no client translator association because the gateway owns their execution and public lifecycle.
- Client-owned tools implement
-
ownership.rs—ToolOwnership::ClientversusToolOwnership::Gateway(Option<GatewayBinding>). AGatewayBindingcombines the resolved executor, its typedExecutionParams, and the optional same-tool semaphore derived from its parallel-safety declaration. A generic adapter checks the executor/parameter pair at construction and erases only that valid bound pair for heterogeneous registry storage. The scheduler therefore never handles untyped JSON configuration or downcasts. Keeping ownership explicit avoids inferring execution policy from whether a handler happens to be present. -
registry.rs—ToolRegistry, a request-scoped map from model-visible tool name toToolEntry { tool_type, server_label, ownership }. Its responsibility is to keep the read-only catalog of every available model-visible tool for the request, including declared client tools, namespace members, built-ins, and discovered MCP tools. A lookup answers which tool type a name identifies and whether its ownership isClientorGateway; a gateway entry may also carry its typedGatewayBinding. Its constructor,pub async fn build_with_handlers( tools: &mut [ResponsesTool], executors: &mut GatewayExecutors, ) -> Result<Self, ToolError>
is the stable entry point every caller (Responses and Messages) uses to build a registry for a request — its signature should not change. It resolves namespace members, inserts one entry per declared/discovered tool, and for
Mcp/WebSearchpulls the actual executor fromGatewayExecutors(discovering live MCP tools viatools/listin the process).ToolRegistry::dispatch(call)is the per-call routing method the Messages loop uses; the ResponsesGatewaySchedulerresolves the same binding into one call plan so execution, self-exclusion, item position, and lifecycle hooks cannot drift apart. After construction, consumers query this catalog for classification, ownership, and gateway bindings.upstream.rssnapshots its classification facts into a per-roundTranslationContext, which translators use when restoring client-owned tool calls.MCP discovery history is also request-scoped registry state:
mcp_list_tools_items: HashMap<String, Vec<McpListTools>>groups records byserver_label. Registry construction puts the current discovery item first; rehydration appends priorInputItem::McpListToolsrecords only for labels already present in that map.mcp_list_tool_items()exposes entries whose vector still has exactly one element—the current item with no history—to both blocking output assembly and streaming lifecycle emission. Streaming clears the map after the first inference round. Consequently a server's list lifecycle is emitted only when no prior list record exists and never repeats across rounds. -
executors.rs—GatewayExecutors, a shared registry built once at startup and reused across requests, specifically for gateway tools that need lazy, per-request connection setup: MCP servers (connects and cachesMcpClients keyed by server URL, falling back to connecting a fresh request-declared server) and the sharedWebSearchHandler. It also has an optional, application-providedShellExecutorslot.GatewayExecutorRegistration::Shellis an explicit execution grant; an unregistered shell declaration remains client-executed.ShellExecutoraccepts a typed call with bounded action limits and cancellation and returns typed command outputs. The adapter binds into the existing gateway scheduler, not a second tool loop. Client-owned tools (function,custom,namespace) never touch this file; their registry entries are inserted withToolOwnership::Clientand noGatewayExecutorsinvolvement.
Shell item history is preserved publicly in storage. At the inference boundary,
ShellHandler::model_input lowers shell calls and outputs into matching function
history, just as declarations and explicit shell selectors are normalized. For an
opt-in gateway executor, storage additionally retains the canonical internal function
call/output pair; rehydration omits that pair's public shell-call projection to avoid
replaying the invocation twice. Client-executed shell history is not omitted.
To add a new tool type:
- Implement
ToolHandler, including its typedToolParams, for it. - If it's gateway-executed, also declare typed
GatewayExecutor::ExecutionParamsand implementexecute— seeweb_search/mod.rs/mcp/handler.rs. - Wire it into
tool/normalize.rs'svalidate/to_function_toolsmatch arms. - Wire it into
tool/registry.rs'sbuild_with_handlers(aninsert_*_entrycall). - For a client-executed function shape, add its
ToolTranslatorand associate the handler inexecutor/translate/client.rs. Gateway-executed tools do not receive a translator association; their public events come from gateway lifecycle plans. - If it needs lazy per-request connection setup, add a slot to
GatewayExecutorsintool/executors.rsand reference it from the registry's match arm for that type.
Currently a placeholder (src/lib.rs is a comment describing intent). Per ADR-03, this
crate will eventually provide HttpFilter implementations, one per
agentic-server-core public function, composed into a Praxis filter chain with branch
support for tool-call looping — an alternative orchestrator to agentic-server's axum
router, reusing the same core logic in-process.
| Task | Where |
|---|---|
| Add a new HTTP or WebSocket route | agentic-server/src/handler/{http,websocket}/, wire it in app.rs's build_router_with_auth |
| Support a new upstream SSE event | events/types.rs → events/normalize.rs → executor/accumulator/ → executor/translate/ when the event needs public tool-shape translation |
| Add a new tool type | tool/handler.rs impl(s) → tool/normalize.rs → tool/registry.rs → executor/translate/client.rs for client-executed function shapes → tool/executors.rs for lazy gateway setup |
| Change gateway-round concurrency or lifecycle ordering | executor/gateway.rs (GatewayScheduler/event plans) + executor/engine.rs (round decision/ordered streaming) + tool/ownership.rs (typed binding and same-tool safety) |
| Change client streaming order, buffering, or backpressure | executor/pipeline/delivery.rs; keep parsing in events/ and response assembly in executor/accumulator/ |
| Move streaming ingestion to a worker | Benchmark the equivalent inline and worker paths under #245 before changing executor placement |
| Feed response output into the next inference round | types/io/output.rs::OutputItem::to_input_item; use executor/gateway.rs::append_output_items_to_input only to append those converted items |
| Change continuation history visibility | storage/types/item.rs::into_input_items → types/io/output.rs::to_input_item (preservation) → types/io/input.rs::model_input (upstream visibility) |
| Add a CRUD operation beyond persist/rehydrate | executor/modes/conversation.rs or modes/response.rs, backed by storage/conversation.rs / storage/response.rs |
| Change how output items are assembled from a stream | executor/accumulator/ — extend the typed slot and transition matches through the existing RoundIngestion path |
| Add a new Responses/Messages wire field | types/io/ or types/messages/ — shape only, no behavior |
- AGENTS.md — module boundaries, lint/format rules, commit and PR conventions
- TERMINOLOGY.md — normative vocabulary for API/state/tool/streaming concepts
- ROADMAP.md — project direction and near-term focus
- docs/adr/ — architecture decision records
- docs/design/ — as-built design docs (tool framework, core public API, MCP integration, Codex integration)