Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion CLAUDE.MD
Original file line number Diff line number Diff line change
Expand Up @@ -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` |
Expand Down
7 changes: 6 additions & 1 deletion backend/src/agents/main_agent/chat_agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
29 changes: 27 additions & 2 deletions backend/src/agents/main_agent/streaming/stream_coordinator.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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]:
Expand Down
11 changes: 8 additions & 3 deletions backend/src/apis/inference_api/chat/routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]:
Expand Down Expand Up @@ -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
Expand Down
60 changes: 49 additions & 11 deletions backend/src/apis/inference_api/chat/service.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
import logging
import hashlib
import os
import threading
from typing import Any, Callable, Dict, Optional, List, Tuple

import boto3
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -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 = {
Expand All @@ -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"],
Expand All @@ -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

Expand Down
Loading