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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions eval/src/comparison.rs
Original file line number Diff line number Diff line change
Expand Up @@ -578,6 +578,7 @@ mod tests {
turns,
function_calls: 0,
function_call_errors: 0,
context_compactions: 0,
input_tokens: input,
output_tokens: output,
cache_read_tokens: None,
Expand Down
5 changes: 5 additions & 0 deletions harness/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,11 @@ with the optional [`approval-gate`](https://github.com/iii-hq/workers/tree/main/
The full function reference (every `harness::*` id and its request/response
schema) lives in the code and `iii worker info harness`.

`harness::metrics` reports `context_compactions` per session and for the full
root/descendant tree. It counts durable `compaction` custom records on each
session's active path; the latest `context.compacted` snapshot remains a
separate latest-generation indicator.

Building a consumer — a chat UI, a Telegram/WhatsApp bridge, a cron worker, or
any event-driven loop on top of the harness? Start with the integration
contract in
Expand Down
5 changes: 3 additions & 2 deletions harness/src/clients/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -370,7 +370,8 @@ impl SessionClient {
}

/// Strict evidence reader for evaluation metrics. Unlike [`Self::messages`],
/// malformed entries fail the read instead of being skipped.
/// malformed entries fail the read instead of being skipped. Custom entries
/// are included so metrics can count durable compaction records.
pub async fn messages_strict(
&self,
session_id: &str,
Expand All @@ -383,7 +384,7 @@ impl SessionClient {
loop {
let mut payload = json!({
"session_id": session_id,
"include_custom": false,
"include_custom": true,
"limit": PAGE_LIMIT,
});
if let Some(value) = &cursor {
Expand Down
65 changes: 61 additions & 4 deletions harness/src/functions/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::json;

use crate::clients::session::LoadedEntry;
use crate::context_snapshot::ContextSnapshotV1;
use crate::deps::Deps;
use crate::error::HarnessError;
Expand Down Expand Up @@ -42,6 +43,9 @@ pub struct SessionUsageTotalsV1 {
pub turns: u64,
pub function_calls: u64,
pub function_call_errors: u64,
/// Durable context compaction records across the active paths of the
/// root session and all descendants.
pub context_compactions: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub input_tokens: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
Expand All @@ -66,6 +70,8 @@ pub struct SessionUsageV1 {
pub turns: u64,
pub function_calls: u64,
pub function_call_errors: u64,
/// Durable context compaction records on this session's active path.
pub context_compactions: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub input_tokens: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
Expand Down Expand Up @@ -166,10 +172,8 @@ pub async fn handle(
let entries = session.messages_strict(&node.session_id).await?;
let mut current = UsageAccumulator::default();
for entry in entries {
if let Some(message) = entry.message.as_ref() {
current.observe(message);
total.observe(message);
}
current.observe_entry(&entry);
total.observe_entry(&entry);
}
// The turn record usually carries the session's latest snapshot for
// free. A freshly seeded turn has not generated yet, so its record has
Expand Down Expand Up @@ -325,6 +329,7 @@ struct UsageAccumulator {
turns: u64,
function_calls: u64,
function_call_errors: u64,
context_compactions: u64,
input: OptionalU64Sum,
output: OptionalU64Sum,
cache_read: OptionalU64Sum,
Expand All @@ -334,6 +339,19 @@ struct UsageAccumulator {
}

impl UsageAccumulator {
fn observe_entry(&mut self, entry: &LoadedEntry) {
if entry
.custom
.as_ref()
.is_some_and(|custom| custom.custom_type == "compaction")
{
self.context_compactions = self.context_compactions.saturating_add(1);
}
if let Some(message) = entry.message.as_ref() {
self.observe(message);
}
}

fn observe(&mut self, message: &AgentMessage) {
match message {
AgentMessage::Assistant(assistant) => {
Expand Down Expand Up @@ -374,6 +392,7 @@ impl UsageAccumulator {
turns: self.turns,
function_calls: self.function_calls,
function_call_errors: self.function_call_errors,
context_compactions: self.context_compactions,
input_tokens: self.input.finish(),
output_tokens: self.output.finish(),
cache_read_tokens: self.cache_read.finish(),
Expand All @@ -390,6 +409,7 @@ impl UsageAccumulator {
turns: self.turns,
function_calls: self.function_calls,
function_call_errors: self.function_call_errors,
context_compactions: self.context_compactions,
input_tokens: self.input.finish(),
output_tokens: self.output.finish(),
cache_read_tokens: self.cache_read.finish(),
Expand Down Expand Up @@ -453,6 +473,7 @@ mod tests {
use serde_json::json;

use super::*;
use crate::clients::session::LoadedCustom;
use crate::types::turn::TurnStatus;

fn message(value: serde_json::Value) -> AgentMessage {
Expand All @@ -468,6 +489,41 @@ mod tests {
}
}

fn custom_entry(custom_type: &str) -> LoadedEntry {
LoadedEntry {
entry_id: format!("entry-{custom_type}"),
message: None,
custom: Some(LoadedCustom {
custom_type: custom_type.into(),
data: json!({}),
}),
}
}

#[test]
fn counts_only_durable_context_compactions() {
let mut root = UsageAccumulator::default();
let mut total = UsageAccumulator::default();
for entry in [
custom_entry("compaction"),
custom_entry("checkpoint"),
custom_entry("compaction"),
] {
root.observe_entry(&entry);
total.observe_entry(&entry);
}
let root = root.finish_session(&session("root", None, 0), None);
assert_eq!(root.context_compactions, 2);

let mut child = UsageAccumulator::default();
let entry = custom_entry("compaction");
child.observe_entry(&entry);
total.observe_entry(&entry);
let child = child.finish_session(&session("child", Some("root"), 1), None);
assert_eq!(child.context_compactions, 1);
assert_eq!(total.finish_totals(2).context_compactions, 3);
}

#[test]
fn counts_generations_calls_errors_and_usage() {
let mut usage = UsageAccumulator::default();
Expand Down Expand Up @@ -499,6 +555,7 @@ mod tests {
assert_eq!(totals.turns, 1);
assert_eq!(totals.function_calls, 1);
assert_eq!(totals.function_call_errors, 1);
assert_eq!(totals.context_compactions, 0);
assert_eq!(totals.input_tokens, Some(10));
assert_eq!(totals.output_tokens, Some(2));
assert_eq!(totals.cost_usd, Some(0.1));
Expand Down
14 changes: 14 additions & 0 deletions harness/tests/golden/schemas/harness.metrics.json
Original file line number Diff line number Diff line change
Expand Up @@ -232,6 +232,12 @@
"null"
]
},
"context_compactions": {
"description": "Durable context compaction records across the active paths of the root session and all descendants.",
"format": "uint64",
"minimum": 0.0,
"type": "integer"
},
"cost_usd": {
"format": "double",
"type": [
Expand Down Expand Up @@ -285,6 +291,7 @@
}
},
"required": [
"context_compactions",
"function_call_errors",
"function_calls",
"sessions",
Expand Down Expand Up @@ -322,6 +329,12 @@
],
"description": "The session's latest per-generation context snapshot (categories, budget, usage) — absent for sessions that have not generated since snapshots landed."
},
"context_compactions": {
"description": "Durable context compaction records on this session's active path.",
"format": "uint64",
"minimum": 0.0,
"type": "integer"
},
"cost_usd": {
"format": "double",
"type": [
Expand Down Expand Up @@ -384,6 +397,7 @@
}
},
"required": [
"context_compactions",
"depth",
"function_call_errors",
"function_calls",
Expand Down
Loading