Serving Agents¶
latent serve turns any agent that satisfies the Agent Protocol
into an HTTP service that Agent Studio (or any client) can stream from — over a
dev tunnel while you iterate, or as a deployed container in production. A served
agent registers an addressable version of itself with Studio; Studio routes
traffic to it and falls back to its in-process runtime when no healthy pod is
registered.
Quick start¶
# Serve locally on :8080 (POST /stream, /healthz, /readyz) — no registration.
latent serve my_agent
# Dev loop: open a cloudflared tunnel and register a `dev` version with Studio.
latent serve --share my_agent
# Deployed: self-register a pinned `deployed` version (container entrypoint).
LATENT_AGENT_VERSION=v3 latent serve --version v3 my_agent
Framework flags (--share, --version, --host, --port) come before the
agent name; anything after the agent name is forwarded to the agent's
constructor, exactly like latent chat.
The readiness contract¶
An agent is ready once constructed — BaseAgent.ready() and
FrameworkAdapter.ready() return True by default, so an agent with no external
dependencies (like the echo example) serves with no ready() boilerplate.
Override ready() to gate registration on an external dependency being
reachable:
class MyAgent(ReActAgent):
async def ready(self) -> bool:
# cheap liveness of the agent's own dependencies
return self._chroma.count() >= 0
Keep the check cheap and side-effect-free: latent serve --share/--version
calls it at startup (before the first registration) and on every heartbeat.
Return False (or raise) while a dependency is unreachable and the pod stays out
of Studio's registry, with /readyz returning 503 — so the consuming runtime
falls back to local instead of routing to a broken agent. Plain latent serve
<agent> (no --share/--version) does not gate on readiness.
Registration & fallback¶
--share/--version run a health-gated heartbeat: on each cycle the pod
evaluates ready(), registers (or re-registers) with Studio at ttl/2, and a
circuit breaker guards the loop. On shutdown (Ctrl-C) a --share process
de-registers, so its dev version clears immediately. A deployed pod does not: the
version's row is shared by every replica and by the pod replacing it during a
rolling restart, so one pod's exit must not expire it. When no pod renews the
lease, the row drops out of routing (routing-plan reads live leases only); a pod
that registers again brings it back. Until the lease lapses (300 s), a version
whose last replica stopped still routes to its now-empty Service.
Studio resolves an agent name to its active registered version (preferring live over lapsed leases, deployed over dev, newest over older) and streams to it over an HMAC-signed hop; if no live version is registered, or the pod is unreachable, it falls back to its local runtime. This applies to both the chat webserver route and the Temporal worker.
Configuration¶
| Variable | Purpose |
|---|---|
AGENT_STUDIO_API_KEY |
Project key (lsk_…) used to register. Required for --share/--version — the CLI fails fast if it's unset. Also settable as [agent_studio].api_key in latent.toml. |
AGENT_STUDIO_URL |
Studio base URL to register against (default http://localhost:8000). |
LATENT_SERVE_PUBLIC_URL |
Public URL a deployed pod advertises (required in --version mode). |
LATENT_AGENT_VERSION |
Deployed version label (equivalent to --version). |
LATENT_SERVE_SECRET |
HMAC secret (K). Set a fixed per-deployment value so every replica of a version signs consistently behind a load balancer / across a rolling update; if unset, an ephemeral per-process secret is generated (single-replica only). |
LATENT_SERVE_HOST / LATENT_SERVE_PORT |
Bind address (default 0.0.0.0:8080). |
LATENT_SERVE_MAX_CONCURRENCY |
uvicorn limit_concurrency — concurrent connections/tasks past which the pod sheds load with 503. Unset = unlimited. See Scaling & tuning. |
LATENT_SERVE_KEEP_ALIVE |
uvicorn timeout_keep_alive in seconds (default 5). |
Deploying¶
infra/Dockerfile.agent builds a container whose entrypoint is exec latent
serve "$LATENT_AGENT", exec'd as PID 1 so SIGTERM reaches Python and the pod
drains in-flight turns on shutdown. The deployed version comes from the
LATENT_AGENT_VERSION env var (bound to --version), which flips latent serve
into deployed self-registration — the entrypoint fails fast if it is unset.
Provide LATENT_AGENT, LATENT_AGENT_VERSION, AGENT_STUDIO_API_KEY,
AGENT_STUDIO_URL, LATENT_SERVE_PUBLIC_URL, and (for multi-replica) a fixed
LATENT_SERVE_SECRET.
Scaling & tuning¶
A served agent pod is I/O-bound on the LLM, not CPU-bound: each turn spends its wall-clock time awaiting completion tokens. So the failure mode under load is socket/FD exhaustion and unbounded memory, not a busy CPU — which changes how you size, shed, and autoscale it.
Concurrency knobs (load-shedding)¶
serve() and start_background_server() expose two uvicorn knobs; each falls
back to its env var, then a default:
| Knob (arg) | uvicorn setting | Env | Default |
|---|---|---|---|
max_concurrency |
limit_concurrency |
LATENT_SERVE_MAX_CONCURRENCY |
unset (unlimited) |
keep_alive_s |
timeout_keep_alive |
LATENT_SERVE_KEEP_ALIVE |
5s |
When in-flight connections exceed max_concurrency, uvicorn returns 503
instead of accepting more work. Studio's RemoteAgentProxy maps 503 to
unavailable → fallback (its _UNAVAILABLE_STATUSES), so an overloaded pod
sheds load back to the local runtime rather than accepting requests it can't
serve and OOMing. Set max_concurrency to a value the pod's memory and FD budget
can sustain.
The LLM connection pool is the real wall¶
The concurrency ceiling that actually bites is the outbound LLM HTTP pool,
not CPU. latent calls the LLM through litellm.acompletion
(src/latent/rate_limiter/wrapper.py) and does not build its own httpx
client — the HTTP connection pool is litellm's process-global default. Two levers:
- Application concurrency cap (the knob latent owns). The rate limiter caps
concurrent in-flight completions per
(provider, model)—max_concurrent, default 50. Override withLATENT_RATE_LIMIT_PROFILES_JSON(a JSON map of per-provider profiles). This bounds how many LLM calls run at once and is the primary serve-pod concurrency lever. - httpx pool size (litellm-global). To raise the underlying socket pool above litellm's default, set litellm's global async client once at process start, before serving:
import httpx, litellm
litellm.aclient_session = httpx.AsyncClient(
limits=httpx.Limits(max_connections=200, max_keepalive_connections=100)
)
latent does not rewire litellm for you; set this in your pod's bootstrap if the
default pool throttles throughput. Keep the pool ≥ the rate-limiter
max_concurrent so the app cap, not the socket pool, is the binding limit.
File descriptors (ulimit)¶
Each in-flight turn holds an inbound socket plus one or more outbound LLM
sockets, so a busy pod can exhaust the default 1024 open-file limit and start
failing to open sockets. A non-root process can only lower its hard limit, so
the image can't raise it — the container runtime must. Raise nofile in the
pod spec / HelmRelease (e.g. ~65536), or docker run --ulimit nofile=65536:65536.
See the comment in infra/Dockerfile.agent.
Autoscaling — target in-flight, not CPU¶
GET /metrics serves Prometheus text exposition:
# TYPE latent_serve_in_flight_turns gauge
latent_serve_in_flight_turns 7
# TYPE latent_serve_requests_total counter
latent_serve_requests_total 128
# TYPE latent_serve_responses_total counter
latent_serve_responses_total{status="200"} 120
latent_serve_responses_total{status="401"} 8
latent_serve_in_flight_turns counts the turns currently being served — the
lifetime of each /stream generator plus every in-progress /invoke. Because
CPU stays idle while sockets saturate, a CPU-target HPA never scales this pod in
time. Point the HPA at latent_serve_in_flight_turns (via
prometheus-adapter / KEDA) or at the ingress connection count — a per-replica
target a little under max_concurrency gives headroom before the pod starts
shedding 503s.
Multi-replica anti-replay — the studio owns the shared store¶
Per-request auth rejects a replayed signature using a nonce store. The
default InMemoryNonceStore is per-pod: behind N replicas, a captured
signature replayed to a different replica within the skew window is not caught,
because that replica never saw the nonce.
- Single pod / standalone / dev → in-memory is fine; no configuration needed.
- Multiple replicas → configure the agent studio (
AGENT_STUDIO_URL+AGENT_STUDIO_API_KEY— the same credential the pod registers with). The pod then usesStudioNonceStore, which POSTs each nonce to the studio'sPOST /api/serving/nonceroute ({"nonce", "ttl_seconds"}→{"first_use": bool}). The studio owns Postgres and does the atomic first-use-wins insert, so every replica shares one view — the first request wins, a replay to any replica returnsfirst_use: false. No Redis, no extra dependency: the store rides the studio the pod already talks to on every turn.
build_nonce_store selects StudioNonceStore automatically whenever the studio
is configured, and InMemoryNonceStore otherwise.
Fails open on a studio blip. The studio is already the per-turn routing dependency, so a brief nonce-route outage must degrade latency, not correctness: on any network error, timeout, or non-200 the pod accepts the request (as if first-use) and logs a loud warning, rather than 401-ing legitimate traffic. Anti-replay is degraded only for the duration of the outage.
If you run latent serve --version (deployed) with more than one replica and no
studio configured, anti-replay silently degrades to per-pod — so the deployed
path logs a loud warning when no studio key is present. Configure the studio
to close the hole.
Wire protocols¶
POST /stream negotiates the response codec via the X-Stream-Protocol header —
ai-sdk (Vercel AI SDK v5) or ag-ui. Both are lossless: every AgentEvent
round-trips, so a consumer reconstructs the exact stream the agent produced. See
latent.wire.
Endpoints¶
| Endpoint | Purpose |
|---|---|
POST /stream |
Streaming turn — SSE, encoded with the negotiated wire protocol. |
POST /invoke |
Non-streaming turn — the buffered InvokeResult as one JSON body ({"text", "steps"}), for callers that don't want to consume SSE. Bounded by invoke_timeout_s (default 600s). |
GET /healthz |
Process is alive. |
GET /readyz |
Registered and able to serve (503 until ready). |
GET /metrics |
Prometheus text exposition — in-flight-turns gauge + request/status counters. Autoscale on this, not CPU. See Scaling & tuning. |
/stream and /invoke share one pipeline: size gate → per-request HMAC auth →
decode → correlation-header merge → fresh agent → context install. The HMAC
signature covers the request path, so a signature minted for /stream does not
authorize /invoke — sign each endpoint distinctly.
/invoke timeout — prefer /stream for long turns¶
/invoke is silent until the whole turn completes, so it is exposed to every
hop's idle/read timeout applied to that one window (httpx's default is 5s;
cloudflared/ingress/LBs are typically 30–100s). It carries a server-side ceiling
(invoke_timeout_s, default 600s, mirroring the Claude API's 10-minute
non-streaming timeout); a turn past it is cancelled and returns 504
timeout_error. For turns that can run longer, use /stream and buffer the
result client-side with segment_stream (each streamed event resets the per-hop
read clock — the same reason the Claude SDK's get_final_message() streams
under the hood). /stream is not bounded by invoke_timeout_s.
Injecting services and context¶
A served agent needs two very different kinds of input, and they arrive by two different mechanisms. Getting the split right is the whole game:
| Kind | Example | How it's injected | Lifetime |
|---|---|---|---|
| Services | a Chroma KB / retriever, an HTTP or DB client, playbook data, model config | the factory closure you pass to serve(make_agent) |
built once at process start, shared by every request |
| Typed context (data) | medium, attachments, conversation_id, audience |
the context payload in the request body → self.context |
validated per request, fresh each turn |
The test for which one applies: is it a live object you build once (service → closure), or a small serializable value that changes per turn (data → context)? A Chroma client is the former; medium: "voice" is the latter. A service can't be context — it isn't JSON-serializable, it's expensive to build, and it's the same for every request.
Services — via the factory closure¶
serve(make_agent, …) takes a zero-arg factory called fresh per request. Build
the heavy, shared dependencies once and close over them; the agent that serves
each turn already has them. Using grow's shape (Deps holding Knowledge — the
Chroma retriever, link registry, escalation — and Playbooks):
from grow_support_agent.deps import Deps
from grow_support_agent.playbooks.containers import Knowledge, Playbooks
from latent.serve import serve
# Built ONCE at startup: expensive, shared across every request, not serializable.
knowledge = Knowledge(retriever=load_chroma_kb(...), link_registry=..., escalation=...)
deps = Deps(knowledge=knowledge, playbooks=load_playbooks(...))
def make_agent() -> GrowSupportAgent:
return create_agent(deps) # fresh agent per request, closing over deps
serve(make_agent, name="grow-support", secret=SECRET)
The Chroma KB lives in that closure — it never touches the request. Every request
gets a new agent already wired to deps.retriever, with no per-request setup.
Typed context — via the request body¶
An agent that binds a context_type (a frozen dataclass; see agents)
is populated per request. Put the raw context object under a top-level context
key in the request body:
@dataclass(frozen=True)
class GrowSupportAgentContext:
medium: str = "chat"
attachments: tuple[AttachmentMeta, ...] = ()
class GrowSupportAgent(ReActAgent[GrowSupportAgentContext]):
context_type = GrowSupportAgentContext
async def stream(self, messages, *, config=None):
medium = self.context.medium if self.context else "chat" # typed read
// request body (ai-sdk shape) — data only, never a service object
{ "messages": [/* … */], "context": { "medium": "voice" } }
The server validates the payload into the agent's context_type and sets
self.context on the fresh instance before the turn — the same contract the
Temporal worker path uses, so an agent behaves identically served or in-process.
Concretely, per request the server: decodes the body (the wire codec surfaces the
top-level context key into config), merges correlation headers, builds a fresh
agent, then runs inject_context, which:
- pops
contextout ofconfig(so it never leaks into behaviour config), - constructs
context_type(**raw)and assignsself.context— the construction is the validation: an unknown/missing field or a failing__post_init__(e.g. grow coercing attachment dicts →AttachmentMeta) is the caller's error → 422, - fills a
conversation_idfield the payload omits from theX-Conversation-Idcorrelation header (matching how the worker fillsConversationContext), - leaves an agent with no
context_typeuntouched (a payload is ignored, not an error).
Fields marked studio_field(redact=True) are masked automatically in the emitted
context Metadata event — no server-side handling. (The standalone server
validates by constructing the frozen dataclass + its __post_init__; it does not
depend on agent_studio_shared, so the worker's extra Literal-enforcement isn't
replicated on this path.)
When the service varies per request (multi-tenant)¶
If the KB is per-partner, you still don't put a KB in the payload — you close over
a registry of services and key it off per-request context data. grow does
exactly this (customer_key_kind in the context + a loader closure):
kbs = {"partner_a": kb_a, "partner_b": kb_b} # built once, or a lazy loader
def make_agent() -> GrowSupportAgent:
agent = create_agent(deps)
# a tool/hook resolves the KB from the per-request context key
agent.resolve_kb = lambda: kbs[agent.context.customer_key_kind]
return agent
The service registry is the closure; the per-request selector is context. The service object itself never crosses the wire.
Run it locally¶
Two levels: a two-command smoke test of a pod on its own, and the full studio +
playground with a turn routed to that pod. Both use the echo example agent, so
they are deterministic and need no LLM key.
A. Pod-only smoke test¶
From this repo:
uv sync --extra serve --extra wire
LATENT_SERVE_SECRET=devsecret LATENT_SERVE_PORT=8091 \
uv run latent serve --version v1 echo # framework flags come BEFORE the agent name
# in another shell:
uv run latent verify-wire http://127.0.0.1:8091 --secret devsecret
verify-wire checks SSE framing, tool-input-before-output ordering, the [DONE]
sentinel, a terminal text step, /healthz + /readyz, an unsigned request → 401,
and the fidelity tier. Non-zero exit on any violation.
B. Full stack: studio + playground + a routed turn¶
Requires the three repos checked out as siblings on the serving branches
(agent-studio, latent-py, and an agent-demo whose worker opts into serving
with run_worker(studio_client=…)), Docker, devbox (it provides
process-compose), and uv on PATH.
Three commands from the agent-demo checkout — no admin signup, no key minting,
no hand-set env:
CLI="bun ../agent-studio/packages/create-agent-studio/bin/cli.ts"
# Postgres + Temporal + migrations. Seeds a dev project key + LATENT_SERVE_SECRET
# and writes AGENT_STUDIO_SERVER_URL + BETTER_AUTH_* into the managed .env.
$CLI services up -d
# Frontend + backend + worker via process-compose. --expose binds 0.0.0.0 so the
# browser over 127.0.0.1 reaches Vite.
$CLI app dev --expose # studio → http://127.0.0.1:5173 (use this host, not
# localhost — BETTER_AUTH_URL is 127.0.0.1, so signing
# in on localhost fails with "Invalid origin")
# In another shell: serve one of agent-demo's own agents. It lives in this repo,
# so no --from is needed; export the model provider key the agent uses.
GEMINI_API_KEY=... $CLI serve Agent_QA_Compliance
agent-studio serve reads the seeded AGENT_STUDIO_API_KEY + LATENT_SERVE_SECRET
from the managed .env, maps the studio URL to AGENT_STUDIO_URL, and assigns a
public URL/port — so nothing is minted or hand-set (--version and --port
default to v1 and 8091). It runs latent serve in the current project, so the
agent must be discoverable there; agent-demo ships several (Agent_QA_Compliance
and its placeholder siblings, stratus, latent-docs).
--from <dir> runs latent serve in another checkout — needed only to serve an
agent that lives elsewhere. The keyless echo example lives in latent-py, so to
serve it (no provider key, deterministic): $CLI serve echo --from ../latent-py.
In the browser: /playground lists the pod with a live presence dot; open its
/try/<slug> page and send a turn — it routes studio → worker → RemoteAgentProxy
→ the pod → back (kill the pod and it flips offline; restart and it recovers); the
agent detail "Serving" tab promotes a version; /project-keys creates and revokes
keys.