fix: preserve compaction metric precision and CAS category

This commit is contained in:
2026-09-16 06:13:32 +09:00
parent c0a73c12ec
commit fef3b6f4a0
3 changed files with 76 additions and 27 deletions
+46 -4
View File
@@ -2,10 +2,11 @@ use std::collections::BTreeMap;
use std::time::Duration; use std::time::Duration;
use agen::token_counter::EstimateSource; use agen::token_counter::EstimateSource;
use agen::usage_record::UsageRecord;
use session_metrics::Metric; use session_metrics::Metric;
use session_store::{SegmentId, SessionId}; 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; const MAX_SAFE_INTEGER: u64 = (1_u64 << 53) - 1;
@@ -112,7 +113,7 @@ impl CompactAttempt {
pub(crate) fn start_metric(&self) -> Metric { pub(crate) fn start_metric(&self) -> Metric {
self.metric("compact.start") 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("occupancy_source", estimate_source(self.pre_context_source))
.with_dimension( .with_dimension(
"retained_token_budget", "retained_token_budget",
@@ -300,13 +301,32 @@ fn metric_with_context(
correlation_id: &str, correlation_id: &str,
) -> Metric { ) -> Metric {
let mut metric = Metric::now(name) let mut metric = Metric::now(name)
.with_value(safe_number(value)) .with_value(safe_metric_number(value))
.with_correlation_id(correlation_id); .with_correlation_id(correlation_id);
metric.dimensions = dimensions.clone(); metric.dimensions = dimensions.clone();
metric 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 value.min(MAX_SAFE_INTEGER) as f64
} }
@@ -346,6 +366,28 @@ mod tests {
assert!(start.dimensions.values().all(|value| value.len() <= 64)); 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] #[test]
fn failure_metrics_never_include_error_text() { fn failure_metrics_never_include_error_text() {
let attempt = CompactAttempt::new( let attempt = CompactAttempt::new(
+10 -21
View File
@@ -42,7 +42,7 @@ use manifest::{
use crate::compact::state::CompactState; use crate::compact::state::CompactState;
use crate::compact::telemetry::{ use crate::compact::telemetry::{
CompactAttempt, CompactFailureCategory, CompactMode, CompactSuccessStats, CompactAttempt, CompactFailureCategory, CompactMode, CompactSuccessStats,
CompactThresholdPolicy, CompactThresholdPolicy, correlated_post_request_metric,
}; };
use crate::compact::usage_tracker::UsageTracker; use crate::compact::usage_tracker::UsageTracker;
use crate::feature::background::{BackgroundTaskRewriteGuard, FeatureBackgroundTaskRegistry}; use crate::feature::background::{BackgroundTaskRewriteGuard, FeatureBackgroundTaskRegistry};
@@ -3094,9 +3094,7 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
WorkerActiveSegmentRef::active_segment(replacement.session_id, replacement.segment_id), WorkerActiveSegmentRef::active_segment(replacement.session_id, replacement.segment_id),
)?; )?;
if !matched { if !matched {
return Err(WorkerError::InvalidState( return Err(WorkerError::CompactActiveSegmentChanged);
"active Segment changed before compaction commit".into(),
));
} }
Ok(()) Ok(())
} }
@@ -4906,22 +4904,8 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
output_tokens: record.output_tokens, output_tokens: record.output_tokens,
})?; })?;
for link in post_requests { for link in post_requests {
let value = match link.metric { let metric =
crate::compact::usage_tracker::PostRequestMetric::Prune => { correlated_post_request_metric(link.metric, &link.correlation_id, &record);
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());
self.try_record_metric(&metric); self.try_record_metric(&metric);
} }
self.usage_history self.usage_history
@@ -7143,7 +7127,9 @@ fn compact_failure_category(error: &WorkerError) -> CompactFailureCategory {
WorkerError::CompactResultContextTooLarge { .. } => { WorkerError::CompactResultContextTooLarge { .. } => {
CompactFailureCategory::ResultContextTooLarge CompactFailureCategory::ResultContextTooLarge
} }
WorkerError::WorkerStore(_) => CompactFailureCategory::ActiveSegmentCommit, WorkerError::WorkerStore(_) | WorkerError::CompactActiveSegmentChanged => {
CompactFailureCategory::ActiveSegmentCommit
}
WorkerError::Store(_) => CompactFailureCategory::Storage, WorkerError::Store(_) => CompactFailureCategory::Storage,
WorkerError::InvalidState(_) | WorkerError::Engine(_) => { WorkerError::InvalidState(_) | WorkerError::Engine(_) => {
CompactFailureCategory::InternalWorker CompactFailureCategory::InternalWorker
@@ -7208,6 +7194,9 @@ pub enum WorkerError {
#[error(transparent)] #[error(transparent)]
Provider(#[from] crate::model_client::ProviderError), Provider(#[from] crate::model_client::ProviderError),
#[error("active Segment changed before compaction commit")]
CompactActiveSegmentChanged,
#[error("compaction thrash: context still exceeds threshold immediately after compact")] #[error("compaction thrash: context still exceeds threshold immediately after compact")]
CompactThrash, CompactThrash,
+20 -2
View File
@@ -475,10 +475,12 @@ async fn active_segment_cas_rejects_stale_compaction_writer() {
let client = MockClient::new(vec![ let client = MockClient::new(vec![
single_text_events("seed response"), single_text_events("seed response"),
write_summary_tool_use_events("summary-1", "replacement summary"), 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(); worker.run_text("seed input").await.unwrap();
let old_segment_id = worker.segment_id(); let old_segment_id = worker.segment_id();
let session_id = worker.session_id();
let competing_segment_id = uuid::Uuid::now_v7(); let competing_segment_id = uuid::Uuid::now_v7();
metadata_store metadata_store
.update_by_name("test-worker", |metadata| { .update_by_name("test-worker", |metadata| {
@@ -486,9 +488,25 @@ async fn active_segment_cas_rejects_stale_compaction_writer() {
}) })
.unwrap(); .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); 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!( assert_eq!(
metadata_store metadata_store
.read_by_name("test-worker") .read_by_name("test-worker")