From 13f6158feb8815f799ca9d06078755c14a82a480 Mon Sep 17 00:00:00 2001 From: Phil Merrell Date: Tue, 29 Sep 2026 10:20:51 -0600 Subject: [PATCH 1/3] docs: map the turn path from send to first token and plan the TTFT work Adds docs/specs/turn-path-ttft.md: every stage of a turn (before the handler, inside it, inside the stream, inside the agent build) with what runs where, what is network, what blocks the event loop and which turn_prelude mark times it; the findings that map exposes (fresh boto3 sessions re-parse service models so warm-up never reaches the build, a full-history ListEvents on the head of every turn just to count it, the serial session_mgr/tools build, agent-turn writes on the critical path, the still-undecomposed agent_build.tools); an assessment of PR #1377; an ordered, measure-first plan that leans toward Runtime V2; and a regression plan a cheaper model can execute. Pointers from CLAUDE.md and the preamble spec. Co-Authored-By: Claude Fable 5.1 --- CLAUDE.MD | 1 + docs/specs/turn-latency-preamble.md | 2 + docs/specs/turn-path-ttft.md | 563 ++++++++++++++++++++++++++++ 3 files changed, 566 insertions(+) create mode 100644 docs/specs/turn-path-ttft.md diff --git a/CLAUDE.MD b/CLAUDE.MD index 2764eed96..2159fbafc 100644 --- a/CLAUDE.MD +++ b/CLAUDE.MD @@ -41,6 +41,7 @@ npx cdk deploy {prefix}-PlatformStack - **Signal-based state** throughout frontend (`signal()`, `computed()`) - **Prompt-cache stability is a contract.** Bedrock prompt caching is exact-prefix-match: any list that reaches the system prompt or `toolConfig` (skills, tools, models, MCP tool listings) must be deterministically ordered at its source, and restored conversation history must be byte-stable between compaction-state changes (see the truncation anchor in `TurnBasedSessionManager`). An order flip or history mutation between turns silently re-writes a 30k–150k-token prefix at the cache-write premium, which is **1.25× the model's own base input rate** (2× at 1h TTL; a cache read is 0.1×) — there is no flat per-MTok figure, so price it against the model actually in play: $1.375/MTok on our default Haiku 4.5, $4.125 on Sonnet 4.6. Our `us.*` ids are **Regional (CRIS)** inference profiles and price ~10% above `global.*` for the same model; rates live in `curated-models.ts`. Per-model AWS **model cards** stay the primary source — they carry context windows, caching support, parameter tables and cutoffs, which no pricing API does — but the **Price List API now corroborates almost all of them**, and a rate worth trusting should agree in both. **Which offer file matters more than which vendor:** `--service-code AmazonBedrockFoundationModels` (372 us-west-2 SKUs) carries Claude, Cohere, Palmyra, TwelveLabs, Luma, Stability; `--service-code AmazonBedrock` (1052) carries Nova, xAI, Google, DeepSeek, Qwen and Moonshot. Querying the wrong one returns nothing and reads exactly like an unpublished model — which is how an earlier revision of this note concluded the API held "10 Claude SKUs, none newer than Claude 3." **That is no longer true** (re-verified 2026-09-21): every Claude model we curate is present and matches its card to the cent — Haiku 4.5 $1.10/$5.50 Regional and $1.00/$5.00 Global, Sonnet 4.6 $3.30/$16.50, Opus 4.7 $5.50/$27.50, Sonnet 5 $2.00/$10.00 Global, Fable 5.1 $10/$50 Global with a cache read of $0.25 that independently confirms its 0.025× multiplier. The 1.100× Regional premium reproduces throughout. **Two traps remain.** The `model` attribute is gone from the product schema — match on `servicename` (FoundationModels) or `usagetype` (AmazonBedrock), or every filter silently returns zero. And the hosted OpenAI family (`us.openai.gpt-5.6-*`, `us.openai.gpt-6-astra`, `openai.gpt-5.4`) is genuinely absent from **both** offer files, so its cards remain the only source; the old explanation that Marketplace billing is uncovered is wrong — Claude's rows are `MP:` usagetypes and are covered. `PROMPT_CACHE_OBSERVABILITY_ENABLED=false` disables the observability layer (fingerprint hook, cacheStatus derivation, EMF metrics) — the caching itself stays on - **Token cost effectiveness is a design tenet — engineer against waste, not against context.** Before merging a change that touches the model call path, answer: what does this add to the prompt, on every turn, for the life of every session? (1) Anything in the cacheable prefix (system prompt, `toolConfig`, restored history) must be deterministic and append-only — see the prompt-cache contract above. (2) Per-turn payloads (tool results, MCP responses, retrieved documents) should be bounded or offloaded, never unbounded pass-through. (3) Don't guess — verify: `cacheStatus` + fingerprint hashes on the session's `C#` rows, `GET /admin/costs/sessions/{id}/calls`, and the `AgentCoreStack/PromptCache` EMF metrics exist to prove a change's cost impact. The balance: "cost effective" means eliminating waste (avoidable cache re-writes, duplicated context, oversized payloads) — it never means stripping context the model needs for a quality answer. When cost and answer quality genuinely conflict, quality wins; look for the cheaper path to the *same* quality, not a cheaper answer +- **The turn path is mapped, stage by stage, in `docs/specs/turn-path-ttft.md`** — every stage from the SPA's send to the first token, with what runs where, what is network, what blocks the event loop, and which `turn_prelude` mark times it. Read it before touching `inference_api/chat/routes.py`, `chat/service.py`, `base_agent.py`, `session_factory.py` or the stream coordinator; the route's phases (`_resolve_turn_attachments`, `_prepare_session_state`, `_check_turn_quota`, `_resolve_turn_model`, `_resolve_effective_tools`, `_build_turn_tools`) are named for that map. It also carries the regression plan for this path (§6): the invariants, the tests per change, and the dev recipes that settle a timing claim - **Time to first token (TTFT) is a hard budget — add no latency before the first token.** That path is everything between a request arriving and the first streamed token: request handling in `inference_api/chat/routes.py`, agent build/cache lookup, MCP pre-flight, and every `BeforeInvocationEvent` / `BeforeModelCallEvent` hook — the latter fires before *every* model call, so its cost repeats on each tool round. Default pattern: a hook captures raw facts only, and anything derived from them is computed after the model answers, memoized (e.g. `get_context_breakdown(agent, itemized=True)` itemizes on read, not in `ContextAttributionHook`). If a change genuinely cannot avoid adding latency there, measure it at realistic scale (dozens of tools, a production-sized system prompt — not a toy fixture), state the number in the PR, and get the developer driving the change to accept it explicitly before merging. A millisecond or two sounds free; the preamble was cut from 455ms to 22–37ms, so each such addition is several percent of that budget, and they compound - **One session can be served by more than one agent — never cache session state on an agent instance.** The agent cache keys on *configuration* (system prompt, tools, model, skills), so an `@`-mention turn builds a second `Agent`, and each `Agent` builds its own `TurnBasedSessionManager`. Both write the same DynamoDB session row and neither knows the other exists, while `initialize()` never re-runs on a cache hit — so anything a manager loads once and holds goes stale silently. This has bitten twice: conversation history (#741, fixed by aliasing the message list in `_adopt_session_conversation`) and compaction state (#751, fixed by re-reading it every turn in `_adopt_persisted_compaction_state`). Per-session state must be aliased across instances or re-read per turn, and must never move backwards — a clobbered checkpoint or truncation anchor is a prompt-cache **cost** bug before it is a correctness one - **All dependencies use exact version pins** — no `^`, `~`, or `>=` diff --git a/docs/specs/turn-latency-preamble.md b/docs/specs/turn-latency-preamble.md index aac0de11e..b61a28d7d 100644 --- a/docs/specs/turn-latency-preamble.md +++ b/docs/specs/turn-latency-preamble.md @@ -1,5 +1,7 @@ # Turn latency: inside the preamble +**Continued by:** `docs/specs/turn-path-ttft.md` (2026-09-29), which maps the whole path from send to first token, assesses PR #1377, and orders what is left — including this spec's unstarted PR-5. + **Status:** COMPLETE for the warm path. Five PRs shipped, merged, deployed and validated on dev (#1184, #1191, #1193, #1198, #1201): diff --git a/docs/specs/turn-path-ttft.md b/docs/specs/turn-path-ttft.md new file mode 100644 index 000000000..0e81fa3ca --- /dev/null +++ b/docs/specs/turn-path-ttft.md @@ -0,0 +1,563 @@ +# The turn path: request to first token + +**Status:** assessment and plan, written 2026-09-29 against `develop` @ `e4646277` +(PR #1378 merged) with PR #1377 (`feature/agent-build-latency`) open. +**Supersedes nothing; it joins three specs that each cover one slice of this path:** +- `docs/specs/turn-latency-preamble.md` — the preamble (455ms → 22–37ms warm) and the + decomposition of `agent_build`. Its PR-5 (split `agent_build.tools`) is still not started. +- `docs/specs/agent-state-feedback.md` — PR-3 deferred the agent build into the stream and + narrates it (`preparing` / `prepared`). +- `docs/specs/agentcore-runtime-v2.md` — the Runtime V2 migration and the prewarm-on-intent + idea (§5a). +**Assesses:** PR #1377, *perf: A/B the first-turn agent build (shared Memory clients, +off-loop build), default off.* + +Labels: **Measured** = a number from dev or prod with a source. **Read** = read off the code +on this date. **Hypothesis** = inferred, not observed. The preamble spec records what +happens when those get confused (its "~53ms per GSI query" turned out to be boto3 client +construction on a throttled CPU), so the labels are load-bearing. + +--- + +## 1. Why this document exists + +A turn is the most complicated process in the app and nobody had one map of it. The handler +in `apis/inference_api/chat/routes.py` is ~2,200 lines in one function with two nested +generators closing over ~40 locals; the build it triggers spans `service.py`, +`base_agent.py`, `chat_agent.py`, `agent_factory.py`, `session_factory.py`, +`turn_based_session_manager.py`, `external_mcp_client.py` and two SDKs; the stream that +follows starts in `stream_coordinator.py` and hands off to Strands. Each spec above opened +one window on it. This document is the whole path, stage by stage, with what runs where, +what is network, what blocks the event loop, and which `turn_prelude` mark (if any) sees it. + +Three things fall out of drawing it: + +1. **The plan** (§5): where the time still is on a cold first turn, ordered so every change + is measured before it is trusted, and leaning toward Runtime V2. +2. **The PR #1377 assessment** (§4): it aims at the right stage, half the fix lands lazily, + and the arm that carries all the complexity buys something other than what it measures. +3. **A regression plan a cheaper model can execute** (§6): invariants that must never move, + named tests per change, and the dev recipes that settle timing claims. + +--- + +## 2. The map + +### 2A. Before the handler (no server-side clock sees any of this) + +| # | Stage | Where | What happens | IO | Cost (warm / cold) | +|---|---|---|---|---|---| +| A1 | Send | SPA | `fetch-event-source` `POST /api/chat/stream` with an `InvocationRequest` body (`session_id`, `message`, `model_id`, `enabled_tools`, `enabled_skills`, `file_upload_ids`, `rag_assistant_id`, …). The SPA mints the session id client-side. | — | — | +| A2 | Edge + BFF | CloudFront → ALB → app-api `chat_stream` (`apis/app_api/chat/proxy_routes.py`) | Cookie auth (`get_current_user_from_session`: BFF session row; user upsert throttled to 300s), CSRF. Reads `session_id` out of the body and sets the **runtime affinity header** `sha256(session_id)` so the turn lands on the conversation's own microVM. Forwards `Authorization: Bearer ` and `OAuth2CallbackUrl`. `httpx` `stream=True`; headers are relayed as soon as inference-api flushes them. | 1 DynamoDB read | **Measured** ~478ms warm for A2+A3 together (agent-state-feedback PR-3) | +| A3 | Data plane | AgentCore Runtime | Routes to the microVM pinned by the affinity header. **No microVM → cold start**: V1 pulls the image and boots the container (**Measured** cold turn 6.7s vs warm 3.75s end to end on 2026-09-19; AWS's V1 curve is 5.4s at 200MB); V2 restores a snapshot (**Asserted** ~2s flat). `/ping` is polled ~every 2s. Only `POST /invocations` and `GET /ping` are proxied. | — | 0 / 1.5–5s | +| A4 | Process start (cold only) | inference-api `main.py` lifespan | Starts `warmup.py` on a daemon thread: imports `strands_tools.calculator` (sympy), the agent factory modules, and builds one client each for `bedrock-runtime`, `bedrock-agentcore`, `dynamodb`, `s3` **on boto3's default session**. `/ping` answers immediately. | none (client construction only) | **Measured** container: `warmup step=boto:bedrock-runtime ms=6` after imports | +| A5 | Middleware | inference-api | Outermost `InvocationActivityMiddleware` (marks the container busy for the whole streamed body), CORS, `AgentCoreContextMiddleware` (workload token, callback URL, session id → contextvars), `StreamSafeGZipMiddleware` (forwards SSE headers immediately; stock GZip withheld them until the first token). Then `get_current_user_trusted` verifies the bearer. | — | ~ms | + +**Two facts about A4 that matter later.** First, the warm-up parses service models into the +*default* boto3 session's loader. **Every fresh `boto3.Session()` re-parses** (**Measured** +on a laptop 2026-09-29: default-session second client 2ms; a client on a fresh session +150–170ms *every time*, for `bedrock-agentcore` and `bedrock-runtime` alike). The AgentCore +Memory SDK, `MemoryClient`, and Strands' `BedrockModel` all construct fresh sessions, so +today's warm-up does not reach them (§3 F2). Second, under V2 whatever the warm-up has +finished before the snapshot is taken is paid once per deploy, not once per session — and +whatever it has not is redone on every restore (`agentcore-runtime-v2.md` B2.2). + +### 2B. Inside `POST /invocations`, before the stream opens + +Every row here runs **before** the `StreamingResponse` returns, i.e. with no channel open +to the client. The `turn_prelude` log line and the `AgentCoreStack/TurnLatency` EMF record +carry one mark per stage; the mark column is the name to filter on. + +Legend for IO: **D** DynamoDB, **S** S3, **M** AgentCore Memory, **B** Bedrock, **C** +tool catalog (DynamoDB through a per-process 60s `config_cache`, or a per-tool 10s +freshness cache), **R** RBAC (per-process TTL cache), **sync** = a synchronous boto3 +call inside `async def`, which blocks the event loop for its round trip. + +| # | Mark | What runs | IO | Warm (**Measured** post-#1198) | +|---|---|---|---|---| +| B1 | `preamble.ownership` | `load_session_meta`: **one** `SessionLookupIndex` query answering both "is this someone else's session" (→ 404) and "the caller's META row" (threaded as `session_meta` into everything below). | D ×1 sync | 6–7ms | +| B2 | `preamble.skills` | `_apply_enabled_skills_filter(await _resolve_accessible_skill_ids(...))` **only when** `enabled_skills` is non-empty (opt-in, D6). | R + skill catalog | ~0 (6–21 with skills) | +| B3 | *(unmarked)* | Model retirement resolve (`resolve_effective_model`, 60s catalog cache). Then the two **non-turn early exits**: `app_tool_call` and `app_context_update` (MCP Apps) build/reuse the agent with `cache_write=False` and return JSON. | C | ~0 | +| B4 | `preamble.files` | `pop_pending_attachments` (snapshot; a write only if a marker exists) → `_select_recovered_attachments` → per-message file cap → `resolve_files` (**one S3 GET per upload id**) → dedupe → `_partition_attachments` (inline / tabular / pptx / oversized) → inline byte budget → marker names. | S ×N | 1ms with no attachments | +| B5 | `preamble.session_state` | `ensure_session_metadata_exists` (a conditional `put_item` on a new session; snapshot short-circuit otherwise) → `clear_paused_turn`, `clear_pending_interrupts`, `clear_truncated_turn`, `clear_interrupted_turn` (all snapshot short-circuits; a write only when a marker exists) → `set_pending_attachments` (write if uploads) → **title task spawned** on a first turn (`generate_conversation_title`, Nova Micro in `to_thread`, persists in the background). | D 0–3 | 3–4ms | +| B6 | `preamble.quota` | `check_quota(session_total_cost=)`: tier (cached), user cost summary (cached), session notice from the snapshot. Early return with a conversational message if exceeded. | cached | 5ms | +| B7 | *(unmarked)* | Retired-model denial (conversational). Model access RBAC check. `_load_user_settings` (**one GetItem**, feeds default model + personal instructions). | R, D ×1 sync | ~6ms | +| B8 | `rag` | **Agent turns only** (`rag_assistant_id`, not resume). In order: `_session_has_messages` (mention turns only); `get_session_metadata` (**a second read of the META row**, outside the PR-2 snapshot); `get_assistant_with_access_check` (assistants table + shares); `resolve_invocation_agent` (version snapshot); `mark_share_as_interacted` (**write**, non-owner); `bump_last_used_at` (**conditional write every turn**, throttled to one real bump/day) + `resume_inactive_policies`; project gate (project row); `resolve_agent_invocation` (bound tools/skills/model/memory, RBAC); **knowledge-base search** (embedding + vector query; the one wait here the answer genuinely needs); prompt composition; memory hydration (`to_thread`); `get_session_metadata` (**third read**) + `store_session_metadata` (**full-row write every bound turn**). Every DynamoDB call in `assistants/service.py` constructs its own `boto3.resource("dynamodb")` (14 sites) rather than using `apis.shared.aws_clients` — the ~47ms-per-construction cost the preamble spec removed from `metadata.py` is still paid here. | D ×4–6 sync, KB | 0 on plain chat; unmeasured on agent turns | +| B9 | *(unmarked)* | Active custom prompt (`resolve_active_prompt_text`: a prompts-table read when the turn or session selects one); personal instructions composed for plain turns. | D 0–1 | ~0 | +| B10 | *(in `tools`)* | **Single-flight lease** `acquire_session_lease` (one conditional update; 409 if held) and, on a resume, `seed_steer_queue`. | D ×1 sync | ~6ms | +| B11 | `tools` | Non-resume: model fallback (`_resolve_fallback_model`: settings → catalog default → RBAC re-check), `_resolve_model_settings` (catalog cache; admin bounds/locks), attachment auto-enable (`_session_has_tabular`: a DynamoDB query unless memoized; RBAC per spreadsheet tool id), always-on union (catalog snapshot), **injected tool builders** (spreadsheet / artifact / word / workspace / excel / powerpoint / account / memory — closures, no IO), `_build_document_tools` (`_session_has_documents`: a query unless memoized), cache-key material. Resume: `get_paused_turn` + an eager `get_agent` from the snapshot. | C, R, D 0–2 | 98–102ms warm, 196–244 cold — **never decomposed** | +| B12 | `stream_setup` | Citations list, the two generators wired, `StreamingResponse` returned. Headers flush here. | — | ~1ms | + +**`agent_build` is not in this table** because with `AGENT_PREPARING_PHASE_ENABLED` +(default on) it runs *inside* the stream (C1 below). With the flag off it runs eagerly +between B11 and B12 and the `preparing`/`prepared` frames are never sent. + +### 2C. Inside the stream, before the first token + +| # | Mark / clock | What runs | IO | Cost | +|---|---|---|---|---| +| C1 | `agent_build.*` | `preparing` frame → `get_agent` → `prepared` frame (`durationMs`) → `turn_prelude` emitted → lease heartbeat task started. **The build is synchronous and runs on the event loop** (`create_agent` is called without `await`), so for its whole duration nothing else in the process runs: the title task cannot deliver, `/ping` cannot answer, and a server-side timer around it cannot fire (agent-state-feedback PR-3 learned this on dev). Sub-stages in §2D. | see 2D | **Measured** 1–49ms warm (cache hit); **~3300ms cold** | +| C2 | *(unmarked)* | `stream_with_quota_warning`: `quota_warning` / `quota_session_notice` / `agent_notice` / `citation` frames; attachment guidance; **`_build_tabular_inventory`** (when spreadsheet tools are enabled: KB file listing + session file listing — DynamoDB/S3 reads); MCP App context drain; recovery / interruption / slash-command notes; `agent.stream_async(...)`. | D/S 0–2 | ~0 unless spreadsheet tools | +| C3 | *(unmarked)* | `ChatAgent.stream_async` → `PromptBuilder.build_prompt` → **`StreamCoordinator.stream_response` head of turn**: env vars, fingerprint/cancel reset, lease stamped, `apply_pending_compaction` (**re-reads compaction state from DynamoDB every turn**, #751), `apply_document_offload`, `reset_stale_interrupt_state`, **`_get_initial_message_count` → `session_manager.list_messages` → `ListEvents` over the whole session history**, just for `len()` (§3 F5), display-text arm, MCP Apps broker subscribe. | D ×1, **M ×1 (full history)** | unmeasured; grows with history | +| C4 | coordinator `stream_start_time` | `agent.stream_async(prompt)` through `process_agent_stream` and the 100ms `_merge_agent_status` poll (which also polls the title). Strands `MessageAddedEvent` → SDK `append_message` (`CreateEvent`, `to_thread`) **and `retrieve_customer_context`**: up to two `RetrieveMemoryRecords` (preferences + facts namespaces) in a `ThreadPoolExecutor`, awaited before the model call, one attempt each with a bounded timeout (`memory_retrieval_timeout_seconds`). `BeforeInvocationEvent` / `BeforeModelCallEvent` hooks (status, attribution capture, fingerprint, ledger, census). `count_tokens` is local (`native_projection=False`). `ConverseStream` opens. | M ×1–3, B | unmeasured server-side except the model call | +| C5 | `first_token_time` | First `contentBlockDelta` with text. The persisted `time_to_first_token` field measures **C4 → C5 only** — it is a *model* TTFT and is blind to everything above it (preamble spec PR-1b). | — | — | +| C6 | — | app-api relays chunks (keepalive comments during silence); the SPA renders. | — | — | + +### 2D. Inside the build (`get_agent` → `create_agent`, a cache miss) + +| Sub-stage | What runs | IO | Cold (**Measured** 2026-09-19) | +|---|---|---|---| +| *(pre-mark)* | `get_freshness_hash(enabled_tools)`: per-tool `updated_at`, 10s TTL, **one `GetItem` per enabled tool on a cold cache**; skills freshness; cache key; hit → return. | C ×N | in `tools` (B11) | +| `prompt` | `ModelConfig.from_params`, `RetryConfig.from_env`, `SystemPromptBuilder`. | — | 160 | +| `registry` | `create_default_registry` (imports; sympy behind `calculator`), `ToolFilter`, `_register_external_mcp_tools`: a `ThreadPoolExecutor` + fresh event loop + `asyncio.run`, **one `get_tool` GetItem per enabled tool** (uncached; `repository.get_tool` is a direct `get_item`). | C ×N | 50 | +| `session_mgr` | `SessionFactory.create_session_manager`: `load_memory_config`; **`_discover_strategy_ids`** (`MemoryClient(region)` → 2 clients on a fresh session, then a control-plane `get_memory_strategies` call; `lru_cache` per process, so **once per fresh process = once per conversation**); retrieval config; `TurnBasedSessionManager(...)` → SDK ctor builds **`MemoryClient()` (fresh session + 2 clients) and then a second fresh session + 2 clients that replace the first pair**; Strands `RepositorySessionManager.__init__` → `read_session` (**1 `ListEvents`, 2 on a new session** via the legacy fallback) → `create_session` (**1 `CreateEvent`** on a new session). | 4 parses, M ×2–3 | 616–830 | +| `tools` | `filter_tools_extended`; gateway: `_expand_gateway_tool_ids` (executor + `asyncio.run`, `get_tool` per gateway id) and `get_gateway_client_if_enabled` (client object; started later by Strands); external MCP: capture contextvars, executor + `asyncio.run` → `load_external_tools`: `get_tool` per server, `create_external_mcp_client`, **`await client.load_tools()` per server, serially** = `MCPClient.start()` (a background thread with its own loop, HTTP `initialize`) + `tools/list`; OAuth pre-flight recovery; `extra_tools` appended. | C ×N, MCP ×servers | **2039** (62% of the build) with **one** server on dev — undecomposed | +| `hooks` | Ten hooks constructed (no IO). | — | 0 | +| `plugins` | Skills runtime (skill records if any), `build_tool_result_offloader`. | 0–1 | 13 | +| `finalize` (#1377 splits it into `strands_agent` + `finalize`) | `AgentFactory.create_agent`: `CountTokensBedrockModel` → Strands `BedrockModel` → **fresh `boto3.Session()` + `bedrock-runtime` client (another parse)**; `Agent(...)`: tool registry `process_tools` (MCP clients register as consumers; `load_tools` is a no-op because the pre-flight primed the cache); `AgentInitializedEvent` → `TurnBasedSessionManager.initialize` → SDK `read_agent` (**1–2 `ListEvents`**, skipped on a new session) → `list_messages` (**1 `ListEvents`, `MAX_FETCH_ALL`**) → converter → document-byte strip → sanitize → compaction (`_load_compaction_state`: **1 DynamoDB GSI read**; possibly `_retrieve_session_summaries`) → pairing repair. New session: `create_agent` `CreateEvent`. Then `_adopt_session_conversation`, snapshot stamping, cache insert. | 1 parse, M ×2–3, D ×1 | 194 | + +Sum of a measured cold build: 3286ms of a 4131ms prelude. Everything in 2D runs +**serially on one thread**, and the two big network stages — the MCP pre-flight and the +Memory session read/restore — do not depend on each other. + +--- + +## 3. What the map exposes + +**F1. The largest cold number is still a black box (Read).** `agent_build.tools` is +2039ms with one external MCP server, and the preamble spec's PR-5 (split it) is not +started. PR #1377 adds `session_mgr_clients` and `strands_agent` marks but leaves `tools` +whole. Nothing below should be built on a guess about what is inside it: it may be the MCP +`initialize` handshake against a cold Lambda, the two thread-pool hops with fresh event +loops, TLS setup, or the per-tool catalog reads. PR-5 comes first. + +**F2. The build parses the same service models four to seven times, and the warm-up +parses them into a session nobody on this path uses (Measured on a laptop, mechanism +Read).** `_discover_strategy_ids` builds a `MemoryClient` (2 parses); the SDK constructor +builds one (2 parses) and throws it away for a second pair (2 parses); `BedrockModel` +builds its own session (1 parse). Each parse is ~150ms on a laptop and client +construction on the container measured ~30× a laptop in the preamble spec. `warmup.py` +builds `bedrock-agentcore` and `bedrock-runtime` clients on the **default** session, which +none of those constructors use. PR #1377's `shared_clients` arm fixes the SDK's pair — +lazily, on the first build, so a first turn still pays two parses plus the strategy-id +call — and does not touch `_discover_strategy_ids` or `BedrockModel`. + +**F3. The build is serial where it could overlap (Read).** `session_mgr` (network: +read/create session; ~0.6–0.8s) and `tools` (network: MCP pre-flight; ~2s) are +independent, and the restore in `finalize` (network: `read_agent` + `list_messages`) is +independent of both. On an agent turn the knowledge-base search (B8) is independent of the +build too. Today each waits for the previous one. PR #1377's off-loop arm moves the whole +serial build to a worker thread; it does not overlap any of these with each other. + +**F4. Agent turns carry writes and duplicate reads on the critical path (Read).** B8 +reads the META row twice more (PR-2's snapshot covers only the preamble), writes it once +in full, does a conditional write for `lastUsedAt` on every turn, and marks shares on +every non-owner turn — all awaited before the stream opens, and each constructing its own +DynamoDB resource. None of them changes what the model is told this turn. + +**F5. Every turn fetches the whole session history a second time to count it (Read).** +`StreamCoordinator._get_initial_message_count` prefers `session_manager.list_messages` +over the maintained `message_count`, "because metadata retrieval uses global indices across +all agents' messages" (voice + text). For a `TurnBasedSessionManager` that is one +`ListEvents` with `MAX_FETCH_ALL` on the head of **every** turn, cache hit or not, and its +cost grows with the conversation. It sits in the gap between `prepared` and the first +`thinking` that the agent-state-feedback epic observed and could not explain. + +**F6. The clock stops at the agent (Read).** `turn_prelude` ends at C1. C2–C4 (inventory, +head-of-turn compaction re-read, the history count, LTM retrieval, the hooks) are between +`prepared` and `thinking` and no server-side number covers them; the persisted +`time_to_first_token` starts after them. A first token that moved because of C3 would be +invisible to every dashboard we have. + +**F7. Per-tool catalog reads are paid once per build per tool, twice (Read).** +`get_freshness_hash` and `_register_external_mcp_tools` each do a `get_item` per enabled +tool (the first on a 10s TTL, the second uncached); `load_external_tools` and +`_expand_gateway_tool_ids` read again per server. The whole catalog is already in the 60s +`config_cache`. Small per read with a cached client, but it is N×4 round trips per cold +build and every one is `sync` inside `async`. + +**F8. Three sync→async bridges, each a thread and a fresh event loop (Read).** +`_register_external_mcp_tools`, `_expand_gateway_tool_ids` and the external MCP load each +spin a `ThreadPoolExecutor` to `asyncio.run` a coroutine, because the constructor is +synchronous and the loop is busy. Correct, and the reason the build cannot be `await`ed +piecemeal today. Any restructuring of the build should collapse these into one deliberate +boundary rather than three incidental ones. + +**F9. What V2 changes about all of the above (Read from the V2 spec; Hypothesis on +timing).** Under V2 the cold start (A3) drops to ~2s flat and *every* first turn is a +restore. Process-start work (A4) is paid once per snapshot. So: anything static that is +computed at start (strategy ids, parsed service models, imported modules) becomes free per +session if it finishes before the snapshot; anything that opens a connection at start is a +dead socket in every restore (§B2.3 of the V2 spec); anything time-bounded pre-populated at +start (TTL caches on a monotonic clock) is a stale read after restore. The idle-clock +origin in `runtime_health.py` is the known hazard. The prewarm-on-intent idea (§5a) is the +one lever that hides the build itself, and it is V2-only for cost reasons. + +**F10. The route handler is where the ambiguity lives (Read).** Most of the "how does a +turn work" questions above were answered by reading a 2,200-line function whose phases are +separated by comments and whose two generators close over its locals. The path is *not* +badly designed — every stage has a reason and most have a spec — but it is not legible. +§7 starts on that. + +--- + +## 4. PR #1377, assessed + +**What it does well.** It measured on dev before proposing, found the real mechanism (the +SDK's double client construction), noticed that every conversation is its own Runtime +process (so laptop warm-process numbers do not transfer), wrote the decision rule before +the data, and shipped everything **off by default** behind one flag with an interleaved, +session-stratified A/B and a script that joins client and server timings. That is exactly +the discipline the preamble spec asked for. The instrumentation (`session_mgr_clients`, +`strands_agent`, `buildArm`, `processBuilds`, the experiment script) has no cost on the +control arm and should merge regardless of the arms' fate. + +**Where it lands short, by component.** + +*Shared clients (`memory_shared_clients_enabled`).* Right target (F2), incomplete fix: + +- The shared session is created **lazily by the first build**, so the first turn of every + conversation — the only turn the A/B is about — still pays two parses plus the + strategy-id round trip. The session, its two clients, and `_discover_strategy_ids` belong + in `warmup.py`, on the same daemon thread, so a first turn finds them built. On V2 that + also puts them in the snapshot. (Build the *clients*; make no calls except the one + strategy-id read, which is static config — see P2 for the V2 caveat.) +- It leaves `_discover_strategy_ids`'s own `MemoryClient` and `BedrockModel`'s fresh + session alone. Both take a `boto3_session` / `boto_session` argument today (verified in + the pinned SDKs), so the same shared session closes them with no monkeypatching. +- The mechanism — rebinding `MemoryClient` inside the SDK module, gated by a contextvar + set for the duration of one constructor — is the only seam the SDK offers for the + throwaway client, and the PR pins it with a test that fails if an upgrade stops going + through it. Acceptable, but it is a workaround for a one-line upstream bug + (`MemoryClient(region_name=region_name, boto3_session=boto_session)`). File that + upstream; the rebinding should have an end date. + +*Off-loop build (`agent_build_off_loop_enabled`).* This arm carries all of the PR's risk +(the integration lock, the singleton lock, consumer pins, process-wide build +serialization) and its stated payoff is the title landing before `prepared`. Two things +to be clear about: + +- **Moving a synchronous build to a thread does not make it faster.** CPU-bound parts + contend on the GIL; network-bound parts take the same time. By itself the arm cannot + reduce TTFT, and the A/B's TTFT comparison will most likely read as noise. +- **What it actually buys is concurrency on the loop during the build**: the title + (measured), `/ping` (never measured — today a 3s build blocks the health poll; nothing + has failed because of it, and V2's platform may or may not be as tolerant), and — the + valuable one — the ability to start an agent turn's knowledge-base search (and anything + else independent) *before* the build and have it finish *during* it. The A/B as designed + measures the title only, and on plain chat turns, so the arm's real value is not on the + scorecard. + +**Recommendation.** + +1. Merge the instrumentation and the `shared_clients` arm now, amended so the shared + session, its clients and the strategy ids are built at warm-up (P2). Keep the arm behind + the flag until the dev A/B reads; the decision rule in the PR stands. +2. Keep the off-loop arm in the A/B, but decide it on the right question. Either extend the + experiment to agent turns with a knowledge base and start the KB search ahead of the + build (P3b, small), or accept the PR's own rule and expect to remove the arm and its + MCP hardening. Do not ship the hardening for a title. +3. Do not let #1377 be the vehicle for F1/F3. Split `tools` first (P1); the overlap that + actually shortens a cold build (P3a) is a different, smaller change inside the + constructor, and it does not need any of the off-loop machinery. + +--- + +## 5. The plan + +Ordered so every change has a measured before and after. Each item names its V2 posture +because the migration is coming and nothing here should have to be undone for it. + +### P0 — Land the measurement in #1377, run the A/B, decide `shared_clients` + +Unchanged from the PR: `AGENT_BUILD_EXPERIMENT=ab` on the dev Runtime out of band, +`experiment_agent_build_arms.py --per-arm 15 --cleanup`, the decision rule in the spec. +Amend the arm per §4 before the run, or the first-turn number will understate it. + +### P1 — Close the two measurement gaps (F1, F6) + +*P1a. Split `agent_build.tools`* (the preamble spec's PR-5) into `tools.filter`, +`tools.catalog` (the per-tool reads), `tools.gateway`, `tools.mcp_preflight` and +`tools.extra`, with `groups.agent_build` still summing. Marks are taken on the calling +thread around each executor hop, never inside it (contextvars do not cross the pool). Add a +per-server `ms` list as a log property of the MCP stage so one slow server is +distinguishable from many. + +*P1b. Extend the clock to the first token.* Add `head_of_turn` (C3: compaction re-read, +history count, offload) and `pre_model` (C4: LTM retrieval + hooks, up to the model call) +marks, recorded by the coordinator on the same `TurnPrelude` (thread it through +`stream_async` like `turn_started_at` already is), and stamp `firstTokenMs` from handler +entry on the emitted line. `PreludeTotalMs` stays what it is; the new number is +`FirstTokenMs`. This is what makes P4/P5 provable. + +*Tests.* `test_turn_timing.py`: new marks render, group sums hold, metric names derive. +Cost: perf_counter calls only. V2: none. + +### P2 — Warm the right things, once (F2) + +- One process-wide `boto3.Session` in `apis/shared/aws_clients.py` (it already owns the + client cache); warm-up builds `bedrock-agentcore`, `bedrock-agentcore-control` and + `bedrock-runtime` clients on **it**. +- `SessionFactory` hands that session to the SDK (`boto_session=`), to + `_discover_strategy_ids` (`MemoryClient(boto3_session=)`), and `ModelConfig.to_bedrock_config` + passes it to `BedrockModel` (`boto_session=`; drop `region_name` when doing so — Strands + rejects both). The `_SharedSessionMemoryClient` rebinding from #1377 covers the SDK's + throwaway client until upstream takes the one-line fix. +- Warm-up calls `_discover_strategy_ids` once. It is a control-plane read of static + configuration; on V2 that puts the ids in the snapshot. **Caveat:** it opens a + connection. Either accept that the pool may hold one idle socket at snapshot time + (botocore retries connection errors on this idempotent call), or take the ids from the + environment instead — the memory construct could export them, but `CfnMemory` exposes no + strategy-id attributes today, so that is a CDK custom resource, not a one-liner. Start + with the warm-up call and watch the first restore's logs. +- Expected: `session_mgr` drops by the four parses it no longer does; `finalize` by one. + Laptop arithmetic says ~0.6–0.9s of a fresh process's first build; the container ratio is + unknown until measured. Warm turns: nothing (the agent is cached). + +*Tests.* Warm-up builds each client once on the shared session (patch `boto3.Session`); +the factory passes the shared session (assert on the SDK ctor kwargs); `BedrockModel` +receives `boto_session` and no `region_name`; `reset_cached_clients` also resets the +session (the moto trap). V2: **positive** — all of it is snapshot-friendly, except the one +connection, which is the thing to watch. + +### P3 — Overlap the build's independent IO (F3) — the real cold-turn lever + +*P3a. Inside the constructor, no async plumbing.* In `BaseAgent.__init__`, `session_mgr` +and `tools` do not depend on each other (hooks, which take the session manager, are built +after both). Run them on two threads of one `ThreadPoolExecutor` and join — the same +mechanism the MCP load already uses, applied one level up. Tool order stays deterministic +(one thread builds one list). Marks: take `agent_build.session_mgr_ms` and +`agent_build.tools_ms` as properties from each thread's own clock and one wall-time stage +`agent_build.parallel` on the calling thread, so the sum-of-stages invariant the +unaccounted-time widget depends on still holds. Expected on the measured cold build: +`max(2039, 830)` instead of `2039 + 830`, roughly **−0.8s**. Then, if P1a shows the MCP +pre-flight is the bulk of `tools`, the per-server loads run on the pool too, merged in +catalog order (the preamble spec's warning: the *merge* must be deterministic because tool +order is the prompt-cache prefix). + +The restore (`initialize`, inside `Agent(...)`) can join the overlap later by prefetching +its three `ListEvents` on a third thread and handing the SDK the results for the one +`initialize` that follows; that is a second step, after P3a is measured. + +*P3b. On agent turns, start the knowledge-base search before the build.* The build closure +does not read `context_chunks`; only the stream does, when it composes `final_message`. +Submit the search to a thread before the deferred build starts and await its future in +`stream_with_quota_warning`. This is the overlap #1377's off-loop arm would enable "for +free"; without that arm it needs the search to be on a thread, which it can be. Keep +citations and augmentation byte-identical (order of chunks, cap) — they land in the +persisted user message, which is the cacheable prefix on every later turn. + +*Tests.* P3a: a fake session factory and a fake tool loader that each sleep; assert the +build's wall time is the max not the sum, that `tools` order is unchanged across runs +(`test_prompt_cache_determinism.py` already pins order), and that a failure in either thread +surfaces as the build's exception with the other thread joined (no orphaned MCP client: +`consumer_pin` semantics from #1377 apply here too). P3b: the KB search is started before +`preparing` is emitted and its result reaches `final_message` unchanged; a failed search +still yields the unaugmented message (today's fail-open). Dev: fetch-tee harness, first +turn with an MCP tool enabled and a KB agent; compare `agent_build` group and +`FirstTokenMs` before/after. V2: neutral — threads, no sockets at start. + +### P4 — Take the bookkeeping off the critical path (F4, F5, F7) + +- **`_get_initial_message_count`:** use the maintained `message_count` and drop the + `ListEvents`. The global-index concern it cites is for mixed voice+text sessions; verify + with `get_messages_from_cloud` whether the index still needs to be global (voice writes + under a different agent id) and, if it does, count once at `initialize` and keep the + count maintained rather than refetching the history every turn. **Measure C3 first + (P1b)** — this is a claim about a cost that has never been timed. +- **B8 writes:** `mark_share_as_interacted`, `bump_last_used_at` + `resume_inactive_policies`, + and the binding persistence `store_session_metadata` become fire-and-forget tasks (strong + references held, like `_pending_title_writes`) or move to the coordinator's post-stream + metadata write. Validation stays where it is. The binding write in particular must still + land before the *next* turn's validation reads it; post-stream is early enough. +- **B8 reads:** thread `session_meta` (PR-2's snapshot) into the assistant block's two + `get_session_metadata` calls; convert `assistants/service.py`'s 14 resource constructions + to `get_dynamodb_table`. +- **Catalog reads:** `get_tool` served from the `config_cache` snapshot inside its TTL, + falling through to `get_item` on a miss. Removes N×4 round trips per cold build. + +*Tests.* Each deferred write still happens exactly once per turn and after `done` at the +latest (patch the repository, assert call order against the frames); a turn that errors +before the stream still binds; the snapshot-threaded reads behave as PR-2's did +(`test_a_snapshot_with_no_row_is_an_answer_not_a_cache_miss` pattern). V2: neutral. + +### P5 — Runtime V2 track (F9) + +In the order the V2 spec gives, with two additions from this map: + +1. Deploy-script fix (B1: `platformVersion` in `ALLOWED`, post-update assert, CLI check). +2. Explicit `PlatformVersion` flag in CDK, always set. +3. Idle-clock re-stamp in `runtime_health.py` (B2.1) — plus **an audit of every + `time.monotonic()`-keyed cache reachable at start** (tool freshness, `config_cache`, + `oauth_token_cache`, RBAC): none may be pre-populated by warm-up. P2 deliberately warms + clients and strategy ids, which are not time-bounded. +4. Dev A/B with `FirstTokenMs` (P1b) and the client-side send→first-delta gap, split cold + vs warm; `get-agent-runtime` after `platform.yml` **and** after the next `backend.yml`. +5. Decide warm-up placement from the restore logs: if warm-up work lands in the snapshot, + hold `/ping` until it completes so every restore inherits it. +6. **Prewarm on intent** (V2 spec §5a): a no-op `/invocations` action that takes no lease, + writes no rows, charges no quota and makes no model call, fired from the SPA on composer + focus with the staged session id. This is the only change that hides A3 + A4 + the first + build from the user, and it is worth more than everything in P2–P4 on a cold turn. A + later step pre-builds the agent for the current selection; keep it separate (selections + change before send). + +### P6 — Keep the room clean + +- `count_tokens` stays local; the LTM retrieval stays bounded and one-attempt; the 100ms + status poll stays. None of them are on this list because each was already measured and + chosen. +- Every hook on `BeforeModelCallEvent` keeps capturing raw facts only (CLAUDE.md). +- Nothing in P2–P5 touches the prompt, `toolConfig` or restored history. If a change here + ever needs to, it is a different spec. + +--- + +## 6. Regression plan — for hand-off + +Written so a less expensive model can run it without this document's context. Three +parts: invariants that must never move, the checklist per plan item, and the dev recipes +that settle timing claims. Run everything from the worktree with the main checkout's +interpreter: + +```bash +cd /backend +PYTHONPATH=$PWD/src AWS_EC2_METADATA_DISABLED=true \ +
/backend/.venv/bin/python -m pytest tests/ -q +``` + +`PYTHONPATH` is not optional: without it the interpreter imports `develop`'s code from the +main checkout and a green run proves nothing about the branch. Verify once with +`python -c "import apis.shared; print(apis.shared.__file__)"`. The suite's socket guard +fails any test that reaches the network; a new test that "passes" while hitting AWS is a +test whose assertion never ran. + +### 6A. Invariants (assert before and after every change on this path) + +| # | Invariant | How to check | +|---|---|---| +| I1 | **SSE frame order.** Head: `quota_warning?` → `quota_session_notice?` → `agent_notice?` → `citation*` → `message_start`. Deferred build: `agent_status{preparing}` first frame, `agent_status{prepared, durationMs}` before `message_start`, always paired. Tail: `message_stop` → interrupt events → `metadata` → `compaction?` → `done`. `session_title` may appear anywhere, at most once, never `"New Conversation"`. | `tests/routes/test_inference.py`, `tests/apis/inference_api/test_preparing_phase.py`, `TestRouteContract` (string pins in `routes.py`) | +| I2 | **Prompt-cache byte stability.** Same session, same configuration, consecutive turns: `toolConfigHash` and `systemPromptHash` identical; `historyHash` extends. Tool order in `toolConfig` is deterministic across builds. | `tests/agents/main_agent/test_prompt_cache_determinism.py`; on dev, `GET /admin/costs/sessions/{id}/calls` — consecutive rows with `cacheStatus=hit` and equal `toolConfigHash` | +| I3 | **One cache key per configuration, shared by all three callers.** The main turn, the resume (`is_resume=True`, `cache_write=False`, `memory_binding=snapshot.memory_binding`) and the MCP App dispatch (`cache_write=False`) compute the same key for the same session state. | `tests/apis/inference_api/test_get_agent_call_sites.py` (AST-pinned), `test_chat_service.py::test_resume_replay_from_snapshot_hits_same_cache_slot` | +| I4 | **Lease discipline.** Acquired exactly once per turn before the stream; released on every exit: stream end, cancellation, pre-stream `HTTPException`, pre-stream error; heartbeat cancelled first. | `tests/apis/inference_api/test_turn_lease_release.py`, `test_carried_steering.py` | +| I5 | **No session state on an agent instance.** `turn_lease`, `cancelled`, display-text arm, fingerprints, interrupt state are reset at the head of every turn; conversation list is aliased across instances; compaction state is re-read per turn. | `test_chat_service.py::test_second_cache_key_for_a_session_shares_the_conversation`, `tests/agents/main_agent/session/test_compaction_deferred_apply.py`, `test_turn_based_session_manager.py` | +| I6 | **Early exits stream, they do not 4xx.** Quota exceeded, retired model without successor, blocked agent binding, archived project: a conversational `message_start … done` with the metadata event, persisted via `persist_synthetic_messages`. Only ownership (404), model access (403), resume validation (400), duplicate turn (409), missing assistant (404/403) raise. | `tests/routes/test_inference.py`, `test_model_retirement_routes.py`, `test_project_harness_invocation.py` | +| I7 | **Attachment semantics.** CSV/XLSX never inline; PPTX never inline; oversized dropped with a note; per-message cap and byte budget applied before S3 reads; marker names follow attachment order; recovered attachments dropped when the turn carries its own. | `test_attachment_turn_guard.py` (37), `test_attachment_recovery.py`, `test_presentation_attachment_carveout.py`, `test_attachment_tool_autoenable.py` | +| I8 | **`turn_prelude` shape.** One line per agent turn with `stages`, `groups` (`preamble`, `agent_build`), `totalMs`, `isResume`, `deferredBuild`, `hasAssistant`; EMF metric names derived from stage names; disabled by `TURN_LATENCY_METRICS_ENABLED=false`. | `test_turn_timing.py` | +| I9 | **No network in tests.** | the socket guard in `tests/conftest.py` (fails at teardown) | +| I10 | **Warm-up never opens a connection** (until P2 deliberately adds the one strategy-id read). | `test_warmup.py` (`boto3.client` patched, assert no call methods invoked) | + +### 6B. Per-change checklist + +For each plan item, the cheaper model should: (1) run the named tests before touching +anything and record the count; (2) make the change; (3) run the same tests plus the new +ones; (4) run the full suite once; (5) run the dev recipe and paste the numbers into the +PR with the *before* numbers from the same recipe on the same day. + +| Item | Unit tests to add | Existing tests to run | Dev recipe | Accept if | +|---|---|---|---|---| +| P1a split `tools` | marks appear in order; `groups.agent_build` sums; per-server ms list is a property | `test_turn_timing.py`, `test_preparing_phase.py` | R1 on a first turn with one MCP tool enabled | sub-stages sum to `tools` ±5ms; a named owner of ≥60% of it | +| P1b first-token clock | `head_of_turn`, `pre_model`, `firstTokenMs` on the line; `FirstTokenMs` metric | `test_turn_timing.py`, `tests/agents/main_agent/streaming/` | R1 + R2 on a warm and a cold turn | `firstTokenMs` ≈ client send→first-delta minus the A2/A3 hop | +| P2 shared session at warm-up | see §5 P2 | `test_warmup.py`, `test_session_factory*.py`, `tests/agents/main_agent/core/` | R1 on three first turns per arm | `session_mgr` + `finalize` down by ≥ the parses removed (≥300ms cold); warm unchanged | +| P3a parallel constructor | see §5 P3a | `test_prompt_cache_determinism.py`, `test_base_agent_external_registration.py`, `tests/agents/main_agent/integrations/` | R1, R3 | `groups.agent_build` cold ≈ max not sum; I2 holds across 5 turns | +| P3b KB search ahead of build | see §5 P3b | `tests/routes/test_inference.py` (assistant cases), `test_project_harness_invocation.py` | R1 on a KB agent's first turn | `rag` mark ≈ 0 and `FirstTokenMs` down by ≈ the old `rag` | +| P4 bookkeeping off the path | see §5 P4 | `tests/routes/test_assistants.py`, `tests/routes/test_sessions.py`, `test_agent_binding_scoped_tools_e2e.py` | R1 + R4 | `head_of_turn` and `rag` down; every write still observed in DynamoDB after the turn | +| P5 V2 | per the V2 spec | `test_runtime_health.py`, infra `jest` | R2 cold vs warm, `get-agent-runtime` twice | no 424s on first turns; `FirstTokenMs` cold down | + +### 6C. Dev recipes + +**R1 — server-side stages.** Log group from SSM `//inference-api/runtime-id` +(the group is `/aws/bedrock-agentcore/runtimes/-DEFAULT`; the prefix-named +group returns zero rows, not an error). Then: + +```bash +aws logs filter-log-events --log-group-name --filter-pattern turn_prelude \ + --profile dev-ai --start-time --query 'events[].message' --output text +``` + +Lines are wrapped in OTEL JSON with escaped quotes. Confirm the image tag in +`//inference-api/image-tag` changed before trusting a number as "after". A window +that spans a deploy mixes two populations: read min or individual lines, never the mean. + +**R2 — client-side first token.** In the in-app browser signed in to dev, wrap +`window.fetch` to `tee()` the `/chat/stream` body and timestamp frames matching +`agent_status|session_title|content_block_delta|done` from the moment of the request (the +recipe in the agent-state-feedback and session-title memories). The gap between request +and the first `content_block_delta` is the number the user feels; subtract `PreludeTotalMs` +(or `FirstTokenMs` after P1b) from the same turn to get the A2/A3 hop. + +**R3 — cache stability.** `GET /admin/costs/sessions/{id}/calls` for the session; on +turns 2..n, `cacheStatus` should be `hit` with equal `toolConfigHash` / `systemPromptHash`. +`partial_miss` or a changed hash after a build change means tool order or prompt bytes +moved. + +**R4 — the writes still land.** After an agent turn: the session's META row carries +`preferences.assistantId`; the assistant's `lastUsedAt` is today (first turn of the day); +the share row is marked interacted (non-owner). One `get-item` each. + +**R5 — the A/B.** `backend/scripts/experiment_agent_build_arms.py` (from #1377) with +`AGENT_BUILD_EXPERIMENT=ab` set on the dev Runtime out of band. Pass env vars explicitly; +zsh does not word-split `$var`. + +### 6D. The end-to-end matrix (run once per release on dev) + +For each row: send the turn, watch R2's frames, confirm the row's "expect", then reload the +page and confirm the conversation restores with the same content. + +| Turn shape | Expect | +|---|---| +| First turn, plain chat | `preparing`/`prepared`, `session_title` before or during the answer, title persisted (R4) | +| Second turn, same session | no `preparing` label visible (build 0–40ms), `cacheStatus=hit` | +| Agent with KB, first turn | `citation` frames before `message_start`; answer uses the corpus; binding persisted | +| `@`-mention of an Agent in a plain thread | that turn runs the Agent; the next plain turn sees the mention's messages (#741) | +| Attach a PDF | inline; `[Attached files: …]` marker; `document_read` appears in tool list; card survives reload | +| Attach a CSV | diverted; Spreadsheet Analysis auto-enabled (granted role); guidance note; card survives reload | +| Attach a PPTX | diverted; PowerPoint tools note; card survives reload | +| Enable a skill and `/invoke` it | `` present; the directive line; the `skills` tool called | +| External MCP tool needing OAuth | `oauth_required` after `message_stop`; consent; resume finishes the same turn (I3) | +| `ask_user_question` | `user_question_required`; answer resumes; "Skip" resumes | +| Stop mid-answer | partial persisted; `interrupted_turn` marker; next turn carries the interruption note; no 409 on resend | +| Steer mid-turn (a tool turn) | `steering_applied` at the tool boundary, or the follow-up sent as the next turn | +| Continue after `max_tokens` | the answer continues, does not restart; Agent tools/skills intact | +| Quota exceeded (test tier) | conversational message, `quota_exceeded` event, persisted, no model call | +| Preview session (`preview-…`) | works; nothing in the sidebar; no lease | +| Duplicate send while streaming | 409 → "already streaming" (the SPA queues instead) | + +--- + +## 7. Clarity refactors + +The route's phases are now named functions with typed results, in the order the map +lists them, so the handler reads as the sequence §2B describes. Behaviour-preserving by +construction: each extraction moves a contiguous block, keeps every module-level seam the +tests patch (`_load_user_settings`, `_resolve_fallback_model`, `get_agent`, +`ensure_session_metadata_exists`, …) and every string the `TestRouteContract` and +call-site AST tests pin. See the PR that introduced this document for the list. + +Next, in this order, each its own PR: + +1. The assistant block (B8) into `resolve_agent_turn(...)` returning either a refusal to + stream or the resolved instructions/overrides/citations, with its writes moved per P4. +2. The build's three sync→async bridges (F8) into one `run_blocking_async(coro)` helper + with the contextvar capture in one place. +3. `stream_with_quota_warning` and `_guarded_stream` into a `TurnStream` object that owns + the head frames, the build, the lease heartbeat and the title poll, so the handler ends + at "return the stream". + +Rules that keep this legible afterwards: a phase function takes what it reads and returns +what it produces (no closure over handler locals); a phase that can end the turn returns +the `StreamingResponse` or raises, it does not set a flag; every new stage gets a mark; and +a comment that explains *why* stays with the code it explains, not in this document. + +--- + +## 8. Open questions + +- What is inside `agent_build.tools` (P1a answers it). +- What C3 costs on a long conversation (P1b answers it); whether the global message index + still needs a per-turn `ListEvents` for mixed voice+text sessions. +- The container's ratio for service-model parsing (laptop ~150ms each). +- Whether AgentCore Runtime tolerates `/ping` being blocked for the length of a cold build + on V2 as it evidently does on V1. +- Whether a `CfnMemory` custom resource is worth it to make strategy ids environment + configuration instead of a start-up read. From 9d51d5186cd628570de2c9ef81db553d46b7bfc2 Mon Sep 17 00:00:00 2001 From: Phil Merrell Date: Tue, 29 Sep 2026 10:20:52 -0600 Subject: [PATCH 2/3] refactor(inference-api): extract the invocation route's turn phases Behaviour-preserving. The invocations handler's contiguous phases become named module-level functions with typed results, in the order the turn-path map lists them: _resolve_turn_attachments, _prepare_session_state, _check_turn_quota, _resolve_turn_model, _resolve_effective_tools, _build_turn_tools, _agent_for_app_dispatch (the two MCP App dispatches shared 60 identical lines), _validate_resume_interrupts, _build_citations, and _refuse_turn / _forbidden_turn for the six conversational early exits. The handler drops from ~2,240 to ~1,640 lines and its docstring is now the phase index keyed to the turn_prelude marks. Every module-level seam the tests patch is kept; the one AST pin on the main get_agent call now reads the binding off turn_tools. Four comments that still described the pre-PR-2 preamble (including the disproved ~53ms-per-GSI-query claim) are corrected. Full backend suite: 10950 passed, 3 skipped. Co-Authored-By: Claude Fable 5.1 --- backend/src/apis/inference_api/chat/routes.py | 1704 ++++++++++------- .../test_get_agent_call_sites.py | 2 +- 2 files changed, 997 insertions(+), 709 deletions(-) diff --git a/backend/src/apis/inference_api/chat/routes.py b/backend/src/apis/inference_api/chat/routes.py index 44a414f42..b5919b1bc 100644 --- a/backend/src/apis/inference_api/chat/routes.py +++ b/backend/src/apis/inference_api/chat/routes.py @@ -11,7 +11,8 @@ import json import logging from collections import OrderedDict -from typing import TYPE_CHECKING, AsyncGenerator, Optional, Tuple, Union +from dataclasses import dataclass, field +from typing import TYPE_CHECKING, Any, AsyncGenerator, Optional, Tuple, Union from fastapi import APIRouter, Depends, HTTPException, status from fastapi.responses import JSONResponse, StreamingResponse @@ -1986,283 +1987,172 @@ def _build_skill_invocation_note(skill_slugs: list[str]) -> str: ) -@router.post("/invocations") -async def invocations(request: InvocationRequest, current_user: User = Depends(get_current_user_trusted)): + +# ============================================================ +# Turn phases +# ============================================================ +# +# `invocations` runs a turn as the sequence docs/specs/turn-path-ttft.md §2B +# lays out. Each function below is one contiguous stretch of that sequence, +# named for what it settles and for the `turn_prelude` mark that times it. A +# phase takes what it reads and returns what it produces; none closes over the +# handler's locals. That is what lets the handler read as the map, and lets a +# change to one phase be reasoned about — and tested — without the rest. + + +def _sse_headers(session_id: str) -> dict: + """Response headers for every SSE stream this route returns. + + ``X-Accel-Buffering: no`` defeats proxy buffering so events after + ``message_stop`` (``oauth_required`` and friends) reach the browser at once. """ - AgentCore Runtime standard invocation endpoint (required) + return {"Cache-Control": "no-cache", "X-Accel-Buffering": "no", "X-Session-ID": session_id} - Supports user-specific tool filtering and SSE streaming. - Creates/caches agent instance per session + tool configuration. - Uses the authenticated user's ID from the JWT token. - Quota enforcement (when enabled via ENABLE_QUOTA_ENFORCEMENT=true): - - Checks user quota before processing - - Streams quota_exceeded as assistant message if quota exceeded (better UX) - - Injects quota_warning event into stream if approaching limit +def _refuse_turn( + input_data: InvocationRequest, + user_id: str, + *, + message: str, + stop_reason: str, + metadata_event: Union[QuotaExceededEvent, ConversationalErrorEvent, None], +) -> StreamingResponse: + """End the turn before any model call, as one short assistant message. + + Errors stream as assistant messages (CLAUDE.md): a quota block, a retired + model, a blocked Agent binding or an archived project each become a + persisted assistant turn carrying its metadata event, rather than an HTTP + status the SPA would render as a generic failure. """ - input_data = request - user_id = current_user.user_id - auth_token = current_user.raw_token + return StreamingResponse( + stream_conversational_message( + message=message, + stop_reason=stop_reason, + metadata_event=metadata_event, + session_id=input_data.session_id, + user_id=user_id, + user_input=input_data.message, + ), + media_type="text/event-stream", + headers=_sse_headers(input_data.session_id), + ) - # Where the pre-stream time goes. Everything between here and the - # `StreamingResponse` return happens with NO channel open to the client — - # measured at 3.75s on a warm turn — so this is the only way to see which - # stage owns it. Pure timing: nothing reaches the model. - # See `turn_timing.py` and docs/specs/agent-state-feedback.md. - prelude = TurnPrelude() - # Whether this turn's agent is built inside the stream (PR-3). Recorded on - # the `turn_prelude` line so the two shapes stay distinguishable in the - # logs once the flag has been on for a while. - deferred_build = False - # Refuse a turn against a session id another user already owns. - # - # Session ids travel in shareable URLs (`/s/{sessionId}`). Opening someone - # else's link 404s on the metadata read, but the SPA then treats the - # session as new and lets the user send — which used to fork the id: a - # SECOND metadata row under the requester, on the same session, invisible - # to both parties. In prod on 2026-08-31 that also left the original - # owner's session resolving non-deterministically between the two rows. - # - # Not a confidentiality fix — conversation content is keyed by actor id in - # AgentCore Memory, so the second user only ever saw an empty thread. This - # stops the id from being forked at all. 404 rather than 403 so the - # response says nothing about whether the session exists, matching what - # `GET /sessions/{id}/metadata` already returns for the same case. - # - # ONE read of the session's META row, shared by everything in the preamble - # that used to fetch it again (PR-2, docs/specs/turn-latency-preamble.md). - # Measured on dev: eight separate reads of this item cost ~445ms of a - # ~455ms stage, because a GSI query from an AgentCore Runtime container is - # ~53ms rather than the ~12ms an in-region figure would suggest. - # - # Deliberately explicit rather than a per-request memo inside - # `_get_session_by_gsi`: CLAUDE.md's "never cache session state" rule has - # been paid for twice (#741, #751), and a snapshot callers opt into cannot - # leak into one that needs a fresh read. - session_meta = await load_session_meta(input_data.session_id, user_id) - if session_meta.owned_by_other: - logger.warning( - "Rejected invocation for session %s — owned by a different user", - _sanitize_log(input_data.session_id), - ) - raise HTTPException(status_code=404, detail="Session not found") - # First of the preamble's five sub-stages (docs/specs/turn-latency-preamble.md). - # The coarse `preamble` number survives as `groups.preamble` in the emitted - # line, so the four-turn baseline in the agent-state-feedback spec stays - # comparable across this split. - prelude.mark("preamble.ownership") - # Resume requests reuse the cached agent and its paused interrupt state; - # they bypass quota, file resolution, and RAG augmentation because those - # already ran on the original turn that got paused. - is_resume = bool(input_data.interrupt_responses) - # Resolve the effective agent type: the client's explicit choice, else the - # compiled-in default ("chat"). Used for the skill resolution below and the - # non-resume get_agent calls (resume reuses the snapshot's type). An Agent - # that binds skills is coerced to "skill" later by the agent-binding - # resolver. - effective_agent_type = input_data.agent_type or DEFAULT_AGENT_TYPE - # Skills feature deferred for this environment: neutralize the legacy - # "skill" agent type, which is a ChatAgent alias since v2 PR-2. Voice and - # other agent types pass through untouched. - if not skills_enabled() and effective_agent_type == "skill": - effective_agent_type = "chat" - # Resolve the user's *effective* skills once for the whole request: the - # accessible set (catalog ∪ own), narrowed by the client's per-turn - # enabled_skills selection. Threaded into every get_agent call below so they - # share one skills_hash cache key (otherwise the app-tool-call / resume paths - # would miss the main turn's cached agent). - # - # Skills v2: this is no longer gated on agent_type == "skill". Skills are a - # plain-chat capability now — the picker in model settings sends - # enabled_skills on an ordinary turn and ChatAgent mounts the AgentSkills - # plugin. The opt-in default (D6) is what keeps this cheap: an absent or - # empty selection short-circuits to [] without touching RBAC or the skill - # table, so every turn that doesn't ask for skills costs exactly what it did - # before. An Agent's skill bindings override this further down. - effective_skill_ids = None - if skills_enabled() and input_data.enabled_skills: - effective_skill_ids = _apply_enabled_skills_filter( - await _resolve_accessible_skill_ids(current_user), - input_data.enabled_skills, +def _forbidden_turn(input_data: InvocationRequest, user_id: str, message: str) -> StreamingResponse: + """A refusal the user cannot recover by retrying (``recoverable=False``).""" + return _refuse_turn( + input_data, + user_id, + message=message, + stop_reason="error", + metadata_event=ConversationalErrorEvent( + code=ErrorCode.FORBIDDEN, message=message, recoverable=False + ), + ) + + +async def _agent_for_app_dispatch( + input_data: InvocationRequest, + current_user: User, + user_id: str, + auth_token: Optional[str], + *, + agent_type: str, + accessible_skill_ids: Optional[list], +): + """The conversation's agent for an MCP App dispatch (map stage B3). + + An App's ``tools/call`` and ``ui/update-model-context`` run no model turn, + but they need the same ``Agent`` the real turns use: the MCP client session, + its auth, and ``agent.state``, where pushed context lives. So the model, its + settings and the effective tool list are resolved exactly as the main turn + resolves them, and ``get_agent`` reads the same cache slot with + ``cache_write=False`` — this path builds no injected tools, so an agent it + built must never seed a slot the real turns would then hit. + """ + request_inference_params = dict(input_data.inference_params or {}) + dispatch_model_id, dispatch_provider = input_data.model_id, input_data.provider + if not dispatch_model_id: + dispatch_model_id, dispatch_provider = await _resolve_fallback_model( + user_id, current_user, dispatch_provider ) - # Near-zero on a turn that selects no skills — the opt-in default (D6) - # short-circuits before touching RBAC or the skill table. A non-trivial - # number here means the RBAC cache missed or the owner-index query is slow. - prelude.mark("preamble.skills") - # A "Continue" after a max_tokens truncation. Like resume, it bypasses - # quota / RAG / file resolution and does NOT clear the turn state; unlike - # resume there is no interrupt to validate — the agent is rebuilt from the - # resent params and re-entered with an empty prompt (assistant-prefill). - is_continuation = bool(input_data.continue_truncated) - # Marketplace D11: the Agent was `@`-mentioned in the composer, so it runs - # this turn only — it does not bind the conversation. Only meaningful - # alongside `rag_assistant_id`; on its own it does nothing. - is_agent_mention = bool(input_data.agent_mention) and bool(input_data.rag_assistant_id) - logger.info( - "Invocation request received (resume=%s, continue_truncated=%s, agent_mention=%s)" - % (is_resume, is_continuation, is_agent_mention) + caching_enabled, inference_params, mantle_api_mode, mantle_region, registry_provider = await _resolve_model_settings( + model_id=dispatch_model_id, + explicit_caching_enabled=input_data.caching_enabled, + request_inference_params=request_inference_params, + ) + return await get_agent( + session_id=input_data.session_id, + user_id=user_id, + auth_token=auth_token, + # Same auto-enable seam as the main turn, so a spreadsheet + # session's dispatch reads the slot the real turns fill — and + # the same always-on union, or the dispatch would compute a + # different effective list and miss into its own agent-cache + # slot on every App call. + enabled_tools=await _apply_admin_always_on_tools( + await _apply_attachment_tool_autoenable( + input_data.enabled_tools, current_user, input_data.session_id, user_id + ), + current_user, + ), + model_id=dispatch_model_id, + system_prompt=await _plain_turn_prompt(input_data, user_id), + caching_enabled=caching_enabled, + provider=dispatch_provider or registry_provider, + inference_params=inference_params, + mantle_api_mode=mantle_api_mode, + mantle_region=mantle_region, + agent_type=agent_type, + is_resume=False, + accessible_skill_ids=accessible_skill_ids, + # This path builds no injected tools, but shares a cache slot + # with the real turns that do. Read the slot; never seed it. + cache_write=False, + assistant_id=input_data.rag_assistant_id, ) - logger.info("Message received") - # Model retirement (docs/specs/model-retirement.md §7). Resolved before anything - # builds an agent from ``input_data.model_id``, so the App tool-call / context / - # continuation paths land in the same agent-cache slot as the turn itself. A - # redirect swaps the provider as well: the request's described the retired - # model, and a successor on another transport misroutes with it. A denial is - # streamed at the access check below, once there is a turn to answer. - retired_model_denial: Optional[str] = None - requested_model = await resolve_effective_model(input_data.model_id) - if requested_model is not None: - if requested_model.denied: - retired_model_denial = retired_model_message(requested_model.retired) - elif requested_model.redirected: - input_data.model_id = requested_model.model_id - input_data.provider = requested_model.provider - # App-initiated tools/call (MCP Apps PR #5). Like resume/continuation it - # bypasses quota / RAG / file resolution / title — there is no model - # turn. We rebuild the conversation agent (so the MCP client session + - # auth are wired exactly as for a model-driven call), dispatch the one - # named tool, publish synthesized tool_use/tool_result into the thread - # via the per-session broker, and return the CallToolResult as JSON for - # app-api to relay back to the iframe. Inert behind the host flag (the - # UIToolCatalog is empty, so dispatch rejects every call as not - # app-visible). - if input_data.app_tool_call is not None: - atc = input_data.app_tool_call - try: - request_inference_params = dict(input_data.inference_params or {}) - dispatch_model_id, dispatch_provider = input_data.model_id, input_data.provider - if not dispatch_model_id: - dispatch_model_id, dispatch_provider = await _resolve_fallback_model( - user_id, current_user, dispatch_provider - ) - caching_enabled, inference_params, mantle_api_mode, mantle_region, registry_provider = await _resolve_model_settings( - model_id=dispatch_model_id, - explicit_caching_enabled=input_data.caching_enabled, - request_inference_params=request_inference_params, - ) - agent = await get_agent( - session_id=input_data.session_id, - user_id=user_id, - auth_token=auth_token, - # Same auto-enable seam as the main turn, so a spreadsheet - # session's dispatch reads the slot the real turns fill — and - # the same always-on union, or the dispatch would compute a - # different effective list and miss into its own agent-cache - # slot on every App call. - enabled_tools=await _apply_admin_always_on_tools( - await _apply_attachment_tool_autoenable( - input_data.enabled_tools, current_user, input_data.session_id, user_id - ), - current_user, - ), - model_id=dispatch_model_id, - system_prompt=await _plain_turn_prompt(input_data, user_id), - caching_enabled=caching_enabled, - provider=dispatch_provider or registry_provider, - inference_params=inference_params, - mantle_api_mode=mantle_api_mode, - mantle_region=mantle_region, - agent_type=effective_agent_type, - is_resume=False, - accessible_skill_ids=effective_skill_ids, - # This path builds no injected tools, but shares a cache slot - # with the real turns that do. Read the slot; never seed it. - cache_write=False, - assistant_id=input_data.rag_assistant_id, - ) - payload = await dispatch_app_tool_call( - agent, - session_id=input_data.session_id, - user_id=user_id, - tool_use_id=atc.tool_use_id, - tool_name=atc.tool_name, - arguments=atc.arguments, - ) - return JSONResponse(payload) - except AppToolCallError as e: - # 200 + envelope, not `status_code=e.code`: AgentCore Runtime - # rewrites any non-2xx to a generic 424 and discards the - # message, so a deliberate 409 ("connect the account") reached - # the SPA as "check your CloudWatch logs". app-api restores the - # real status. See `mcp_apps.error_envelope`. - return app_tool_error_response(e.message, e.code) - except HTTPException: - raise - except Exception: - logger.error("app tools/call invocation failed", exc_info=True) - return JSONResponse({"error": "Internal error"}, status_code=500) +@dataclass +class TurnAttachments: + """What this turn's attachments resolved to (map stage B4, ``preamble.files``). - # App-pushed model context (MCP Apps PR #6, `ui/update-model-context`). - # Like app_tool_call it bypasses quota / RAG / file resolution / title - # and runs NO model turn — we rebuild the conversation agent (so the - # same cached `agent.state` is reused) and stash the payload under - # `mcp_apps.context[resource_uri]`. The next real user turn merges and - # clears it. Inert behind the host flag (no live App ever calls this). - if input_data.app_context_update is not None: - acu = input_data.app_context_update - try: - request_inference_params = dict(input_data.inference_params or {}) - dispatch_model_id, dispatch_provider = input_data.model_id, input_data.provider - if not dispatch_model_id: - dispatch_model_id, dispatch_provider = await _resolve_fallback_model( - user_id, current_user, dispatch_provider - ) - caching_enabled, inference_params, mantle_api_mode, mantle_region, registry_provider = await _resolve_model_settings( - model_id=dispatch_model_id, - explicit_caching_enabled=input_data.caching_enabled, - request_inference_params=request_inference_params, - ) - agent = await get_agent( - session_id=input_data.session_id, - user_id=user_id, - auth_token=auth_token, - # Same auto-enable seam as the main turn, so a spreadsheet - # session's dispatch reads the slot the real turns fill — and - # the same always-on union, or the dispatch would compute a - # different effective list and miss into its own agent-cache - # slot on every App call. - enabled_tools=await _apply_admin_always_on_tools( - await _apply_attachment_tool_autoenable( - input_data.enabled_tools, current_user, input_data.session_id, user_id - ), - current_user, - ), - model_id=dispatch_model_id, - system_prompt=await _plain_turn_prompt(input_data, user_id), - caching_enabled=caching_enabled, - provider=dispatch_provider or registry_provider, - inference_params=inference_params, - mantle_api_mode=mantle_api_mode, - mantle_region=mantle_region, - agent_type=effective_agent_type, - is_resume=False, - accessible_skill_ids=effective_skill_ids, - # Same partial-toolset hazard as app_tool_call above. - cache_write=False, - assistant_id=input_data.rag_assistant_id, - ) - payload = dispatch_app_context_update( - agent, - resource_uri=acu.resource_uri, - content=acu.content, - structured_content=acu.structured_content, - ) - return JSONResponse(payload) - except AppContextUpdateError as e: - # Same AgentCore flattening as the app_tool_call path above. - return app_tool_error_response(e.message, e.code) - except HTTPException: - raise - except Exception: - logger.error("app context update invocation failed", exc_info=True) - return JSONResponse({"error": "Internal error"}, status_code=500) + ``files_to_send`` is the inline set the model receives as document/image + blocks. The three diverted lists never go inline (spreadsheets and decks + route through their tools; oversized and over-budget files are dropped with + a note). ``marker_names`` is every filename that still exists in the + session, in attachment order, for the ``[Attached files: …]`` marker the + SPA replays on reload. ``turn_has_document`` is the *classified* answer to + "did this turn attach something ``document_read`` can read", which feeds + the tool-injection gate and therefore ``toolConfig``. + """ - if input_data.enabled_tools: - logger.info(f"Enabled tools ({len(input_data.enabled_tools)})") + recovered_upload_ids: list = field(default_factory=list) + files_to_send: list = field(default_factory=list) + diverted_tabular: list = field(default_factory=list) + diverted_presentations: list = field(default_factory=list) + oversized_inline: list = field(default_factory=list) + over_budget_inline: list = field(default_factory=list) + dropped_over_count_names: list = field(default_factory=list) + dropped_over_count_total: int = 0 + turn_has_document: bool = False + marker_names: list = field(default_factory=list) + +async def _resolve_turn_attachments( + input_data: InvocationRequest, + user_id: str, + session_meta, + *, + is_resume: bool, + is_continuation: bool, +) -> TurnAttachments: + """Map stage B4 (``preamble.files``): recover, fetch, dedupe, partition, budget. + + Mutates ``input_data.file_upload_ids`` when it re-attaches a previous + turn's unanswered uploads, because the write-ahead marker recorded in the + session-state phase reads the ids off the request. + """ # Recover attachments the PREVIOUS turn sent but never got an answer for. # # Inline document bytes are one-shot: they are stripped out of restored @@ -2460,12 +2350,41 @@ async def invocations(request: InvocationRequest, current_user: User = Depends(g all_files, oversized_inline + over_budget_inline ) - # Covers the unconsumed-attachment recovery read, the S3 fetch behind - # `resolve_files`, and the inline/tabular/oversized partitioning. Expected - # to be ~0 on a turn with no attachments; if it is not, the hypothesis in - # docs/specs/turn-latency-preamble.md is wrong about where the time is. - prelude.mark("preamble.files") + return TurnAttachments( + recovered_upload_ids=recovered_upload_ids, + files_to_send=files_to_send, + diverted_tabular=diverted_tabular, + diverted_presentations=diverted_presentations, + oversized_inline=oversized_inline, + over_budget_inline=over_budget_inline, + dropped_over_count_names=dropped_over_count_names, + dropped_over_count_total=dropped_over_count_total, + turn_has_document=turn_has_document, + marker_names=attachment_marker_names, + ) + +async def _prepare_session_state( + input_data: InvocationRequest, + user_id: str, + session_meta, + *, + is_resume: bool, + is_continuation: bool, +) -> Tuple[bool, Optional[str]]: + """Map stage B5 (``preamble.session_state``). + + Pre-creates the session row on a first turn, clears the markers a previous + turn may have left (paused turn, pending interrupts, truncation, + interruption) and records the write-ahead attachment marker. Every read is + answered from ``session_meta`` (PR-2 of the preamble spec); only a marker + that actually exists costs a write. + + Returns ``(is_new_session, interrupted_turn_reason)``: whether this is the + session's first turn (it drives title generation), and the settled reason + the previous turn was interrupted, if it was, for the note prepended to + this turn's prompt. + """ # Pre-create session metadata so OAuth interrupts and other state can # attach to the session row from turn one. Best-effort; on failure the # post-stream lazy-create in StreamCoordinator still covers it. @@ -2481,67 +2400,737 @@ async def invocations(request: InvocationRequest, current_user: User = Depends(g input_data.session_id, user_id, snapshot=session_meta ) try: - from apis.shared.sessions.metadata import ( - clear_paused_turn, - clear_pending_interrupts, + from apis.shared.sessions.metadata import ( + clear_paused_turn, + clear_pending_interrupts, + ) + await clear_paused_turn(input_data.session_id, user_id, snapshot=session_meta) + # The snapshot's breadcrumbs go with it. They are the other half of + # the same record, and a breadcrumb that outlives the snapshot + # re-renders a prompt the user can no longer answer: the resume + # route 400s on an interrupt id the rebuilt agent never saw. Safe + # here specifically because this runs at the *head* of a non-resume + # turn — any breadcrumb this turn goes on to write lands later, on + # its own `done` event. + await clear_pending_interrupts( + input_data.session_id, user_id, snapshot=session_meta + ) + except Exception as e: + logger.error("Failed to clear stale paused_turn on new turn: %s", e, exc_info=True) + + # Invalidate any prior max_tokens "Continue" marker on every new model + # turn that isn't an interrupt-resume — both a fresh turn and a + # continuation supersede it. If a continuation itself re-truncates, the + # stream_coordinator intercept re-sets the marker. + interrupted_turn_reason: Optional[str] = None + if not is_resume: + try: + from apis.shared.sessions.metadata import clear_truncated_turn + await clear_truncated_turn(input_data.session_id, user_id, snapshot=session_meta) + except Exception as e: + logger.error("Failed to clear stale truncated_turn on new turn: %s", e, exc_info=True) + + # Same lifecycle for the interrupted-turn marker: any new non-resume + # turn supersedes a prior interruption, so a stale marker can't + # resurrect the "response interrupted" state against a turn the user + # has moved past. The pop returns the settled reason (user_stopped + # beats connection_lost via write precedence) so this same read+write + # also drives the one-turn interruption note prepended to the prompt + # in the stream generator below. + try: + from apis.shared.sessions.metadata import clear_interrupted_turn + interrupted_turn_reason = await clear_interrupted_turn( + input_data.session_id, user_id, snapshot=session_meta + ) + except Exception as e: + logger.error("Failed to clear stale interrupted_turn on new turn: %s", e, exc_info=True) + + # Write-ahead half of the unconsumed-attachment pair (the pop lives + # up with the file-resolution block). Recorded BEFORE the model call + # so it survives every way a turn can die — including the ones no + # error handler sees, like the AgentCore data plane dropping the + # stream. Cleared by StreamCoordinator the moment the turn produces + # assistant content, so a turn that succeeds never leaves a marker + # behind. Runs after ensure_session_metadata_exists above, which is + # what guarantees the session row exists to update. + if input_data.file_upload_ids: + try: + from apis.shared.sessions.metadata import set_pending_attachments + await set_pending_attachments( + input_data.session_id, user_id, input_data.file_upload_ids + ) + except Exception as e: + logger.error("Failed to record pending attachments: %s", e, exc_info=True) + + return is_new_session, interrupted_turn_reason + + +@dataclass +class TurnQuota: + """Map stage B6 (``preamble.quota``): what the quota check decided.""" + + warning_event: Any = None + session_notice_event: Any = None + exceeded_event: Any = None + + +async def _check_turn_quota( + current_user: User, + input_data: InvocationRequest, + session_meta, + *, + is_resume: bool, + is_continuation: bool, +) -> TurnQuota: + """Map stage B6 (``preamble.quota``). Fails open: a quota-service error + lets the turn run rather than blocking it.""" + # Check quota if enforcement is enabled + quota_warning_event = None + quota_session_notice_event = None + quota_exceeded_event = None + if is_quota_enforcement_enabled() and not is_resume and not is_continuation: + try: + quota_checker = get_quota_checker() + # Hand the quota checker the session cost we already read (PR-2b). + # ONLY when the row actually carries `totalCost`: absent means a + # legacy row that `get_session_metadata` still needs to backfill, + # and `None` routes the checker back to that read. Passing 0.0 for + # a missing attribute would silence the notice on exactly the + # long-lived conversations it exists to catch. + session_total_cost = None + if session_meta.row is not None and "totalCost" in session_meta.row: + try: + session_total_cost = float(session_meta.row["totalCost"]) + except (TypeError, ValueError): + session_total_cost = None + quota_result = await quota_checker.check_quota( + user=current_user, + session_id=input_data.session_id, + session_total_cost=session_total_cost, + ) + + if not quota_result.allowed: + # Quota blocked - stream as SSE instead of 429 for better UX + logger.warning("Quota blocked for user") + if quota_result.tier is None: + # No quota tier configured for this user + quota_exceeded_event = build_no_quota_configured_event(quota_result) + else: + # Quota limit exceeded + quota_exceeded_event = build_quota_exceeded_event(quota_result) + else: + # Check for warning level + quota_warning_event = build_quota_warning_event(quota_result) + if quota_warning_event: + logger.info("Quota warning for user") + + # Independent of the per-user ladder: is THIS conversation + # eating the month? (#833 PR-5 — the incident session spent + # 90% of a user's quota while every per-user warning stayed + # quiet until the day the block landed.) + quota_session_notice_event = build_quota_session_notice_event(quota_result) + if quota_session_notice_event: + logger.info("Quota session notice for user") + + except Exception as e: + # Log error but don't block request - fail open for quota errors + logger.error("Error checking quota for user", exc_info=True) + + return TurnQuota( + warning_event=quota_warning_event, + session_notice_event=quota_session_notice_event, + exceeded_event=quota_exceeded_event, + ) + + +@dataclass +class TurnModel: + """Map stage B11: the model this turn runs on and the knobs resolved for it.""" + + model_id: Optional[str] + provider: Optional[str] + caching_enabled: Optional[bool] + inference_params: dict + mantle_api_mode: Optional[str] + mantle_region: Optional[str] + + +async def _resolve_turn_model( + input_data: InvocationRequest, + current_user: User, + user_id: str, + user_settings: dict, + agent_model_override, +) -> TurnModel: + """Map stage B11 (in ``tools``): which model, and with what settings. + + Precedence for the id: the Agent's governed ``modelConfig`` (already + access-checked against the invoker), else the request's, else the user's + saved default, else the catalog's default. The registry then supplies + caching, admin bounds and locks on the inference params, the Mantle + transport fields, and the provider when nothing else named one. + """ + # Build the canonical request inference-params dict. The frontend + # sends ``inference_params`` directly; legacy ``temperature`` / + # ``max_tokens`` fields are folded in for older clients and + # treated as defaults that lose to anything in ``inference_params``. + request_inference_params: dict = dict(input_data.inference_params or {}) + if input_data.temperature is not None: + request_inference_params.setdefault("temperature", input_data.temperature) + if input_data.max_tokens is not None: + request_inference_params.setdefault("max_tokens", input_data.max_tokens) + + # Resolve the user's persisted default when the request does + # not pin a model. Without this, a "no default selected" client + # always lands on the hardcoded factory default and the user's + # saved preference is silently ignored at chat time (#161). + effective_model_id = input_data.model_id + effective_provider = input_data.provider + if agent_model_override is not None: + # The Agent's governed modelConfig wins over the request / user-default + # chain. Already access-checked against the invoker in the resolver (R2), + # so the earlier request-only gate at the top doesn't leave a hole. + effective_model_id = agent_model_override.model_id + effective_provider = agent_model_override.provider or effective_provider + if not effective_model_id: + effective_model_id, effective_provider = await _resolve_fallback_model( + user_id, current_user, effective_provider, settings=user_settings + ) + + # Agent-authored params sit as defaults BENEATH explicit request params, + # then flow through _resolve_model_settings' admin bounds/locks like any + # other request params — an author can't smuggle out-of-bounds values. + if agent_model_override is not None and agent_model_override.params: + request_inference_params = {**agent_model_override.params, **request_inference_params} + + # Single registry lookup resolves caching + inference params + + # the Mantle endpoint path + provider, merging admin defaults with + # request overrides. + caching_enabled, inference_params, mantle_api_mode, mantle_region, registry_provider = await _resolve_model_settings( + model_id=effective_model_id, + explicit_caching_enabled=input_data.caching_enabled, + request_inference_params=request_inference_params, + ) + + # Recover the provider from the registry when neither the request nor + # the Agent's model binding carried one. Agent bindings persist only + # ``model_id`` (no provider), so without this a Mantle model like + # ``openai.gpt-5.4`` resolves to provider=None → Bedrock and blows up + # in ConverseStream with "invalid model identifier" — even though the + # same model works from the normal chat path, which always sends + # ``provider`` alongside ``model_id``. + if not effective_provider and registry_provider: + effective_provider = registry_provider + + if caching_enabled is False: + logger.info("Prompt caching disabled for model") + + return TurnModel( + model_id=effective_model_id, + provider=effective_provider, + caching_enabled=caching_enabled, + inference_params=inference_params, + mantle_api_mode=mantle_api_mode, + mantle_region=mantle_region, + ) + + +async def _resolve_effective_tools( + input_data: InvocationRequest, + current_user: User, + user_id: str, + agent_tools_override, + *, + turn_has_tabular: bool, +) -> Optional[list]: + """Map stage B11 (in ``tools``): the tool ids this turn actually carries. + + The Agent's bindings if it has any, else the request's picker; then the + spreadsheet auto-enable for a session holding a spreadsheet, then the + admin's always-on set and the platform's system tools. One value, so the + cache key, every builder, the attachment guidance and the paused-turn + snapshot all see the same list. + """ + # Get agent instance with user-specific configuration + # AgentCore Memory tracks preferences across sessions per user_id + # Supports multiple LLM providers: AWS Bedrock, OpenAI, and Google Gemini + # Use augmented message and assistant system prompt if assistant RAG was applied + + # Spreadsheet tools scoped to the assistant's document corpus, + # when an assistant is attached to this request. The frontend + # keeps the assistant id in the URL for the whole session's + # lifetime, so we can trust `input_data.rag_assistant_id` + # directly; no preferences fallback needed. + # An Agent's tool bindings replace the request's enabled_tools for this + # turn (D5, resolved per invoker above). None ⇒ no tool binding ⇒ the + # request drives the toolset exactly as today. Drives both the built-in + # extra tools (spreadsheet/artifact gate on specific ids) and get_agent. + effective_enabled_tools = ( + agent_tools_override.tool_ids + if agent_tools_override is not None + else input_data.enabled_tools + ) + # A session holding a spreadsheet gets the Spreadsheet Analysis + # tools whether or not the picker has them on, gated on the + # caller's RBAC grant. Applied to the *effective* list so it + # flows into the cache key, every builder below, the attachment + # guidance and the paused-turn snapshot as one value. Sticky + # across the session (see `_session_has_tabular`), so the key + # does not flip between the attach turn and the follow-up. + effective_enabled_tools = await _apply_attachment_tool_autoenable( + effective_enabled_tools, + current_user, + input_data.session_id, + user_id, + turn_has_tabular=turn_has_tabular, + ) + + # Tools an admin pinned are unioned in for users whose roles grant + # them, unless this Agent binds its own toolset (D4). Applied to + # the same *effective* list for the same reason as the line above: + # one value flows into the cache key, every builder, and the + # paused-turn snapshot. The set depends only on the catalog and the + # user's roles, so it is constant across a session and does not + # flip the key turn to turn. + effective_enabled_tools = await _apply_admin_always_on_tools( + effective_enabled_tools, + current_user, + agent_bound_tools=agent_tools_override is not None, + ) + + return effective_enabled_tools + + +@dataclass +class TurnTools: + """Map stage B11: the context-bound tools built for this turn. + + ``extra_tools`` is appended to the registry's filtered tools by + ``BaseAgent``. ``document_tools`` is the ``document_read`` subset, whose + presence is a cache-key element of its own. ``extra_tools_key_described`` + says whether every builder that fired closes over values the cache key + already carries (otherwise the agent is not cached). ``memory_binding_key`` + is what the memory tools close over, or None. + """ + + extra_tools: list + document_tools: list + extra_tools_key_described: bool + memory_binding_key: Optional[dict] + + +async def _build_turn_tools( + input_data: InvocationRequest, + current_user: User, + user_id: str, + effective_enabled_tools: Optional[list], + *, + agent_memory, + project_memory, + turn_has_document: bool, +) -> TurnTools: + """Map stage B11 (in ``tools``): every injected tool, and the key material + that describes them. Closures only — no IO except the ``document_read`` + gate, which is memoized per session once it answers yes.""" + extra_tools = _build_spreadsheet_tools( + enabled_tools=effective_enabled_tools, + assistant_id=input_data.rag_assistant_id, + session_id=input_data.session_id, + user_id=user_id, + ) + _build_artifact_tools( + enabled_tools=effective_enabled_tools, + session_id=input_data.session_id, + user_id=user_id, + ) + _build_word_document_tools( + enabled_tools=effective_enabled_tools, + session_id=input_data.session_id, + user_id=user_id, + ) + _build_workspace_tools( + enabled_tools=effective_enabled_tools, + session_id=input_data.session_id, + user_id=user_id, + ) + _build_excel_spreadsheet_tools( + enabled_tools=effective_enabled_tools, + session_id=input_data.session_id, + user_id=user_id, + ) + _build_powerpoint_presentation_tools( + enabled_tools=effective_enabled_tools, + session_id=input_data.session_id, + user_id=user_id, + ) + _build_account_tools(effective_enabled_tools, current_user) + + memory_tools = _build_memory_tools( + agent_memory=agent_memory, + user_id=user_id, + user_email=current_user.email, + ) + # A project harness addresses its spaces by scope instead (2.4b). + if project_memory is not None: + from apis.inference_api.chat.project_memory import build_project_memory_tools + + memory_tools = build_project_memory_tools(project_memory, current_user) + extra_tools = extra_tools + memory_tools + + # document_read for any session that carries a readable attachment + # (this turn's uploads count). Gated on session state, not the + # picker; its presence goes into the cache key below rather than + # vetoing the cache, so an attachment session that could keep a + # warm agent still does. + document_tools = await _build_document_tools( + session_id=input_data.session_id, + user_id=user_id, + turn_has_document=turn_has_document, + ) + extra_tools = extra_tools + document_tools + + # Can this turn's agent be cached despite carrying injected tools? + # Only when every builder that fired closes over values the cache + # key already carries (session, user, enabled_tools). Derived from + # the same `effective_enabled_tools` that goes into the key below — + # passing the request's list here instead would let the predicate + # and the key disagree about which builders ran. + extra_tools_key_described = injected_tools_are_key_described( + enabled_tools=effective_enabled_tools, + ) + # Memory tools close over the resolved binding, so it is a cache-key + # element (Shared Projects 2.1). Only set when the tools were built, + # so the key and the toolset cannot disagree. + if project_memory is not None: + memory_binding_key = project_memory.binding_key() + elif memory_tools: + memory_binding_key = { + "spaceId": agent_memory.space_id, + "spaceName": agent_memory.space_name, + "access": agent_memory.access, + } + else: + memory_binding_key = None + + return TurnTools( + extra_tools=extra_tools, + document_tools=document_tools, + extra_tools_key_described=extra_tools_key_described, + memory_binding_key=memory_binding_key, + ) + + +def _validate_resume_interrupts(agent, input_data: InvocationRequest) -> None: + """Reject a resume that names interrupts the cached agent has not paused. + + Cache eviction, a process restart, or a forged request would otherwise be + silently accepted by Strands and drop the client's response. Raising up + front lets the client see a 400 and restart the turn cleanly. + """ + strands_agent = getattr(agent, "agent", None) + interrupt_state = getattr(strands_agent, "_interrupt_state", None) if strands_agent else None + known_ids: set[str] = set() + if interrupt_state and getattr(interrupt_state, "activated", False): + interrupts = getattr(interrupt_state, "interrupts", None) or {} + known_ids = set(interrupts.keys()) + submitted_ids = [entry.interruptId for entry in (input_data.interrupt_responses or [])] + unknown_ids = [iid for iid in submitted_ids if iid not in known_ids] + if unknown_ids: + logger.warning( + "Resume rejected: submitted interrupt ids not in paused state" + ) + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail="Unknown or expired interrupt ids; restart the turn.", + ) + + +def _build_citations(assistant, context_chunks, rag_assistant_id: Optional[str]) -> list: + """Citations to persist and to stream ahead of the answer.""" + # #111: when the agent's ``show_citations`` flag is off, suppress citations + # entirely — leaving this list empty is a single choke point that turns off all + # three downstream consumers at once: the ``event: citation`` SSE below, the + # ``citations=...`` persisted on the stored message, and the copy handed to + # ``agent.stream_async``. RAG retrieval and prompt augmentation above are + # deliberately untouched: the model still receives the context chunks, the user + # just is not shown (or able to download) the sources. + show_citations = getattr(assistant, "show_citations", True) + citations_for_storage = [] + if context_chunks and show_citations: + for chunk in context_chunks: + citations_for_storage.append( + { + "assistantId": rag_assistant_id, + "documentId": chunk.get("metadata", {}).get("document_id", ""), + # Managed KBs carry the filename under ``filename`` (set at ingest, + # managed_backend.py); legacy S3-Vectors used ``source``. Read + # managed first, fall back to legacy, then the placeholder — before + # this, every managed-KB citation rendered "Unknown Source". + "fileName": ( + chunk.get("metadata", {}).get("filename") + or chunk.get("metadata", {}).get("source") + or "Unknown Source" + ), + "text": chunk.get("text", "")[:500], # Limit excerpt length + } + ) + return citations_for_storage + + +@router.post("/invocations") +async def invocations(request: InvocationRequest, current_user: User = Depends(get_current_user_trusted)): + """The AgentCore Runtime invocation endpoint: one conversation turn. + + Runs the sequence docs/specs/turn-path-ttft.md maps, in this order. The + marks are the `turn_prelude` stage names; everything up to `stream_setup` + happens before the response opens, with no channel to the client. + + 1. `preamble.ownership` — one read of the session's META row; 404 if it is + another user's. + 2. `preamble.skills` — the turn's effective skills (opt-in; ~0 otherwise). + 3. *(no mark)* — model retirement; the two MCP App dispatches + (`app_tool_call`, `app_context_update`) return here, JSON not SSE. + 4. `preamble.files` — `_resolve_turn_attachments`. + 5. `preamble.session_state` — `_prepare_session_state`; the title task on a + first turn. + 6. `preamble.quota` — `_check_turn_quota`; a block streams as an + assistant message. + 7. *(no mark)* — retired-model denial, model access, user settings. + 8. `rag` — Agent turns: access, version, bindings, the + knowledge-base search, prompt composition, memory, binding persistence. + 9. *(no mark)* — custom prompt, personal instructions, the + single-flight lease. + 10. `tools` — resume: rebuild from the paused snapshot; else + `_resolve_turn_model`, `_resolve_effective_tools`, `_build_turn_tools`. + 11. `agent_build.*` — `get_agent`, deferred into the stream and + narrated (`preparing` / `prepared`) unless the preparing phase is off. + 12. `stream_setup` — citations, the generators, the response. + + Then `_guarded_stream` opens the response, builds the agent, emits + `turn_prelude`, starts the lease heartbeat and yields + `stream_with_quota_warning`, which composes the final message and streams + the agent. The SSE contract is in CLAUDE.md § SSE Event Types. + """ + input_data = request + user_id = current_user.user_id + auth_token = current_user.raw_token + + # Where the pre-stream time goes. Everything between here and the + # `StreamingResponse` return happens with NO channel open to the client — + # measured at 3.75s on a warm turn — so this is the only way to see which + # stage owns it. Pure timing: nothing reaches the model. + # See `turn_timing.py` and docs/specs/agent-state-feedback.md. + prelude = TurnPrelude() + # Whether this turn's agent is built inside the stream (PR-3). Recorded on + # the `turn_prelude` line so the two shapes stay distinguishable in the + # logs once the flag has been on for a while. + deferred_build = False + + # Refuse a turn against a session id another user already owns. + # + # Session ids travel in shareable URLs (`/s/{sessionId}`). Opening someone + # else's link 404s on the metadata read, but the SPA then treats the + # session as new and lets the user send — which used to fork the id: a + # SECOND metadata row under the requester, on the same session, invisible + # to both parties. In prod on 2026-08-31 that also left the original + # owner's session resolving non-deterministically between the two rows. + # + # Not a confidentiality fix — conversation content is keyed by actor id in + # AgentCore Memory, so the second user only ever saw an empty thread. This + # stops the id from being forked at all. 404 rather than 403 so the + # response says nothing about whether the session exists, matching what + # `GET /sessions/{id}/metadata` already returns for the same case. + # + # ONE read of the session's META row, shared by everything in the preamble + # that used to fetch it again (PR-2, docs/specs/turn-latency-preamble.md). + # Measured on dev: eight separate reads of this item cost ~445ms of a + # ~455ms stage. That was later traced (PR-3 of the same spec) to boto3 + # client construction on a throttled container CPU — ~47ms per read site — + # not to the DynamoDB round trip, which is ~6ms. Both fixes stand: one read, + # on a cached client. + # + # Deliberately explicit rather than a per-request memo inside + # `_get_session_by_gsi`: CLAUDE.md's "never cache session state" rule has + # been paid for twice (#741, #751), and a snapshot callers opt into cannot + # leak into one that needs a fresh read. + session_meta = await load_session_meta(input_data.session_id, user_id) + if session_meta.owned_by_other: + logger.warning( + "Rejected invocation for session %s — owned by a different user", + _sanitize_log(input_data.session_id), + ) + raise HTTPException(status_code=404, detail="Session not found") + # First of the preamble's five sub-stages (docs/specs/turn-latency-preamble.md). + # The coarse `preamble` number survives as `groups.preamble` in the emitted + # line, so the four-turn baseline in the agent-state-feedback spec stays + # comparable across this split. + prelude.mark("preamble.ownership") + # Resume requests reuse the cached agent and its paused interrupt state; + # they bypass quota, file resolution, and RAG augmentation because those + # already ran on the original turn that got paused. + is_resume = bool(input_data.interrupt_responses) + # Resolve the effective agent type: the client's explicit choice, else the + # compiled-in default ("chat"). Used for the skill resolution below and the + # non-resume get_agent calls (resume reuses the snapshot's type). An Agent + # that binds skills is coerced to "skill" later by the agent-binding + # resolver. + effective_agent_type = input_data.agent_type or DEFAULT_AGENT_TYPE + # Skills feature deferred for this environment: neutralize the legacy + # "skill" agent type, which is a ChatAgent alias since v2 PR-2. Voice and + # other agent types pass through untouched. + if not skills_enabled() and effective_agent_type == "skill": + effective_agent_type = "chat" + # Resolve the user's *effective* skills once for the whole request: the + # accessible set (catalog ∪ own), narrowed by the client's per-turn + # enabled_skills selection. Threaded into every get_agent call below so they + # share one skills_hash cache key (otherwise the app-tool-call / resume paths + # would miss the main turn's cached agent). + # + # Skills v2: this is no longer gated on agent_type == "skill". Skills are a + # plain-chat capability now — the picker in model settings sends + # enabled_skills on an ordinary turn and ChatAgent mounts the AgentSkills + # plugin. The opt-in default (D6) is what keeps this cheap: an absent or + # empty selection short-circuits to [] without touching RBAC or the skill + # table, so every turn that doesn't ask for skills costs exactly what it did + # before. An Agent's skill bindings override this further down. + effective_skill_ids = None + if skills_enabled() and input_data.enabled_skills: + effective_skill_ids = _apply_enabled_skills_filter( + await _resolve_accessible_skill_ids(current_user), + input_data.enabled_skills, + ) + # Near-zero on a turn that selects no skills — the opt-in default (D6) + # short-circuits before touching RBAC or the skill table. A non-trivial + # number here means the RBAC cache missed or the owner-index query is slow. + prelude.mark("preamble.skills") + # A "Continue" after a max_tokens truncation. Like resume, it bypasses + # quota / RAG / file resolution and does NOT clear the turn state; unlike + # resume there is no interrupt to validate — the agent is rebuilt from the + # resent params and re-entered with an empty prompt (assistant-prefill). + is_continuation = bool(input_data.continue_truncated) + # Marketplace D11: the Agent was `@`-mentioned in the composer, so it runs + # this turn only — it does not bind the conversation. Only meaningful + # alongside `rag_assistant_id`; on its own it does nothing. + is_agent_mention = bool(input_data.agent_mention) and bool(input_data.rag_assistant_id) + logger.info( + "Invocation request received (resume=%s, continue_truncated=%s, agent_mention=%s)" + % (is_resume, is_continuation, is_agent_mention) + ) + logger.info("Message received") + + # Model retirement (docs/specs/model-retirement.md §7). Resolved before anything + # builds an agent from ``input_data.model_id``, so the App tool-call / context / + # continuation paths land in the same agent-cache slot as the turn itself. A + # redirect swaps the provider as well: the request's described the retired + # model, and a successor on another transport misroutes with it. A denial is + # streamed at the access check below, once there is a turn to answer. + retired_model_denial: Optional[str] = None + requested_model = await resolve_effective_model(input_data.model_id) + if requested_model is not None: + if requested_model.denied: + retired_model_denial = retired_model_message(requested_model.retired) + elif requested_model.redirected: + input_data.model_id = requested_model.model_id + input_data.provider = requested_model.provider + + # App-initiated tools/call (MCP Apps PR #5). Like resume/continuation it + # bypasses quota / RAG / file resolution / title — there is no model + # turn. We rebuild the conversation agent (so the MCP client session + + # auth are wired exactly as for a model-driven call), dispatch the one + # named tool, publish synthesized tool_use/tool_result into the thread + # via the per-session broker, and return the CallToolResult as JSON for + # app-api to relay back to the iframe. Inert behind the host flag (the + # UIToolCatalog is empty, so dispatch rejects every call as not + # app-visible). + if input_data.app_tool_call is not None: + atc = input_data.app_tool_call + try: + agent = await _agent_for_app_dispatch( + input_data, + current_user, + user_id, + auth_token, + agent_type=effective_agent_type, + accessible_skill_ids=effective_skill_ids, ) - await clear_paused_turn(input_data.session_id, user_id, snapshot=session_meta) - # The snapshot's breadcrumbs go with it. They are the other half of - # the same record, and a breadcrumb that outlives the snapshot - # re-renders a prompt the user can no longer answer: the resume - # route 400s on an interrupt id the rebuilt agent never saw. Safe - # here specifically because this runs at the *head* of a non-resume - # turn — any breadcrumb this turn goes on to write lands later, on - # its own `done` event. - await clear_pending_interrupts( - input_data.session_id, user_id, snapshot=session_meta + payload = await dispatch_app_tool_call( + agent, + session_id=input_data.session_id, + user_id=user_id, + tool_use_id=atc.tool_use_id, + tool_name=atc.tool_name, + arguments=atc.arguments, ) - except Exception as e: - logger.error("Failed to clear stale paused_turn on new turn: %s", e, exc_info=True) - - # Invalidate any prior max_tokens "Continue" marker on every new model - # turn that isn't an interrupt-resume — both a fresh turn and a - # continuation supersede it. If a continuation itself re-truncates, the - # stream_coordinator intercept re-sets the marker. - interrupted_turn_reason: Optional[str] = None - if not is_resume: - try: - from apis.shared.sessions.metadata import clear_truncated_turn - await clear_truncated_turn(input_data.session_id, user_id, snapshot=session_meta) - except Exception as e: - logger.error("Failed to clear stale truncated_turn on new turn: %s", e, exc_info=True) + return JSONResponse(payload) + except AppToolCallError as e: + # 200 + envelope, not `status_code=e.code`: AgentCore Runtime + # rewrites any non-2xx to a generic 424 and discards the + # message, so a deliberate 409 ("connect the account") reached + # the SPA as "check your CloudWatch logs". app-api restores the + # real status. See `mcp_apps.error_envelope`. + return app_tool_error_response(e.message, e.code) + except HTTPException: + raise + except Exception: + logger.error("app tools/call invocation failed", exc_info=True) + return JSONResponse({"error": "Internal error"}, status_code=500) - # Same lifecycle for the interrupted-turn marker: any new non-resume - # turn supersedes a prior interruption, so a stale marker can't - # resurrect the "response interrupted" state against a turn the user - # has moved past. The pop returns the settled reason (user_stopped - # beats connection_lost via write precedence) so this same read+write - # also drives the one-turn interruption note prepended to the prompt - # in the stream generator below. + # App-pushed model context (MCP Apps PR #6, `ui/update-model-context`). + # Like app_tool_call it bypasses quota / RAG / file resolution / title + # and runs NO model turn — we rebuild the conversation agent (so the + # same cached `agent.state` is reused) and stash the payload under + # `mcp_apps.context[resource_uri]`. The next real user turn merges and + # clears it. Inert behind the host flag (no live App ever calls this). + if input_data.app_context_update is not None: + acu = input_data.app_context_update try: - from apis.shared.sessions.metadata import clear_interrupted_turn - interrupted_turn_reason = await clear_interrupted_turn( - input_data.session_id, user_id, snapshot=session_meta + agent = await _agent_for_app_dispatch( + input_data, + current_user, + user_id, + auth_token, + agent_type=effective_agent_type, + accessible_skill_ids=effective_skill_ids, ) - except Exception as e: - logger.error("Failed to clear stale interrupted_turn on new turn: %s", e, exc_info=True) + payload = dispatch_app_context_update( + agent, + resource_uri=acu.resource_uri, + content=acu.content, + structured_content=acu.structured_content, + ) + return JSONResponse(payload) + except AppContextUpdateError as e: + # Same AgentCore flattening as the app_tool_call path above. + return app_tool_error_response(e.message, e.code) + except HTTPException: + raise + except Exception: + logger.error("app context update invocation failed", exc_info=True) + return JSONResponse({"error": "Internal error"}, status_code=500) - # Write-ahead half of the unconsumed-attachment pair (the pop lives - # up with the file-resolution block). Recorded BEFORE the model call - # so it survives every way a turn can die — including the ones no - # error handler sees, like the AgentCore data plane dropping the - # stream. Cleared by StreamCoordinator the moment the turn produces - # assistant content, so a turn that succeeds never leaves a marker - # behind. Runs after ensure_session_metadata_exists above, which is - # what guarantees the session row exists to update. - if input_data.file_upload_ids: - try: - from apis.shared.sessions.metadata import set_pending_attachments - await set_pending_attachments( - input_data.session_id, user_id, input_data.file_upload_ids - ) - except Exception as e: - logger.error("Failed to record pending attachments: %s", e, exc_info=True) + + if input_data.enabled_tools: + logger.info(f"Enabled tools ({len(input_data.enabled_tools)})") + + # Attachments: recover the previous turn's unanswered uploads, fetch this + # turn's from S3, dedupe, partition (inline / spreadsheet / deck / + # oversized) and hold the inline set to one message's byte budget. + attachments = await _resolve_turn_attachments( + input_data, + user_id, + session_meta, + is_resume=is_resume, + is_continuation=is_continuation, + ) + + # Covers the unconsumed-attachment recovery read (answered from the + # snapshot), the S3 fetch behind `resolve_files`, and the partitioning. + # ~1ms on a turn with no attachments; otherwise the S3 GETs are the cost. + prelude.mark("preamble.files") + + # Session state: the row pre-create, the four stale-marker clears and the + # write-ahead attachment marker, all answered from `session_meta`. + is_new_session, interrupted_turn_reason = await _prepare_session_state( + input_data, + user_id, + session_meta, + is_resume=is_resume, + is_continuation=is_continuation, + ) # First turn → kick off title generation concurrently with the stream. # Runs as a background task so it doesn't add latency to TTFT. The @@ -2559,84 +3148,38 @@ async def invocations(request: InvocationRequest, current_user: User = Depends(g ) ) - # The prime suspect: five of the preamble's eight reads of the session's - # META row live in this stage (pre-create plus the four stale-marker - # clears), each on its own round trip, and four of them short-circuit - # without writing anything. The title task is inside the boundary because - # spawning it is first-turn session state; it is an `asyncio.create_task`, - # so it contributes nothing to the number. + # Five reads of the META row used to live here (pre-create plus the four + # stale-marker clears), each on its own round trip; since PR-2 all five are + # answered from `session_meta`, so a non-trivial number means a marker + # write or a first turn's pre-create. The title task is inside the boundary + # because spawning it is first-turn session state; it is an + # `asyncio.create_task`, so it contributes nothing to the number. prelude.mark("preamble.session_state") - # Check quota if enforcement is enabled - quota_warning_event = None - quota_session_notice_event = None - quota_exceeded_event = None - if is_quota_enforcement_enabled() and not is_resume and not is_continuation: - try: - quota_checker = get_quota_checker() - # Hand the quota checker the session cost we already read (PR-2b). - # ONLY when the row actually carries `totalCost`: absent means a - # legacy row that `get_session_metadata` still needs to backfill, - # and `None` routes the checker back to that read. Passing 0.0 for - # a missing attribute would silence the notice on exactly the - # long-lived conversations it exists to catch. - session_total_cost = None - if session_meta.row is not None and "totalCost" in session_meta.row: - try: - session_total_cost = float(session_meta.row["totalCost"]) - except (TypeError, ValueError): - session_total_cost = None - quota_result = await quota_checker.check_quota( - user=current_user, - session_id=input_data.session_id, - session_total_cost=session_total_cost, - ) - - if not quota_result.allowed: - # Quota blocked - stream as SSE instead of 429 for better UX - logger.warning("Quota blocked for user") - if quota_result.tier is None: - # No quota tier configured for this user - quota_exceeded_event = build_no_quota_configured_event(quota_result) - else: - # Quota limit exceeded - quota_exceeded_event = build_quota_exceeded_event(quota_result) - else: - # Check for warning level - quota_warning_event = build_quota_warning_event(quota_result) - if quota_warning_event: - logger.info("Quota warning for user") - - # Independent of the per-user ladder: is THIS conversation - # eating the month? (#833 PR-5 — the incident session spent - # 90% of a user's quota while every per-user warning stayed - # quiet until the day the block landed.) - quota_session_notice_event = build_quota_session_notice_event(quota_result) - if quota_session_notice_event: - logger.info("Quota session notice for user") - - except Exception as e: - # Log error but don't block request - fail open for quota errors - logger.error("Error checking quota for user", exc_info=True) + quota = await _check_turn_quota( + current_user, + input_data, + session_meta, + is_resume=is_resume, + is_continuation=is_continuation, + ) + quota_warning_event = quota.warning_event + quota_session_notice_event = quota.session_notice_event + quota_exceeded_event = quota.exceeded_event - # The quota round trip: a cached tier resolve, a cached O(1) cost-summary - # GetItem, and — the uncached one — the per-session notice, which reads the - # META row for the eighth time this turn. + # The quota check: a cached tier resolve, a cached O(1) cost-summary + # GetItem, and the per-session notice, answered from the snapshot's + # `totalCost` since PR-2b (or by its own read on a legacy row without one). prelude.mark("preamble.quota") # If quota exceeded, stream the quota exceeded message instead of agent response if quota_exceeded_event: - return StreamingResponse( - stream_conversational_message( - message=quota_exceeded_event.message, - stop_reason="quota_exceeded", - metadata_event=quota_exceeded_event, - session_id=input_data.session_id, - user_id=user_id, - user_input=input_data.message, - ), - media_type="text/event-stream", - headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no", "X-Session-ID": input_data.session_id}, + return _refuse_turn( + input_data, + user_id, + message=quota_exceeded_event.message, + stop_reason="quota_exceeded", + metadata_event=quota_exceeded_event, ) # A retired model with no successor is denied for everyone, wildcard holders @@ -2644,21 +3187,7 @@ async def invocations(request: InvocationRequest, current_user: User = Depends(g # resume: that turn finishes on its paused snapshot's model, whatever the # request carries. if retired_model_denial and not is_resume: - retired_event = ConversationalErrorEvent( - code=ErrorCode.FORBIDDEN, message=retired_model_denial, recoverable=False - ) - return StreamingResponse( - stream_conversational_message( - message=retired_model_denial, - stop_reason="error", - metadata_event=retired_event, - session_id=input_data.session_id, - user_id=user_id, - user_input=input_data.message, - ), - media_type="text/event-stream", - headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no", "X-Session-ID": input_data.session_id}, - ) + return _forbidden_turn(input_data, user_id, retired_model_denial) # Check model access if a specific model_id is requested if input_data.model_id: @@ -2917,22 +3446,7 @@ async def invocations(request: InvocationRequest, current_user: User = Depends(g from apis.shared.assistants.service import is_disabled_project_harness if await is_disabled_project_harness(input_data.rag_assistant_id): - refusal = PROJECTS_DISABLED_MESSAGE - refused_event = ConversationalErrorEvent( - code=ErrorCode.FORBIDDEN, message=refusal, recoverable=False - ) - return StreamingResponse( - stream_conversational_message( - message=refusal, - stop_reason="error", - metadata_event=refused_event, - session_id=input_data.session_id, - user_id=user_id, - user_input=input_data.message, - ), - media_type="text/event-stream", - headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no", "X-Session-ID": input_data.session_id}, - ) + return _forbidden_turn(input_data, user_id, PROJECTS_DISABLED_MESSAGE) # Check if assistant exists at all to provide better error message from apis.shared.assistants.service import assistant_exists @@ -3035,21 +3549,7 @@ async def invocations(request: InvocationRequest, current_user: User = Depends(g if runs_project_harness: refusal, turn_project = await _project_turn_gate(assistant.project_id) if refusal: - refused_event = ConversationalErrorEvent( - code=ErrorCode.FORBIDDEN, message=refusal, recoverable=False - ) - return StreamingResponse( - stream_conversational_message( - message=refusal, - stop_reason="error", - metadata_event=refused_event, - session_id=input_data.session_id, - user_id=user_id, - user_input=input_data.message, - ), - media_type="text/event-stream", - headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no", "X-Session-ID": input_data.session_id}, - ) + return _forbidden_turn(input_data, user_id, refusal) turn_project_id = assistant.project_id # Shared Projects 2.4b: the project's memory spaces, read now and awaited at # prompt assembly (5b), so the reads overlap binding resolution and the @@ -3083,21 +3583,7 @@ async def invocations(request: InvocationRequest, current_user: User = Depends(g project_id=turn_project_id, ) except AgentBindingBlockedError as block: - blocked_event = ConversationalErrorEvent( - code=ErrorCode.FORBIDDEN, message=block.message, recoverable=False - ) - return StreamingResponse( - stream_conversational_message( - message=block.message, - stop_reason="error", - metadata_event=blocked_event, - session_id=input_data.session_id, - user_id=user_id, - user_input=input_data.message, - ), - media_type="text/event-stream", - headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no", "X-Session-ID": input_data.session_id}, - ) + return _forbidden_turn(input_data, user_id, block.message) # Skipped on a continuation: the turn carries an empty message, so a # knowledge-base search would spend a query on "" and augment nothing. The @@ -3499,106 +3985,20 @@ async def invocations(request: InvocationRequest, current_user: User = Depends(g # snapshot's toolset, the same source the resume `get_agent` used. effective_enabled_tools = snapshot.enabled_tools else: - # Build the canonical request inference-params dict. The frontend - # sends ``inference_params`` directly; legacy ``temperature`` / - # ``max_tokens`` fields are folded in for older clients and - # treated as defaults that lose to anything in ``inference_params``. - request_inference_params: dict = dict(input_data.inference_params or {}) - if input_data.temperature is not None: - request_inference_params.setdefault("temperature", input_data.temperature) - if input_data.max_tokens is not None: - request_inference_params.setdefault("max_tokens", input_data.max_tokens) - - # Resolve the user's persisted default when the request does - # not pin a model. Without this, a "no default selected" client - # always lands on the hardcoded factory default and the user's - # saved preference is silently ignored at chat time (#161). - effective_model_id = input_data.model_id - effective_provider = input_data.provider - if agent_model_override is not None: - # The Agent's governed modelConfig wins over the request / user-default - # chain. Already access-checked against the invoker in the resolver (R2), - # so the earlier request-only gate at the top doesn't leave a hole. - effective_model_id = agent_model_override.model_id - effective_provider = agent_model_override.provider or effective_provider - if not effective_model_id: - effective_model_id, effective_provider = await _resolve_fallback_model( - user_id, current_user, effective_provider, settings=user_settings - ) - - # Agent-authored params sit as defaults BENEATH explicit request params, - # then flow through _resolve_model_settings' admin bounds/locks like any - # other request params — an author can't smuggle out-of-bounds values. - if agent_model_override is not None and agent_model_override.params: - request_inference_params = {**agent_model_override.params, **request_inference_params} - - # Single registry lookup resolves caching + inference params + - # the Mantle endpoint path + provider, merging admin defaults with - # request overrides. - caching_enabled, inference_params, mantle_api_mode, mantle_region, registry_provider = await _resolve_model_settings( - model_id=effective_model_id, - explicit_caching_enabled=input_data.caching_enabled, - request_inference_params=request_inference_params, + # The model: the Agent's binding, else the request's, else the + # user's default, else the catalog's; then the registry's settings. + turn_model = await _resolve_turn_model( + input_data, current_user, user_id, user_settings, agent_model_override ) - # Recover the provider from the registry when neither the request nor - # the Agent's model binding carried one. Agent bindings persist only - # ``model_id`` (no provider), so without this a Mantle model like - # ``openai.gpt-5.4`` resolves to provider=None → Bedrock and blows up - # in ConverseStream with "invalid model identifier" — even though the - # same model works from the normal chat path, which always sends - # ``provider`` alongside ``model_id``. - if not effective_provider and registry_provider: - effective_provider = registry_provider - - if caching_enabled is False: - logger.info("Prompt caching disabled for model") - - # Get agent instance with user-specific configuration - # AgentCore Memory tracks preferences across sessions per user_id - # Supports multiple LLM providers: AWS Bedrock, OpenAI, and Google Gemini - # Use augmented message and assistant system prompt if assistant RAG was applied - - # Spreadsheet tools scoped to the assistant's document corpus, - # when an assistant is attached to this request. The frontend - # keeps the assistant id in the URL for the whole session's - # lifetime, so we can trust `input_data.rag_assistant_id` - # directly; no preferences fallback needed. - # An Agent's tool bindings replace the request's enabled_tools for this - # turn (D5, resolved per invoker above). None ⇒ no tool binding ⇒ the - # request drives the toolset exactly as today. Drives both the built-in - # extra tools (spreadsheet/artifact gate on specific ids) and get_agent. - effective_enabled_tools = ( - agent_tools_override.tool_ids - if agent_tools_override is not None - else input_data.enabled_tools - ) - # A session holding a spreadsheet gets the Spreadsheet Analysis - # tools whether or not the picker has them on, gated on the - # caller's RBAC grant. Applied to the *effective* list so it - # flows into the cache key, every builder below, the attachment - # guidance and the paused-turn snapshot as one value. Sticky - # across the session (see `_session_has_tabular`), so the key - # does not flip between the attach turn and the follow-up. - effective_enabled_tools = await _apply_attachment_tool_autoenable( - effective_enabled_tools, + # The tool ids this turn carries — one value for the cache key, + # every builder, the attachment guidance and the paused snapshot. + effective_enabled_tools = await _resolve_effective_tools( + input_data, current_user, - input_data.session_id, user_id, - turn_has_tabular=bool(diverted_tabular), - ) - - # Tools an admin pinned are unioned in for users whose roles grant - # them, unless this Agent binds its own toolset (D4). Applied to - # the same *effective* list for the same reason as the line above: - # one value flows into the cache key, every builder, and the - # paused-turn snapshot. The set depends only on the catalog and the - # user's roles, so it is constant across a session and does not - # flip the key turn to turn. - effective_enabled_tools = await _apply_admin_always_on_tools( - effective_enabled_tools, - current_user, - agent_bound_tools=agent_tools_override is not None, + agent_tools_override, + turn_has_tabular=bool(attachments.diverted_tabular), ) # An Agent's skill bindings replace the request's skills for this turn so @@ -3611,79 +4011,15 @@ async def invocations(request: InvocationRequest, current_user: User = Depends(g effective_agent_type = "skill" effective_skill_ids = agent_skills_override.skill_ids - extra_tools = _build_spreadsheet_tools( - enabled_tools=effective_enabled_tools, - assistant_id=input_data.rag_assistant_id, - session_id=input_data.session_id, - user_id=user_id, - ) + _build_artifact_tools( - enabled_tools=effective_enabled_tools, - session_id=input_data.session_id, - user_id=user_id, - ) + _build_word_document_tools( - enabled_tools=effective_enabled_tools, - session_id=input_data.session_id, - user_id=user_id, - ) + _build_workspace_tools( - enabled_tools=effective_enabled_tools, - session_id=input_data.session_id, - user_id=user_id, - ) + _build_excel_spreadsheet_tools( - enabled_tools=effective_enabled_tools, - session_id=input_data.session_id, - user_id=user_id, - ) + _build_powerpoint_presentation_tools( - enabled_tools=effective_enabled_tools, - session_id=input_data.session_id, - user_id=user_id, - ) + _build_account_tools(effective_enabled_tools, current_user) - - memory_tools = _build_memory_tools( + turn_tools = await _build_turn_tools( + input_data, + current_user, + user_id, + effective_enabled_tools, agent_memory=agent_memory, - user_id=user_id, - user_email=current_user.email, - ) - # A project harness addresses its spaces by scope instead (2.4b). - if project_memory is not None: - from apis.inference_api.chat.project_memory import build_project_memory_tools - - memory_tools = build_project_memory_tools(project_memory, current_user) - extra_tools = extra_tools + memory_tools - - # document_read for any session that carries a readable attachment - # (this turn's uploads count). Gated on session state, not the - # picker; its presence goes into the cache key below rather than - # vetoing the cache, so an attachment session that could keep a - # warm agent still does. - document_tools = await _build_document_tools( - session_id=input_data.session_id, - user_id=user_id, - turn_has_document=turn_has_document, - ) - extra_tools = extra_tools + document_tools - - # Can this turn's agent be cached despite carrying injected tools? - # Only when every builder that fired closes over values the cache - # key already carries (session, user, enabled_tools). Derived from - # the same `effective_enabled_tools` that goes into the key below — - # passing the request's list here instead would let the predicate - # and the key disagree about which builders ran. - extra_tools_key_described = injected_tools_are_key_described( - enabled_tools=effective_enabled_tools, + project_memory=project_memory, + turn_has_document=attachments.turn_has_document, ) - # Memory tools close over the resolved binding, so it is a cache-key - # element (Shared Projects 2.1). Only set when the tools were built, - # so the key and the toolset cannot disagree. - if project_memory is not None: - memory_binding_key = project_memory.binding_key() - elif memory_tools: - memory_binding_key = { - "spaceId": agent_memory.space_id, - "spaceName": agent_memory.space_name, - "access": agent_memory.access, - } - else: - memory_binding_key = None # System-prompt assembly, the single-flight lease, skill # resolution and every tool builder (documents, attachments, @@ -3707,22 +4043,22 @@ async def _build_main_agent(): user_id=user_id, auth_token=auth_token, enabled_tools=effective_enabled_tools, - model_id=effective_model_id, + model_id=turn_model.model_id, system_prompt=system_prompt, # Use assistant's instructions if available - caching_enabled=caching_enabled, - provider=effective_provider, - inference_params=inference_params, - mantle_api_mode=mantle_api_mode, - mantle_region=mantle_region, + caching_enabled=turn_model.caching_enabled, + provider=turn_model.provider, + inference_params=turn_model.inference_params, + mantle_api_mode=turn_model.mantle_api_mode, + mantle_region=turn_model.mantle_region, agent_type=effective_agent_type, - extra_tools=extra_tools, + extra_tools=turn_tools.extra_tools, is_resume=False, accessible_skill_ids=effective_skill_ids, - extra_tools_key_described=extra_tools_key_described, - has_document_tools=bool(document_tools), + extra_tools_key_described=turn_tools.extra_tools_key_described, + has_document_tools=bool(turn_tools.document_tools), assistant_id=input_data.rag_assistant_id, build_stage_recorder=_mark_build_stage, - memory_binding=memory_binding_key, + memory_binding=turn_tools.memory_binding_key, memory_context=memory_context, ) @@ -3749,58 +4085,15 @@ async def _build_main_agent(): prelude.mark("agent_build.rest") # Resume requests must target interrupts that the cached agent - # actually has paused. Cache eviction, a process restart, or a - # forged request will otherwise be silently accepted by Strands - # and drop the client's response. Reject up front so the client - # sees a 400 and can restart the turn cleanly. + # actually has paused; see `_validate_resume_interrupts`. if is_resume: - strands_agent = getattr(agent, "agent", None) - interrupt_state = getattr(strands_agent, "_interrupt_state", None) if strands_agent else None - known_ids: set[str] = set() - if interrupt_state and getattr(interrupt_state, "activated", False): - interrupts = getattr(interrupt_state, "interrupts", None) or {} - known_ids = set(interrupts.keys()) - submitted_ids = [entry.interruptId for entry in (input_data.interrupt_responses or [])] - unknown_ids = [iid for iid in submitted_ids if iid not in known_ids] - if unknown_ids: - logger.warning( - "Resume rejected: submitted interrupt ids not in paused state" - ) - raise HTTPException( - status_code=status.HTTP_400_BAD_REQUEST, - detail="Unknown or expired interrupt ids; restart the turn.", - ) + _validate_resume_interrupts(agent, input_data) - # Build citations list for persistence (convert context chunks to citation format) - # Build citations list for persistence (convert context chunks to citation format) - # - # #111: when the agent's ``show_citations`` flag is off, suppress citations - # entirely — leaving this list empty is a single choke point that turns off all - # three downstream consumers at once: the ``event: citation`` SSE below, the - # ``citations=...`` persisted on the stored message, and the copy handed to - # ``agent.stream_async``. RAG retrieval and prompt augmentation above are - # deliberately untouched: the model still receives the context chunks, the user - # just is not shown (or able to download) the sources. - show_citations = getattr(assistant, "show_citations", True) - citations_for_storage = [] - if context_chunks and show_citations: - for chunk in context_chunks: - citations_for_storage.append( - { - "assistantId": input_data.rag_assistant_id, - "documentId": chunk.get("metadata", {}).get("document_id", ""), - # Managed KBs carry the filename under ``filename`` (set at ingest, - # managed_backend.py); legacy S3-Vectors used ``source``. Read - # managed first, fall back to legacy, then the placeholder — before - # this, every managed-KB citation rendered "Unknown Source". - "fileName": ( - chunk.get("metadata", {}).get("filename") - or chunk.get("metadata", {}).get("source") - or "Unknown Source" - ), - "text": chunk.get("text", "")[:500], # Limit excerpt length - } - ) + # Citations to persist and to stream ahead of the answer. Empty when + # the Agent hides them, which switches off every downstream consumer. + citations_for_storage = _build_citations( + assistant, context_chunks, input_data.rag_assistant_id + ) # Create stream with optional quota warning injection async def stream_with_quota_warning() -> AsyncGenerator[str, None]: @@ -3871,13 +4164,13 @@ def _session_title_sse() -> Optional[str]: # The original text becomes the single source of truth for UI display, # while the full augmented prompt stays in AgentCore Memory for the LLM. attachment_guidance = _build_attachment_guidance( - diverted_tabular, - diverted_presentations, - oversized_inline, + attachments.diverted_tabular, + attachments.diverted_presentations, + attachments.oversized_inline, effective_enabled_tools, - over_budget=over_budget_inline, - dropped_over_count_names=dropped_over_count_names, - dropped_over_count_total=dropped_over_count_total, + over_budget=attachments.over_budget_inline, + dropped_over_count_names=attachments.dropped_over_count_names, + dropped_over_count_total=attachments.dropped_over_count_total, max_files=MAX_FILES_PER_MESSAGE, ) # When multiple spreadsheets are visible, ship the full inventory @@ -3923,9 +4216,9 @@ def _session_title_sse() -> Optional[str]: # files the previous turn never got an answer for. Say so, or # the model has to guess why documents it was not asked about # are attached to the message. - if recovered_upload_ids and attachment_marker_names: + if attachments.recovered_upload_ids and attachments.marker_names: final_message = ( - f"{_build_attachment_recovery_note(attachment_marker_names)}\n\n{final_message}" + f"{_build_attachment_recovery_note(attachments.marker_names)}\n\n{final_message}" ) if interrupted_turn_reason: @@ -3950,11 +4243,11 @@ def _session_title_sse() -> Optional[str]: message_will_be_modified = ( final_message != input_data.message # RAG augmentation / attachment guidance / inventory - or bool(files_to_send) # File attachments + or bool(attachments.files_to_send) # File attachments # The `[Attached files: …]` marker is appended for diverted # attachments too, so the persisted text differs from what the # user typed even when nothing went inline (a lone .pptx). - or bool(attachment_marker_names) + or bool(attachments.marker_names) ) # Strands' resume protocol wants each entry wrapped as # {"interruptResponse": {...}}. The InvocationRequest schema @@ -3969,8 +4262,8 @@ def _session_title_sse() -> Optional[str]: async for event in agent.stream_async( final_message, session_id=input_data.session_id, - files=files_to_send if files_to_send else None, - attachment_names=attachment_marker_names or None, + files=attachments.files_to_send if attachments.files_to_send else None, + attachment_names=attachments.marker_names or None, citations=citations_for_storage if citations_for_storage else None, original_message=input_data.message if message_will_be_modified else None, interrupt_responses=interrupt_responses_payload, @@ -4193,7 +4486,7 @@ async def _guarded_stream() -> AsyncGenerator[str, None]: return StreamingResponse( _guarded_stream(), media_type="text/event-stream", - headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no", "X-Session-ID": input_data.session_id}, + headers=_sse_headers(input_data.session_id), ) except HTTPException: @@ -4213,15 +4506,10 @@ async def _guarded_stream() -> AsyncGenerator[str, None]: error_event = build_conversational_error_event(code=ErrorCode.AGENT_ERROR, error=e, session_id=input_data.session_id, recoverable=True) - return StreamingResponse( - stream_conversational_message( - message=error_event.message, - stop_reason="error", - metadata_event=error_event, - session_id=input_data.session_id, - user_id=user_id, - user_input=input_data.message, - ), - media_type="text/event-stream", - headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no", "X-Session-ID": input_data.session_id}, + return _refuse_turn( + input_data, + user_id, + message=error_event.message, + stop_reason="error", + metadata_event=error_event, ) diff --git a/backend/tests/apis/inference_api/test_get_agent_call_sites.py b/backend/tests/apis/inference_api/test_get_agent_call_sites.py index 811a1c2e9..2dd97cffb 100644 --- a/backend/tests/apis/inference_api/test_get_agent_call_sites.py +++ b/backend/tests/apis/inference_api/test_get_agent_call_sites.py @@ -40,7 +40,7 @@ def test_the_main_turn_passes_the_binding_it_built_tools_from(): if _is_false(kw.get("is_resume")) and "extra_tools_key_described" in kw ] assert len(mains) == 1 - assert ast.unparse(mains[0]["memory_binding"]) == "memory_binding_key" + assert ast.unparse(mains[0]["memory_binding"]) == "turn_tools.memory_binding_key" def test_memory_context_is_passed_on_the_main_turn_and_replayed_on_resume(): From 719a19085fbb7c6629a060071e1d278399adbdec Mon Sep 17 00:00:00 2001 From: Phil Merrell Date: Wed, 30 Sep 2026 09:16:01 -0600 Subject: [PATCH 3/3] docs(turn-path): record that P2 landed in the narrowed PR #1377 Co-Authored-By: Claude Fable 5.1 --- docs/specs/turn-path-ttft.md | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/docs/specs/turn-path-ttft.md b/docs/specs/turn-path-ttft.md index 0e81fa3ca..6ef673ba2 100644 --- a/docs/specs/turn-path-ttft.md +++ b/docs/specs/turn-path-ttft.md @@ -248,7 +248,7 @@ to be clear about: measures the title only, and on plain chat turns, so the arm's real value is not on the scorecard. -**Recommendation.** +**Recommendation** (applied to #1377 on 2026-09-30; kept here as the record of why)**.** 1. Merge the instrumentation and the `shared_clients` arm now, amended so the shared session, its clients and the strategy ids are built at warm-up (P2). Keep the arm behind @@ -295,6 +295,8 @@ Cost: perf_counter calls only. V2: none. ### P2 — Warm the right things, once (F2) +**Status (2026-09-30): built in PR #1377 as narrowed** — the off-loop arm was withdrawn, the shared session lives in `apis/shared/aws_clients.py`, warm-up builds its three clients and primes the strategy ids unconditionally, and the SDK session manager, `_discover_strategy_ids` and `BedrockModel` all take it on the `shared_clients` arm. The strategy-id cache is keyed by arm so control's first turn is unchanged. What remains of P2 is P0: the dev A/B decides whether the arm ships default-on. + - One process-wide `boto3.Session` in `apis/shared/aws_clients.py` (it already owns the client cache); warm-up builds `bedrock-agentcore`, `bedrock-agentcore-control` and `bedrock-runtime` clients on **it**.