fix: align compaction metric schema

This commit is contained in:
2026-09-16 06:24:45 +09:00
parent fef3b6f4a0
commit b6960878a6
5 changed files with 74 additions and 54 deletions
+10 -5
View File
@@ -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)?;
+30 -27
View File
@@ -130,13 +130,8 @@ 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")
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()),
@@ -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"));
+3 -2
View File
@@ -5043,12 +5043,13 @@ impl<C: LlmClient + 'static, St: Store> Worker<C, St> {
}
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);
}
lifecycle.revision = lifecycle.revision.saturating_add(1);
lifecycle.state = if matches!(error, WorkerError::CompactCancelled) {
CompactionLifecycleState::Interrupted
+25 -14
View File
@@ -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"
+2 -2
View File
@@ -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.