Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
97 changes: 82 additions & 15 deletions codex-rs/rollout/src/persistence_metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,13 +14,31 @@ use crate::ResponseItemEnvelope;
use crate::RolloutItem;
use crate::policy::persisted_rollout_item;

const ITEM_BYTES_METRIC: &str = "codex.rollout.persistence.item_bytes";
const ITEM_BYTES_METRIC: &str = "codex.rollout.persistence.item_bytes_v2";
const BYTES_REMOVED_METRIC: &str = "codex.rollout.persistence.bytes_removed";
const APPEND_METRIC: &str = "codex.rollout.persistence.append";
const TURN_BYTES_METRIC: &str = "codex.rollout.persistence.turn_bytes";
const TURN_BYTES_METRIC: &str = "codex.rollout.persistence.turn_bytes_v2";
const MEASUREMENT_ERROR_METRIC: &str = "codex.rollout.persistence.measurement_error";
const SAMPLE_DENOMINATOR: u64 = 100;
const SAMPLE_RATE_LABEL: &str = "0.01";

const KIB: f64 = 1_024.0;
const MIB: f64 = 1_024.0 * KIB;

// Cover the persistence cap while retaining enough range for large items and turns. Values above
// 64 MiB remain visible in the overflow bucket and histogram sum.
const ROLLOUT_BYTES_BOUNDARIES: &[f64] = &[
KIB,
4.0 * KIB,
16.0 * KIB,
64.0 * KIB,
256.0 * KIB,
MIB,
4.0 * MIB,
16.0 * MIB,
64.0 * MIB,
];

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PersistenceDecision {
Kept,
Expand Down Expand Up @@ -315,6 +333,26 @@ fn response_item_type(item: &ResponseItemEnvelope) -> &'static str {
}
}

fn record_item_bytes(
metrics: &MetricsClient,
stage: &'static str,
item: &RolloutItemMeasurement,
payload_bytes: u64,
) {
let _ = metrics.histogram_with_boundaries(
ITEM_BYTES_METRIC,
saturating_i64(payload_bytes),
ROLLOUT_BYTES_BOUNDARIES,
&[
("stage", stage),
("decision", item.decision.as_str()),
("rollout_item_type", item.rollout_item_type.as_str()),
("encoding", "rollout_item_json_v1"),
("sample_rate", SAMPLE_RATE_LABEL),
],
);
}

#[derive(Clone)]
pub struct RolloutPersistenceTelemetry {
metrics: Option<MetricsClient>,
Expand Down Expand Up @@ -348,16 +386,44 @@ impl RolloutPersistenceTelemetry {

for item in &measurement.items {
if let Some(payload_bytes) = item.payload_bytes {
let _ = metrics.histogram(
ITEM_BYTES_METRIC,
saturating_i64(payload_bytes),
&[
("decision", item.decision.as_str()),
("rollout_item_type", item.rollout_item_type.as_str()),
("encoding", "rollout_item_json_v1"),
("sample_rate", SAMPLE_RATE_LABEL),
],
);
record_item_bytes(metrics, "before_projection", item, payload_bytes);
let persisted_payload_bytes = match item.persisted_payload_bytes {
Some(persisted_payload_bytes) => {
record_item_bytes(
metrics,
"after_projection",
item,
persisted_payload_bytes,
);
persisted_payload_bytes
}
None if item.decision == PersistenceDecision::Dropped => 0,
None => {
let _ = metrics.counter(
MEASUREMENT_ERROR_METRIC,
/*inc*/ 1,
&[
("rollout_item_type", item.rollout_item_type.as_str()),
("phase", "serialize_projection"),
],
);
continue;
}
};
let bytes_removed = payload_bytes.saturating_sub(persisted_payload_bytes);
if bytes_removed > 0 {
let _ = metrics.histogram_with_boundaries(
BYTES_REMOVED_METRIC,
saturating_i64(bytes_removed),
ROLLOUT_BYTES_BOUNDARIES,
&[
("decision", item.decision.as_str()),
("rollout_item_type", item.rollout_item_type.as_str()),
("encoding", "rollout_item_json_v1"),
("sample_rate", SAMPLE_RATE_LABEL),
],
);
}
} else {
let _ = metrics.counter(
MEASUREMENT_ERROR_METRIC,
Expand Down Expand Up @@ -406,12 +472,13 @@ impl RolloutPersistenceTelemetry {
}
for turn in turn_update.completed {
for (stage, totals) in [
("pre_filter", turn.totals.pre_filter),
("post_filter", turn.totals.post_filter),
("before_projection", turn.totals.pre_filter),
("after_projection", turn.totals.post_filter),
] {
let _ = metrics.histogram(
let _ = metrics.histogram_with_boundaries(
TURN_BYTES_METRIC,
saturating_i64(totals.payload_bytes),
ROLLOUT_BYTES_BOUNDARIES,
&[
("stage", stage),
("outcome", turn.outcome.as_str()),
Expand Down
Loading