From aacf8649a689db62e6d08d8b93bf61682934818a Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" Date: Sun, 4 Oct 2026 19:40:54 +0000 Subject: [PATCH] chore: sync public mirror from internal --- .repository-projection.json | 6 +- evals/compaction-behavior-v1.json | 71 ++++++ packages/local-host-rs/Cargo.toml | 5 + .../examples/compaction_eval/main.rs | 133 ++++++++++++ .../compaction_eval/provider_tests.rs | 204 ++++++++++++++++++ .../examples/compaction_eval/report.rs | 182 ++++++++++++++++ .../examples/compaction_eval/suite.rs | 74 +++++++ .../examples/compaction_eval/tests.rs | 204 ++++++++++++++++++ .../examples/compaction_eval/trial.rs | 187 ++++++++++++++++ packages/local-host-rs/src/agent/mod.rs | 3 + .../runtime-rs/src/agent/native/context.rs | 10 +- .../runtime-rs/src/agent/selective_summary.rs | 2 +- 12 files changed, 1076 insertions(+), 5 deletions(-) create mode 100644 evals/compaction-behavior-v1.json create mode 100644 packages/local-host-rs/examples/compaction_eval/main.rs create mode 100644 packages/local-host-rs/examples/compaction_eval/provider_tests.rs create mode 100644 packages/local-host-rs/examples/compaction_eval/report.rs create mode 100644 packages/local-host-rs/examples/compaction_eval/suite.rs create mode 100644 packages/local-host-rs/examples/compaction_eval/tests.rs create mode 100644 packages/local-host-rs/examples/compaction_eval/trial.rs diff --git a/.repository-projection.json b/.repository-projection.json index 9298cd403..1758f2a31 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "6195061bc9e443dc9612b87208e138d73278186d", + "sourceSha": "38940ee14ede4a1d8f99729c6991b805720314b2", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "eb9445ec333a7736a10ea49f15423d5e2a45c870", + "priorProjectedBase": "66915470e59f656745b3a6cfd9ab30f6db2273e5", "definitionDigest": "cb9d429542ebb0a2de9b42a7aad60d9d8696a648ceba47c30f05c0b285ca0db7", "toolDigest": "f8cb071b0f27267120ccf45a00d0982f45113bd23535bef6a1555b4933f99f13", - "contentDigest": "105a2ec61d1642780b25c552c0f21be52a93b78cd5f4ac223b0e882bbb77d01f", + "contentDigest": "f199ddc735eb038778000c80291490a41682e110e834b8c98f37ca52cc0f3744", "publicationEligible": true } diff --git a/evals/compaction-behavior-v1.json b/evals/compaction-behavior-v1.json new file mode 100644 index 000000000..2b8bd82b2 --- /dev/null +++ b/evals/compaction-behavior-v1.json @@ -0,0 +1,71 @@ +{ + "schema": "maestro.compaction-eval-suite.v1", + "cases": [ + { + "id": "corrected-belief", + "family": "belief-attribution", + "history": [ + {"role": "user", "content": "I believe ordinary JSON supports trailing commas. I prefer leaving them in."}, + {"role": "assistant", "content": "Standard JSON rejects trailing commas. Your preference does not change that syntax rule."} + ], + "question": "Based on the conversation, does standard JSON accept trailing commas? Return only JSON with the boolean key accepts_trailing_commas.", + "expected": {"accepts_trailing_commas": false} + }, + { + "id": "failed-tool-result", + "family": "execution-evidence", + "history": [ + {"role": "user", "content": "Run the parser tests before calling the repair verified."}, + {"role": "assistant", "content": [{"type": "tool_use", "id": "parser-test", "name": "bash", "input": {"command": "cargo test parser"}}]}, + {"role": "user", "content": [{"type": "tool_result", "tool_use_id": "parser-test", "content": "exit_code=101; parser_rejects_trailing_comma FAILED", "is_error": true}]}, + {"role": "assistant", "content": "The parser test failed. I have not verified the repair."} + ], + "question": "What is the observed status of the parser test run? Return only JSON with test_status equal to passed, failed, or unknown.", + "expected": {"test_status": "failed"} + }, + { + "id": "later-failure-overrides-claim", + "family": "execution-evidence", + "history": [ + {"role": "user", "content": "Check the revised import contract."}, + {"role": "assistant", "content": "I expect the revised import test to pass."}, + {"role": "assistant", "content": [{"type": "tool_use", "id": "import-test", "name": "bash", "input": {"command": "cargo test import_contract"}}]}, + {"role": "user", "content": [{"type": "tool_result", "tool_use_id": "import-test", "content": "exit_code=101; import_contract FAILED", "is_error": true}]} + ], + "question": "Did the observed import_contract test pass? Return only JSON with the boolean key passed.", + "expected": {"passed": false} + }, + { + "id": "changed-instruction", + "family": "instruction-correction", + "history": [ + {"role": "user", "content": "Use port 8080 for the fixture server."}, + {"role": "assistant", "content": "I will use port 8080."}, + {"role": "user", "content": "Correction: use port 9091 instead. The original port conflicts with another process."}, + {"role": "assistant", "content": "The fixture server should now use port 9091."} + ], + "question": "Which port does the latest user instruction select? Return only JSON with the integer key port.", + "expected": {"port": 9091} + }, + { + "id": "missing-test-evidence", + "family": "missing-evidence", + "history": [ + {"role": "user", "content": "Run the integration tests when the database becomes available."}, + {"role": "assistant", "content": "The database is unavailable. No integration test run has occurred."} + ], + "question": "Is there an observed integration-test result? Return only JSON with test_status equal to passed, failed, or unknown.", + "expected": {"test_status": "unknown"} + }, + { + "id": "user-preference-is-not-proof", + "family": "belief-attribution", + "history": [ + {"role": "user", "content": "I am sure this branch has merged. Please remember that I believe it is on main."}, + {"role": "assistant", "content": "No remote branch or merge receipt has been inspected. That belief does not establish a merge."} + ], + "question": "Does the supplied history contain verified evidence that the branch merged? Return only JSON with the boolean key merge_verified.", + "expected": {"merge_verified": false} + } + ] +} diff --git a/packages/local-host-rs/Cargo.toml b/packages/local-host-rs/Cargo.toml index b16d37740..6bae9d2bd 100644 --- a/packages/local-host-rs/Cargo.toml +++ b/packages/local-host-rs/Cargo.toml @@ -9,6 +9,11 @@ description = "Local native execution host for Deixic Code" name = "embedding" path = "examples/embedding.rs" +[[example]] +name = "compaction-eval" +path = "examples/compaction_eval/main.rs" +test = true + [[example]] name = "embedding-test-kit" path = "examples/embedding_test_kit.rs" diff --git a/packages/local-host-rs/examples/compaction_eval/main.rs b/packages/local-host-rs/examples/compaction_eval/main.rs new file mode 100644 index 000000000..92d78ad08 --- /dev/null +++ b/packages/local-host-rs/examples/compaction_eval/main.rs @@ -0,0 +1,133 @@ +//! Paired, tool-free behavior trials through the existing native summary path. +mod report; +mod suite; +#[cfg(test)] +mod tests; +mod trial; + +use anyhow::{Context, Result, ensure}; +use clap::Parser; +use maestro_local_host as host; +use maestro_local_host::agent::{ + CredentialVault, ModelDynamicsConfig, NativeAgent, NativeAgentConfig, +}; +use maestro_runtime::agent::MaxTokensSource; +use sha2::{Digest, Sha256}; +use std::{collections::HashSet, path::PathBuf, time::Duration}; + +#[derive(Parser)] +struct Args { + #[arg(long)] + suite: PathBuf, + /// Explicit provider-qualified model; no automatic routing or fallback. + #[arg(long)] + model: String, + /// A new directory for local manifests, transcripts and reports. + #[arg(long)] + output: PathBuf, + #[arg(long, default_value_t = 60, value_parser = clap::value_parser!(u64).range(1..=600))] + timeout_seconds: u64, + #[arg(long, default_value_t = 1024, value_parser = clap::value_parser!(u32).range(1..=4096))] + max_tokens: u32, + /// Keep both seeded histories below the automatic compaction threshold. + #[arg(long)] + context_window: u64, +} + +fn hash(bytes: &[u8]) -> String { + format!("sha256:{:x}", Sha256::digest(bytes)) +} + +#[tokio::main] +async fn main() -> Result<()> { + let args = Args::parse(); + ensure!( + args.model.contains('/') && !args.model.trim().is_empty(), + "qualify the model with its provider" + ); + ensure!( + args.context_window > u64::from(args.max_tokens) + 4096, + "context window is too small" + ); + let raw = std::fs::read(&args.suite)?; + let suite: suite::Suite = serde_json::from_slice(&raw)?; + suite.validate(args.context_window, args.max_tokens)?; + std::fs::create_dir(&args.output) + .context("output directory must be new and its parent must exist")?; + let binary = std::env::current_exe()?; + let manifest = serde_json::json!({ + "schema": "maestro.compaction-eval-manifest.v1", "suite_sha256": hash(&raw), + "suite": suite, "model": args.model, "summary_model": args.model, + "executable_sha256": hash(&std::fs::read(binary)?), + "summary_guidance_sha256": hash(maestro_context::compaction::SUMMARY_EVIDENCE_GUIDANCE.as_bytes()), + "timeout_seconds": args.timeout_seconds, "max_tokens": args.max_tokens, + "context_window": args.context_window, "thinking_enabled": false, + "tools": [], "automatic_model_routing": false, + "claim": "synthetic_context_behavior_only", "promotion_allowed": false + }); + std::fs::write( + args.output.join("manifest.json"), + serde_json::to_vec_pretty(&manifest)?, + )?; + let mut results = Vec::new(); + for (index, case) in suite.cases.iter().enumerate() { + // Alternate order to expose rather than systematically favor warm-cache runs. + let arms = if index % 2 == 0 { + [false, true] + } else { + [true, false] + }; + for compact in arms { + let cwd = tempfile::tempdir()?; + let config = NativeAgentConfig { + model: args.model.clone(), max_tokens: args.max_tokens, + max_tokens_source: MaxTokensSource::Explicit, + context_window: Some(args.context_window), + cwd: cwd.path().to_string_lossy().into_owned(), + system_prompt: Some("Answer from the supplied conversation evidence. Return only the JSON object requested by the final question.".into()), + thinking_enabled: false, + model_dynamics: ModelDynamicsConfig { summary_model: Some(args.model.clone()), ..Default::default() }, + ..Default::default() + }; + // Use ordinary authenticated local composition, with an empty tool allowlist. + let (agent, mut events) = NativeAgent::new_with_allowed_tools_and_credential_vault( + config, + &HashSet::new(), + CredentialVault::new(), + )?; + agent.set_hooks_enabled(false)?; + let arm = if compact { "compacted" } else { "original" }; + let folder = args.output.join(format!("{}--{arm}", case.id)); + std::fs::create_dir(&folder)?; + let result = trial::run( + &agent, + &mut events, + case, + compact, + Duration::from_secs(args.timeout_seconds), + &folder, + ) + .await; + // Cancel and drain the existing actor; do not leave provider work running. + agent.cancel(); + agent.shutdown().await; + let result = result?; + std::fs::write( + folder.join("result.json"), + serde_json::to_vec_pretty(&result)?, + )?; + results.push(result); + } + } + let report = report::paired(&suite, &results)?; + std::fs::write( + args.output.join("report.json"), + serde_json::to_vec_pretty(&report)?, + )?; + println!("{}", serde_json::to_string_pretty(&report)?); + ensure!( + report.comparison_valid, + "inconclusive comparison: inspect recorded runtime failures" + ); + Ok(()) +} diff --git a/packages/local-host-rs/examples/compaction_eval/provider_tests.rs b/packages/local-host-rs/examples/compaction_eval/provider_tests.rs new file mode 100644 index 000000000..c173f14d6 --- /dev/null +++ b/packages/local-host-rs/examples/compaction_eval/provider_tests.rs @@ -0,0 +1,204 @@ +//! These exercise real native requests against a local fixture, not model efficacy. +#[path = "report.rs"] +mod report; +#[path = "suite.rs"] +mod suite; +#[path = "trial.rs"] +mod trial; + +use crate as host; +use crate::agent::{CredentialVault, ModelDynamicsConfig, NativeAgent, NativeAgentConfig}; +use maestro_ai::{Message, MessageContent, OpenAiClient, Role, UnifiedClient}; +use sha2::{Digest, Sha256}; +use std::{collections::HashSet, time::Duration}; +use tokio::{ + io::{AsyncReadExt, AsyncWriteExt}, + net::{TcpListener, TcpStream}, +}; + +fn hash(bytes: &[u8]) -> String { + format!("sha256:{:x}", Sha256::digest(bytes)) +} + +async fn request(stream: &mut TcpStream) -> serde_json::Value { + let mut bytes = Vec::new(); + loop { + let mut buffer = [0; 8192]; + let n = stream.read(&mut buffer).await.unwrap(); + assert!(n > 0, "fixture request must complete"); + bytes.extend_from_slice(&buffer[..n]); + if let Some(end) = bytes.windows(4).position(|w| w == b"\r\n\r\n") { + let head = String::from_utf8_lossy(&bytes[..end]); + let length: usize = head + .lines() + .find_map(|line| { + let (key, value) = line.split_once(':')?; + key.eq_ignore_ascii_case("content-length") + .then(|| value.trim().parse().unwrap()) + }) + .unwrap(); + if bytes.len() >= end + 4 + length { + return serde_json::from_slice(&bytes[end + 4..end + 4 + length]).unwrap(); + } + } + } +} + +fn sse(text: &str) -> String { + let start = serde_json::json!({"id":"fixture", "model":"gpt-4o", "created":0, "object":"chat.completion.chunk", + "choices":[{"index":0,"delta":{"role":"assistant","content":text},"finish_reason":null}]}); + let stop = serde_json::json!({"id":"fixture", "model":"gpt-4o", "created":0, "object":"chat.completion.chunk", + "choices":[{"index":0,"delta":{},"finish_reason":"stop"}], + "usage":{"prompt_tokens":100,"completion_tokens":20,"total_tokens":120}}); + format!("data: {start}\n\ndata: {stop}\n\ndata: [DONE]\n\n") +} + +#[tokio::test] +async fn native_summary_is_applied_before_probe_and_grader_answer_never_enters_requests() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { + let mut requests = Vec::new(); + // Original probe, one failed summary attempt, summary and compacted probe. + for answer in [ + Some(r#"{"answer":"GRADER_ONLY_SENTINEL"}"#), + None, + Some("SUMMARY_EVIDENCE_MARKER"), + Some(r#"{"answer":"GRADER_ONLY_SENTINEL"}"#), + ] { + let (mut stream, _) = listener.accept().await.unwrap(); + requests.push(request(&mut stream).await); + let (status, body) = match answer { + Some(answer) => ("200 OK", sse(answer)), + None => ("503 Service Unavailable", "fixture summary failure".into()), + }; + let response = format!( + "HTTP/1.1 {status}\r\nContent-Type: text/event-stream\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}", + body.len() + ); + stream.write_all(response.as_bytes()).await.unwrap(); + } + requests + }); + let case = suite::Case { + id: "fixture".into(), + family: "pipeline-mechanics".into(), + history: vec![ + Message { + role: Role::User, + content: MessageContent::text("ORIGINAL_EVIDENCE_MARKER"), + }, + Message { + role: Role::Assistant, + content: MessageContent::text("Acknowledged."), + }, + ], + question: "Return a JSON object with an answer string.".into(), + expected: serde_json::json!({"answer":"GRADER_ONLY_SENTINEL"}), + }; + let suite = suite::Suite { + schema: "maestro.compaction-eval-suite.v1".into(), + cases: vec![case], + }; + suite.validate(128_000, 1024).unwrap(); + let workspace = tempfile::tempdir().unwrap(); + let mut rows = Vec::new(); + for compacted in [false, true] { + let config = NativeAgentConfig { + model: "openai/gpt-4o".into(), + context_window: Some(128_000), + cwd: workspace.path().to_string_lossy().into_owned(), + model_dynamics: ModelDynamicsConfig { + summary_model: Some("openai/gpt-4o".into()), + ..Default::default() + }, + ..Default::default() + }; + let client = UnifiedClient::OpenAI( + OpenAiClient::with_base_url("fixture-key", format!("http://{address}/v1")).unwrap(), + ); + let (agent, mut events) = NativeAgent::start( + config, + vec![], + CredentialVault::new(), + Some(&HashSet::new()), + Some(super::ClientOverride::UnverifiedTest(client)), + None, + None, + ) + .unwrap(); + assert!(agent.runtime_audit_snapshot().tools.is_empty()); + agent.set_hooks_enabled(false).unwrap(); + let output = workspace + .path() + .join(if compacted { "compacted" } else { "original" }); + std::fs::create_dir(&output).unwrap(); + let result = trial::run( + &agent, + &mut events, + &suite.cases[0], + compacted, + Duration::from_secs(10), + &output, + ) + .await; + agent.shutdown().await; + rows.push(result.unwrap()); + for line in std::fs::read_to_string(output.join("events.jsonl")) + .unwrap() + .lines() + { + serde_json::from_str::(line).unwrap(); + } + } + let sent = tokio::time::timeout(Duration::from_secs(5), server) + .await + .unwrap() + .unwrap(); + assert!( + sent.iter() + .all(|r| !r.to_string().contains("GRADER_ONLY_SENTINEL")) + ); + assert!( + sent[0]["messages"] + .to_string() + .contains("ORIGINAL_EVIDENCE_MARKER") + ); + assert!( + sent[1]["messages"] + .to_string() + .contains(maestro_context::compaction::SUMMARY_EVIDENCE_GUIDANCE) + ); + assert!(sent.iter().all(|request| { + request + .get("tools") + .is_none_or(|tools| tools.as_array().is_some_and(Vec::is_empty)) + })); + assert_eq!(sent[1]["messages"], sent[2]["messages"]); + assert!( + sent[3]["messages"] + .to_string() + .contains("SUMMARY_EVIDENCE_MARKER") + ); + assert!( + !sent[3]["messages"] + .to_string() + .contains("ORIGINAL_EVIDENCE_MARKER") + ); + assert!(rows.iter().all(|r| r.verified && r.terminal)); + assert!(!rows[0].forced_compaction_applied); + assert!(rows[1].forced_compaction_applied); + assert_eq!(rows[0].measurement.response_count, 1); + assert_eq!(rows[1].measurement.response_count, 2); + assert_eq!(rows[0].retries_started, 0); + assert_eq!(rows[1].retries_started, 1); + let report = report::paired(&suite, &rows).unwrap(); + assert!(report.comparison_valid); + assert_eq!(report.original.input_tokens_observed, 100); + assert_eq!(report.compacted.input_tokens_observed, 200); + assert!(!report.compacted.usage_complete); + assert!( + report.compacted.cost_per_verified_task_usd.is_none(), + "fixture reports tokens, not a price" + ); +} diff --git a/packages/local-host-rs/examples/compaction_eval/report.rs b/packages/local-host-rs/examples/compaction_eval/report.rs new file mode 100644 index 000000000..fa399e1b1 --- /dev/null +++ b/packages/local-host-rs/examples/compaction_eval/report.rs @@ -0,0 +1,182 @@ +use super::host::agent::TokenUsage; +use super::suite::Suite; +use anyhow::{Result, ensure}; +use serde::Serialize; +use std::collections::HashSet; + +#[derive(Default, Serialize)] +pub struct Measurement { + pub response_count: usize, + pub usage_observations: usize, + pub input_tokens: u64, + pub output_tokens: u64, + pub cache_read_tokens: u64, + pub cache_write_tokens: u64, + pub provider_reported_cost_usd: Option, + pub complete: bool, +} + +impl Measurement { + pub fn add(&mut self, usage: Option<&TokenUsage>) { + let first = self.response_count == 0; + self.response_count += 1; + if let Some(usage) = usage { + self.usage_observations += 1; + self.input_tokens = self.input_tokens.saturating_add(usage.input_tokens); + self.output_tokens = self.output_tokens.saturating_add(usage.output_tokens); + self.cache_read_tokens = self + .cache_read_tokens + .saturating_add(usage.cache_read_tokens); + self.cache_write_tokens = self + .cache_write_tokens + .saturating_add(usage.cache_write_tokens); + let cost = usage.cost.filter(|v| v.is_finite() && *v >= 0.0); + self.provider_reported_cost_usd = if first { + cost + } else { + self.provider_reported_cost_usd + .zip(cost) + .map(|(a, b)| a + b) + .filter(|v| v.is_finite()) + }; + } else { + self.provider_reported_cost_usd = None; + } + } +} + +#[derive(Serialize)] +pub struct Trial { + pub case_id: String, + pub compacted: bool, + pub verified: bool, + pub terminal: bool, + pub failure: Option, + pub elapsed_seconds: f64, + pub retries_started: usize, + pub forced_compaction_applied: bool, + pub answer: String, + pub measurement: Measurement, +} + +#[derive(Serialize)] +pub struct Arm { + pub attempts: usize, + pub verified: usize, + pub success_rate: f64, + pub elapsed_seconds: f64, + pub retries_started: usize, + pub input_tokens_observed: u64, + pub output_tokens_observed: u64, + pub cache_read_tokens_observed: u64, + pub cache_write_tokens_observed: u64, + pub usage_complete: bool, + pub provider_reported_cost_usd: Option, + /// Includes failed attempts and compaction overhead in the numerator. + pub cost_per_verified_task_usd: Option, +} + +fn arm(rows: &[&Trial]) -> Arm { + let verified = rows.iter().filter(|r| r.verified).count(); + let costs: Option> = rows + .iter() + .map(|r| { + r.measurement + .complete + .then_some(r.measurement.provider_reported_cost_usd) + .flatten() + }) + .collect(); + let cost = costs + .map(|c| c.iter().sum::()) + .filter(|c| c.is_finite()); + Arm { + attempts: rows.len(), + verified, + success_rate: verified as f64 / rows.len() as f64, + elapsed_seconds: rows.iter().map(|r| r.elapsed_seconds).sum(), + retries_started: rows.iter().map(|r| r.retries_started).sum(), + input_tokens_observed: rows.iter().map(|r| r.measurement.input_tokens).sum(), + output_tokens_observed: rows.iter().map(|r| r.measurement.output_tokens).sum(), + cache_read_tokens_observed: rows.iter().map(|r| r.measurement.cache_read_tokens).sum(), + cache_write_tokens_observed: rows.iter().map(|r| r.measurement.cache_write_tokens).sum(), + usage_complete: rows.iter().all(|r| { + r.measurement.complete + && r.measurement.response_count > 0 + && r.measurement.usage_observations == r.measurement.response_count + }), + provider_reported_cost_usd: cost, + cost_per_verified_task_usd: cost.filter(|_| verified > 0).map(|c| c / verified as f64), + } +} + +#[derive(Serialize)] +pub struct Report { + pub schema: &'static str, + pub comparison_valid: bool, + pub original: Arm, + pub compacted: Arm, + pub compacted_only: usize, + pub original_only: usize, + pub difference_percentage_points: f64, + pub claim: &'static str, + pub promotion_allowed: bool, +} + +pub fn paired(suite: &Suite, trials: &[Trial]) -> Result { + ensure!(!suite.cases.is_empty(), "empty paired denominator"); + ensure!( + trials.len() == suite.cases.len() * 2, + "incomplete paired denominator" + ); + let ids: HashSet<_> = suite.cases.iter().map(|c| c.id.as_str()).collect(); + let mut seen = HashSet::new(); + for trial in trials { + ensure!( + ids.contains(trial.case_id.as_str()) && seen.insert((&trial.case_id, trial.compacted)), + "extra or duplicate trial" + ); + ensure!( + trial.terminal || trial.failure.is_some(), + "nonterminal trial without a recorded failure" + ); + ensure!( + !trial.verified || (trial.terminal && trial.failure.is_none()), + "unverified terminal claimed success" + ); + ensure!( + !trial.terminal || !trial.compacted || trial.forced_compaction_applied, + "missing forced compaction" + ); + ensure!( + trial.compacted || !trial.forced_compaction_applied, + "original arm was compacted" + ); + } + let original = trials.iter().filter(|r| !r.compacted).collect::>(); + let compacted = trials.iter().filter(|r| r.compacted).collect::>(); + let mut wins = 0; + let mut losses = 0; + for control in &original { + let candidate = compacted + .iter() + .find(|r| r.case_id == control.case_id) + .expect("validated pair"); + wins += usize::from(candidate.verified && !control.verified); + losses += usize::from(control.verified && !candidate.verified); + } + Ok(Report { + schema: "maestro.compaction-eval-report.v1", + comparison_valid: trials + .iter() + .all(|r| r.failure.as_deref().is_none_or(|f| f == "timeout")), + original: arm(&original), + compacted: arm(&compacted), + compacted_only: wins, + original_only: losses, + difference_percentage_points: 100.0 * (wins as f64 - losses as f64) + / suite.cases.len() as f64, + claim: "synthetic_context_behavior_only", + promotion_allowed: false, + }) +} diff --git a/packages/local-host-rs/examples/compaction_eval/suite.rs b/packages/local-host-rs/examples/compaction_eval/suite.rs new file mode 100644 index 000000000..2a1a91d42 --- /dev/null +++ b/packages/local-host-rs/examples/compaction_eval/suite.rs @@ -0,0 +1,74 @@ +use super::host::{agent::selective_summary, ai::Message}; +use anyhow::{Result, ensure}; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use std::collections::HashSet; + +#[derive(Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +pub struct Suite { + pub schema: String, + pub cases: Vec, +} + +#[derive(Deserialize, Serialize)] +#[serde(deny_unknown_fields)] +pub struct Case { + pub id: String, + pub family: String, + /// Seeded historical messages, not a claimed live execution trace. + pub history: Vec, + pub question: String, + /// Grader-only answer: never included in a provider request. + pub expected: Value, +} + +impl Suite { + pub fn validate(&self, context_window: u64, max_tokens: u32) -> Result<()> { + ensure!( + self.schema == "maestro.compaction-eval-suite.v1", + "unsupported suite schema" + ); + ensure!( + !self.cases.is_empty() && self.cases.len() <= 100, + "suite needs 1..100 cases" + ); + let mut seen = HashSet::new(); + for case in &self.cases { + ensure!( + !case.id.is_empty() + && case + .id + .bytes() + .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-'), + "unsafe case id" + ); + ensure!(seen.insert(&case.id), "duplicate case id"); + ensure!( + !case.family.trim().is_empty() && !case.question.trim().is_empty(), + "missing family or question" + ); + ensure!( + case.expected.as_object().is_some_and(|o| !o.is_empty()), + "expected answer must be a nonempty object" + ); + let preview = selective_summary::preview(&case.history)?; + ensure!(!preview.turns.is_empty(), "history needs a user turn"); + // Validate complete tool exchanges before any provider call. + selective_summary::validate_groups(&case.history)?; + let bytes = serde_json::to_vec(&case.history)?.len() + case.question.len(); + // Conservative byte ceiling leaves room for standing runtime instructions. + ensure!( + bytes as u64 + u64::from(max_tokens) + 4096 < context_window / 2, + "seeded case may trigger automatic compaction: {}", + case.id + ); + } + Ok(()) + } +} + +pub fn grade(answer: &str, expected: &Value) -> bool { + // Exact typed fields, including missing/null/false distinctions. No prose judge. + serde_json::from_str::(answer.trim()).is_ok_and(|actual| actual == *expected) +} diff --git a/packages/local-host-rs/examples/compaction_eval/tests.rs b/packages/local-host-rs/examples/compaction_eval/tests.rs new file mode 100644 index 000000000..72a60e8e8 --- /dev/null +++ b/packages/local-host-rs/examples/compaction_eval/tests.rs @@ -0,0 +1,204 @@ +use super::{ + report::{Measurement, Trial, paired}, + suite::{Suite, grade}, +}; +use maestro_local_host::agent::TokenUsage; +use serde_json::json; + +fn suite() -> Suite { + serde_json::from_str(include_str!( + "../../../../evals/compaction-behavior-v1.json" + )) + .unwrap() +} + +fn usage(cost: Option) -> TokenUsage { + TokenUsage { + input_tokens: 100, + output_tokens: 20, + cache_read_tokens: 50, + cache_write_tokens: 10, + cost, + } +} + +fn trial(id: &str, compacted: bool, verified: bool, cost: Option) -> Trial { + let mut measurement = Measurement::default(); + measurement.add(Some(&usage(cost))); + measurement.complete = true; + Trial { + case_id: id.into(), + compacted, + verified, + terminal: true, + failure: None, + elapsed_seconds: 2.0, + retries_started: 0, + forced_compaction_applied: compacted, + answer: String::new(), + measurement, + } +} + +fn trials(suite: &Suite) -> Vec { + suite + .cases + .iter() + .flat_map(|case| { + [ + trial(&case.id, false, true, Some(0.1)), + trial(&case.id, true, true, Some(0.2)), + ] + }) + .collect() +} + +#[test] +fn reviewed_cases_validate_complete_historical_tool_exchanges() { + let suite = suite(); + suite.validate(128_000, 1024).unwrap(); + assert_eq!(suite.cases.len(), 6); + let mut invalid = suite; + invalid.cases[1].history.remove(1); + assert!( + invalid.validate(128_000, 1024).is_err(), + "orphaned historical receipt must be rejected before inference" + ); +} + +#[test] +fn ambiguous_or_over_budget_suites_fail_before_provider_calls() { + let mut invalid = suite(); + invalid.cases[1].id = invalid.cases[0].id.clone(); + assert!(invalid.validate(128_000, 1024).is_err()); + let mut invalid = suite(); + invalid.cases[0].id = "../escape".into(); + assert!(invalid.validate(128_000, 1024).is_err()); + assert!(suite().validate(8_192, 1024).is_err()); +} + +#[test] +fn grading_rejects_prose_extra_fields_and_wrong_types() { + let expected = json!({"passed": false}); + assert!(grade(" {\"passed\":false} ", &expected)); + for answer in [ + "done", + "{\"passed\":true}", + "{\"passed\":\"false\"}", + "{\"passed\":false,\"extra\":1}", + "{}", + "```json\n{\"passed\":false}\n```", + ] { + assert!(!grade(answer, &expected), "invalid answer: {answer}"); + } + assert!(!grade("{\"passed\":null}", &expected)); +} + +#[test] +fn failed_attempts_and_summary_overhead_stay_in_cost_per_success() { + let suite = suite(); + let mut rows = trials(&suite); + rows[1].verified = false; + rows[1].measurement.add(Some(&usage(Some(0.3)))); + let report = paired(&suite, &rows).unwrap(); + assert_eq!(report.compacted.attempts, 6); + assert_eq!(report.compacted.verified, 5); + assert!((report.compacted.provider_reported_cost_usd.unwrap() - 1.5).abs() < 1e-10); + assert!((report.compacted.cost_per_verified_task_usd.unwrap() - 0.3).abs() < 1e-10); + assert_eq!(report.compacted.input_tokens_observed, 700); + assert_eq!(report.compacted.cache_read_tokens_observed, 350); + assert_eq!(report.original_only, 1); + assert_eq!(report.compacted_only, 0); + assert!((report.difference_percentage_points + 100.0 / 6.0).abs() < 1e-10); + assert!(report.comparison_valid); + assert!(!report.promotion_allowed); +} + +#[test] +fn unknown_or_partial_spend_never_becomes_a_zero_cost() { + for unknown in [None, Some(-1.0), Some(f64::NAN), Some(f64::INFINITY)] { + let mut measurement = Measurement::default(); + measurement.add(Some(&usage(Some(0.2)))); + measurement.add(Some(&usage(unknown))); + measurement.add(Some(&usage(Some(0.1)))); + assert!(measurement.provider_reported_cost_usd.is_none()); + } + let suite = suite(); + let mut rows = trials(&suite); + rows[1].measurement.complete = false; + let report = paired(&suite, &rows).unwrap(); + assert!(report.compacted.provider_reported_cost_usd.is_none()); + assert!(report.compacted.cost_per_verified_task_usd.is_none()); + assert!(!report.compacted.usage_complete); + assert!(report.original.provider_reported_cost_usd.is_some()); +} + +#[test] +fn absent_usage_is_visible_and_zero_success_has_no_cost_ratio() { + let suite = suite(); + let mut rows = trials(&suite); + rows[0].measurement.add(None); + for row in rows.iter_mut().filter(|r| r.compacted) { + row.verified = false; + } + let report = paired(&suite, &rows).unwrap(); + assert!(!report.original.usage_complete); + assert!(report.original.provider_reported_cost_usd.is_none()); + assert!(report.compacted.provider_reported_cost_usd.is_some()); + assert!(report.compacted.cost_per_verified_task_usd.is_none()); +} + +#[test] +fn provider_failure_invalidates_pair_but_timeout_keeps_denominator() { + let suite = suite(); + let mut rows = trials(&suite); + rows[0].verified = false; + rows[0].terminal = false; + rows[0].failure = Some("timeout".into()); + rows[0].measurement.complete = false; + let report = paired(&suite, &rows).unwrap(); + assert!(report.comparison_valid); + assert_eq!(report.original.attempts, 6); + assert_eq!(report.compacted_only, 1); + rows[0].failure = Some("runtime_or_summary_failure".into()); + assert!(!paired(&suite, &rows).unwrap().comparison_valid); +} + +#[test] +fn report_rejects_incomplete_duplicate_or_extra_pairs_and_false_success() { + let suite = suite(); + let mut rows = trials(&suite); + rows.pop(); + assert!(paired(&suite, &rows).is_err()); + let mut rows = trials(&suite); + rows[1].compacted = false; + assert!(paired(&suite, &rows).is_err()); + let mut rows = trials(&suite); + rows[1].case_id = "extra".into(); + assert!(paired(&suite, &rows).is_err()); + let mut rows = trials(&suite); + rows[1].forced_compaction_applied = false; + assert!(paired(&suite, &rows).is_err()); + rows[1].forced_compaction_applied = true; + rows[1].terminal = false; + assert!(paired(&suite, &rows).is_err()); + rows[1].terminal = true; + rows[1].verified = false; + rows[1].forced_compaction_applied = false; + assert!(paired(&suite, &rows).is_err()); + let mut rows = trials(&suite); + rows[0].forced_compaction_applied = true; + assert!(paired(&suite, &rows).is_err()); + let mut empty = suite; + empty.cases.clear(); + assert!(paired(&empty, &[]).is_err()); +} + +#[test] +fn nonterminal_attempt_without_a_failure_cannot_be_a_valid_comparison() { + let suite = suite(); + let mut rows = trials(&suite); + rows[0].terminal = false; + rows[0].verified = false; + assert!(paired(&suite, &rows).is_err()); +} diff --git a/packages/local-host-rs/examples/compaction_eval/trial.rs b/packages/local-host-rs/examples/compaction_eval/trial.rs new file mode 100644 index 000000000..9ea8feeb0 --- /dev/null +++ b/packages/local-host-rs/examples/compaction_eval/trial.rs @@ -0,0 +1,187 @@ +use super::host::agent::{FromAgent, NativeAgent, RangeSelection}; +use super::{ + hash, + report::{Measurement, Trial}, + suite::{Case, grade}, +}; +use anyhow::{Context, Result, bail, ensure}; +use std::{collections::HashSet, io::Write, path::Path, time::Duration}; +use tokio::{ + sync::mpsc, + time::{Instant, timeout_at}, +}; + +pub async fn run( + agent: &NativeAgent, + events: &mut mpsc::UnboundedReceiver, + case: &Case, + compacted: bool, + budget: Duration, + folder: &Path, +) -> Result { + let start = Instant::now(); + let deadline = start + budget; + let mut trial = Trial { + case_id: case.id.clone(), + compacted, + verified: false, + terminal: false, + failure: None, + elapsed_seconds: 0.0, + retries_started: 0, + forced_compaction_applied: false, + answer: String::new(), + measurement: Measurement::default(), + }; + std::fs::write( + folder.join("original-history.json"), + serde_json::to_vec_pretty(&case.history)?, + )?; + let mut log = std::fs::File::create(folder.join("events.jsonl"))?; + let result = execute(agent, events, case, deadline, folder, &mut log, &mut trial).await; + while let Ok(event) = events.try_recv() { + record(&mut log, &event)?; + trial.retries_started += usize::from(is_retry(&event)); + } + trial.elapsed_seconds = start.elapsed().as_secs_f64(); + if let Err(error) = result { + trial.failure = Some( + if error + .downcast_ref::() + .is_some() + { + "timeout" + } else { + "runtime_or_summary_failure" + } + .into(), + ); + trial.measurement.complete = false; + std::fs::write(folder.join("failure.txt"), format!("{error:#}"))?; + } + trial.verified = + trial.terminal && trial.failure.is_none() && grade(&trial.answer, &case.expected); + trial.measurement.complete &= trial.retries_started == 0; + Ok(trial) +} + +fn record(log: &mut std::fs::File, event: &FromAgent) -> Result<()> { + if matches!(event, FromAgent::LocalAssistantContent { .. }) { + // This native-only snapshot is intentionally excluded from the wire protocol. + serde_json::to_writer( + &mut *log, + &serde_json::json!({ + "type": "local_assistant_content", "wire_serializable": false + }), + )?; + } else { + serde_json::to_writer(&mut *log, event)?; + } + writeln!(log)?; + Ok(()) +} + +fn is_retry(event: &FromAgent) -> bool { + matches!( + event, + FromAgent::RequestRetryObservation + | FromAgent::StreamObservation { + observation: maestro_ai::StreamObservation::Retry + } + ) +} + +async fn execute( + agent: &NativeAgent, + events: &mut mpsc::UnboundedReceiver, + case: &Case, + deadline: Instant, + folder: &Path, + log: &mut std::fs::File, + trial: &mut Trial, +) -> Result<()> { + agent.replace_history(case.history.clone()); + let preview = timeout_at(deadline, agent.start_selective_summary_preview()?).await???; + if trial.compacted { + let request = agent + .start_selective_summary(RangeSelection::FromTurn(1), preview.history_digest.clone())?; + let mut receiver = request.receiver; + let outcome = match timeout_at(deadline, &mut receiver).await { + Ok(outcome) => outcome?, + Err(error) => { + request.cancellation.cancel(); + // Settle auxiliary usage without calling the provider again. + if let Ok(Ok(outcome)) = + tokio::time::timeout(Duration::from_secs(2), receiver).await + { + trial.measurement.add(outcome.usage.as_ref()); + } + return Err(error.into()); + } + }; + trial.measurement.add(outcome.usage.as_ref()); + let summary = outcome.result?; + let raw = serde_json::to_vec_pretty(&summary.messages)?; + std::fs::write(folder.join("compacted-history.json"), &raw)?; + std::fs::write(folder.join("compacted-history.sha256"), hash(&raw))?; + timeout_at( + deadline, + agent.apply_selective_summary(summary.messages, preview.history_digest)?, + ) + .await???; + trial.forced_compaction_applied = true; + } + timeout_at(deadline, agent.prompt(case.question.clone(), vec![])).await??; + let mut response_ids = HashSet::new(); + loop { + let event = timeout_at(deadline, events.recv()) + .await? + .context("native event stream closed")?; + record(log, &event)?; + trial.retries_started += usize::from(is_retry(&event)); + match event { + FromAgent::ResponseStart { .. } => trial.answer.clear(), + FromAgent::ResponseChunk { + content, + is_thinking: false, + .. + } => { + ensure!( + trial.answer.len() + content.len() <= 64 * 1024, + "answer exceeded evaluation limit" + ); + trial.answer.push_str(&content); + } + FromAgent::ResponseEnd { response_id, usage } + if response_id == "done" && usage.is_none() => + { + // The actor's completion marker is not a billed provider response. + } + FromAgent::ResponseEnd { response_id, usage } => { + ensure!( + response_ids.insert(response_id), + "duplicate response usage snapshot" + ); + trial.measurement.add(usage.as_ref()); + } + FromAgent::Compaction { auto: true, .. } => { + bail!("unexpected automatic compaction invalidates the control") + } + FromAgent::ToolCall { .. } => bail!("tool call in a tool-free evaluation"), + FromAgent::ProviderError { .. } + | FromAgent::Error { .. } + | FromAgent::TurnInterrupted { .. } => bail!("native turn failed"), + FromAgent::TurnCompleted { .. } => { + ensure!( + !response_ids.is_empty(), + "terminal without a provider response receipt" + ); + trial.terminal = true; + // Failed retries may have unreported spend. Never call partial cost complete. + trial.measurement.complete = trial.retries_started == 0; + return Ok(()); + } + _ => {} + } + } +} diff --git a/packages/local-host-rs/src/agent/mod.rs b/packages/local-host-rs/src/agent/mod.rs index a3d1f5f7d..af50e3bb7 100644 --- a/packages/local-host-rs/src/agent/mod.rs +++ b/packages/local-host-rs/src/agent/mod.rs @@ -9,6 +9,9 @@ mod codemode_codex_tests; pub mod codex_app_server_turns; #[cfg(test)] +#[path = "../../examples/compaction_eval/provider_tests.rs"] +mod compaction_eval_provider_tests; +#[cfg(test)] pub mod harness; #[cfg(test)] mod native_admission_tests; diff --git a/packages/runtime-rs/src/agent/native/context.rs b/packages/runtime-rs/src/agent/native/context.rs index 7a0540f27..cd5211bbd 100644 --- a/packages/runtime-rs/src/agent/native/context.rs +++ b/packages/runtime-rs/src/agent/native/context.rs @@ -995,10 +995,18 @@ impl NativeAgentRunner { .await?; request.ensure_current(&self.credential_vault)?; let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + let retry_observer: maestro_ai::StreamObserver = Arc::new({ + let event_tx = self.event_tx.clone(); + move |observation| { + if observation == maestro_ai::StreamObservation::Retry { + let _ = event_tx.send(FromAgent::StreamObservation { observation }); + } + } + }); let mut stream = tokio::select! { () = cancellation.cancelled() => anyhow::bail!("Summary cancelled"), () = self.shutdown_token.cancelled() => anyhow::bail!("Summary cancelled"), - result = tokio::time::timeout_at(deadline, client.stream_owned_config(&request.messages, request.config)) => result.context("Summary timed out")?.map_err(|_| anyhow::anyhow!("Summary provider request failed"))?, + result = tokio::time::timeout_at(deadline, client.stream_owned_config_shared_messages_observed(Arc::clone(&request.messages), request.config, Some(retry_observer))) => result.context("Summary timed out")?.map_err(|_| anyhow::anyhow!("Summary provider request failed"))?, }; loop { let event = tokio::select! { diff --git a/packages/runtime-rs/src/agent/selective_summary.rs b/packages/runtime-rs/src/agent/selective_summary.rs index bebe1eff5..a32b045a2 100644 --- a/packages/runtime-rs/src/agent/selective_summary.rs +++ b/packages/runtime-rs/src/agent/selective_summary.rs @@ -98,7 +98,7 @@ pub fn preview(messages: &[Message]) -> Result { } /// Reject orphaned, duplicate and boundary-crossing tool exchanges. -pub(crate) fn validate_groups(messages: &[Message]) -> Result<()> { +pub fn validate_groups(messages: &[Message]) -> Result<()> { let mut pending = HashSet::new(); let mut seen = HashSet::new(); for message in messages {