From b6960878a643cfd9c9b1a1e1a92db4ec6ed2616c Mon Sep 17 00:00:00 2001 From: Hare Date: Wed, 16 Sep 2026 06:24:45 +0900 Subject: [PATCH] fix: align compaction metric schema --- crates/session-metrics/README.md | 15 ++++-- crates/worker/src/compact/telemetry.rs | 63 +++++++++++----------- crates/worker/src/worker.rs | 7 +-- crates/worker/tests/compact_events_test.rs | 39 +++++++++----- docs/design/compaction.md | 4 +- 5 files changed, 74 insertions(+), 54 deletions(-) diff --git a/crates/session-metrics/README.md b/crates/session-metrics/README.md index 96f83586..5bef618f 100644 --- a/crates/session-metrics/README.md +++ b/crates/session-metrics/README.md @@ -4,9 +4,9 @@ Session 単位の append-only な観測値を既存 session-log に記録し、 metrics 読取 / JSONL export 経路で取り出すための小さなヘルパークレートです。 - 保存先は `session-store` の `LogEntry::Extension` -- domain は `session.metrics` +- extension domain は `metrics` - metric は `name / ts / dimensions / value / correlation_id` の最小 envelope -- `record_metric` / `record_metric_at` で現在の Segment に append する +- `record_metric` で指定した Session / Segment に append する - `read_segment_metrics` は 1 Segment、`read_session_metrics` は Session 内の全 Segment を読み、各 metric に `segment_id` と `compacted_from` を付ける - `export_metrics_jsonl` はその located metric を newline-delimited JSON にする @@ -19,14 +19,19 @@ request を結合できる。 ```rust,ignore use session_metrics::{ - Metric, export_metrics_jsonl, read_session_metrics, record_metric_at, + Metric, export_metrics_jsonl, read_session_metrics, record_metric, }; let metric = Metric::now("compact.start") .with_value(12_345.0) - .with_dimension("mode", "automatic") + .with_dimension("trigger", "automatic") .with_correlation_id("018f6f8a-9822-7b11-8b35-706f30313700"); -record_metric_at(&store, location, &metric)?; +record_metric( + &store, + location.session_id, + location.segment_id, + &metric, +)?; let records = read_session_metrics(&store, location.session_id)?; let jsonl = export_metrics_jsonl(&records)?; diff --git a/crates/worker/src/compact/telemetry.rs b/crates/worker/src/compact/telemetry.rs index 2c3c85a4..799530d0 100644 --- a/crates/worker/src/compact/telemetry.rs +++ b/crates/worker/src/compact/telemetry.rs @@ -130,16 +130,11 @@ impl CompactAttempt { let dimensions = self.base_dimensions(); let correlation_id = self.correlation_id.clone(); let mut metrics = vec![ - metric_with_context( - "compact.finish", - elapsed.as_millis().min(u128::from(MAX_SAFE_INTEGER)) as u64, - &dimensions, - &correlation_id, - ) - .with_dimension("outcome", "success") - .with_dimension("result_segment_id", result_segment_id.to_string()) - .with_dimension("retained_items", stats.retained_items.to_string()) - .with_dimension("summarized_items", stats.summarized_items.to_string()), + metric_with_context("compact.finish", 1, &dimensions, &correlation_id) + .with_dimension("outcome", "succeeded") + .with_dimension("result_segment_id", result_segment_id.to_string()) + .with_dimension("retained_items", stats.retained_items.to_string()) + .with_dimension("summarized_items", stats.summarized_items.to_string()), metric_with_context( "compact.retained_tokens", stats.retained_tokens, @@ -173,43 +168,38 @@ impl CompactAttempt { ) .with_dimension("source", estimate_source(stats.result_context_source)), metric_with_context( - "compact.compactor.input_tokens", + "compact.input_tokens", stats.usage.input_total_tokens, &dimensions, &correlation_id, ), metric_with_context( - "compact.compactor.output_tokens", + "compact.output_tokens", stats.usage.output_tokens, &dimensions, &correlation_id, ), metric_with_context( - "compact.compactor.cache_read_tokens", + "compact.cache_read_tokens", stats.usage.cache_read_tokens, &dimensions, &correlation_id, ), metric_with_context( - "compact.compactor.cache_write_tokens", + "compact.cache_creation_tokens", stats.usage.cache_write_tokens, &dimensions, &correlation_id, ), metric_with_context( - "compact.compactor.requests", + "compact.requests", stats.requests, &dimensions, &correlation_id, ), + metric_with_context("compact.turns", stats.turns, &dimensions, &correlation_id), metric_with_context( - "compact.compactor.turns", - stats.turns, - &dimensions, - &correlation_id, - ), - metric_with_context( - "compact.compactor.tool_calls", + "compact.tool_calls", stats.tool_calls, &dimensions, &correlation_id, @@ -229,7 +219,7 @@ impl CompactAttempt { // Provider UsageEvent currently carries tokens but no price or cost. Keep // the field explicit and valueless rather than fabricating a zero cost. metrics.push( - self.metric("compact.compactor.cost_usd") + self.metric("compact.cost_usd") .with_dimension("status", "unavailable") .with_dimension("reason", "provider_usage_unpriced") .with_dimension("result_segment_id", result_segment_id.to_string()), @@ -237,22 +227,29 @@ impl CompactAttempt { metrics } - pub(crate) fn failure_metric( + pub(crate) fn failure_metrics( &self, observed_segment_id: SegmentId, elapsed: Duration, category: CompactFailureCategory, - ) -> Metric { + ) -> [Metric; 2] { let outcome = if category == CompactFailureCategory::Cancelled { "cancelled" } else { - "failure" + "failed" }; - self.metric("compact.finish") - .with_value(elapsed.as_millis().min(u128::from(MAX_SAFE_INTEGER)) as f64) + let outcome_metric = self + .metric("compact.finish") + .with_value(1.0) .with_dimension("outcome", outcome) .with_dimension("failure_category", category.as_str()) - .with_dimension("observed_segment_id", observed_segment_id.to_string()) + .with_dimension("observed_segment_id", observed_segment_id.to_string()); + let duration_metric = self + .metric("compact.duration_ms") + .with_value(elapsed.as_millis().min(u128::from(MAX_SAFE_INTEGER)) as f64) + .with_dimension("outcome", outcome) + .with_dimension("observed_segment_id", observed_segment_id.to_string()); + [outcome_metric, duration_metric] } fn metric(&self, name: &'static str) -> Metric { @@ -269,6 +266,7 @@ impl CompactAttempt { self.source_segment_id.to_string(), ), ("mode".into(), self.mode.as_str().into()), + ("trigger".into(), self.mode.as_str().into()), ( "threshold_policy".into(), self.threshold_policy.as_str().into(), @@ -359,6 +357,7 @@ mod tests { assert_eq!(start.name, "compact.start"); assert_eq!(start.value, Some(MAX_SAFE_INTEGER as f64)); assert_eq!(start.dimensions["mode"], "automatic"); + assert_eq!(start.dimensions["trigger"], "automatic"); assert_eq!(start.dimensions["threshold_policy"], "request_threshold"); assert_eq!(start.dimensions["occupancy_source"], "measured"); assert!(start.correlation_id.is_some()); @@ -400,12 +399,16 @@ mod tests { EstimateSource::NoData, 1, ); - let metric = attempt.failure_metric( + let [metric, duration] = attempt.failure_metrics( uuid::Uuid::now_v7(), Duration::from_millis(7), CompactFailureCategory::InternalWorker, ); let encoded = serde_json::to_string(&metric).unwrap(); + assert_eq!(metric.value, Some(1.0)); + assert_eq!(metric.dimensions["outcome"], "failed"); + assert_eq!(duration.name, "compact.duration_ms"); + assert_eq!(duration.value, Some(7.0)); assert!(encoded.contains("internal_worker")); assert!(!encoded.contains("error")); assert!(!encoded.contains("path")); diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 3a4a6957..98787a98 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -5043,12 +5043,13 @@ impl Worker { } Err(error) => { let observed_segment_id = self.segment_state.location().segment_id; - let metric = attempt.failure_metric( + for metric in attempt.failure_metrics( observed_segment_id, started.elapsed(), compact_failure_category(&error), - ); - self.try_record_metric(&metric); + ) { + self.try_record_metric(&metric); + } lifecycle.revision = lifecycle.revision.saturating_add(1); lifecycle.state = if matches!(error, WorkerError::CompactCancelled) { CompactionLifecycleState::Interrupted diff --git a/crates/worker/tests/compact_events_test.rs b/crates/worker/tests/compact_events_test.rs index 5e8c9599..a0aa908c 100644 --- a/crates/worker/tests/compact_events_test.rs +++ b/crates/worker/tests/compact_events_test.rs @@ -501,7 +501,8 @@ async fn active_segment_cas_rejects_stale_compaction_writer() { .iter() .find(|record| record.metric.name == "compact.finish") .unwrap(); - assert_eq!(finish.metric.dimensions["outcome"], "failure"); + assert_eq!(finish.metric.dimensions["outcome"], "failed"); + assert_eq!(finish.metric.value, Some(1.0)); assert_eq!( finish.metric.dimensions["failure_category"], "active_segment_commit" @@ -550,7 +551,8 @@ async fn failed_active_segment_commit_keeps_live_and_durable_history_on_old_segm .iter() .find(|record| record.metric.name == "compact.finish") .unwrap(); - assert_eq!(finish.metric.dimensions["outcome"], "failure"); + assert_eq!(finish.metric.dimensions["outcome"], "failed"); + assert_eq!(finish.metric.value, Some(1.0)); assert_eq!( finish.metric.dimensions["failure_category"], "active_segment_commit" @@ -758,6 +760,7 @@ async fn compact_emits_session_start_carrying_summary_and_task_snapshot() { assert_eq!(starts.len(), 1); assert_eq!(starts[0].segment_id, source_segment_id); assert_eq!(starts[0].metric.dimensions["mode"], "automatic"); + assert_eq!(starts[0].metric.dimensions["trigger"], "automatic"); assert_eq!( starts[0].metric.dimensions["threshold_policy"], "request_threshold" @@ -772,7 +775,8 @@ async fn compact_emits_session_start_carrying_summary_and_task_snapshot() { .find(|record| record.metric.name == "compact.finish") .unwrap(); assert_eq!(finish.segment_id, compacted_segment_id); - assert_eq!(finish.metric.dimensions["outcome"], "success"); + assert_eq!(finish.metric.dimensions["outcome"], "succeeded"); + assert_eq!(finish.metric.value, Some(1.0)); assert_eq!( finish.metric.correlation_id.as_deref(), Some(correlation_id) @@ -789,17 +793,17 @@ async fn compact_emits_session_start_carrying_summary_and_task_snapshot() { .and_then(|record| record.metric.value) .unwrap() as u64 }; - assert_eq!(value("compact.compactor.input_tokens"), 150); - assert_eq!(value("compact.compactor.cache_read_tokens"), 13); - assert_eq!(value("compact.compactor.cache_write_tokens"), 7); - assert_eq!(value("compact.compactor.output_tokens"), 30); - assert_eq!(value("compact.compactor.requests"), 2); - assert!(value("compact.compactor.tool_calls") >= 1); - assert!(value("compact.compactor.turns") >= 2); + assert_eq!(value("compact.input_tokens"), 150); + assert_eq!(value("compact.cache_read_tokens"), 13); + assert_eq!(value("compact.cache_creation_tokens"), 7); + assert_eq!(value("compact.output_tokens"), 30); + assert_eq!(value("compact.requests"), 2); + assert!(value("compact.tool_calls") >= 1); + assert!(value("compact.turns") >= 2); assert!(value("compact.duration_ms") <= u64::MAX); let cost = metrics .iter() - .find(|record| record.metric.name == "compact.compactor.cost_usd") + .find(|record| record.metric.name == "compact.cost_usd") .unwrap(); assert_eq!(cost.metric.value, None); assert_eq!(cost.metric.dimensions["status"], "unavailable"); @@ -834,12 +838,14 @@ async fn manual_compact_metrics_identify_manual_mode() { .find(|record| record.metric.name == "compact.start") .unwrap(); assert_eq!(start.metric.dimensions["mode"], "manual"); + assert_eq!(start.metric.dimensions["trigger"], "manual"); assert_eq!(start.metric.dimensions["threshold_policy"], "manual"); let finish = metrics .iter() .find(|record| record.metric.name == "compact.finish") .unwrap(); - assert_eq!(finish.metric.dimensions["outcome"], "success"); + assert_eq!(finish.metric.dimensions["outcome"], "succeeded"); + assert_eq!(finish.metric.value, Some(1.0)); assert_eq!(finish.metric.correlation_id, start.metric.correlation_id); } @@ -862,10 +868,11 @@ async fn compact_failure_and_cancellation_emit_bounded_categories() { .iter() .find(|record| { record.metric.name == "compact.finish" - && record.metric.dimensions["outcome"] == "failure" + && record.metric.dimensions["outcome"] == "failed" }) .unwrap(); assert_eq!(failure.segment_id, source_segment_id); + assert_eq!(failure.metric.value, Some(1.0)); assert_eq!( failure.metric.dimensions["failure_category"], "summary_missing" @@ -889,6 +896,7 @@ async fn compact_failure_and_cancellation_emit_bounded_categories() { .last() .unwrap(); assert_eq!(cancelled.segment_id, source_segment_id); + assert_eq!(cancelled.metric.value, Some(1.0)); assert_eq!(cancelled.metric.dimensions["failure_category"], "cancelled"); } @@ -962,12 +970,14 @@ async fn pre_run_compact_publishes_runtime_progress_phases() { .unwrap(); assert_eq!(start.segment_id, segment_before); assert_eq!(start.metric.dimensions["mode"], "automatic"); + assert_eq!(start.metric.dimensions["trigger"], "automatic"); assert_eq!(start.metric.dimensions["threshold_policy"], "pre_run"); let finish = metrics .iter() .find(|record| record.metric.name == "compact.finish") .unwrap(); - assert_eq!(finish.metric.dimensions["outcome"], "success"); + assert_eq!(finish.metric.dimensions["outcome"], "succeeded"); + assert_eq!(finish.metric.value, Some(1.0)); assert_eq!(finish.metric.correlation_id, start.metric.correlation_id); } @@ -1017,6 +1027,7 @@ async fn request_threshold_compact_publishes_runtime_progress() { .iter() .find(|record| record.metric.name == "compact.start") .unwrap(); + assert_eq!(start.metric.dimensions["trigger"], "automatic"); assert_eq!( start.metric.dimensions["threshold_policy"], "request_threshold" diff --git a/docs/design/compaction.md b/docs/design/compaction.md index 1583a676..70479ba7 100644 --- a/docs/design/compaction.md +++ b/docs/design/compaction.md @@ -48,9 +48,9 @@ Compaction measurements stay out of the ordinary transcript. They are appended a 3. Compare `compact.retained_tokens`, `compact.overview_tokens`, `compact.summary_tokens`, `compact.auto_read_tokens`, and `compact.result_context_tokens` to explain the context-size change. -4. Compare the Compactor's input/output/cache-read/cache-write tokens, request, +4. Compare the Compactor's input/output/cache-read/cache-creation tokens, request, turn, tool-call, and duration metrics. The current provider `UsageEvent` has no - pricing authority, so `compact.compactor.cost_usd` is valueless with + pricing authority, so `compact.cost_usd` is valueless with `status=unavailable` and `reason=provider_usage_unpriced`; do not fabricate a zero cost. Record a numeric value only after a price authority exists. 5. Aggregate failure and cancellation using the fixed `failure_category` values.