mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-11 03:40:05 +00:00
Fold context-window data into agent.message, remove snapshot event (#390)
## 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<StageContextWindowProjection>`. 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<StageContextWindowProjection>` 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
<details>
<summary>Ran 8 stages in 60m 3s for $55.78</summary>
| 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** |
</details>
<details>
<summary>Ran <code>ImplementPlan.fabro</code> (11 nodes and 14
edges)</summary>
```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
}
```
</details>
⚒️ Generated with [Fabro](https://fabro.sh)
---------
Co-authored-by: Fabro <noreply@fabro.sh>
Co-authored-by: Bryan Helmkamp <bryan@brynary.com>
This commit is contained in:
parent
b8cf39ab7b
commit
98c26d5370
20 changed files with 381 additions and 412 deletions
|
|
@ -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([]);
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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([
|
||||
|
|
|
|||
|
|
@ -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
|
||||
? [
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<StageContextWindowWarning> {
|
||||
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<StageContextWindowWarning> {
|
||||
warnings
|
||||
.iter()
|
||||
.map(|warning| StageContextWindowWarning {
|
||||
|
|
@ -131,18 +152,10 @@ pub(crate) fn warnings_from_llm(warnings: &[LlmWarning]) -> Vec<StageContextWind
|
|||
.collect()
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub(crate) fn warning(code: &str, message: &str) -> StageContextWindowWarning {
|
||||
StageContextWindowWarning {
|
||||
code: code.to_string(),
|
||||
message: message.to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
fn add_message_breakdown(
|
||||
builder: &mut BreakdownBuilder,
|
||||
warnings: &mut Vec<StageContextWindowWarning>,
|
||||
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,
|
||||
|
|
|
|||
|
|
@ -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<ToolDefinitionWithSource>,
|
||||
}
|
||||
|
||||
struct EmittedContextWindowSnapshot {
|
||||
local_snapshot: StageContextWindowProjection,
|
||||
fingerprint: Option<u64>,
|
||||
request: Request,
|
||||
context_window: StageContextWindowProjection,
|
||||
}
|
||||
|
||||
pub struct Session {
|
||||
|
|
@ -333,7 +325,6 @@ pub struct Session {
|
|||
control_notify: Arc<Notify>,
|
||||
followup_queue: Arc<Mutex<VecDeque<String>>>,
|
||||
cancel_token: CancellationToken,
|
||||
close_token: CancellationToken,
|
||||
round_token: Arc<RwLock<CancellationToken>>,
|
||||
interrupt_reason: Arc<Mutex<Option<InterruptReason>>>,
|
||||
memory: Vec<MemoryDocument>,
|
||||
|
|
@ -341,8 +332,6 @@ pub struct Session {
|
|||
skills: Vec<Skill>,
|
||||
system_prompt: String,
|
||||
activated_skill_context_observed: bool,
|
||||
context_window_counted_fingerprints: HashSet<u64>,
|
||||
context_window_response_usage_fingerprints: Arc<Mutex<HashSet<u64>>>,
|
||||
file_tracker: FileTracker,
|
||||
tool_env_provider: Option<Arc<dyn ToolEnvProvider>>,
|
||||
subagent_manager: Option<Arc<AsyncMutex<SubAgentManager>>>,
|
||||
|
|
@ -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<u64> {
|
||||
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<Mutex<Option<ObservedSpan>>>,
|
||||
notify: Arc<Notify>,
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl ProviderAdapter for SpanCheckingTokenCountProvider {
|
||||
fn name(&self) -> &'static str {
|
||||
"mock"
|
||||
}
|
||||
|
||||
async fn complete(&self, _request: &Request) -> Result<Response, LlmError> {
|
||||
Err(LlmError::Configuration {
|
||||
message: "SpanCheckingTokenCountProvider does not implement complete()".into(),
|
||||
source: None,
|
||||
})
|
||||
}
|
||||
|
||||
async fn stream(&self, _request: &Request) -> Result<StreamEventStream, LlmError> {
|
||||
Ok(Box::pin(stream::iter(
|
||||
ScriptedStreamProvider::events_for_response(text_response("OK")),
|
||||
)))
|
||||
}
|
||||
|
||||
async fn count_input_tokens(
|
||||
&self,
|
||||
request: &Request,
|
||||
) -> Result<Option<InputTokenCount>, 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<dyn ProviderAdapter>) -> 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")];
|
||||
|
|
|
|||
|
|
@ -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<StageContextWindowProjection>,
|
||||
},
|
||||
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 {
|
||||
|
|
|
|||
|
|
@ -578,6 +578,7 @@ mod tests {
|
|||
},
|
||||
usage: TokenCounts::default(),
|
||||
tool_call_count: 0,
|
||||
context_window: None,
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
}),
|
||||
),
|
||||
]
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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<TranscriptMessage>,
|
||||
/// 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<StageContextWindowProjection>,
|
||||
}
|
||||
|
||||
#[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");
|
||||
|
|
|
|||
|
|
@ -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]
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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": {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
37
lib/packages/fabro-api-client/src/models/agent-message-props.ts
generated
Normal file
37
lib/packages/fabro-api-client/src/models/agent-message-props.ts
generated
Normal file
|
|
@ -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;
|
||||
}
|
||||
|
|
@ -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';
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue