From 98c26d5370909f670b59a49ac1a5ffa49b079c1a Mon Sep 17 00:00:00 2001 From: "fabro-sh-0530[bot]" <281434857+fabro-sh-0530[bot]@users.noreply.github.com> Date: Sun, 24 May 2026 15:48:21 -0400 Subject: [PATCH] Fold context-window data into `agent.message`, remove snapshot event (#390) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Summary Removes the standalone `agent.context_window.snapshot` event and instead attaches the context-window projection directly to `agent.message`. This eliminates the async provider token-count API calls that the old approach required, and simplifies the event log to a single event type carrying all post-response agent data. ## What Changed and Why **Before:** After each LLM turn, the agent emitted a separate `agent.context_window.snapshot` event — first a local estimate, then potentially a second one after an async `count_input_tokens` call resolved (or after response usage arrived). This required fingerprint deduplication state, a `close_token` to cancel in-flight counts, and frontend handling for the extra event type. **After:** The `AgentEvent::AssistantMessage` variant carries an `Option`. The projection is computed locally at request-build time and then refined using response token usage when available (`ResponseUsageScaledBreakdown`), or kept as a `LocalEstimate` when response usage is absent. No provider API calls are made. ### Plan Summary - **Task 1:** Added `context_window: Option` to `AgentMessageProps` (Rust types + OpenAPI), removed `AgentContextWindowSnapshotProps` and `EventBody::AgentContextWindowSnapshot`. - **Task 2:** Removed the spawned `count_input_tokens` task, `close_token`, fingerprint sets, and both snapshot-emit methods from `Session`. Added `context_window_from_response_usage` to `context_window.rs`; `BuiltRequest` now holds the local projection instead of the tool list. - **Task 3:** Workflow conversion copies `context_window` from `AgentEvent::AssistantMessage` into `AgentMessageProps`; store reducer reads it from `AgentMessage` instead of the removed snapshot variant and stamps `event_seq`. - **Task 4:** GET endpoint tests updated to seed data via `agent.message` with embedded context-window; endpoint behavior unchanged. - **Task 5:** Frontend constant and tests for `agent.context_window.snapshot` removed; `agent.message` already invalidates `stageContextWindow` through existing stage-activity handling. TypeScript client regenerated with the new `AgentMessageProps` model. ### Key Design Decisions - **No provider token-count API calls** during normal execution — context-window accuracy relies on local estimates scaled by response usage, which is always available for successful turns. - **Failed-before-response turns** emit no context-window data (`context_window: None`), matching the old behavior where a snapshot would have been emitted but response-usage scaling would never arrive. - `BuiltRequest` drops the `tools` field (only needed for the now-removed snapshot emission path); the local projection is computed at build time and stored directly. ### Fabro Details
Ran 8 stages in 60m 3s for $55.78 | Stage | Duration | Cost | Retries | |---|---|---|---| | start | 0s | – | 0 | | toolchain | 1s | – | 0 | | preflight_compile | 2m 1s | – | 0 | | preflight_lint | 2m 16s | – | 0 | | implement | 30m 39s | $44.82 | 0 | | simplify_opus | 10m 48s | $4.03 | 0 | | simplify_gpt | 5m 1s | $6.93 | 0 | | verify | 8m 47s | – | 0 | | **Total** | **60m 3s** | **$55.78** | **0** |
Ran ImplementPlan.fabro (11 nodes and 14 edges) ```dot digraph ImplementPlan { graph [ goal="Implement and simplify", model_stylesheet=" * { model: claude-opus-4-7; } " ] rankdir=LR start [shape=Mdiamond, label="Start"] exit [shape=Msquare, label="Exit"] toolchain [label="Toolchain", shape=parallelogram, script="command -v cargo >/dev/null || { curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y && sudo ln -sf $HOME/.cargo/bin/* /usr/local/bin/; }; cargo --version 2>&1", max_retries=0] preflight_compile [label="Preflight Compile", shape=parallelogram, script="cargo check -q --workspace 2>&1", max_retries=0] preflight_lint [label="Preflight Lint", shape=parallelogram, script="cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1", max_retries=0] fix_lints [label="Fix Lints", prompt="The preflight lint step failed. Read the build output from context and fix all clippy lint warnings.", max_visits=3] implement [label="Implement", prompt="Read the plan file referenced in the goal and implement every step. Make all the code changes described in the plan. Use red/green TDD.", model="gpt-55", reasoning_effort="xhigh"] simplify_opus [label="Simplify (Opus)", prompt="@prompts/simplify.md"] simplify_gpt [label="Simplify (GPT-55)", prompt="@prompts/simplify.md", model="gpt-55"] verify [label="Verify", shape=parallelogram, script="git fetch origin main 2>&1 && git merge --no-edit --no-stat origin/main 2>&1 && cargo +nightly-2026-04-14 fmt --all 2>&1 && cargo dev docs refresh 2>&1 && cargo +nightly-2026-04-14 fmt --check --all 2>&1 && { command -v rg >/dev/null 2>&1 || { echo 'rg is required for verify'; exit 127; }; } && ! rg -n 'AuthMode::Disabled|RunAuthMethod|RunSubjectProvenance|\bActorRef\b|\bActorKind\b|AuthenticatedSubject|AuthenticatedService|AuthorizeRunScoped|AuthorizeRunBlob|AuthorizeStageArtifact|AuthorizeCommandLog|auth_method\s*==\s*\"disabled\"' lib/crates apps lib/packages docs/public/api-reference/fabro-api.yaml 2>&1 && cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings 2>&1 && cargo nextest run --workspace --status-level slow --profile ci 2>&1 && cargo dev docs check 2>&1 && bun install --frozen-lockfile 2>&1 && (cd apps/fabro-web && bun run typecheck) 2>&1 && (cd apps/fabro-web && bun run test) 2>&1 && (cd lib/packages/fabro-api-client && bun run typecheck) 2>&1 && cargo dev build -- -p fabro-cli --release 2>&1", goal_gate=true, retry_target="fixup"] fixup [label="Fixup", prompt="The verify step failed. Read the build output from context and fix all format, clippy, Rust test, docs, TypeScript typecheck/test, and build failures.", max_visits=3] start -> toolchain toolchain -> preflight_compile [condition="outcome=succeeded"] toolchain -> exit preflight_compile -> preflight_lint [condition="outcome=succeeded"] preflight_compile -> exit preflight_lint -> implement [condition="outcome=succeeded"] preflight_lint -> fix_lints fix_lints -> preflight_lint implement -> simplify_opus -> simplify_gpt -> verify verify -> exit [condition="outcome=succeeded"] verify -> fixup fixup -> verify } ```
⚒️ Generated with [Fabro](https://fabro.sh) --------- Co-authored-by: Fabro Co-authored-by: Bryan Helmkamp --- apps/fabro-web/app/lib/query-keys.test.ts | 10 - apps/fabro-web/app/lib/run-events.test.tsx | 10 - apps/fabro-web/app/lib/run-events.ts | 10 - docs/public/api-reference/fabro-api.yaml | 34 ++ lib/crates/fabro-agent/src/context_window.rs | 45 ++- lib/crates/fabro-agent/src/session.rs | 318 ++++-------------- lib/crates/fabro-agent/src/types.rs | 15 +- .../src/commands/run/run_progress/mod.rs | 1 + lib/crates/fabro-server/src/demo/mod.rs | 2 + .../fabro-server/src/server/handler/pair.rs | 2 + lib/crates/fabro-server/src/server/tests.rs | 20 +- lib/crates/fabro-store/src/run_state.rs | 86 ++--- lib/crates/fabro-types/src/run_event/agent.rs | 16 +- lib/crates/fabro-types/src/run_event/mod.rs | 120 +++++-- .../fabro-workflow/src/event/convert.rs | 62 +++- lib/crates/fabro-workflow/src/event/names.rs | 1 - lib/packages/fabro-api-client/package.json | 2 +- .../src/.openapi-generator/FILES | 1 + .../src/models/agent-message-props.ts | 37 ++ .../fabro-api-client/src/models/index.ts | 1 + 20 files changed, 381 insertions(+), 412 deletions(-) create mode 100644 lib/packages/fabro-api-client/src/models/agent-message-props.ts diff --git a/apps/fabro-web/app/lib/query-keys.test.ts b/apps/fabro-web/app/lib/query-keys.test.ts index 914f3b11f..8ad4a4a3a 100644 --- a/apps/fabro-web/app/lib/query-keys.test.ts +++ b/apps/fabro-web/app/lib/query-keys.test.ts @@ -92,16 +92,6 @@ describe("queryKeys", () => { } }); - test("context-window snapshot invalidates context window, run events, and stage events", () => { - expect( - queryKeysForRunEvent("run-1", "agent.context_window.snapshot", "stage-1"), - ).toEqual([ - queryKeys.runs.events("run-1", 1000), - queryKeys.runs.stageEvents("run-1", "stage-1"), - queryKeys.runs.stageContextWindow("run-1", "stage-1"), - ]); - }); - test("agent activity events without a node_id invalidate nothing", () => { expect(queryKeysForRunEvent("run-1", "agent.message")).toEqual([]); }); diff --git a/apps/fabro-web/app/lib/run-events.test.tsx b/apps/fabro-web/app/lib/run-events.test.tsx index 0a33cc3d5..b73d272f1 100644 --- a/apps/fabro-web/app/lib/run-events.test.tsx +++ b/apps/fabro-web/app/lib/run-events.test.tsx @@ -92,16 +92,6 @@ describe("queryKeysForRunEvent", () => { ]); }); - test("context-window snapshots invalidate context window, run events, and stage events", () => { - expect( - queryKeysForRunEvent("run-1", "agent.context_window.snapshot", "agent@1"), - ).toEqual([ - queryKeys.runs.events("run-1", 1000), - queryKeys.runs.stageEvents("run-1", "agent@1"), - queryKeys.runs.stageContextWindow("run-1", "agent@1"), - ]); - }); - test("todo events invalidate run state and run events", () => { for (const event of ["todo.created", "todo.updated", "todo.deleted"]) { expect(queryKeysForRunEvent("run-1", event)).toEqual([ diff --git a/apps/fabro-web/app/lib/run-events.ts b/apps/fabro-web/app/lib/run-events.ts index c88c76341..8efd40547 100644 --- a/apps/fabro-web/app/lib/run-events.ts +++ b/apps/fabro-web/app/lib/run-events.ts @@ -107,7 +107,6 @@ const TODO_EVENTS = new Set([ "todo.updated", "todo.deleted", ]); -const CONTEXT_WINDOW_SNAPSHOT_EVENT = "agent.context_window.snapshot"; export function queryKeysForRunEvent( runId: string, @@ -169,15 +168,6 @@ export function queryKeysForRunEvent( return keys; } - if (event === CONTEXT_WINDOW_SNAPSHOT_EVENT) { - const keys: Key[] = [queryKeys.runs.events(runId, 1000)]; - if (stageId) { - keys.push(queryKeys.runs.stageEvents(runId, stageId)); - keys.push(queryKeys.runs.stageContextWindow(runId, stageId)); - } - return keys; - } - if (STAGE_ACTIVITY_EVENTS.has(event)) { return stageId ? [ diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 58a1e02e9..e9ba9c616 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -7971,6 +7971,40 @@ components: type: integer minimum: 1 + AgentMessageProps: + description: Properties for the `agent.message` event. + type: object + required: + - text + - model + - billing + - tool_call_count + - visit + properties: + text: + type: string + model: + $ref: "#/components/schemas/BillingModelRef" + billing: + $ref: "#/components/schemas/BilledTokenCounts" + tool_call_count: + type: integer + minimum: 0 + visit: + type: integer + minimum: 1 + message: + oneOf: + - type: object + additionalProperties: true + - type: "null" + description: Canonical replay-authoritative transcript message, when present. + context_window: + oneOf: + - $ref: "#/components/schemas/StageContextWindowProjection" + - type: "null" + description: Latest content-free context-window projection for this agent stage. + RunSupersededByProps: description: Properties for the `run.superseded_by` audit event emitted on a rewound source run after archive succeeds. type: object diff --git a/lib/crates/fabro-agent/src/context_window.rs b/lib/crates/fabro-agent/src/context_window.rs index da8da9d72..ae2cdd343 100644 --- a/lib/crates/fabro-agent/src/context_window.rs +++ b/lib/crates/fabro-agent/src/context_window.rs @@ -5,7 +5,7 @@ use fabro_llm::token_count::{ estimate_message_tokens, estimate_request_control_tokens, estimate_text_tokens, estimate_tool_definition_tokens, is_local_estimator_warning, }; -use fabro_llm::types::{Request, Role, Warning as LlmWarning}; +use fabro_llm::types::{Request, Role, TokenCounts, Warning as LlmWarning}; use fabro_types::{ StageContextWindowBreakdownItem, StageContextWindowCategory, StageContextWindowCountMethod, StageContextWindowProjection, StageContextWindowStaleness, StageContextWindowWarning, @@ -16,7 +16,7 @@ use crate::skills::{Skill, format_skills_prompt_section}; use crate::tool_registry::{ToolDefinitionWithSource, ToolSource}; #[derive(Clone, Copy)] -pub(crate) struct ContextWindowSnapshotInput<'a> { +pub(crate) struct ContextWindowInput<'a> { pub request: &'a Request, pub tools: &'a [ToolDefinitionWithSource], pub system_prompt: &'a str, @@ -29,9 +29,7 @@ pub(crate) struct ContextWindowSnapshotInput<'a> { } #[must_use] -pub(crate) fn build_local_snapshot( - input: ContextWindowSnapshotInput<'_>, -) -> StageContextWindowProjection { +pub(crate) fn build_local_snapshot(input: ContextWindowInput<'_>) -> StageContextWindowProjection { let mut builder = BreakdownBuilder::default(); let mut warnings = Vec::new(); @@ -117,8 +115,31 @@ const fn total_is_provider_authoritative(method: StageContextWindowCountMethod) ) } +/// Build a projection from a previously-computed local snapshot and the +/// token usage returned by the LLM response. If the response carried no +/// usable input tokens, fall back to the local estimate unchanged. #[must_use] -pub(crate) fn warnings_from_llm(warnings: &[LlmWarning]) -> Vec { +pub(crate) fn context_window_from_response_usage( + local_snapshot: &StageContextWindowProjection, + usage: &TokenCounts, +) -> StageContextWindowProjection { + let input_tokens = usage + .input_tokens + .saturating_add(usage.cache_read_tokens) + .saturating_add(usage.cache_write_tokens); + if input_tokens <= 0 { + return local_snapshot.clone(); + } + scaled_snapshot( + local_snapshot, + u64::try_from(input_tokens).unwrap_or(u64::MAX), + StageContextWindowCountMethod::ResponseUsageScaledBreakdown, + local_snapshot.warnings.clone(), + ) +} + +#[must_use] +fn warnings_from_llm(warnings: &[LlmWarning]) -> Vec { warnings .iter() .map(|warning| StageContextWindowWarning { @@ -131,18 +152,10 @@ pub(crate) fn warnings_from_llm(warnings: &[LlmWarning]) -> Vec StageContextWindowWarning { - StageContextWindowWarning { - code: code.to_string(), - message: message.to_string(), - } -} - fn add_message_breakdown( builder: &mut BreakdownBuilder, warnings: &mut Vec, - input: &ContextWindowSnapshotInput<'_>, + input: &ContextWindowInput<'_>, ) { let memory_text = memory_prompt_suffix(input.memory); let skills_text = skills_prompt_suffix(input.skills); @@ -394,7 +407,7 @@ mod tests { tools.iter().map(|tool| tool.definition.clone()).collect(), ); - let snapshot = build_local_snapshot(ContextWindowSnapshotInput { + let snapshot = build_local_snapshot(ContextWindowInput { request: &req, tools: &tools, system_prompt: &system_prompt, diff --git a/lib/crates/fabro-agent/src/session.rs b/lib/crates/fabro-agent/src/session.rs index 5c98a9123..ee77ad219 100644 --- a/lib/crates/fabro-agent/src/session.rs +++ b/lib/crates/fabro-agent/src/session.rs @@ -1,5 +1,4 @@ -use std::collections::{HashMap, HashSet, VecDeque}; -use std::hash::{Hash, Hasher}; +use std::collections::{HashMap, VecDeque}; use std::sync::{Arc, Mutex, RwLock}; use std::time::SystemTime; @@ -8,10 +7,9 @@ use fabro_llm::client::Client; use fabro_llm::error::ProviderErrorKind; use fabro_llm::generate::StreamAccumulator; use fabro_llm::provider::StreamEventStream; -use fabro_llm::token_count::{InputTokenCountMethod, InputTokenCountPreference}; use fabro_llm::types::{ ContentPart, Message as LlmMessage, ReasoningEffort, Request, RetryPolicy, StreamEvent, - TokenCounts, ToolChoice, + ToolChoice, }; use fabro_llm::{Error as LlmError, retry}; use fabro_mcp::config::{McpServerSettings, McpTransport}; @@ -26,13 +24,13 @@ use futures::StreamExt; use tokio::sync::{Mutex as AsyncMutex, Notify, broadcast}; use tokio::time; use tokio_util::sync::CancellationToken; -use tracing::{Instrument as _, Span, debug, info, warn}; +use tracing::{debug, info, warn}; use crate::agent_profile::AgentProfile; use crate::compaction::{check_context_usage, compact_context}; use crate::config::SessionOptions; use crate::context_window::{ - ContextWindowSnapshotInput, build_local_snapshot, scaled_snapshot, warning, warnings_from_llm, + ContextWindowInput, build_local_snapshot, context_window_from_response_usage, }; use crate::error::{Error, InterruptReason}; use crate::event::Emitter; @@ -48,7 +46,6 @@ use crate::skills::{ }; use crate::subagent::{SubAgentCallbackEvent, SubAgentEventCallback, SubAgentManager}; use crate::tool_execution::execute_tool_calls; -use crate::tool_registry::ToolDefinitionWithSource; use crate::types::{ AgentEvent, McpToolSummary, MemoryFileSummary, Message, SessionEvent, SessionState, SkillActivationSource, SkillSummary, @@ -306,13 +303,8 @@ impl ToolEnvProvider for StaticEnvProvider { } struct BuiltRequest { - request: Request, - tools: Vec, -} - -struct EmittedContextWindowSnapshot { - local_snapshot: StageContextWindowProjection, - fingerprint: Option, + request: Request, + context_window: StageContextWindowProjection, } pub struct Session { @@ -333,7 +325,6 @@ pub struct Session { control_notify: Arc, followup_queue: Arc>>, cancel_token: CancellationToken, - close_token: CancellationToken, round_token: Arc>, interrupt_reason: Arc>>, memory: Vec, @@ -341,8 +332,6 @@ pub struct Session { skills: Vec, system_prompt: String, activated_skill_context_observed: bool, - context_window_counted_fingerprints: HashSet, - context_window_response_usage_fingerprints: Arc>>, file_tracker: FileTracker, tool_env_provider: Option>, subagent_manager: Option>>, @@ -373,7 +362,6 @@ impl Session { control_notify: Arc::new(Notify::new()), followup_queue: Arc::new(Mutex::new(VecDeque::new())), cancel_token: CancellationToken::new(), - close_token: CancellationToken::new(), round_token: Arc::new(RwLock::new(CancellationToken::new())), interrupt_reason: Arc::new(Mutex::new(None)), memory: Vec::new(), @@ -381,8 +369,6 @@ impl Session { skills: Vec::new(), system_prompt: String::new(), activated_skill_context_observed: false, - context_window_counted_fingerprints: HashSet::new(), - context_window_response_usage_fingerprints: Arc::new(Mutex::new(HashSet::new())), file_tracker: FileTracker::default(), tool_env_provider: None, subagent_manager, @@ -1147,7 +1133,6 @@ impl Session { pub fn close(&mut self) -> bool { let was_open = self.state != SessionState::Closed; - self.close_token.cancel(); self.transition(SessionState::Closed); was_open } @@ -1347,7 +1332,7 @@ impl Session { // Build request let built_request = self.build_request(); - let context_window_snapshot = self.emit_context_window_snapshots(&built_request); + let local_context_window = built_request.context_window.clone(); let request = built_request.request; // Emit AssistantTextStart before LLM call @@ -1618,7 +1603,10 @@ impl Session { .cloned() .collect(); let usage = response.usage.clone(); - self.emit_response_usage_context_window_snapshot(&context_window_snapshot, &usage); + let context_window = Some(context_window_from_response_usage( + &local_context_window, + &usage, + )); self.history.push(Message::Assistant { content: text.clone(), @@ -1645,6 +1633,7 @@ impl Session { model, usage: response.usage.clone(), tool_call_count: tool_calls.len(), + context_window, }); // Post-response compaction: trim context after appending assistant turn @@ -1749,136 +1738,6 @@ impl Session { Ok(()) } - fn emit_context_window_snapshots( - &mut self, - built_request: &BuiltRequest, - ) -> EmittedContextWindowSnapshot { - let provider = self.provider_profile.provider_id().to_string(); - let model = self.provider_profile.model().to_string(); - let local_snapshot = build_local_snapshot(ContextWindowSnapshotInput { - request: &built_request.request, - tools: &built_request.tools, - system_prompt: &self.system_prompt, - memory: &self.memory, - skills: &self.skills, - activated_skill_context_observed: self.activated_skill_context_observed, - provider: &provider, - model: &model, - context_window_tokens: self.provider_profile.context_window_size(), - }); - self.event_emitter.emit( - self.id.clone(), - AgentEvent::ContextWindowSnapshot(local_snapshot.clone()), - ); - - let Some(fingerprint) = request_fingerprint(&built_request.request) else { - return EmittedContextWindowSnapshot { - local_snapshot, - fingerprint: None, - }; - }; - if !self.context_window_counted_fingerprints.insert(fingerprint) { - return EmittedContextWindowSnapshot { - local_snapshot, - fingerprint: Some(fingerprint), - }; - } - - let client = self.llm_client.clone(); - let request = built_request.request.clone(); - let session_id = self.id.clone(); - let emitter = self.event_emitter.clone(); - let local_for_count = local_snapshot.clone(); - let close_token = self.close_token.clone(); - let response_usage_fingerprints = - Arc::clone(&self.context_window_response_usage_fingerprints); - let span = Span::current(); - let count_task = async move { - let count_result = tokio::select! { - biased; - () = close_token.cancelled() => return, - result = client.count_input_tokens(&request, InputTokenCountPreference::PreferProvider) => result, - }; - if close_token.is_cancelled() { - return; - } - if response_usage_fingerprints - .lock() - .expect("context window response-usage fingerprint lock poisoned") - .contains(&fingerprint) - { - return; - } - let snapshot = match count_result { - Ok(count) if count.method == InputTokenCountMethod::ProviderApi => { - let input_tokens = u64::try_from(count.input_tokens.max(0)).unwrap_or(u64::MAX); - scaled_snapshot( - &local_for_count, - input_tokens, - fabro_types::StageContextWindowCountMethod::ProviderApiScaledBreakdown, - warnings_from_llm(&count.warnings), - ) - } - Ok(count) => { - let mut warnings = local_for_count.warnings.clone(); - warnings.extend(warnings_from_llm(&count.warnings)); - let input_tokens = u64::try_from(count.input_tokens.max(0)).unwrap_or(u64::MAX); - scaled_snapshot( - &local_for_count, - input_tokens, - fabro_types::StageContextWindowCountMethod::LocalEstimate, - warnings, - ) - } - Err(_) => { - let mut warnings = local_for_count.warnings.clone(); - warnings.push(warning( - "provider_token_count_unavailable", - "provider input token counting was unavailable; retained local estimate", - )); - let mut snapshot = local_for_count.clone(); - snapshot.warnings = warnings; - snapshot - } - }; - emitter.emit(session_id, AgentEvent::ContextWindowSnapshot(snapshot)); - }; - tokio::spawn(count_task.instrument(span)); - - EmittedContextWindowSnapshot { - local_snapshot, - fingerprint: Some(fingerprint), - } - } - - fn emit_response_usage_context_window_snapshot( - &self, - context_window_snapshot: &EmittedContextWindowSnapshot, - usage: &TokenCounts, - ) { - let input_tokens = usage - .input_tokens - .saturating_add(usage.cache_read_tokens) - .saturating_add(usage.cache_write_tokens); - if input_tokens <= 0 { - return; - } - if let Some(fingerprint) = context_window_snapshot.fingerprint { - self.context_window_response_usage_fingerprints - .lock() - .expect("context window response-usage fingerprint lock poisoned") - .insert(fingerprint); - } - let snapshot = scaled_snapshot( - &context_window_snapshot.local_snapshot, - u64::try_from(input_tokens).unwrap_or(u64::MAX), - fabro_types::StageContextWindowCountMethod::ResponseUsageScaledBreakdown, - context_window_snapshot.local_snapshot.warnings.clone(), - ); - self.event_emitter - .emit(self.id.clone(), AgentEvent::ContextWindowSnapshot(snapshot)); - } - async fn compact_if_needed(&mut self) { let Some(estimate) = check_context_usage( &self.system_prompt, @@ -2016,9 +1875,22 @@ impl Session { metadata: None, provider_options: None, }; + let provider = self.provider_profile.provider_id().to_string(); + let model = self.provider_profile.model().to_string(); + let context_window = build_local_snapshot(ContextWindowInput { + request: &request, + tools: &tools_with_source, + system_prompt: &self.system_prompt, + memory: &self.memory, + skills: &self.skills, + activated_skill_context_observed: self.activated_skill_context_observed, + provider: &provider, + model: &model, + context_window_tokens: self.provider_profile.context_window_size(), + }); BuiltRequest { request, - tools: tools_with_source, + context_window, } } @@ -2047,13 +1919,6 @@ const fn is_auth_error(err: &LlmError) -> bool { ) } -fn request_fingerprint(request: &Request) -> Option { - let bytes = serde_json::to_vec(request).ok()?; - let mut hasher = std::collections::hash_map::DefaultHasher::new(); - bytes.hash(&mut hasher); - Some(hasher.finish()) -} - /// Best-effort kill of a sandbox MCP server process group. Used when /// `start_sandbox_mcp_server` is cancelled after spawning a detached /// `setsid` child but before reporting readiness. Errors from the sandbox @@ -2080,15 +1945,13 @@ mod tests { use anyhow::Context as _; use fabro_llm::error::{ProviderErrorDetail, ProviderErrorKind}; use fabro_llm::provider::{ProviderAdapter, StreamEventStream}; - use fabro_llm::token_count::{InputTokenCount, InputTokenCountMethod}; use fabro_llm::types::{ ContentPart, ReasoningEffort, Request, Response, Role, StreamEvent, TokenCounts, ToolCall, ToolDefinition, }; + use fabro_types::StageContextWindowCountMethod; use futures::stream; - use tokio::sync::Notify; use tokio::time::{sleep, timeout}; - use tracing::{Instrument as _, subscriber}; use super::*; use crate::config::{ToolAccess, ToolAccessPolicy, ToolApprovalAdapter, ToolExposureMode}; @@ -2207,57 +2070,6 @@ mod tests { } } - #[derive(Clone, Debug, PartialEq, Eq)] - enum ObservedSpan { - Missing, - Name(String), - } - - struct SpanCheckingTokenCountProvider { - observed_span: Arc>>, - notify: Arc, - } - - #[async_trait::async_trait] - impl ProviderAdapter for SpanCheckingTokenCountProvider { - fn name(&self) -> &'static str { - "mock" - } - - async fn complete(&self, _request: &Request) -> Result { - Err(LlmError::Configuration { - message: "SpanCheckingTokenCountProvider does not implement complete()".into(), - source: None, - }) - } - - async fn stream(&self, _request: &Request) -> Result { - Ok(Box::pin(stream::iter( - ScriptedStreamProvider::events_for_response(text_response("OK")), - ))) - } - - async fn count_input_tokens( - &self, - request: &Request, - ) -> Result, LlmError> { - let span = Span::current() - .metadata() - .map_or(ObservedSpan::Missing, |metadata| { - ObservedSpan::Name(metadata.name().to_string()) - }); - *self.observed_span.lock().unwrap() = Some(span); - self.notify.notify_waiters(); - Ok(Some(InputTokenCount { - input_tokens: 12, - method: InputTokenCountMethod::ProviderApi, - provider: "mock".to_string(), - model: request.model.clone(), - warnings: vec![], - })) - } - } - async fn make_session_with_provider(provider: Arc) -> Session { make_session_with_provider_and_manager(provider, None).await } @@ -2642,10 +2454,15 @@ mod tests { .iter() .any(|e| matches!(e.event, AgentEvent::UserInput { .. })) ); - assert!( - events - .iter() - .any(|e| matches!(e.event, AgentEvent::AssistantMessage { .. })) + let assistant_context_window = events.iter().find_map(|e| match &e.event { + AgentEvent::AssistantMessage { context_window, .. } => context_window.as_ref(), + _ => None, + }); + let context_window = + assistant_context_window.expect("assistant message should carry context window data"); + assert_eq!( + context_window.count_method, + StageContextWindowCountMethod::ResponseUsageScaledBreakdown ); assert!( events @@ -2654,6 +2471,33 @@ mod tests { ); } + #[tokio::test] + async fn assistant_message_context_window_uses_local_estimate_without_response_usage() { + let mut session = make_session(vec![response_with_usage( + text_response("Hello"), + TokenCounts::default(), + )]) + .await; + let mut rx = session.subscribe(); + + session.process_input("Hi").await.unwrap(); + + let context_window = std::iter::from_fn(|| rx.try_recv().ok()).find_map(|event| { + if let AgentEvent::AssistantMessage { context_window, .. } = event.event { + context_window + } else { + None + } + }); + + let context_window = context_window.expect("assistant message should carry context window"); + assert_eq!( + context_window.count_method, + StageContextWindowCountMethod::LocalEstimate + ); + assert!(context_window.input_tokens > 0); + } + #[tokio::test] async fn tool_call_end_has_untruncated_output() { let mut registry = ToolRegistry::new(); @@ -3059,42 +2903,6 @@ mod tests { assert_eq!(request.reasoning_effort, Some(ReasoningEffort::High)); } - #[tokio::test] - async fn context_window_token_count_task_inherits_current_run_span() { - let subscriber = tracing_subscriber::fmt().with_test_writer().finish(); - let _guard = subscriber::set_default(subscriber); - - let observed_span = Arc::new(Mutex::new(None)); - let notify = Arc::new(Notify::new()); - let provider = Arc::new(SpanCheckingTokenCountProvider { - observed_span: Arc::clone(&observed_span), - notify: Arc::clone(¬ify), - }); - let client = make_client(provider).await; - let registry = ToolRegistry::new(); - let profile = Arc::new(TestProfile::with_context_window(registry, 200_000)); - let env = Arc::new(MockSandbox::default()); - let mut session = Session::new(client, profile, env, SessionOptions::default(), None); - - let run_span = tracing::info_span!("run", id = %"run_context_window"); - session - .process_input("Hi") - .instrument(run_span) - .await - .unwrap(); - - if observed_span.lock().unwrap().is_none() { - timeout(Duration::from_secs(1), notify.notified()) - .await - .expect("provider token count should run"); - } - - assert_eq!( - observed_span.lock().unwrap().clone(), - Some(ObservedSpan::Name("run".to_string())) - ); - } - #[tokio::test] async fn context_window_no_warning_under_threshold() { let responses = vec![text_response("OK")]; diff --git a/lib/crates/fabro-agent/src/types.rs b/lib/crates/fabro-agent/src/types.rs index 9aa1a014c..d6907880f 100644 --- a/lib/crates/fabro-agent/src/types.rs +++ b/lib/crates/fabro-agent/src/types.rs @@ -246,6 +246,8 @@ pub enum AgentEvent { model: ModelRef, usage: TokenCounts, tool_call_count: usize, + #[serde(default, skip_serializing_if = "Option::is_none")] + context_window: Option, }, TextDelta { delta: String, @@ -304,7 +306,6 @@ pub enum AgentEvent { delay_secs: f64, error: LlmError, }, - ContextWindowSnapshot(StageContextWindowProjection), SubAgentSpawned { agent_id: String, depth: usize, @@ -504,17 +505,6 @@ impl AgentEvent { "LLM request failed, retrying" ); } - Self::ContextWindowSnapshot(snapshot) => { - debug!( - session_id, - provider = snapshot.provider.as_str(), - model = snapshot.model.as_str(), - input_tokens = snapshot.input_tokens, - context_window_tokens = snapshot.context_window_tokens, - count_method = %snapshot.count_method, - "Context window snapshot" - ); - } Self::SubAgentSpawned { agent_id, depth, @@ -891,6 +881,7 @@ mod tests { }, usage: usage.clone(), tool_call_count: 2, + context_window: None, }; match &event { AgentEvent::AssistantMessage { diff --git a/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs b/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs index eee7c3789..b13f34dea 100644 --- a/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs +++ b/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs @@ -578,6 +578,7 @@ mod tests { }, usage: TokenCounts::default(), tool_call_count: 0, + context_window: None, }) } diff --git a/lib/crates/fabro-server/src/demo/mod.rs b/lib/crates/fabro-server/src/demo/mod.rs index f13adafd5..e2c9b1885 100644 --- a/lib/crates/fabro-server/src/demo/mod.rs +++ b/lib/crates/fabro-server/src/demo/mod.rs @@ -1430,6 +1430,7 @@ mod runs { tool_call_count: 0, visit: 1, message: None, + context_window: None, }), ), make_envelope( @@ -1498,6 +1499,7 @@ mod runs { tool_call_count: 0, visit: 1, message: None, + context_window: None, }), ), ] diff --git a/lib/crates/fabro-server/src/server/handler/pair.rs b/lib/crates/fabro-server/src/server/handler/pair.rs index f5170aa86..f60228376 100644 --- a/lib/crates/fabro-server/src/server/handler/pair.rs +++ b/lib/crates/fabro-server/src/server/handler/pair.rs @@ -926,6 +926,7 @@ mod tests { tool_call_count: 0, visit: 1, message: None, + context_window: None, }), ), ) @@ -957,6 +958,7 @@ mod tests { tool_call_count: 0, visit: 1, message: None, + context_window: None, }), ), ) diff --git a/lib/crates/fabro-server/src/server/tests.rs b/lib/crates/fabro-server/src/server/tests.rs index 6488ab364..bd9759c43 100644 --- a/lib/crates/fabro-server/src/server/tests.rs +++ b/lib/crates/fabro-server/src/server/tests.rs @@ -14,7 +14,7 @@ use fabro_config::bind::Bind; use fabro_interview::{ AnswerValue, ControlInterviewer, Interviewer, Question, WorkerControlMessage, }; -use fabro_llm::types::{Message as LlmMessage, Request as LlmRequest}; +use fabro_llm::types::{Message as LlmMessage, Request as LlmRequest, TokenCounts}; use fabro_model::catalog::LlmCatalogSettings; use fabro_model::{Catalog, ModelRef, ProviderId, ReasoningEffort, Speed}; use fabro_types::settings::ServerAuthMethod; @@ -3023,12 +3023,22 @@ fn stage_completed_event(node_id: &str) -> workflow_event::Event { fn context_window_event( stage: &str, visit: u32, - snapshot: StageContextWindowProjection, + context_window: StageContextWindowProjection, ) -> workflow_event::Event { workflow_event::Event::Agent { stage: stage.to_string(), visit, - event: fabro_agent::AgentEvent::ContextWindowSnapshot(snapshot), + event: fabro_agent::AgentEvent::AssistantMessage { + text: "assistant response".to_string(), + model: ModelRef { + provider: ProviderId::openai(), + model_id: "gpt-5.4".to_string(), + speed: None, + }, + usage: TokenCounts::default(), + tool_call_count: 0, + context_window: Some(context_window), + }, session_id: Some("session-1".to_string()), parent_session_id: None, tool_call_id: None, @@ -3045,7 +3055,7 @@ fn context_window_snapshot( context_window_tokens: 400_000, input_tokens, usage_percent: input_tokens as f64 * 100.0 / 400_000.0, - count_method: StageContextWindowCountMethod::ProviderApiScaledBreakdown, + count_method: StageContextWindowCountMethod::ResponseUsageScaledBreakdown, staleness: StageContextWindowStaleness::Live, generated_at: Utc::now(), event_seq: None, @@ -6602,7 +6612,7 @@ async fn get_run_stage_context_window_returns_live_projected_snapshot() { assert_eq!(body["stage_id"], "agent_node@1"); assert_eq!(body["available"], true); assert_eq!(body["provider"], "openai"); - assert_eq!(body["count_method"], "provider_api_scaled_breakdown"); + assert_eq!(body["count_method"], "response_usage_scaled_breakdown"); assert_eq!(body["staleness"], "live"); assert_eq!(body["input_tokens"], 123_456); assert_eq!(body["breakdown"][0]["category"], "conversation"); diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index d5854291e..e2ce31e07 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -402,6 +402,11 @@ impl RunProjectionReducer for RunProjection { }; stage.usage.add_counts(&props.billing); stage.model = Some(props.model.clone()); + if let Some(context_window) = &props.context_window { + let mut context_window = context_window.clone(); + context_window.event_seq = Some(event.seq); + stage.context_window = Some(context_window); + } } EventBody::AgentSessionActivated(props) => { let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq) @@ -615,15 +620,6 @@ impl RunProjectionReducer for RunProjection { } } } - EventBody::AgentContextWindowSnapshot(props) => { - let Some(stage) = stage_at_stored_or_visit(self, stored, props.visit, event.seq) - else { - return Ok(()); - }; - let mut snapshot = props.snapshot.clone(); - snapshot.event_seq = Some(event.seq); - stage.context_window = Some(snapshot); - } _ => {} } @@ -1227,14 +1223,14 @@ mod tests { use fabro_types::run_event::run::RunFailedProps; use fabro_types::run_event::{ AgentAcpCancelledProps, AgentAcpCompletedProps, AgentAcpStartedProps, - AgentAcpTimedOutProps, AgentContextWindowSnapshotProps, AgentMcpFailedProps, - AgentMcpReadyProps, AgentMcpToolSummary, AgentMessageProps, AgentSessionActivatedProps, - AgentSessionEndedProps, AgentSessionStartedProps, AgentSkillActivatedProps, - AgentSkillActivationSource, AgentSkillSummary, AgentSkillsDiscoveredProps, - AgentSubClosedProps, AgentSubCompletedProps, AgentSubFailedProps, AgentSubSpawnedProps, - AgentToolStartedProps, CheckpointCompletedProps, InterviewCompletedProps, InterviewOption, - InterviewStartedProps, RunCompletedProps, RunControlEffectProps, StageCompletedProps, - StageFailedProps, StagePromptProps, StageRetryingProps, StageStartedProps, + AgentAcpTimedOutProps, AgentMcpFailedProps, AgentMcpReadyProps, AgentMcpToolSummary, + AgentMessageProps, AgentSessionActivatedProps, AgentSessionEndedProps, + AgentSessionStartedProps, AgentSkillActivatedProps, AgentSkillActivationSource, + AgentSkillSummary, AgentSkillsDiscoveredProps, AgentSubClosedProps, AgentSubCompletedProps, + AgentSubFailedProps, AgentSubSpawnedProps, AgentToolStartedProps, CheckpointCompletedProps, + InterviewCompletedProps, InterviewOption, InterviewStartedProps, RunCompletedProps, + RunControlEffectProps, StageCompletedProps, StageFailedProps, StagePromptProps, + StageRetryingProps, StageStartedProps, }; use fabro_types::{ AgentBackend, BilledModelUsage, BilledTokenCounts, BlockedReason, Checkpoint, @@ -3462,6 +3458,7 @@ mod tests { tool_call_count: 0, visit: 1, message: None, + context_window: None, } } @@ -4816,7 +4813,7 @@ mod tests { } #[test] - fn context_window_snapshots_replace_latest_for_matching_stage() { + fn agent_messages_replace_latest_context_window_for_matching_stage() { let mut state = initialized_projection(); let stage_id = stage_id(); let first = context_window_snapshot(10); @@ -4825,22 +4822,14 @@ mod tests { state .apply_event(&test_stage_event( 7, - EventBody::AgentContextWindowSnapshot(AgentContextWindowSnapshotProps { - stage_id: stage_id.clone(), - visit: 1, - snapshot: first, - }), + EventBody::AgentMessage(agent_message_with_context_window(first)), stage_id.clone(), )) .unwrap(); state .apply_event(&test_stage_event( 8, - EventBody::AgentContextWindowSnapshot(AgentContextWindowSnapshotProps { - stage_id: stage_id.clone(), - visit: 1, - snapshot: second, - }), + EventBody::AgentMessage(agent_message_with_context_window(second)), stage_id.clone(), )) .unwrap(); @@ -4852,25 +4841,44 @@ mod tests { } #[test] - fn context_window_snapshot_does_not_update_other_stage() { + fn agent_message_without_context_window_preserves_existing_context_window() { let mut state = initialized_projection(); - let target = stage_id(); - let other = StageId::new("review", 1); + let stage_id = stage_id(); state .apply_event(&test_stage_event( 7, - EventBody::AgentContextWindowSnapshot(AgentContextWindowSnapshotProps { - stage_id: target.clone(), - visit: 1, - snapshot: context_window_snapshot(10), - }), - target.clone(), + EventBody::AgentMessage(agent_message_with_context_window( + context_window_snapshot(10), + )), + stage_id.clone(), + )) + .unwrap(); + state + .apply_event(&test_stage_event( + 8, + EventBody::AgentMessage(live_agent_message_props(live_counts(1, 1))), + stage_id.clone(), )) .unwrap(); - assert!(state.stage(&target).unwrap().context_window.is_some()); - assert!(state.stage(&other).is_none()); + let snapshot = state + .stage(&stage_id) + .unwrap() + .context_window + .as_ref() + .unwrap(); + assert_eq!(snapshot.input_tokens, 10); + assert_eq!(snapshot.event_seq, Some(7)); + } + + fn agent_message_with_context_window( + context_window: StageContextWindowProjection, + ) -> AgentMessageProps { + AgentMessageProps { + context_window: Some(context_window), + ..live_agent_message_props(live_counts(1, 1)) + } } fn context_window_snapshot(input_tokens: u64) -> StageContextWindowProjection { diff --git a/lib/crates/fabro-types/src/run_event/agent.rs b/lib/crates/fabro-types/src/run_event/agent.rs index d172b2c9e..b1e567e6e 100644 --- a/lib/crates/fabro-types/src/run_event/agent.rs +++ b/lib/crates/fabro-types/src/run_event/agent.rs @@ -6,7 +6,7 @@ use super::BilledTokenCounts; use crate::transcript::{ToolCall, ToolResult, TranscriptMessage}; use crate::{ MessageId, ModelRef, PairId, PairMessageId, PairSystemMessageKind, PermissionLevel, - StageContextWindowProjection, StageId, TurnId, + StageContextWindowProjection, TurnId, }; #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -77,6 +77,10 @@ pub struct AgentMessageProps { /// payloads so older events still deserialize. #[serde(default, skip_serializing_if = "Option::is_none")] pub message: Option, + /// Latest content-free context-window projection for this agent stage, + /// computed from the request that produced this assistant response. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub context_window: Option, } #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] @@ -218,14 +222,6 @@ pub struct AgentLlmRetryProps { pub visit: u32, } -#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] -pub struct AgentContextWindowSnapshotProps { - pub stage_id: StageId, - pub visit: u32, - #[serde(flatten)] - pub snapshot: StageContextWindowProjection, -} - #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] pub struct AgentSubSpawnedProps { pub agent_id: String, @@ -357,6 +353,7 @@ mod tests { let props: AgentMessageProps = serde_json::from_value(v).unwrap(); assert_eq!(props.text, "hello"); assert!(props.message.is_none()); + assert!(props.context_window.is_none()); } #[test] @@ -371,6 +368,7 @@ mod tests { tool_call_count: 0, visit: 1, message: Some(msg.clone()), + context_window: None, }; let v = serde_json::to_value(&props).unwrap(); assert_eq!(v["message"]["kind"], "agent"); diff --git a/lib/crates/fabro-types/src/run_event/mod.rs b/lib/crates/fabro-types/src/run_event/mod.rs index 84a58ef16..571a3d860 100644 --- a/lib/crates/fabro-types/src/run_event/mod.rs +++ b/lib/crates/fabro-types/src/run_event/mod.rs @@ -238,8 +238,6 @@ pub enum EventBody { AgentCompactionCompleted(AgentCompactionCompletedProps), #[serde(rename = "agent.llm.retry")] AgentLlmRetry(AgentLlmRetryProps), - #[serde(rename = "agent.context_window.snapshot")] - AgentContextWindowSnapshot(AgentContextWindowSnapshotProps), #[serde(rename = "agent.sub.spawned")] AgentSubSpawned(AgentSubSpawnedProps), #[serde(rename = "agent.sub.completed")] @@ -520,7 +518,6 @@ impl EventBody { Self::AgentCompactionStarted(_) => "agent.compaction.started", Self::AgentCompactionCompleted(_) => "agent.compaction.completed", Self::AgentLlmRetry(_) => "agent.llm.retry", - Self::AgentContextWindowSnapshot(_) => "agent.context_window.snapshot", Self::AgentSubSpawned(_) => "agent.sub.spawned", Self::AgentSubCompleted(_) => "agent.sub.completed", Self::AgentSubFailed(_) => "agent.sub.failed", @@ -702,7 +699,6 @@ fn is_known_event_name(event: &str) -> bool { | "agent.compaction.started" | "agent.compaction.completed" | "agent.llm.retry" - | "agent.context_window.snapshot" | "agent.sub.spawned" | "agent.sub.completed" | "agent.sub.failed" @@ -2156,43 +2152,95 @@ mod tests { } #[test] - fn agent_context_window_snapshot_serializes_with_canonical_name() { - let body = EventBody::AgentContextWindowSnapshot(AgentContextWindowSnapshotProps { - stage_id: crate::StageId::new("implement", 1), - visit: 1, - snapshot: crate::StageContextWindowProjection { - provider: "openai".to_string(), - model: "gpt-5.4".to_string(), - context_window_tokens: 400_000, - input_tokens: 123_456, - usage_percent: 30.864, - count_method: - crate::StageContextWindowCountMethod::ProviderApiScaledBreakdown, - staleness: crate::StageContextWindowStaleness::Live, - generated_at: DateTime::parse_from_rfc3339("2026-05-23T12:34:56Z") - .unwrap() - .with_timezone(&Utc), - event_seq: None, - breakdown: vec![crate::StageContextWindowBreakdownItem { - category: crate::StageContextWindowCategory::SystemPrompt, - tokens: 30_000, - usage_percent: 7.5, - }], - warnings: vec![crate::StageContextWindowWarning { - code: "local_token_estimate".to_string(), - message: "input token count is a local estimate".to_string(), - }], + fn agent_message_omits_context_window_when_absent() { + let body = EventBody::AgentMessage(AgentMessageProps { + text: "ok".to_string(), + model: crate::ModelRef { + provider: fabro_model::ProviderId::openai(), + model_id: "gpt-5.4".to_string(), + speed: None, }, + billing: BilledTokenCounts::default(), + tool_call_count: 0, + visit: 1, + message: None, + context_window: None, }); + let value = serde_json::to_value(&body).unwrap(); - assert_eq!(value["event"], "agent.context_window.snapshot"); - assert_eq!(value["properties"]["stage_id"], "implement@1"); - assert_eq!( - value["properties"]["breakdown"][0]["category"], - "system_prompt" + assert_eq!(value["event"], "agent.message"); + assert!( + value["properties"] + .as_object() + .unwrap() + .get("context_window") + .is_none() ); let parsed: EventBody = serde_json::from_value(value).unwrap(); - assert_eq!(parsed.event_name(), "agent.context_window.snapshot"); + assert_eq!(parsed.event_name(), "agent.message"); + } + + #[test] + fn agent_message_round_trips_optional_context_window() { + let context_window = crate::StageContextWindowProjection { + provider: "openai".to_string(), + model: "gpt-5.4".to_string(), + context_window_tokens: 400_000, + input_tokens: 123_456, + usage_percent: 30.864, + count_method: + crate::StageContextWindowCountMethod::ResponseUsageScaledBreakdown, + staleness: crate::StageContextWindowStaleness::Live, + generated_at: DateTime::parse_from_rfc3339("2026-05-23T12:34:56Z") + .unwrap() + .with_timezone(&Utc), + event_seq: None, + breakdown: vec![crate::StageContextWindowBreakdownItem { + category: crate::StageContextWindowCategory::SystemPrompt, + tokens: 30_000, + usage_percent: 7.5, + }], + warnings: vec![crate::StageContextWindowWarning { + code: "local_token_estimate".to_string(), + message: "input token count is a local estimate".to_string(), + }], + }; + let body = EventBody::AgentMessage(AgentMessageProps { + text: "ok".to_string(), + model: crate::ModelRef { + provider: fabro_model::ProviderId::openai(), + model_id: "gpt-5.4".to_string(), + speed: None, + }, + billing: BilledTokenCounts::default(), + tool_call_count: 0, + visit: 1, + message: None, + context_window: Some(context_window), + }); + + let value = serde_json::to_value(&body).unwrap(); + assert_eq!(value["event"], "agent.message"); + assert_eq!( + value["properties"]["context_window"]["breakdown"][0]["category"], + "system_prompt" + ); + assert_eq!( + value["properties"]["context_window"]["count_method"], + "response_usage_scaled_breakdown" + ); + let parsed: EventBody = serde_json::from_value(value).unwrap(); + match parsed { + EventBody::AgentMessage(props) => { + let context_window = props.context_window.expect("context window present"); + assert_eq!(context_window.input_tokens, 123_456); + assert_eq!( + context_window.count_method, + crate::StageContextWindowCountMethod::ResponseUsageScaledBreakdown + ); + } + other => panic!("expected AgentMessage body, got {other:?}"), + } } #[test] diff --git a/lib/crates/fabro-workflow/src/event/convert.rs b/lib/crates/fabro-workflow/src/event/convert.rs index ff098e6f0..3f4b90098 100644 --- a/lib/crates/fabro-workflow/src/event/convert.rs +++ b/lib/crates/fabro-workflow/src/event/convert.rs @@ -591,7 +591,7 @@ fn event_body_from_event(event: &Event) -> EventBody { billing: billing.clone(), }), Event::Agent { - stage, + stage: _, visit, event, .. @@ -610,6 +610,7 @@ fn event_body_from_event(event: &Event) -> EventBody { model, usage, tool_call_count, + context_window, } => { let billing = billed_token_counts_from_llm(usage); EventBody::AgentMessage(fabro_types::AgentMessageProps { @@ -619,6 +620,7 @@ fn event_body_from_event(event: &Event) -> EventBody { tool_call_count: *tool_call_count, visit: *visit, message: None, + context_window: context_window.clone(), }) } AgentEvent::ToolCallStarted { @@ -711,13 +713,6 @@ fn event_body_from_event(event: &Event) -> EventBody { error: serde_json::to_value(error).expect("serializable sdk error"), visit: *visit, }), - AgentEvent::ContextWindowSnapshot(snapshot) => EventBody::AgentContextWindowSnapshot( - fabro_types::AgentContextWindowSnapshotProps { - stage_id: ::fabro_types::StageId::new(stage.clone(), *visit), - visit: *visit, - snapshot: snapshot.clone(), - }, - ), AgentEvent::SubAgentSpawned { agent_id, depth, @@ -2172,6 +2167,7 @@ mod tests { }, usage: LlmTokenCounts::default(), tool_call_count: 0, + context_window: None, }, session_id: Some("ses_agent".to_string()), parent_session_id: None, @@ -2203,6 +2199,7 @@ mod tests { ..LlmTokenCounts::default() }, tool_call_count: 0, + context_window: None, }, session_id: Some("ses_agent".to_string()), parent_session_id: None, @@ -2219,6 +2216,55 @@ mod tests { assert_eq!(message.billing.total_usd_micros, None); } + #[test] + fn agent_assistant_message_copies_context_window_to_props() { + let context_window = ::fabro_types::StageContextWindowProjection { + provider: "openai".to_string(), + model: "gpt-5.4".to_string(), + context_window_tokens: 400_000, + input_tokens: 123, + usage_percent: 0.03075, + count_method: ::fabro_types::StageContextWindowCountMethod::LocalEstimate, + staleness: ::fabro_types::StageContextWindowStaleness::Live, + generated_at: Utc::now(), + event_seq: None, + breakdown: vec![::fabro_types::StageContextWindowBreakdownItem { + category: ::fabro_types::StageContextWindowCategory::Conversation, + tokens: 123, + usage_percent: 0.03075, + }], + warnings: Vec::new(), + }; + let stored = to_run_event(&fixtures::RUN_1, &Event::Agent { + stage: "code".to_string(), + visit: 1, + event: AgentEvent::AssistantMessage { + text: "ok".to_string(), + model: ModelRef { + provider: ProviderId::openai(), + model_id: "gpt-5.4".to_string(), + speed: None, + }, + usage: LlmTokenCounts::default(), + tool_call_count: 0, + context_window: Some(context_window), + }, + session_id: Some("ses_agent".to_string()), + parent_session_id: None, + tool_call_id: None, + }); + + let EventBody::AgentMessage(message) = stored.body else { + panic!("expected agent message body"); + }; + let context_window = message.context_window.expect("context window copied"); + assert_eq!(context_window.input_tokens, 123); + assert_eq!( + context_window.count_method, + ::fabro_types::StageContextWindowCountMethod::LocalEstimate + ); + } + #[test] fn agent_acp_events_map_to_event_bodies_with_stage_scope() { let scope = StageScope { diff --git a/lib/crates/fabro-workflow/src/event/names.rs b/lib/crates/fabro-workflow/src/event/names.rs index 958f8fdeb..c4d19d832 100644 --- a/lib/crates/fabro-workflow/src/event/names.rs +++ b/lib/crates/fabro-workflow/src/event/names.rs @@ -86,7 +86,6 @@ pub fn event_name(event: &Event) -> &'static str { AgentEvent::CompactionStarted { .. } => "agent.compaction.started", AgentEvent::CompactionCompleted { .. } => "agent.compaction.completed", AgentEvent::LlmRetry { .. } => "agent.llm.retry", - AgentEvent::ContextWindowSnapshot(_) => "agent.context_window.snapshot", AgentEvent::SubAgentSpawned { .. } => "agent.sub.spawned", AgentEvent::SubAgentCompleted { .. } => "agent.sub.completed", AgentEvent::SubAgentFailed { .. } => "agent.sub.failed", diff --git a/lib/packages/fabro-api-client/package.json b/lib/packages/fabro-api-client/package.json index 7709ed33a..de8ac9883 100644 --- a/lib/packages/fabro-api-client/package.json +++ b/lib/packages/fabro-api-client/package.json @@ -4,7 +4,7 @@ "private": true, "type": "module", "scripts": { - "generate": "bunx @openapitools/openapi-generator-cli generate -i ../../../docs/public/api-reference/fabro-api.yaml -g typescript-axios --additional-properties=supportsES6=true,typescriptThreePlus=true,withSeparateModelsAndApi=true,apiPackage=api,modelPackage=models,useTags=true,enumPropertyNaming=UPPERCASE -o src && bun run scripts/normalize-generated.ts", + "generate": "bunx @openapitools/openapi-generator-cli@2.20.2 generate -i ../../../docs/public/api-reference/fabro-api.yaml -g typescript-axios --additional-properties=supportsES6=true,typescriptThreePlus=true,withSeparateModelsAndApi=true,apiPackage=api,modelPackage=models,useTags=true,enumPropertyNaming=UPPERCASE -o src && bun run scripts/normalize-generated.ts", "typecheck": "tsc" }, "devDependencies": { diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index f8d333e14..61f64e490 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -23,6 +23,7 @@ configuration.ts index.ts models/activated-skill.ts models/agent-mcp-tool-summary.ts +models/agent-message-props.ts models/agent-permissions.ts models/agent-session-activated-props.ts models/agent-skill-activation-source.ts diff --git a/lib/packages/fabro-api-client/src/models/agent-message-props.ts b/lib/packages/fabro-api-client/src/models/agent-message-props.ts new file mode 100644 index 000000000..a0c268de2 --- /dev/null +++ b/lib/packages/fabro-api-client/src/models/agent-message-props.ts @@ -0,0 +1,37 @@ +/* tslint:disable */ +/* eslint-disable */ +/** + * Fabro Run API + * HTTP API for managing Fabro workflow run executions. + * + * The version of the OpenAPI document: 0.1.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + + +// May contain unused imports in some cases +// @ts-ignore +import type { BilledTokenCounts } from './billed-token-counts'; +// May contain unused imports in some cases +// @ts-ignore +import type { BillingModelRef } from './billing-model-ref'; +// May contain unused imports in some cases +// @ts-ignore +import type { StageContextWindowProjection } from './stage-context-window-projection'; + +/** + * Properties for the `agent.message` event. + */ +export interface AgentMessageProps { + 'text': string; + 'model': BillingModelRef; + 'billing': BilledTokenCounts; + 'tool_call_count': number; + 'visit': number; + 'message'?: { [key: string]: any; } | null; + 'context_window'?: StageContextWindowProjection | null; +} diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index 78d106d70..387306be1 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -1,5 +1,6 @@ export * from './activated-skill'; export * from './agent-mcp-tool-summary'; +export * from './agent-message-props'; export * from './agent-permissions'; export * from './agent-session-activated-props'; export * from './agent-skill-activation-source';