From 8b19ca9408df25d02ce7e7e25b62e6726b4a8a19 Mon Sep 17 00:00:00 2001 From: Phil Merrell Date: Sun, 27 Sep 2026 11:48:47 -0600 Subject: [PATCH 1/2] perf(backend): push the first-turn session title as soon as it exists Four delays sat between Nova Micro answering and the user seeing the title: - The title was only checked between agent events, so it waited out the model's time-to-first-token and every long tool call. It now also rides the coordinator's 100ms live status merge (`poll_side_frame`). - The title task awaited its DynamoDB write before finishing, and the stream only emits a finished task. The write is now scheduled in the background. - That write took two round trips (GSI lookup, then update). Rows are born with the static SK, so it is now one keyed update guarded by attribute_exists, falling back to the GSI only for a legacy SK. - A new bedrock-runtime client per title: ~250ms on first use in a process (on the event loop) and a fresh TLS handshake every time. It is now built once, lazily, inside the worker thread. Co-Authored-By: Claude Opus 5.5 --- CLAUDE.MD | 2 +- backend/src/agents/main_agent/chat_agent.py | 7 +- .../streaming/stream_coordinator.py | 29 +++++- backend/src/apis/inference_api/chat/routes.py | 11 ++- .../src/apis/inference_api/chat/service.py | 60 +++++++++--- backend/src/apis/shared/sessions/metadata.py | 23 ++++- .../streaming/test_agent_status_live_drain.py | 92 +++++++++++++++++++ .../tests/shared/test_sessions_metadata.py | 39 ++++++++ .../test_side_channel_inference_config.py | 53 +++++++++++ 9 files changed, 296 insertions(+), 20 deletions(-) diff --git a/CLAUDE.MD b/CLAUDE.MD index 0bacb1ad3..2764eed96 100644 --- a/CLAUDE.MD +++ b/CLAUDE.MD @@ -75,7 +75,7 @@ Current in-development flags: **Shared Projects** (`PROJECTS_ENABLED` / `CDK_PRO | `ui_tool_input_partial` | Streamed partial tool input for a UI tool (SEP-1865 `ui/notifications/tool-input-partial`) — payload `{type, toolUseId, arguments}`. Emitted repeatedly while the model is still streaming a UI tool's arguments (after the early `ui_resource` mount); `arguments` is the streamed prefix server-side "healed" into a valid object (`apis/shared/mcp_apps/partial_json.py`). The SPA relays each to the App via `ui/notifications/tool-input-partial` so a progressively-rendering App (e.g. Excalidraw's guided camera tour) animates as args arrive; the complete `tool-input` follows once the input is final. Same gating as `ui_resource` | | `agent_status` | What the agent is doing right now — payload `{type, sessionId, phase, cycle, toolName?, toolUseId?, durationMs?, ok?}`. Emitted from `AgentStatusHook` (`BeforeModelCall` / `Before`+`AfterToolCall`) and drained in `stream_coordinator` exactly like `steering_applied`, before the event it precedes, so "Using list_assignments" reaches the client while that tool is running rather than after its result. Phases: `thinking` (one per event-loop cycle — a three-tool turn reports it four times, and `cycle` distinguishes them), `tool_start`, `tool_end` (carries Strands' own measured `durationMs`, and `ok=false` for BOTH a raised exception and a result with `status: "error"`). There is deliberately **no "responding" phase** — the SPA already knows text is streaming from the deltas, and a backend-derived duplicate of a fact the client holds first-hand would only disagree at the edges. Durations are live-only and deliberately NOT persisted: a reloaded conversation shows summaries without timings, where a client-invented number would be one the user could not trust. Costs nothing against the model — nothing it produces reaches the prompt, so the cacheable prefix is untouched. Gated by `AGENT_STATUS_ENABLED` (default on with a kill switch); while off the hook is registered but every callback returns immediately and the SPA stays on its cycling phrases | | `tool_group_summary` | Model-generated one-line summary of a finished tool batch — payload `{type, sessionId, batchId, toolUseIds, summary}`, e.g. "Found the Syllabus Acknowledgment assignment in BIO 101". Produced by a Nova Micro **side-channel** task (`apis/shared/tool_summaries/summarizer.py`) structured exactly like `session_title`: its own Bedrock call on its own messages, concurrent with the agent stream, so it **never appends to the conversation** and adds nothing to the cacheable prefix. Spend is one bounded call per batch (inputs/results truncated at capture in the hook AND again in the summarizer). Lands mid-turn, out of band with the content stream, so it can arrive after the rail that shows it has rendered; the SPA keys it by `toolUseId` (not batch — client-side grouping need not match backend batches, and a group spanning two batches shows the first batch's line). Persisted as `TSUM#` rows in sessions-metadata, reusing the `SessionLookupIndex` GSI — zero new infra — and replayed on `GET /messages` as `toolSummaries`, because the event never re-streams. Deliberately NOT written onto the message content blocks: that is the Converse payload and the cacheable prefix, so a display string there would be paid at model rates every subsequent turn. Gated by `TOOL_SUMMARIES_ENABLED` (default on with a kill switch); while off the SPA's deterministic client-side formatter still renders ("Listed 4 assignments"), so absence is a downgrade in specificity, never a blank | -| `session_title` | Server-generated conversation title on a session's FIRST turn — payload `{type, sessionId, title}`. Title generation (Nova Micro) runs as an asyncio task concurrent with the agent stream; the finished title is interleaved between agent events (non-blocking done-check in `stream_with_quota_warning`), so the sidebar/top-nav rename while the response is still pending. Emitted at most once per stream, possibly after `done` (the SPA parser allowlists it past Completed-state gating); never carries the "New Conversation" placeholder. Best-effort: a stream that finishes before generation emits nothing — the SPA's post-close metadata refresh (`refreshTitleFromServer`) is the fallback, reading the title the task also persisted via `update_session_title` | +| `session_title` | Server-generated conversation title on a session's FIRST turn — payload `{type, sessionId, title}`. Title generation (Nova Micro) runs as an asyncio task concurrent with the agent stream, so the sidebar/top-nav rename while the response is still pending. The finished title is picked up by a non-blocking one-shot check (`_session_title_sse` in `stream_with_quota_warning`) that is polled from two places: the coordinator's live status merge (`poll_side_frame`, every 100ms, so a title that lands during the model's time-to-first-token or a long tool call goes out at once rather than waiting for the agent's next event) and between agent events (its only route while `AGENT_STATUS_LIVE_DRAIN_ENABLED=false`). The task returns the moment Nova answers and persists the title in the background (one keyed `update_item` on the static SK), so the DynamoDB write never sits between generation and the user seeing it. Emitted at most once per stream, possibly after `done` (the SPA parser allowlists it past Completed-state gating); never carries the "New Conversation" placeholder. Best-effort: a stream that finishes before generation emits nothing, and deliberately does not hold the connection open to wait (the SPA clears loading on close, not on `done`) — the SPA's post-close metadata refresh (`refreshTitleFromServer`, up to 5 retries at 400ms) is the fallback, reading the persisted title | | `quota_session_notice` | **This conversation** has reached the tier's session-notice share of the monthly limit — payload `{type, sessionId, sessionCost, quotaLimit, sessionPercentageOfLimit, thresholdPercentage, message}`. Emitted at the head of the stream right after `quota_warning`, and re-emitted every turn while over the share (dismissal is client-side, same contract as `quota_warning`). `sessionCost` is the session's **lifetime** cost — the `totalCost` aggregate on its metadata row — deliberately not period-scoped: a conversation that opened last month and is spending this month's budget is exactly the one worth surfacing. Share is tier-configurable (`sessionNoticePercentage`, default 25%, 0 disables); the whole runway rides the `QUOTA_RUNWAY_ENABLED` kill switch (default on), which also gates the 50%/75% `quota_warning` rungs. The SPA scopes it to the conversation it names — never shown above another thread's composer | | `agent_notice` | A Shared Project's harness is running this turn **without** some of its setup — payload `{type, sessionId, agentId, projectId?, message, unavailableModelId?, unavailableTools, unavailableSkills, unavailableMemory?}`. Only a project harness degrades (shared-projects §9.6): a member missing a bound tool, skill, model or memory space still gets the turn, with that piece dropped — tools keyed by *server* so no later scoped ref of a denied server slips through, an unavailable model falling through to the member's default, and an all-dropped toolset still meaning "no tools", never the request's. Every other shared agent keeps block-with-message (D5, `stream_error`). Emitted before `message_start`, beside `quota_session_notice`. **SSE only, never the prompt** — the model is not told what was dropped, so the cacheable prefix is unaffected — and not persisted, because the next turn re-derives it. The SPA ignores it until the Projects UI (PR-1.8) renders it | | `model_retry` | Backend is retrying a failed model call instead of surfacing it — payload `{type, attempt, delaySeconds}`. Emitted from Strands' `EventLoopThrottleEvent`; `attempt` is 1-based and counted per turn in `stream_processor` (the raw event carries only the delay). **Timing caveat:** Strands sleeps the backoff *inside* its hook and yields the event afterwards, so it lands as the next attempt BEGINS, not when the wait starts — `delaySeconds` describes the gap just endured, it is not a countdown. It also cannot cover the failing model call itself, which is indistinguishable from a slow healthy one. The SPA swaps the loading indicator's cycling phrases for a fixed amber notice, cleared on `message_start`/`done` | diff --git a/backend/src/agents/main_agent/chat_agent.py b/backend/src/agents/main_agent/chat_agent.py index 8b3c13a77..4b71e045b 100644 --- a/backend/src/agents/main_agent/chat_agent.py +++ b/backend/src/agents/main_agent/chat_agent.py @@ -7,7 +7,7 @@ import logging import os -from typing import Any, AsyncGenerator, Dict, List, Optional +from typing import Any, AsyncGenerator, Callable, Dict, List, Optional from agents.main_agent.base_agent import BaseAgent from agents.main_agent.core import AgentFactory @@ -164,6 +164,7 @@ async def stream_async( turn_project_id: Optional[str] = None, turn_lease: Any = None, turn_started_at: Optional[float] = None, + poll_side_frame: Optional[Callable[[], Optional[str]]] = None, ) -> AsyncGenerator[str, None]: """ Stream agent responses. @@ -198,6 +199,9 @@ async def stream_async( off the agent for the same reason as `turn_agent_id`: the agent instance is cached across turns, so per-turn state must never live on it (#741/#751). + poll_side_frame: Non-blocking check for a frame produced outside the + agent stream (the first turn's `session_title`), polled by the + coordinator's live status merge. Per turn, like `turn_lease`. Yields: str: SSE formatted events @@ -234,5 +238,6 @@ async def stream_async( turn_project_id=turn_project_id, turn_lease=turn_lease, turn_started_at=turn_started_at, + poll_side_frame=poll_side_frame, ): yield event diff --git a/backend/src/agents/main_agent/streaming/stream_coordinator.py b/backend/src/agents/main_agent/streaming/stream_coordinator.py index 06d594ab2..ff7053e8f 100644 --- a/backend/src/agents/main_agent/streaming/stream_coordinator.py +++ b/backend/src/agents/main_agent/streaming/stream_coordinator.py @@ -9,7 +9,7 @@ import os import time from datetime import datetime, timezone -from typing import Any, AsyncGenerator, Dict, List, Optional, Union +from typing import Any, AsyncGenerator, Callable, Dict, List, Optional, Union from agents.main_agent.config.constants import EnvVars from agents.main_agent.session.hooks.prefix_fingerprint import ( @@ -256,6 +256,7 @@ async def stream_response( turn_project_id: Optional[str] = None, turn_lease: Any = None, turn_started_at: Optional[float] = None, + poll_side_frame: Optional[Callable[[], Optional[str]]] = None, ) -> AsyncGenerator[str, None]: """ Stream agent responses with proper lifecycle management @@ -286,6 +287,12 @@ async def stream_response( stamped *unconditionally*, including to None, for the same reason ``reset_cancellation_state`` exists: a lease left behind by a previous turn on a cached agent would be read against a row that no longer names us. + poll_side_frame: Non-blocking check for one formatted SSE frame produced + outside the agent stream (the first turn's ``session_title``). Polled + on every pass of the live status merge, so the frame goes out within + one poll interval of being ready instead of waiting for the agent + stream's next event. Must be idempotent: the caller keeps checking it + between events too, which is its only route while the merge is off. Yields: str: SSE formatted events @@ -472,7 +479,8 @@ async def stream_response( ) if agent_status_live_drain_enabled(): processed_stream = self._merge_agent_status( - processed_stream, main_agent_wrapper, session_id + processed_stream, main_agent_wrapper, session_id, + poll_side_frame=poll_side_frame, ) async for event in processed_stream: @@ -2686,6 +2694,7 @@ async def _merge_agent_status( events: AsyncGenerator[Dict[str, Any], None], main_agent_wrapper: Any, session_id: str, + poll_side_frame: Optional[Callable[[], Optional[str]]] = None, ) -> AsyncGenerator[Any, None]: """Yield the agent's events, interleaved with status transitions as they happen. @@ -2758,6 +2767,9 @@ async def _merge_agent_status( main_agent_wrapper, session_id ): yield _StatusFrame(sse) + side_frame = self._poll_side_frame(poll_side_frame) + if side_frame: + yield _StatusFrame(side_frame) if not done: continue @@ -2789,6 +2801,19 @@ async def _merge_agent_status( exc_info=True, ) + @staticmethod + def _poll_side_frame( + poll_side_frame: Optional[Callable[[], Optional[str]]], + ) -> Optional[str]: + """Best-effort: a failing poll must never break the stream it rides on.""" + if poll_side_frame is None: + return None + try: + return poll_side_frame() + except Exception: # noqa: BLE001 - side-channel frames are never load-bearing + logger.warning("Side-frame poll failed", exc_info=True) + return None + def _drain_agent_status_events( self, main_agent_wrapper: Any, session_id: str ) -> List[str]: diff --git a/backend/src/apis/inference_api/chat/routes.py b/backend/src/apis/inference_api/chat/routes.py index 4727deef2..44a414f42 100644 --- a/backend/src/apis/inference_api/chat/routes.py +++ b/backend/src/apis/inference_api/chat/routes.py @@ -3809,9 +3809,13 @@ async def stream_with_quota_warning() -> AsyncGenerator[str, None]: # (kicked off before the quota check on first turns) finishes, # push the title to the client so the sidebar/header rename in # parallel with the pending response instead of at stream end. - # Checked between agent events — never awaited, so it adds no - # latency; a stream that outruns Nova Micro simply never emits - # and the SPA's post-close metadata refresh covers it. + # Never awaited, so it adds no latency. Polled from two places: + # the coordinator's live status merge (every 100ms, so a title + # that lands during the model's time-to-first-token or a long + # tool call goes out right away) and between agent events below + # (the only route while that merge is switched off). A stream + # that outruns Nova Micro simply never emits and the SPA's + # post-close metadata refresh covers it. title_emitted = False def _session_title_sse() -> Optional[str]: @@ -3993,6 +3997,7 @@ def _session_title_sse() -> Optional[str]: # is not. None for preview sessions and the local # no-DynamoDB path, where steering is simply inert. turn_lease=session_lease, + poll_side_frame=_session_title_sse if title_task is not None else None, ): yield event # Interleave the finished title between agent events (same diff --git a/backend/src/apis/inference_api/chat/service.py b/backend/src/apis/inference_api/chat/service.py index cb5ad0f27..8d90145e4 100644 --- a/backend/src/apis/inference_api/chat/service.py +++ b/backend/src/apis/inference_api/chat/service.py @@ -8,6 +8,7 @@ import logging import hashlib import os +import threading from typing import Any, Callable, Dict, Optional, List, Tuple import boto3 @@ -596,6 +597,44 @@ def clear_agent_cache(): Focus on being informative and scannable. The title should allow users to quickly identify this conversation in a list.""" +# One Bedrock client for every title, built on first use. A client per call +# cost ~250ms the first time in a process (botocore loads the service model) +# and, every time, a fresh connection pool — a new TCP+TLS handshake to +# bedrock-runtime on the path of a title the user is watching for. boto3 +# clients are thread-safe, so the worker threads below can share it. +_title_bedrock_client: Any = None +# Guards the first build: boto3's default session is not thread-safe, and two +# first-turn titles can reach it from separate worker threads at once. +_title_client_lock = threading.Lock() + +# Strong references to in-flight title writes. The event loop only holds weak +# references to tasks, so a fire-and-forget write could otherwise be +# garbage-collected before it lands. +_pending_title_writes: set = set() + + +def _get_title_bedrock_client() -> Any: + global _title_bedrock_client + if _title_bedrock_client is None: + with _title_client_lock: + if _title_bedrock_client is None: + bedrock_region = os.environ.get('AWS_REGION', 'us-east-1') + _title_bedrock_client = boto3.client('bedrock-runtime', region_name=bedrock_region) + return _title_bedrock_client + + +def _converse_for_title(**kwargs: Any) -> Dict[str, Any]: + """Run in a worker thread, so a first-use client build stays off the event loop.""" + return _get_title_bedrock_client().converse(**kwargs) + + +def _schedule_title_write(session_id: str, user_id: str, title: str) -> None: + task = asyncio.create_task( + update_session_title(session_id=session_id, user_id=user_id, title=title) + ) + _pending_title_writes.add(task) + task.add_done_callback(_pending_title_writes.discard) + async def generate_conversation_title( session_id: str, @@ -608,8 +647,8 @@ async def generate_conversation_title( This function: 1. Truncates user input to ~500 tokens (2000 chars as rough approximation) 2. Calls Nova Micro with optimized system prompt - 3. Updates session metadata both locally and in cloud - 4. Returns generated title or fallback on error + 3. Returns the title as soon as Nova does, persisting it in the background + 4. Returns fallback on error Args: session_id: Session identifier @@ -628,10 +667,6 @@ async def generate_conversation_title( logger.debug(f"Truncated input from {len(user_input)} to {MAX_INPUT_LENGTH} chars") try: - # Initialize Bedrock Runtime client - bedrock_region = os.environ.get('AWS_REGION', 'us-east-1') - bedrock_client = boto3.client('bedrock-runtime', region_name=bedrock_region) - # Prepare request for Nova Micro # us.amazon.nova-micro-v1:0 is the fastest, most cost-effective model request_body = { @@ -657,7 +692,7 @@ async def generate_conversation_title( # whole Nova round-trip, stalling the agent stream this task runs # concurrently with. response = await asyncio.to_thread( - bedrock_client.converse, + _converse_for_title, modelId="us.amazon.nova-micro-v1:0", messages=request_body["messages"], system=request_body["system"], @@ -682,10 +717,13 @@ async def generate_conversation_title( logger.info("✅ Generated title successfully") - # Targeted update — only writes the title attribute. The post-stream - # update_session_activity write is also targeted and disjoint, so the - # two cannot clobber each other on overlapping turns. - await update_session_title(session_id=session_id, user_id=user_id, title=title) + # Persist without holding the title back. This task finishing is what + # lets the stream push `session_title`, so awaiting the DynamoDB write + # here would put its round trips between Nova answering and the user + # seeing the name. Targeted update — only writes the title attribute. + # The post-stream update_session_activity write is also targeted and + # disjoint, so the two cannot clobber each other on overlapping turns. + _schedule_title_write(session_id, user_id, title) return title diff --git a/backend/src/apis/shared/sessions/metadata.py b/backend/src/apis/shared/sessions/metadata.py index 1ef7d990f..268f97ff9 100644 --- a/backend/src/apis/shared/sessions/metadata.py +++ b/backend/src/apis/shared/sessions/metadata.py @@ -1322,8 +1322,13 @@ async def update_session_title(session_id: str, user_id: str, title: str) -> Non Uses a targeted ``UpdateExpression`` so it can run concurrently with ``store_session_metadata`` (which does a full-row merge) without racing - on other fields like ``messageCount`` or ``lastMessageAt``. Looks up the - current SK via the GSI because the SK contains a timestamp. + on other fields like ``messageCount`` or ``lastMessageAt``. + + One round trip in the common case: every row is born with the static + ``S#{session_id}`` SK (issue #175), and this runs on a session's first + turn, so the key is known without a lookup. A row still on a legacy + timestamped SK fails the ``attribute_exists`` guard and falls back to + resolving the SK through the GSI. No-op when the session row doesn't exist (preview sessions, sessions deleted mid-turn). @@ -1336,9 +1341,23 @@ async def update_session_title(session_id: str, user_id: str, title: str) -> Non raise RuntimeError("DYNAMODB_SESSIONS_METADATA_TABLE_NAME environment variable is required") try: + from botocore.exceptions import ClientError table = get_dynamodb_table(sessions_metadata_table) + try: + table.update_item( + Key={"PK": f"USER#{user_id}", "SK": _static_session_sk(session_id)}, + UpdateExpression="SET title = :t", + ConditionExpression="attribute_exists(PK)", + ExpressionAttributeValues={":t": title}, + ) + logger.info(f"💾 Updated title for session {session_id}") + return + except ClientError as e: + if e.response.get("Error", {}).get("Code") != "ConditionalCheckFailedException": + raise + existing = await _get_session_by_gsi(session_id, user_id, table) if not existing: logger.info(f"update_session_title: session {session_id} not found, skipping") diff --git a/backend/tests/agents/main_agent/streaming/test_agent_status_live_drain.py b/backend/tests/agents/main_agent/streaming/test_agent_status_live_drain.py index 1700295db..a3fe26cf0 100644 --- a/backend/tests/agents/main_agent/streaming/test_agent_status_live_drain.py +++ b/backend/tests/agents/main_agent/streaming/test_agent_status_live_drain.py @@ -206,6 +206,98 @@ def drain_statuses(self): assert any(f.startswith("event: done") for f, _ in frames) +_TITLE_FRAME = 'event: session_title\ndata: {"type": "session_title", "title": "T"}\n\n' + + +class _OneShotTitle: + """The routes' `_session_title_sse` contract: ready once, emitted once. + + It becomes ready when the agent stream goes silent (the hook's first + record), which models a title that lands during a tool call or the + model's time-to-first-token. + """ + + def __init__(self, hook: _Hook) -> None: + self.ready = False + self.emitted = 0 + record = hook.record + + def _record(status: dict) -> None: + self.ready = True + record(status) + + hook.record = _record + + def __call__(self) -> Optional[str]: + if self.emitted or not self.ready: + return None + self.emitted += 1 + return _TITLE_FRAME + + +async def _collect_with_side_frame(agent, wrapper, poll) -> List[tuple]: + coordinator = StreamCoordinator() + seen: List[tuple] = [] + async for sse in coordinator.stream_response( + agent=agent, + prompt="hi", + session_manager=_SessionManager(), + session_id="sess-1", + user_id="user-1", + main_agent_wrapper=wrapper, + poll_side_frame=poll, + ): + seen.append((sse, agent.spoke_again)) + return seen + + +class TestSideFrame: + """The first turn's `session_title` rides the same poll as the statuses.""" + + @pytest.mark.asyncio + async def test_a_ready_title_goes_out_during_the_silence(self, monkeypatch): + monkeypatch.delenv("AGENT_STATUS_LIVE_DRAIN_ENABLED", raising=False) + hook = _Hook() + agent = _StallingAgent(hook) + poll = _OneShotTitle(hook) + + frames = await _collect_with_side_frame(agent, _Wrapper(hook), poll) + + arrivals = [spoke for sse, spoke in frames if sse == _TITLE_FRAME] + assert arrivals == [False] + assert poll.emitted == 1 + + @pytest.mark.asyncio + async def test_it_is_polled_even_without_a_status_hook(self, monkeypatch): + """No hook means no statuses, not no title.""" + monkeypatch.delenv("AGENT_STATUS_LIVE_DRAIN_ENABLED", raising=False) + agent = _StallingAgent(_Hook(), silence=0.25) + emitted: List[str] = [] + + def poll() -> Optional[str]: + if emitted: + return None + emitted.append(_TITLE_FRAME) + return _TITLE_FRAME + + frames = await _collect_with_side_frame(agent, _Wrapper(None), poll) + + assert [sse for sse, _ in frames].count(_TITLE_FRAME) == 1 + + @pytest.mark.asyncio + async def test_a_failing_poll_does_not_break_the_turn(self, monkeypatch): + monkeypatch.delenv("AGENT_STATUS_LIVE_DRAIN_ENABLED", raising=False) + hook = _Hook() + agent = _StallingAgent(hook, silence=0.05) + + def poll() -> Optional[str]: + raise RuntimeError("title task exploded") + + frames = await _collect_with_side_frame(agent, _Wrapper(hook), poll) + + assert any(f.startswith("event: done") for f, _ in frames) + + class TestEarlyExit: @pytest.mark.asyncio async def test_a_cooperative_stop_mid_silence_ends_cleanly(self, monkeypatch): diff --git a/backend/tests/shared/test_sessions_metadata.py b/backend/tests/shared/test_sessions_metadata.py index cf659cf2e..60b2d23f8 100644 --- a/backend/tests/shared/test_sessions_metadata.py +++ b/backend/tests/shared/test_sessions_metadata.py @@ -1590,3 +1590,42 @@ async def test_marker_survives_the_metadata_read_model(self, sessions_metadata_t meta = await get_session_metadata("s1", "u1") assert meta.pending_attachment_upload_ids == ["up-1", "up-2"] assert meta.pending_attachments_at is not None + + +class TestUpdateSessionTitle: + """One round trip on the born-static row; the GSI only for legacy SKs.""" + + @pytest.mark.asyncio + async def test_writes_a_born_static_row(self, sessions_metadata_table): + from apis.shared.sessions.metadata import ( + ensure_session_metadata_exists, + get_session_metadata, + update_session_title, + ) + await ensure_session_metadata_exists("s1", "u1") + + await update_session_title("s1", "u1", "Biology Syllabus") + + assert (await get_session_metadata("s1", "u1")).title == "Biology Syllabus" + + @pytest.mark.asyncio + async def test_falls_back_to_the_gsi_for_a_legacy_row(self, sessions_metadata_table): + from apis.shared.sessions.metadata import update_session_title + + _put_legacy_row(sessions_metadata_table, "s1", "2026-01-01T00:00:00Z") + + await update_session_title("s1", "u1", "Biology Syllabus") + + items = sessions_metadata_table.scan()["Items"] + assert [(i["SK"], i["title"]) for i in items] == [ + ("S#ACTIVE#2026-01-01T00:00:00Z#s1", "Biology Syllabus") + ] + + @pytest.mark.asyncio + async def test_a_missing_row_is_not_created(self, sessions_metadata_table): + """The guard is what keeps the keyed write from upserting a ghost row.""" + from apis.shared.sessions.metadata import update_session_title + + await update_session_title("s1", "u1", "Biology Syllabus") + + assert sessions_metadata_table.scan()["Items"] == [] diff --git a/backend/tests/shared/test_side_channel_inference_config.py b/backend/tests/shared/test_side_channel_inference_config.py index 5d241fe7c..72430133b 100644 --- a/backend/tests/shared/test_side_channel_inference_config.py +++ b/backend/tests/shared/test_side_channel_inference_config.py @@ -8,6 +8,7 @@ other side channels to the same shape. """ +import asyncio from unittest.mock import AsyncMock, MagicMock import pytest @@ -29,6 +30,16 @@ def converse(**kwargs): return client +@pytest.fixture(autouse=True) +def _fresh_title_client(monkeypatch): + """The title client is cached per process; each test builds its own.""" + monkeypatch.setattr(chat_service, "_title_bedrock_client", None) + + +async def _drain_title_writes() -> None: + await asyncio.gather(*list(chat_service._pending_title_writes)) + + def _patch_boto3_module(monkeypatch, client: MagicMock) -> None: module = MagicMock() module.client.return_value = client @@ -64,6 +75,7 @@ async def test_conversation_title(monkeypatch): title = await chat_service.generate_conversation_title(session_id="s", user_id="u", user_input="hi") assert title == "Planning a biology syllabus" assert "topP" not in client.converse.call_args.kwargs["inferenceConfig"] + await _drain_title_writes() @pytest.mark.asyncio @@ -86,3 +98,44 @@ async def test_conversation_title_at_the_token_ceiling_is_still_clipped(monkeypa monkeypatch.setattr(chat_service, "update_session_title", AsyncMock()) title = await chat_service.generate_conversation_title(session_id="s", user_id="u", user_input="hi") assert title == "A" * 47 + "..." + await _drain_title_writes() + + +@pytest.mark.asyncio +async def test_conversation_title_returns_before_its_write_lands(monkeypatch): + """The stream pushes `session_title` when this task finishes, so the + DynamoDB write must not sit between Nova answering and the user seeing it.""" + client = _client("Planning a biology syllabus") + monkeypatch.setattr(chat_service.boto3, "client", MagicMock(return_value=client)) + write_may_finish = asyncio.Event() + written: list = [] + + async def slow_write(session_id, user_id, title): + await write_may_finish.wait() + written.append(title) + + monkeypatch.setattr(chat_service, "update_session_title", slow_write) + + title = await chat_service.generate_conversation_title(session_id="s", user_id="u", user_input="hi") + + assert title == "Planning a biology syllabus" + assert written == [] + write_may_finish.set() + await _drain_title_writes() + assert written == ["Planning a biology syllabus"] + + +@pytest.mark.asyncio +async def test_conversation_title_reuses_one_bedrock_client(monkeypatch): + """A client per title paid botocore's model load and a fresh TLS handshake.""" + client = _client("Planning a biology syllabus") + factory = MagicMock(return_value=client) + monkeypatch.setattr(chat_service.boto3, "client", factory) + monkeypatch.setattr(chat_service, "update_session_title", AsyncMock()) + + for _ in range(3): + await chat_service.generate_conversation_title(session_id="s", user_id="u", user_input="hi") + await _drain_title_writes() + + assert factory.call_count == 1 + assert client.converse.call_count == 3 From 0a48f3209559cb118118251308eb08fc76b2e8bb Mon Sep 17 00:00:00 2001 From: Phil Merrell Date: Sun, 27 Sep 2026 11:48:47 -0600 Subject: [PATCH 2/2] perf(spa): poll for a late session title instead of one 1.5s retry When the stream outruns title generation, the post-close metadata read sees the placeholder. Retry up to 5 times at 400ms (the same ~1.5s+ window) and stop early once the title arrived another way, so it appears within ~400ms of being written instead of after a fixed 1.5s wait. Co-Authored-By: Claude Opus 5.5 --- .../services/chat/chat-http.service.spec.ts | 56 +++++++++++++++++++ .../services/chat/chat-http.service.ts | 23 ++++++-- 2 files changed, 74 insertions(+), 5 deletions(-) diff --git a/frontend/ai.client/src/app/session/services/chat/chat-http.service.spec.ts b/frontend/ai.client/src/app/session/services/chat/chat-http.service.spec.ts index 890ea4f8e..94711fadc 100644 --- a/frontend/ai.client/src/app/session/services/chat/chat-http.service.spec.ts +++ b/frontend/ai.client/src/app/session/services/chat/chat-http.service.spec.ts @@ -334,6 +334,62 @@ describe('ChatHttpService', () => { }); }); + describe('title fallback polling', () => { + // A stream that outran title generation closes before Nova Micro answers, + // so the first read sees the placeholder. Short retries pick the title up + // as soon as it is written instead of after one fixed 1.5s wait. + let sessionSvc: any; + + beforeEach(() => { + vi.useFakeTimers(); + sessionSvc = TestBed.inject(SessionService) as any; + sessionSvc.isNewSession.mockReturnValue(true); + sessionSvc.applyServerTitle = vi.fn(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it('applies the title on the first retry that sees it', async () => { + sessionSvc.getSessionMetadata = vi + .fn() + .mockResolvedValueOnce({ title: 'New Conversation' }) + .mockResolvedValueOnce({ title: 'New Conversation' }) + .mockResolvedValue({ title: 'Biology Syllabus' }); + + void (service as any).refreshTitleFromServer('s1'); + await vi.advanceTimersByTimeAsync(800); + + expect(sessionSvc.getSessionMetadata).toHaveBeenCalledTimes(3); + expect(sessionSvc.applyServerTitle).toHaveBeenCalledWith('s1', 'Biology Syllabus'); + await vi.runAllTimersAsync(); + expect(sessionSvc.getSessionMetadata).toHaveBeenCalledTimes(3); + }); + + it('gives up after a bounded number of retries', async () => { + sessionSvc.getSessionMetadata = vi.fn().mockResolvedValue({ title: 'New Conversation' }); + + void (service as any).refreshTitleFromServer('s1'); + await vi.runAllTimersAsync(); + + expect(sessionSvc.getSessionMetadata).toHaveBeenCalledTimes(6); + expect(sessionSvc.applyServerTitle).not.toHaveBeenCalled(); + }); + + it('stops polling once the title arrived another way', async () => { + sessionSvc.getSessionMetadata = vi.fn().mockResolvedValue({ title: 'New Conversation' }); + + void (service as any).refreshTitleFromServer('s1'); + await vi.advanceTimersByTimeAsync(0); + // A late `session_title` event applied the title and cleared the flag. + sessionSvc.isNewSession.mockReturnValue(false); + await vi.runAllTimersAsync(); + + expect(sessionSvc.getSessionMetadata).toHaveBeenCalledTimes(1); + }); + }); + describe('page-departure attribution', () => { // A refresh / tab close / navigation is the one interruption cause only // the browser witnesses. Unattested it lands server-side as diff --git a/frontend/ai.client/src/app/session/services/chat/chat-http.service.ts b/frontend/ai.client/src/app/session/services/chat/chat-http.service.ts index c8895ed5d..48b923e09 100644 --- a/frontend/ai.client/src/app/session/services/chat/chat-http.service.ts +++ b/frontend/ai.client/src/app/session/services/chat/chat-http.service.ts @@ -11,6 +11,15 @@ import { SessionService } from '../session/session.service'; import { ErrorService } from '../../../services/error/error.service'; import { isPreviewSession } from '../../../shared/constants/session.constants'; +/** + * Title fallback polling. A stream that outruns title generation closes a few + * hundred ms before Nova Micro answers, so one early read usually sees the + * placeholder; short, bounded retries pick the title up as soon as it lands + * instead of after a single fixed wait. 5 x 400ms keeps the old ~1.5s window. + */ +const TITLE_REFRESH_RETRY_MS = 400; +const TITLE_REFRESH_MAX_RETRIES = 5; + class RetriableError extends Error { constructor(message?: string) { super(message); @@ -509,18 +518,22 @@ export class ChatHttpService { * and is normally PUSHED mid-stream as a `session_title` SSE event; this * fetch only runs when the stream closed while the session still looked * new (see onclose). On a "New Conversation" placeholder — generation - * still in flight or failed — we retry once after a short delay before - * giving up. + * still in flight or failed — we poll a few more times at a short interval + * before giving up, stopping early if the title arrived another way. */ - private async refreshTitleFromServer(sessionId: string, retried = false): Promise { + private async refreshTitleFromServer(sessionId: string, attempt = 0): Promise { try { const metadata = await this.sessionService.getSessionMetadata(sessionId); if (metadata.title && metadata.title !== 'New Conversation') { this.sessionService.applyServerTitle(sessionId, metadata.title); return; } - if (!retried) { - setTimeout(() => this.refreshTitleFromServer(sessionId, true), 1500); + if (attempt < TITLE_REFRESH_MAX_RETRIES) { + setTimeout(() => { + if (this.sessionService.isNewSession(sessionId)) { + void this.refreshTitleFromServer(sessionId, attempt + 1); + } + }, TITLE_REFRESH_RETRY_MS); } } catch (error) { console.error('Failed to refresh session title:', error);