From fef3b6f4a0f530c569d35d4bf6b734a1dffe576d Mon Sep 17 00:00:00 2001 From: Hare Date: Wed, 16 Sep 2026 06:13:32 +0900 Subject: [PATCH] fix: preserve compaction metric precision and CAS category --- crates/worker/src/compact/telemetry.rs | 50 ++++++++++++++++++++-- crates/worker/src/worker.rs | 31 +++++--------- crates/worker/tests/compact_events_test.rs | 22 +++++++++- 3 files changed, 76 insertions(+), 27 deletions(-) diff --git a/crates/worker/src/compact/telemetry.rs b/crates/worker/src/compact/telemetry.rs index 93fe52c0..2c3c85a4 100644 --- a/crates/worker/src/compact/telemetry.rs +++ b/crates/worker/src/compact/telemetry.rs @@ -2,10 +2,11 @@ use std::collections::BTreeMap; use std::time::Duration; use agen::token_counter::EstimateSource; +use agen::usage_record::UsageRecord; use session_metrics::Metric; use session_store::{SegmentId, SessionId}; -use super::usage_tracker::UsageSnapshot; +use super::usage_tracker::{PostRequestMetric, UsageSnapshot}; const MAX_SAFE_INTEGER: u64 = (1_u64 << 53) - 1; @@ -112,7 +113,7 @@ impl CompactAttempt { pub(crate) fn start_metric(&self) -> Metric { self.metric("compact.start") - .with_value(safe_number(self.pre_context_tokens)) + .with_value(safe_metric_number(self.pre_context_tokens)) .with_dimension("occupancy_source", estimate_source(self.pre_context_source)) .with_dimension( "retained_token_budget", @@ -300,13 +301,32 @@ fn metric_with_context( correlation_id: &str, ) -> Metric { let mut metric = Metric::now(name) - .with_value(safe_number(value)) + .with_value(safe_metric_number(value)) .with_correlation_id(correlation_id); metric.dimensions = dimensions.clone(); metric } -fn safe_number(value: u64) -> f64 { +pub(crate) fn correlated_post_request_metric( + kind: PostRequestMetric, + correlation_id: &str, + record: &UsageRecord, +) -> Metric { + let value = match kind { + PostRequestMetric::Prune => record.cache_read_tokens, + PostRequestMetric::Compaction => record.input_total_tokens, + }; + Metric::now(kind.name()) + .with_correlation_id(correlation_id) + .with_value(safe_metric_number(value)) + .with_dimension("history_len", record.history_len.to_string()) + .with_dimension("input_total_tokens", record.input_total_tokens.to_string()) + .with_dimension("cache_read_tokens", record.cache_read_tokens.to_string()) + .with_dimension("cache_write_tokens", record.cache_write_tokens.to_string()) + .with_dimension("output_tokens", record.output_tokens.to_string()) +} + +pub(crate) fn safe_metric_number(value: u64) -> f64 { value.min(MAX_SAFE_INTEGER) as f64 } @@ -346,6 +366,28 @@ mod tests { assert!(start.dimensions.values().all(|value| value.len() <= 64)); } + #[test] + fn post_request_metric_saturates_values_above_json_safe_integer() { + let record = UsageRecord { + history_len: 1, + input_total_tokens: u64::MAX, + cache_read_tokens: 0, + cache_write_tokens: 0, + output_tokens: 1, + }; + let metric = correlated_post_request_metric( + PostRequestMetric::Compaction, + "018f6f8a-9822-7b11-8b35-706f30313700", + &record, + ); + assert_eq!(metric.name, "compact.post_request"); + assert_eq!(metric.value, Some(MAX_SAFE_INTEGER as f64)); + assert_eq!( + metric.dimensions["input_total_tokens"], + u64::MAX.to_string() + ); + } + #[test] fn failure_metrics_never_include_error_text() { let attempt = CompactAttempt::new( diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 10c5aa43..3a4a6957 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -42,7 +42,7 @@ use manifest::{ use crate::compact::state::CompactState; use crate::compact::telemetry::{ CompactAttempt, CompactFailureCategory, CompactMode, CompactSuccessStats, - CompactThresholdPolicy, + CompactThresholdPolicy, correlated_post_request_metric, }; use crate::compact::usage_tracker::UsageTracker; use crate::feature::background::{BackgroundTaskRewriteGuard, FeatureBackgroundTaskRegistry}; @@ -3094,9 +3094,7 @@ impl Worker { WorkerActiveSegmentRef::active_segment(replacement.session_id, replacement.segment_id), )?; if !matched { - return Err(WorkerError::InvalidState( - "active Segment changed before compaction commit".into(), - )); + return Err(WorkerError::CompactActiveSegmentChanged); } Ok(()) } @@ -4906,22 +4904,8 @@ impl Worker { output_tokens: record.output_tokens, })?; for link in post_requests { - let value = match link.metric { - crate::compact::usage_tracker::PostRequestMetric::Prune => { - record.cache_read_tokens - } - crate::compact::usage_tracker::PostRequestMetric::Compaction => { - record.input_total_tokens - } - }; - let metric = session_metrics::Metric::now(link.metric.name()) - .with_correlation_id(&link.correlation_id) - .with_value(value as f64) - .with_dimension("history_len", record.history_len.to_string()) - .with_dimension("input_total_tokens", record.input_total_tokens.to_string()) - .with_dimension("cache_read_tokens", record.cache_read_tokens.to_string()) - .with_dimension("cache_write_tokens", record.cache_write_tokens.to_string()) - .with_dimension("output_tokens", record.output_tokens.to_string()); + let metric = + correlated_post_request_metric(link.metric, &link.correlation_id, &record); self.try_record_metric(&metric); } self.usage_history @@ -7143,7 +7127,9 @@ fn compact_failure_category(error: &WorkerError) -> CompactFailureCategory { WorkerError::CompactResultContextTooLarge { .. } => { CompactFailureCategory::ResultContextTooLarge } - WorkerError::WorkerStore(_) => CompactFailureCategory::ActiveSegmentCommit, + WorkerError::WorkerStore(_) | WorkerError::CompactActiveSegmentChanged => { + CompactFailureCategory::ActiveSegmentCommit + } WorkerError::Store(_) => CompactFailureCategory::Storage, WorkerError::InvalidState(_) | WorkerError::Engine(_) => { CompactFailureCategory::InternalWorker @@ -7208,6 +7194,9 @@ pub enum WorkerError { #[error(transparent)] Provider(#[from] crate::model_client::ProviderError), + #[error("active Segment changed before compaction commit")] + CompactActiveSegmentChanged, + #[error("compaction thrash: context still exceeds threshold immediately after compact")] CompactThrash, diff --git a/crates/worker/tests/compact_events_test.rs b/crates/worker/tests/compact_events_test.rs index b44e0d51..5e8c9599 100644 --- a/crates/worker/tests/compact_events_test.rs +++ b/crates/worker/tests/compact_events_test.rs @@ -475,10 +475,12 @@ async fn active_segment_cas_rejects_stale_compaction_writer() { let client = MockClient::new(vec![ single_text_events("seed response"), write_summary_tool_use_events("summary-1", "replacement summary"), + single_text_events("done"), ]); - let (mut worker, metadata_store, _segment_store) = make_faulting_worker(client).await; + let (mut worker, metadata_store, segment_store) = make_faulting_worker(client).await; worker.run_text("seed input").await.unwrap(); let old_segment_id = worker.segment_id(); + let session_id = worker.session_id(); let competing_segment_id = uuid::Uuid::now_v7(); metadata_store .update_by_name("test-worker", |metadata| { @@ -486,9 +488,25 @@ async fn active_segment_cas_rejects_stale_compaction_writer() { }) .unwrap(); - let _error = worker.compact(0).await.unwrap_err(); + let error = worker.compact(0).await.unwrap_err(); + assert!( + matches!(error, worker::WorkerError::CompactActiveSegmentChanged), + "unexpected stale CAS error: {error:?}" + ); assert_eq!(worker.segment_id(), old_segment_id); + let failure_metrics = + session_metrics::read_segment_metrics(&segment_store, session_id, old_segment_id).unwrap(); + let finish = failure_metrics + .iter() + .find(|record| record.metric.name == "compact.finish") + .unwrap(); + assert_eq!(finish.metric.dimensions["outcome"], "failure"); + assert_eq!( + finish.metric.dimensions["failure_category"], + "active_segment_commit" + ); + assert!(!finish.metric.dimensions.contains_key("error")); assert_eq!( metadata_store .read_by_name("test-worker")