diff --git a/Cargo.lock b/Cargo.lock index c4ff7ed8..7322234f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -304,7 +304,7 @@ dependencies = [ [[package]] name = "cntryl-midge" version = "0.1.0" -source = "git+https://github.com/cntryl/midge?branch=main#6764d9c3a9929f560dd59f6521252c81dca67482" +source = "git+https://github.com/cntryl/midge?branch=main#b7e20515e87ecf036398823c27a02fb339582fdb" dependencies = [ "arc-swap", "base64 0.23.1", diff --git a/docs/admin/api/overview-probes-metrics.md b/docs/admin/api/overview-probes-metrics.md index 99373fd7..2497c109 100644 --- a/docs/admin/api/overview-probes-metrics.md +++ b/docs/admin/api/overview-probes-metrics.md @@ -220,6 +220,12 @@ GET /metrics **Listener**: `FITZ_METRICS_BIND_ADDR:FITZ_METRICS_PORT` **Authentication**: None; keep this listener private to the scrape network. **Response**: Prometheus text format + +Scrapes read an in-process Stream metrics projection initialized during startup +and advanced by successful commits and persisted watermark updates. They do not +scan durable Stream inventory or enqueue admin work on a family actor, so a slow +storage backend cannot turn observability polling into data-plane backpressure. + ``` # HELP fitz_connections_total Total number of active connections # TYPE fitz_connections_total gauge diff --git a/docs/development/architecture.md b/docs/development/architecture.md index 3590deae..d245e679 100644 --- a/docs/development/architecture.md +++ b/docs/development/architecture.md @@ -109,6 +109,9 @@ Raw Prometheus is served by the dedicated unauthenticated listener returns `404` for `/metrics`. Admin consumers use structured JSON at `/api/v1/{family}/metrics`; broker-global samples are available only at `/api/v1/all/metrics` with wildcard authority. +Prometheus rendering reads Stream counts and durable progress from in-process +metric projections initialized during startup and advanced on committed work; a +scrape must not scan storage or enqueue admin work onto a domain actor. ### Critical Invariant: Ephemeral Sessions > **Fitz sessions are ephemeral. The broker never restores session state after disconnect. Clients are responsible for rebuilding all state including subscriptions, transactions, workers, leases, and stream resume position.** diff --git a/docs/development/routing-design.md b/docs/development/routing-design.md index a8c613be..524f8d76 100644 --- a/docs/development/routing-design.md +++ b/docs/development/routing-design.md @@ -505,7 +505,8 @@ Posting indexes preserve the order of their parent scope: contending on a mutable tail page. - Background compaction may merge adjacent fragments into larger pages after they are below the governing watermark. Readers accept both representations. -- One synchronous maintenance slice examines at most eight buckets or 4 MiB. +- One synchronous maintenance slice examines at most one bucket or 4 MiB so + strict storage commits yield to client commands between buckets. Successful commits enqueue only their touched bucket prefixes. The first maintenance slice after restart rebuilds pending work with one lazy family scan; later slices consume the queue without rescanning the family history. diff --git a/src/api/admin/metrics/domains/stream.rs b/src/api/admin/metrics/domains/stream.rs index ad8b4ee0..1956e7b6 100644 --- a/src/api/admin/metrics/domains/stream.rs +++ b/src/api/admin/metrics/domains/stream.rs @@ -4,15 +4,25 @@ use std::fmt::Write as _; use super::super::rendering::encode_prometheus_label_value; pub(super) fn append_metrics(output: &mut String, runtime: &Runtime) { - append_core_metrics(output, runtime); - append_lag_bucket_metrics(output, runtime); - append_watermark_metrics(output, runtime); + let durable = runtime.stream_durable_metrics_snapshot(); + append_core_metrics(output, runtime, durable.as_ref()); + append_lag_bucket_metrics(output, runtime, durable.as_ref()); + append_watermark_metrics(output, runtime, durable.as_ref()); } -fn append_core_metrics(output: &mut String, runtime: &Runtime) { +fn append_core_metrics( + output: &mut String, + runtime: &Runtime, + durable: Option<&crate::domains::stream::metrics::StreamDurableMetricsSnapshot>, +) { + let metrics = crate::observability::metrics(); output.push_str("# HELP fitz_stream_active Active streams\n"); output.push_str("# TYPE fitz_stream_active gauge\n"); - let _ = writeln!(output, "fitz_stream_active {}", runtime.stream_active()); + let _ = writeln!( + output, + "fitz_stream_active {}", + metrics.gauge_get(crate::domains::stream::metrics::METRIC_ACTIVE_GAUGE) + ); output.push('\n'); output.push_str("# HELP fitz_stream_response_drops_total Total Stream responses dropped by this broker process\n# TYPE fitz_stream_response_drops_total counter\n"); @@ -30,7 +40,7 @@ fn append_core_metrics(output: &mut String, runtime: &Runtime) { let _ = writeln!( output, "fitz_stream_append_sessions_active {}", - runtime.stream_append_sessions_active() + metrics.gauge_get(crate::domains::stream::metrics::METRIC_APPEND_SESSIONS_GAUGE) ); output.push('\n'); @@ -39,7 +49,10 @@ fn append_core_metrics(output: &mut String, runtime: &Runtime) { let _ = writeln!( output, "fitz_stream_events_total {}", - runtime.stream_events_total() + durable.map_or_else( + || runtime.admin_read_model().stream_events_total(), + |snapshot| snapshot.events_total, + ) ); output.push('\n'); @@ -81,7 +94,7 @@ fn append_core_metrics(output: &mut String, runtime: &Runtime) { let _ = writeln!( output, "fitz_stream_subscriptions_active {}", - runtime.stream_subscriptions_active() + metrics.gauge_get(crate::domains::stream::metrics::METRIC_SUBSCRIPTIONS_GAUGE) ); output.push('\n'); @@ -95,8 +108,15 @@ fn append_core_metrics(output: &mut String, runtime: &Runtime) { output.push('\n'); } -fn append_lag_bucket_metrics(output: &mut String, runtime: &Runtime) { - let watermark_lag_buckets = runtime.stream_watermark_lag_buckets(); +fn append_lag_bucket_metrics( + output: &mut String, + runtime: &Runtime, + durable: Option<&crate::domains::stream::metrics::StreamDurableMetricsSnapshot>, +) { + let watermark_lag_buckets = durable.map_or_else( + || runtime.stream_watermark_lag_buckets(), + crate::domains::stream::metrics::StreamDurableMetricsSnapshot::watermark_lag_buckets, + ); output.push_str("# HELP fitz_stream_watermark_lag_bucket_caught_up Stream family watermarks aligned with the fastest family in their area\n"); output.push_str("# TYPE fitz_stream_watermark_lag_bucket_caught_up gauge\n"); let _ = writeln!( @@ -134,20 +154,27 @@ fn append_lag_bucket_metrics(output: &mut String, runtime: &Runtime) { output.push('\n'); } -fn append_watermark_metrics(output: &mut String, runtime: &Runtime) { +fn append_watermark_metrics( + output: &mut String, + runtime: &Runtime, + durable: Option<&crate::domains::stream::metrics::StreamDurableMetricsSnapshot>, +) { output.push_str( "# HELP fitz_stream_realm_watermark Highest committed realm watermark per Stream route family and realm\n", ); output.push_str("# TYPE fitz_stream_realm_watermark gauge\n"); - for detail in runtime.stream_list_realm_watermark_details() { - let realm = encode_prometheus_label_value(&detail.realm); - for watermark in detail.family_watermarks { + if let Some(snapshot) = durable { + for metric in &snapshot.realm_watermarks { let _ = writeln!( output, "fitz_stream_realm_watermark{{realm=\"{}\",family=\"{}\"}} {}", - realm, watermark.family, watermark.watermark + encode_prometheus_label_value(&metric.realm), + metric.family, + metric.watermark ); } + } else { + append_cached_realm_watermarks(output, runtime); } output.push('\n'); @@ -155,7 +182,38 @@ fn append_watermark_metrics(output: &mut String, runtime: &Runtime) { "# HELP fitz_stream_area_watermark Highest committed area watermark per Stream route family, realm, and area\n", ); output.push_str("# TYPE fitz_stream_area_watermark gauge\n"); - for detail in runtime.stream_list_area_watermark_details() { + if let Some(snapshot) = durable { + for metric in &snapshot.area_watermarks { + let _ = writeln!( + output, + "fitz_stream_area_watermark{{realm=\"{}\",area=\"{}\",family=\"{}\"}} {}", + encode_prometheus_label_value(&metric.realm), + encode_prometheus_label_value(&metric.area), + metric.family, + metric.watermark + ); + } + } else { + append_cached_area_watermarks(output, runtime); + } + output.push('\n'); +} + +fn append_cached_realm_watermarks(output: &mut String, runtime: &Runtime) { + for detail in runtime.admin_read_model().stream_realm_watermarks() { + let realm = encode_prometheus_label_value(&detail.realm); + for watermark in detail.family_watermarks { + let _ = writeln!( + output, + "fitz_stream_realm_watermark{{realm=\"{}\",family=\"{}\"}} {}", + realm, watermark.family, watermark.watermark + ); + } + } +} + +fn append_cached_area_watermarks(output: &mut String, runtime: &Runtime) { + for detail in runtime.admin_read_model().stream_area_watermarks() { let realm = encode_prometheus_label_value(&detail.realm); let area = encode_prometheus_label_value(&detail.area); for watermark in detail.family_watermarks { @@ -166,5 +224,4 @@ fn append_watermark_metrics(output: &mut String, runtime: &Runtime) { ); } } - output.push('\n'); } diff --git a/src/api/admin/metrics/mod.rs b/src/api/admin/metrics/mod.rs index 9dcbd4fb..18122233 100644 --- a/src/api/admin/metrics/mod.rs +++ b/src/api/admin/metrics/mod.rs @@ -44,6 +44,7 @@ pub(crate) struct StructuredMetricSample { /// Handle the authenticated structured metrics contract. pub(crate) fn handle_structured_metrics(runtime: &Runtime, family: Option) -> Response { + runtime.refresh_stream_admin_snapshot(); let mut samples = structured_samples(&generate_prometheus_metrics(runtime), family); if let Some(family) = family { samples.extend(family_attributable_samples(runtime, family)); diff --git a/src/api/admin/metrics/tests.rs b/src/api/admin/metrics/tests.rs index cfab679a..338f2587 100644 --- a/src/api/admin/metrics/tests.rs +++ b/src/api/admin/metrics/tests.rs @@ -218,6 +218,36 @@ fn should_export_schedule_metrics_given_preloaded_schedule_runtime() { assert_metric_exported(&metrics, "fitz_notice_wildcard_limit_rejects_total"); } +#[test] +fn should_render_stream_metrics_from_cached_observability_state() { + // Arrange + let metrics = crate::observability::metrics(); + metrics.gauge_set(crate::domains::stream::metrics::METRIC_ACTIVE_GAUGE, 7); + metrics.gauge_set( + crate::domains::stream::metrics::METRIC_APPEND_SESSIONS_GAUGE, + 5, + ); + metrics.gauge_set( + crate::domains::stream::metrics::METRIC_SUBSCRIPTIONS_GAUGE, + 3, + ); + let read_model = crate::control::admin::read_model::AdminReadModel::new(); + read_model.replace_stream_events_total(11); + let runtime = Arc::new(Runtime::with_admin_read_model( + Arc::new(Router::new()), + read_model, + )); + + // Act + let payload = generate_prometheus_metrics(&runtime); + + // Assert + assert!(payload.contains("fitz_stream_active 7")); + assert!(payload.contains("fitz_stream_append_sessions_active 5")); + assert!(payload.contains("fitz_stream_subscriptions_active 3")); + assert!(payload.contains("fitz_stream_events_total 11")); +} + #[test] fn should_export_type_metadata_for_every_metric_family() { // Arrange diff --git a/src/boot/domains.rs b/src/boot/domains.rs index ca93f07f..0d9eb9ae 100644 --- a/src/boot/domains.rs +++ b/src/boot/domains.rs @@ -261,6 +261,12 @@ impl DomainHandles { self.stream.refresh_admin_snapshot_if_dirty(); } + pub(crate) fn stream_durable_metrics_snapshot( + &self, + ) -> crate::domains::stream::metrics::StreamDurableMetricsSnapshot { + self.stream.durable_metrics_snapshot() + } + pub(crate) fn kv_active_transaction_count(&self) -> usize { self.kv.active_transaction_count() } @@ -560,6 +566,7 @@ pub fn setup( &route_families, &metrics, )?; + stream_sink.initialize_admin_snapshot(); register_domain_sink(DomainKind::Stream, router, stream_sink.clone()); let rpc_sink = Arc::new( diff --git a/src/boot/stats/admin_queries.rs b/src/boot/stats/admin_queries.rs index 3e71caf9..12431996 100644 --- a/src/boot/stats/admin_queries.rs +++ b/src/boot/stats/admin_queries.rs @@ -41,6 +41,15 @@ impl Runtime { } } + pub(crate) fn stream_durable_metrics_snapshot( + &self, + ) -> Option { + self.domains + .read() + .as_ref() + .map(|domains| domains.stream_durable_metrics_snapshot()) + } + #[must_use] pub fn kv_list_transactions( &self, @@ -201,13 +210,6 @@ impl Runtime { domains.stream_admin_read_resource_records(request) } - pub(crate) fn stream_list_realm_watermark_details( - &self, - ) -> Vec { - self.refresh_stream_admin_snapshot(); - self.admin_read_model.stream_realm_watermarks() - } - pub(crate) fn stream_realm_watermark_detail( &self, realm: &str, diff --git a/src/boot/storage/tests.rs b/src/boot/storage/tests.rs index 2387bf6c..7c1efdd4 100644 --- a/src/boot/storage/tests.rs +++ b/src/boot/storage/tests.rs @@ -137,8 +137,13 @@ fn should_apply_cloud_throughput_defaults_when_memtable_is_auto() { "tests", ) .memory_budget(MemoryBudget::Bytes(512 * 1024 * 1024)); - let expected_memtable_bytes = - (512 * 1024 * 1024usize).saturating_sub((512 * 1024 * 1024usize) / 10) / 2; + let memory_budget_bytes = 512 * 1024 * 1024usize; + let transaction_pool_bytes = memory_budget_bytes / 10; + let compaction_pool_bytes = memory_budget_bytes / 10; + let expected_memtable_bytes = memory_budget_bytes + .saturating_sub(transaction_pool_bytes) + .saturating_sub(compaction_pool_bytes) + / 2; // Act let tuned = build_midge_open_options(open_options, &config).expect("build cloud options"); diff --git a/src/domains/stream/area_actor.rs b/src/domains/stream/area_actor.rs index 3c5f50fe..c6023560 100644 --- a/src/domains/stream/area_actor.rs +++ b/src/domains/stream/area_actor.rs @@ -31,6 +31,8 @@ pub struct AreaActor { /// Storage layer for watermark persistence store: Arc, + durable_metrics: Arc, + /// Area watermark (highest contiguous committed offset). /// /// `None` means no offset has committed yet. Offset 0 is a valid committed @@ -57,6 +59,7 @@ impl AreaActor { realm: String, area: String, store: Arc, + durable_metrics: Arc, ) -> Self { let (area_watermark, watermark_initialized) = store .get_persisted_area_watermark(family_id.as_u64(), &realm, &area) @@ -80,6 +83,7 @@ impl AreaActor { realm, area, store, + durable_metrics, area_watermark, watermark_initialized, committed_ranges: BTreeMap::new(), @@ -188,6 +192,12 @@ impl AreaActor { fn apply_persisted_watermark(&mut self, current_watermark: u64, ctx: &mut Context) { let previous_watermark = self.area_watermark.unwrap_or(0); self.area_watermark = Some(current_watermark); + self.durable_metrics.set_area_watermark( + self.family_id.as_u64(), + &self.realm, + &self.area, + current_watermark, + ); self.committed_ranges .retain(|_, last_offset| *last_offset > current_watermark); @@ -301,7 +311,13 @@ mod tests { .set_watermark(family.as_u64(), "realm1", "area1", watermark) .expect("persist area watermark"); } - let actor = AreaActor::new(family, "realm1".to_string(), "area1".to_string(), store); + let actor = AreaActor::new( + family, + "realm1".to_string(), + "area1".to_string(), + store, + Arc::new(crate::domains::stream::metrics::StreamDurableMetrics::default()), + ); let ctx = Context::new(addr, router); (actor, ctx) } @@ -325,7 +341,13 @@ mod tests { .expect("Failed to open store"), ); let store = Arc::new(StreamStore::new(db)); - let actor = AreaActor::new(family, "realm1".to_string(), "area1".to_string(), store); + let actor = AreaActor::new( + family, + "realm1".to_string(), + "area1".to_string(), + store, + Arc::new(crate::domains::stream::metrics::StreamDurableMetrics::default()), + ); let ctx = Context::new(addr, router); (actor, ctx, stream_mailbox) } @@ -340,6 +362,10 @@ mod tests { // Assert assert_eq!(actor.watermark(), 3); + assert_eq!( + actor.durable_metrics.snapshot().area_watermarks[0].watermark, + 3 + ); } #[test] diff --git a/src/domains/stream/metrics.rs b/src/domains/stream/metrics.rs index 958823b9..053eaf2f 100644 --- a/src/domains/stream/metrics.rs +++ b/src/domains/stream/metrics.rs @@ -1,4 +1,7 @@ use crate::observability::metrics::{DomainMetricSet, MetricsCollector}; +use parking_lot::RwLock; +use std::collections::BTreeMap; +use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::Instant; pub const METRIC_REQUESTS_TOTAL: &str = "fitz_stream_requests_total"; @@ -26,6 +29,140 @@ pub const METRIC_MAINTENANCE_BUCKETS_COMPACTED_TOTAL: &str = pub const METRIC_ADMIN_PROJECTION_FAILURES_TOTAL: &str = "fitz_stream_admin_projection_failures_total"; +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct StreamRealmWatermarkMetric { + pub(crate) family: u64, + pub(crate) realm: String, + pub(crate) watermark: u64, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct StreamAreaWatermarkMetric { + pub(crate) family: u64, + pub(crate) realm: String, + pub(crate) area: String, + pub(crate) watermark: u64, +} + +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub(crate) struct StreamDurableMetricsSnapshot { + pub(crate) events_total: usize, + pub(crate) realm_watermarks: Vec, + pub(crate) area_watermarks: Vec, +} + +impl StreamDurableMetricsSnapshot { + pub(crate) fn watermark_lag_buckets(&self) -> crate::control::admin::StreamLagBuckets { + let mut watermarks_by_area: BTreeMap<(&str, &str), Vec> = BTreeMap::new(); + for metric in &self.area_watermarks { + watermarks_by_area + .entry((&metric.realm, &metric.area)) + .or_default() + .push(metric.watermark); + } + + watermarks_by_area.values().fold( + crate::control::admin::StreamLagBuckets::default(), + |mut buckets, watermarks| { + let max_watermark = watermarks.iter().copied().max().unwrap_or(0); + for watermark in watermarks { + buckets.record_lag_events(max_watermark.saturating_sub(*watermark)); + } + buckets + }, + ) + } +} + +#[derive(Default)] +pub(crate) struct StreamDurableMetrics { + events_total: AtomicUsize, + realm_watermarks: RwLock>, + area_watermarks: RwLock>, +} + +impl StreamDurableMetrics { + pub(crate) fn observe_snapshot( + &self, + events_total: usize, + realm_details: &[crate::control::admin::StreamRealmWatermarkDetail], + area_details: &[crate::control::admin::StreamAreaWatermarkDetail], + ) { + self.events_total.fetch_max(events_total, Ordering::Relaxed); + let mut realm_watermarks = self.realm_watermarks.write(); + for detail in realm_details { + for watermark in &detail.family_watermarks { + realm_watermarks + .entry((watermark.family, detail.realm.clone())) + .and_modify(|current| *current = (*current).max(watermark.watermark)) + .or_insert(watermark.watermark); + } + } + drop(realm_watermarks); + + let mut area_watermarks = self.area_watermarks.write(); + for detail in area_details { + for watermark in &detail.family_watermarks { + area_watermarks + .entry((watermark.family, detail.realm.clone(), detail.area.clone())) + .and_modify(|current| *current = (*current).max(watermark.watermark)) + .or_insert(watermark.watermark); + } + } + } + + pub(crate) fn record_events(&self, count: usize) { + let _ = self + .events_total + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| { + Some(current.saturating_add(count)) + }); + } + + pub(crate) fn set_realm_watermark(&self, family: u64, realm: &str, watermark: u64) { + self.realm_watermarks + .write() + .insert((family, realm.to_string()), watermark); + } + + pub(crate) fn set_area_watermark(&self, family: u64, realm: &str, area: &str, watermark: u64) { + self.area_watermarks + .write() + .insert((family, realm.to_string(), area.to_string()), watermark); + } + + pub(crate) fn snapshot(&self) -> StreamDurableMetricsSnapshot { + let realm_watermarks = self + .realm_watermarks + .read() + .iter() + .map(|((family, realm), watermark)| StreamRealmWatermarkMetric { + family: *family, + realm: realm.clone(), + watermark: *watermark, + }) + .collect(); + let area_watermarks = self + .area_watermarks + .read() + .iter() + .map( + |((family, realm, area), watermark)| StreamAreaWatermarkMetric { + family: *family, + realm: realm.clone(), + area: area.clone(), + watermark: *watermark, + }, + ) + .collect(); + StreamDurableMetricsSnapshot { + events_total: self.events_total.load(Ordering::Relaxed), + realm_watermarks, + area_watermarks, + } + } +} + #[derive(Clone)] pub struct StreamMetrics { metrics: DomainMetricSet, @@ -84,3 +221,38 @@ impl StreamMetrics { .gauge_set(METRIC_APPEND_SESSIONS_GAUGE, count as u64); } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn should_not_regress_durable_metrics_given_an_older_admin_snapshot() { + // Arrange + let metrics = StreamDurableMetrics::default(); + metrics.record_events(5); + metrics.set_realm_watermark(1, "prod", 4); + metrics.set_area_watermark(1, "prod", "audit", 4); + let realm = crate::control::admin::StreamRealmWatermarkDetail::snapshot( + "prod", + 1, + 1, + vec![crate::control::admin::StreamRealmWatermark::snapshot(1, 2)], + ); + let area = crate::control::admin::StreamAreaWatermarkDetail::snapshot( + "prod", + "audit", + 1, + vec![crate::control::admin::StreamAreaWatermark::snapshot(1, 2)], + ); + + // Act + metrics.observe_snapshot(3, &[realm], &[area]); + let snapshot = metrics.snapshot(); + + // Assert + assert_eq!(snapshot.events_total, 5); + assert_eq!(snapshot.realm_watermarks[0].watermark, 4); + assert_eq!(snapshot.area_watermarks[0].watermark, 4); + } +} diff --git a/src/domains/stream/realm_actor.rs b/src/domains/stream/realm_actor.rs index 395a7b25..ffa1b60a 100644 --- a/src/domains/stream/realm_actor.rs +++ b/src/domains/stream/realm_actor.rs @@ -36,6 +36,8 @@ pub struct RealmActor { /// Storage layer for watermark persistence store: Arc, + durable_metrics: Arc, + /// Realm watermark (highest contiguous committed realm-wide offset). /// /// `None` means no offset has committed yet. Offset 0 is a valid @@ -57,7 +59,12 @@ pub struct RealmActor { } impl RealmActor { - pub fn new(family_id: RouteFamily, realm: String, store: Arc) -> Self { + pub fn new( + family_id: RouteFamily, + realm: String, + store: Arc, + durable_metrics: Arc, + ) -> Self { let (realm_watermark, watermark_initialized) = match store.get_persisted_realm_watermark(family_id.as_u64(), &realm) { Ok(watermark) => (watermark, true), @@ -77,6 +84,7 @@ impl RealmActor { family_id, realm, store, + durable_metrics, realm_watermark, watermark_initialized, committed_ranges: BTreeMap::new(), @@ -180,6 +188,11 @@ impl RealmActor { fn apply_persisted_watermark(&mut self, current_watermark: u64, ctx: &mut Context) { let previous_watermark = self.realm_watermark.unwrap_or(0); self.realm_watermark = Some(current_watermark); + self.durable_metrics.set_realm_watermark( + self.family_id.as_u64(), + &self.realm, + current_watermark, + ); self.committed_ranges .retain(|_, last_offset| *last_offset > current_watermark); @@ -295,7 +308,12 @@ mod tests { .set_realm_watermark(family.as_u64(), "realm1", watermark) .expect("persist realm watermark"); } - let actor = RealmActor::new(family, "realm1".to_string(), store); + let actor = RealmActor::new( + family, + "realm1".to_string(), + store, + Arc::new(crate::domains::stream::metrics::StreamDurableMetrics::default()), + ); let ctx = Context::new(addr, router); (actor, ctx) } @@ -319,7 +337,12 @@ mod tests { .expect("Failed to open store"), ); let store = Arc::new(StreamStore::new(db)); - let actor = RealmActor::new(family, "realm1".to_string(), store); + let actor = RealmActor::new( + family, + "realm1".to_string(), + store, + Arc::new(crate::domains::stream::metrics::StreamDurableMetrics::default()), + ); let ctx = Context::new(addr, router); (actor, ctx, stream_mailbox) } @@ -334,6 +357,10 @@ mod tests { // Assert assert_eq!(actor.watermark(), 3); + assert_eq!( + actor.durable_metrics.snapshot().realm_watermarks[0].watermark, + 3 + ); } #[test] diff --git a/src/domains/stream/sink/delivery/watermark_coordination.rs b/src/domains/stream/sink/delivery/watermark_coordination.rs index d4f8f882..a478a0be 100644 --- a/src/domains/stream/sink/delivery/watermark_coordination.rs +++ b/src/domains/stream/sink/delivery/watermark_coordination.rs @@ -33,6 +33,7 @@ impl StreamDomainCore { ); let realm_spawned = { let store = self.stream_store.clone(); + let durable_metrics = self.durable_metrics.clone(); let realm_owned = realm.to_string(); self.watermark_coordinators.realm.ensure_spawned( StreamRealmScope { @@ -45,6 +46,7 @@ impl StreamDomainCore { family_id, realm_owned.clone(), store.clone(), + durable_metrics.clone(), ) }, ) @@ -59,6 +61,7 @@ impl StreamDomainCore { ); let area_spawned = { let store = self.stream_store.clone(); + let durable_metrics = self.durable_metrics.clone(); let realm_owned = realm.to_string(); let area_owned = area.to_string(); self.watermark_coordinators.area.ensure_spawned( @@ -74,6 +77,7 @@ impl StreamDomainCore { realm_owned.clone(), area_owned.clone(), store.clone(), + durable_metrics.clone(), ) }, ) diff --git a/src/domains/stream/sink/facade.rs b/src/domains/stream/sink/facade.rs index 976d219b..1f1d0c4e 100644 --- a/src/domains/stream/sink/facade.rs +++ b/src/domains/stream/sink/facade.rs @@ -5,8 +5,8 @@ use super::model::{ AdminStreamReadRequestOwned, Arc, AtomicBool, AtomicU64, AtomicUsize, BTreeMap, CleanedUpSessions, HashMap, Mutex, Ordering, Route, RouteFamily, Router, StreamAdminReadCommand, StreamDomainActor, StreamDomainCommand, StreamDomainCore, - StreamDomainSink, StreamLiveCounts, StreamMetrics, StreamReadItem, StreamStorageLayout, - StreamStore, StreamWorkKey, SubscriptionRegistry, WatermarkCoordinators, + StreamDomainSink, StreamDurableMetrics, StreamLiveCounts, StreamMetrics, StreamReadItem, + StreamStorageLayout, StreamStore, StreamWorkKey, SubscriptionRegistry, WatermarkCoordinators, }; use crate::runtime::routing::RouteAddress; use crate::runtime::DeliveryError; @@ -135,6 +135,7 @@ impl StreamDomainSink { ), sync_write_mode: crate::domains::stream::protocol::StreamWriteMode::Sync, metrics: None, + durable_metrics: Arc::new(StreamDurableMetrics::default()), active: Arc::new(AtomicBool::new(true)), family_cores: Arc::new(Mutex::new(BTreeMap::new())), watermark_coordinators: WatermarkCoordinators { @@ -288,6 +289,7 @@ impl StreamDomainSink { ), sync_write_mode: shared.sync_write_mode, metrics: shared.metrics.clone(), + durable_metrics: shared.durable_metrics.clone(), active: shared.active.clone(), family_cores: shared.family_cores.clone(), watermark_coordinators: WatermarkCoordinators { @@ -683,6 +685,16 @@ impl StreamDomainSink { ); } + pub(crate) fn durable_metrics_snapshot( + &self, + ) -> crate::domains::stream::metrics::StreamDurableMetricsSnapshot { + self.core.durable_metrics.snapshot() + } + + pub(crate) fn initialize_admin_snapshot(&self) { + self.core.refresh_admin_snapshot_if_dirty(); + } + #[cfg(test)] pub(super) fn sync_admin_snapshot(&self) { self.send_admin_snapshot_command(StreamDomainCommand::SyncAdminSnapshot, "sync"); diff --git a/src/domains/stream/sink/mailbox_sink_impl/session_operations.rs b/src/domains/stream/sink/mailbox_sink_impl/session_operations.rs index c832b4d6..5b81799a 100644 --- a/src/domains/stream/sink/mailbox_sink_impl/session_operations.rs +++ b/src/domains/stream/sink/mailbox_sink_impl/session_operations.rs @@ -315,6 +315,7 @@ impl StreamDomainCore { Ok(commit) => { self.session_owners.lock().remove(&session_id); self.counter_inc("fitz_stream_append_sessions_ended_total"); + self.durable_metrics.record_events(commit.batch_size); self.notify_area_batch_committed( owner.key.family, &owner.key.realm, diff --git a/src/domains/stream/sink/model.rs b/src/domains/stream/sink/model.rs index 164315af..07ed4141 100644 --- a/src/domains/stream/sink/model.rs +++ b/src/domains/stream/sink/model.rs @@ -1,4 +1,5 @@ pub(super) use crate::dispatch::protocol::payload_codec::PayloadEncoder; +pub(super) use crate::domains::stream::metrics::StreamDurableMetrics; pub(super) use crate::domains::stream::StreamMetrics; pub(super) use crate::domains::stream::{ StreamActor, StreamClientFrame, StreamClientRequest, StreamClientResponseBody, @@ -378,6 +379,7 @@ pub(super) struct StreamDomainCore { pub(super) admin_snapshot: AdminSnapshotState, pub(super) sync_write_mode: crate::domains::stream::protocol::StreamWriteMode, pub(super) metrics: Option, + pub(super) durable_metrics: Arc, pub(super) active: Arc, /// Weak family-core registry used only to aggregate live/admin views. /// Mutable delivery state itself remains owned by each family core. diff --git a/src/domains/stream/sink/observability.rs b/src/domains/stream/sink/observability.rs index 1805b4b6..b66d0fa1 100644 --- a/src/domains/stream/sink/observability.rs +++ b/src/domains/stream/sink/observability.rs @@ -329,6 +329,11 @@ impl StreamDomainCore { stream_area_watermarks: Vec, committed_events_total: usize, ) { + self.durable_metrics.observe_snapshot( + committed_events_total, + &stream_realm_watermarks, + &stream_area_watermarks, + ); self.admin_snapshot .read_model .replace_streams(streams.into_values().collect()); diff --git a/src/domains/stream/sink/tests/sink_dispatch.rs b/src/domains/stream/sink/tests/sink_dispatch.rs index b4b9143a..75b95f90 100644 --- a/src/domains/stream/sink/tests/sink_dispatch.rs +++ b/src/domains/stream/sink/tests/sink_dispatch.rs @@ -55,7 +55,7 @@ fn should_confirm_stream_family_cleanup_before_reporting_delivery() { } #[test] -fn should_run_bounded_stream_maintenance_through_internal_actor_command() { +fn should_yield_bounded_stream_maintenance_through_internal_actor_command() { // Arrange let context = setup_test_context(); for offset in 0..9 { @@ -100,7 +100,7 @@ fn should_run_bounded_stream_maintenance_through_internal_actor_command() { // Assert assert_eq!(records.len(), 9); - assert!(!context + assert!(context .sink .core .stream_store @@ -466,6 +466,7 @@ fn should_preserve_append_session_without_notify_given_commit_failure() { let (retry_commit_type, retry_commit_payload) = extract_single_tlv_field(&retry_commit_frame); let retry_commit_response = request(&context, route, retry_commit_type, retry_commit_payload); let read_after_retry = stream_read_response(&context, route, 0, 10); + let durable_metrics = context.sink.durable_metrics_snapshot(); context.sink.sync_admin_snapshot(); let stream = context .admin_read_model @@ -490,6 +491,7 @@ fn should_preserve_append_session_without_notify_given_commit_failure() { assert_eq!(stream.offset, 0); assert_eq!(stream.sessions_active, 0); assert_eq!(context.admin_read_model.stream_events_total(), 1); + assert_eq!(durable_metrics.events_total, 1); } #[test] diff --git a/src/domains/stream/store/maintenance.rs b/src/domains/stream/store/maintenance.rs index 298e24d0..8d5b5f5a 100644 --- a/src/domains/stream/store/maintenance.rs +++ b/src/domains/stream/store/maintenance.rs @@ -7,7 +7,10 @@ use super::{ use std::collections::BTreeMap; const FRAGMENT_COMPACTION_THRESHOLD: usize = 8; -const MAX_BUCKETS_PER_INVOCATION: usize = 8; +// Maintenance executes on the same synchronous family actor as client Stream +// commands. Yield after one bucket so a series of strict storage commits cannot +// monopolize that actor beyond the client liveness budget. +const MAX_BUCKETS_PER_INVOCATION: usize = 1; const MAX_BYTES_PER_INVOCATION: usize = 4 * 1024 * 1024; #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] @@ -579,8 +582,8 @@ impl StreamStore { /// Runs one bounded synchronous D4 maintenance slice. /// /// Maintenance is separate from append so discovery never adds historical - /// payload reads to commits. One invocation handles at most eight buckets - /// and four MiB across resource, area, realm, global, and posting planes. + /// payload reads to commits. One invocation handles at most one bucket and + /// four MiB across resource, area, realm, global, and posting planes. /// Absolute record expirations remain authoritative after replacement; /// Midge reclaims the replacement at the latest contained deadline. /// diff --git a/src/domains/stream/store/tests.rs b/src/domains/stream/store/tests.rs index 46d20dfb..fce3c023 100644 --- a/src/domains/stream/store/tests.rs +++ b/src/domains/stream/store/tests.rs @@ -27,6 +27,26 @@ impl crate::runtime::clock::Clock for TestStreamClock { } } +fn drain_maintenance(store: &StreamStore, family: u64) -> StreamMaintenanceResult { + let mut total = StreamMaintenanceResult::default(); + for _ in 0..128 { + let slice = store + .run_maintenance(family) + .expect("run bounded Stream maintenance slice"); + total.buckets_compacted = total + .buckets_compacted + .saturating_add(slice.buckets_compacted); + total.records_compacted = total + .records_compacted + .saturating_add(slice.records_compacted); + total.bytes_examined = total.bytes_examined.saturating_add(slice.bytes_examined); + if !store.has_pending_maintenance(family) { + return total; + } + } + panic!("Stream maintenance did not drain within the test slice bound"); +} + mod sessions_layout_and_watermarks; use sessions_layout_and_watermarks::*; mod filters_ttl_and_metadata; diff --git a/src/domains/stream/store/tests/filters_ttl_and_metadata.rs b/src/domains/stream/store/tests/filters_ttl_and_metadata.rs index 8e663c49..37573427 100644 --- a/src/domains/stream/store/tests/filters_ttl_and_metadata.rs +++ b/src/domains/stream/store/tests/filters_ttl_and_metadata.rs @@ -24,9 +24,7 @@ fn should_compact_zero_ttl_fragments_without_positional_gaps() { }) .expect("commit zero-TTL fragment"); } - store - .run_maintenance(1) - .expect("compact zero-TTL fragment round"); + drain_maintenance(&store, 1); } // Act @@ -138,7 +136,7 @@ fn should_preserve_absolute_expiration_before_plus_after_compaction() { }) .expect("read before TTL compaction") .0; - let maintenance = store.run_maintenance(1).expect("compact TTL fragments"); + let maintenance = drain_maintenance(&store, 1); let after = store .read_resource(&ReadResourceParams { family: 1, diff --git a/src/domains/stream/store/tests/global_ordering.rs b/src/domains/stream/store/tests/global_ordering.rs index 88754071..48a233b1 100644 --- a/src/domains/stream/store/tests/global_ordering.rs +++ b/src/domains/stream/store/tests/global_ordering.rs @@ -59,7 +59,7 @@ fn should_append_resource_history_with_immutable_fragments() { } #[test] -fn should_compact_over_fragmented_resource_bucket_in_one_atomic_slice() { +fn should_compact_over_fragmented_resource_bucket_in_one_atomic_replacement() { // Arrange let db = create_test_engine_with_cfs(vec![1]); let store = StreamStore::new(db.clone()); @@ -79,7 +79,7 @@ fn should_compact_over_fragmented_resource_bucket_in_one_atomic_slice() { } // Act - let result = store.run_maintenance(1).expect("run D4 maintenance"); + let result = drain_maintenance(&store, 1); let replay = store .read_resource(&ReadResourceParams { family: 1, @@ -125,7 +125,7 @@ fn should_compact_over_fragmented_resource_bucket_in_one_atomic_slice() { .expect("collect compacted resource"); // Assert - assert!((1..=8).contains(&result.buckets_compacted)); + assert!(result.buckets_compacted > 0); assert!(result.records_compacted >= 9); assert_eq!(rows.len(), 1); assert_eq!(event_records(replay.0).len(), 9); @@ -255,7 +255,7 @@ fn should_fail_closed_when_large_payload_blob_is_missing() { } #[test] -fn should_bound_one_maintenance_slice_to_eight_buckets() { +fn should_yield_stream_maintenance_after_one_bucket() { // Arrange let store = StreamStore::new(create_test_engine_with_cfs(vec![1])); for resource_index in 0..9 { @@ -279,15 +279,17 @@ fn should_bound_one_maintenance_slice_to_eight_buckets() { let first = store .run_maintenance(1) .expect("run first maintenance slice"); + let pending_after_first = store.has_pending_maintenance(1); let second = store .run_maintenance(1) .expect("run second maintenance slice"); // Assert - assert_eq!(first.buckets_compacted, 8); - assert!(first.records_compacted <= 8 * 64); + assert_eq!(first.buckets_compacted, 1); + assert!(first.records_compacted <= 64); assert!(first.bytes_examined <= 4 * 1024 * 1024); - assert!((1..=8).contains(&second.buckets_compacted)); + assert!(pending_after_first); + assert_eq!(second.buckets_compacted, 1); assert_eq!(store.maintenance_full_scan_count_for_tests(), 1); } @@ -422,7 +424,7 @@ fn should_keep_read_snapshot_stable_across_atomic_compaction() { .expect("collect pre-compaction snapshot"); // Act - store.run_maintenance(1).expect("compact snapshot bucket"); + drain_maintenance(&store, 1); let held: Vec<_> = snapshot .scan(&cntryl_midge::Query::new().prefix(prefix.clone())) .expect("rescan held snapshot") diff --git a/src/domains/stream/store/tests/maintenance_and_payloads.rs b/src/domains/stream/store/tests/maintenance_and_payloads.rs index 66324938..55ea4b22 100644 --- a/src/domains/stream/store/tests/maintenance_and_payloads.rs +++ b/src/domains/stream/store/tests/maintenance_and_payloads.rs @@ -56,14 +56,14 @@ fn should_bound_non_compactable_maintenance_groups_by_buckets_examined() { // Assert assert_eq!(first.buckets_compacted, 0); - assert_eq!(attempts_after_first, 8); + assert_eq!(attempts_after_first, 1); assert!(pending_after_first); assert_eq!(second.buckets_compacted, 0); assert_eq!( metrics.counter_get(crate::domains::stream::metrics::METRIC_MAINTENANCE_ATTEMPTS_TOTAL), - 9 + 2 ); - assert!(!store.has_pending_maintenance(1)); + assert!(store.has_pending_maintenance(1)); } #[test] @@ -107,7 +107,7 @@ fn should_count_bytes_examined_for_non_compactable_group() { } #[test] -fn should_count_plus_requeue_over_budget_maintenance_group() { +fn should_preserve_byte_accounting_across_one_bucket_slices() { // Arrange let db = create_test_engine_with_cfs(vec![1]); let store = StreamStore::new(db.clone()); @@ -142,18 +142,25 @@ fn should_count_plus_requeue_over_budget_maintenance_group() { // Act let first = store .run_maintenance(1) - .expect("run byte-bounded maintenance slice"); + .expect("run first maintenance slice"); let pending_after_first = store.has_pending_maintenance(1); - let second = store - .run_maintenance(1) - .expect("run requeued maintenance group"); + let mut remaining = Vec::new(); + for _ in 0..4 { + remaining.push( + store + .run_maintenance(1) + .expect("run remaining maintenance slice"), + ); + } // Assert - assert_eq!(first.buckets_compacted, 4); - assert!(first.bytes_examined > 4 * 1024 * 1024); + assert_eq!(first.buckets_compacted, 1); + assert!(first.bytes_examined > 0); + assert!(first.bytes_examined <= 4 * 1024 * 1024); assert!(pending_after_first); - assert_eq!(second.buckets_compacted, 1); - assert!(second.bytes_examined > 0); + assert!(remaining + .iter() + .all(|slice| slice.buckets_compacted == 1 && slice.bytes_examined > 0)); assert!(!store.has_pending_maintenance(1)); }