From 50b160e90afe8e4c71a01038c3fe99795ae4e4f5 Mon Sep 17 00:00:00 2001 From: zz_y Date: Tue, 12 May 2026 22:23:39 -0600 Subject: [PATCH] refactor(http): drop Arc from HttpServer + runtime-info adapter MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Phase 5 M2.3.6g step 3 — `HttpServer` and the protocol-adapter trait's `handle_runtime_info` family switch from `Arc` to `Arc`. The Prometheus adapter's earliest-timestamp field is renamed from `earliest_timestamp_per_aggregation_id` to `earliest_timestamp_per_sid` (wire-format change; the data shape is otherwise identical: u64 → u64 map). Adds `SketchIndex::earliest_timestamps_per_sid() -> HashMap` that returns each sid's `first_seen_unix_ms` from instance metadata. Replaces the legacy `Store::get_earliest_timestamp_per_aggregation_id` call sites (HTTP runtime-info endpoint, `/api/v1/store/metrics`). Test fixtures that constructed an HttpServer with a `SketchStore` now bind it to `_store` and pass a fresh empty `SketchIndex` to the constructor — they don't exercise the runtime-info path so the empty index is fine. 811/811 lib tests pass. Co-Authored-By: Claude Opus 4.7 (1M context) --- .../drivers/query/adapters/prometheus_http.rs | 19 +++---- .../src/drivers/query/adapters/traits.rs | 11 ++-- data_plane/src/drivers/query/servers/http.rs | 57 +++++++++---------- data_plane/src/main.rs | 2 +- data_plane/src/stores/sketch_db/index/mod.rs | 13 +++++ .../tests/capability_miss_http_e2e_tests.rs | 4 +- .../src/tests/prometheus_forwarding_tests.rs | 12 +++- 7 files changed, 67 insertions(+), 51 deletions(-) diff --git a/data_plane/src/drivers/query/adapters/prometheus_http.rs b/data_plane/src/drivers/query/adapters/prometheus_http.rs index fabebcab..323623d4 100644 --- a/data_plane/src/drivers/query/adapters/prometheus_http.rs +++ b/data_plane/src/drivers/query/adapters/prometheus_http.rs @@ -360,18 +360,14 @@ impl HttpProtocolAdapter for PrometheusHttpAdapter { async fn handle_runtime_info( &self, - store: Arc, + sketch_index: Arc, ) -> Result, StatusCode> { debug!("Handling runtime info request in Prometheus adapter"); - // Get earliest timestamp per aggregation ID from store - let earliest_timestamps = match store.get_earliest_timestamp_per_aggregation_id() { - Ok(timestamps) => timestamps, - Err(e) => { - error!("Error getting earliest timestamps: {}", e); - HashMap::new() - } - }; + // M2.3.6g — earliest timestamps now come from SketchIndex's + // per-sid `first_seen_unix_ms` metadata. Wire field renamed + // accordingly below. + let earliest_timestamps = sketch_index.earliest_timestamps_per_sid(); // Get runtime info from fallback if available let mut runtime_data = if let Some(fallback) = &self.config.fallback { @@ -390,13 +386,12 @@ impl HttpProtocolAdapter for PrometheusHttpAdapter { // Merge local data with fallback data if let Some(data_obj) = runtime_data.as_object_mut() { data_obj.insert( - "earliest_timestamp_per_aggregation_id".to_string(), + "earliest_timestamp_per_sid".to_string(), serde_json::to_value(earliest_timestamps).unwrap_or(json!({})), ); } else { - // If runtime_data is not an object, just create a new one with local data runtime_data = json!({ - "earliest_timestamp_per_aggregation_id": earliest_timestamps + "earliest_timestamp_per_sid": earliest_timestamps }); } diff --git a/data_plane/src/drivers/query/adapters/traits.rs b/data_plane/src/drivers/query/adapters/traits.rs index 28b7dba0..a8c21275 100644 --- a/data_plane/src/drivers/query/adapters/traits.rs +++ b/data_plane/src/drivers/query/adapters/traits.rs @@ -147,20 +147,23 @@ pub trait HttpProtocolAdapter: QueryRequestAdapter + QueryResponseAdapter + Send /// Handle runtime info request /// - /// The adapter can query the store for internal metrics and + /// The adapter can query the SketchIndex for internal metrics and /// optionally forward to fallback backend for additional info. + /// M2.3.6g — switched from `Arc` to + /// `Arc` now that SketchIndex is the only data + /// backend. async fn handle_runtime_info( &self, - store: std::sync::Arc, + sketch_index: std::sync::Arc, ) -> Result, StatusCode>; async fn handle_runtime_info_with_headers( &self, - store: std::sync::Arc, + sketch_index: std::sync::Arc, headers: HashMap, ) -> Result, StatusCode> { // Default implementation ignores headers and calls the old method let _ = headers; - self.handle_runtime_info(store).await + self.handle_runtime_info(sketch_index).await } } diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 5ea6199b..04d50114 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -156,7 +156,8 @@ pub struct HttpServer { /// `KeyByLabelNames` Prometheus needs to populate the `metric` /// map. See `docs/design-gorilla-s3-cold-engine.md` §8. query_router: Arc, - store: Arc, + /// M2.3.6g — SketchIndex replaces `Arc`. + sketch_index: Arc, /// Hot-reloadable `StreamingConfig` source. `None` when hot-reload /// is not wired up by the caller (unit tests, legacy binaries). hot_reload_config: Option, @@ -220,7 +221,10 @@ struct AppState { query_engine: Arc, /// See [`HttpServer::query_router`]. query_router: Arc, - store: Arc, + /// Phase 5 M2.3.6g — SketchIndex replaces `Arc` as the + /// only data backend HTTP-side endpoints consult. Today the only + /// consumer is the runtime-info handler. + sketch_index: Arc, adapter: Arc, fallback: Option>, hot_reload_config: Option, @@ -246,7 +250,7 @@ impl HttpServer { pub fn new( config: HttpServerConfig, query_engine: Arc, - store: Arc, + sketch_index: Arc, ) -> Self { // Bootstrap the capability router with `ASAPQueryEngine` // registered under its canonical query-engine id. @@ -257,7 +261,7 @@ impl HttpServer { config, query_engine, query_router, - store, + sketch_index, hot_reload_config: None, backend_storage_routing: None, schemas: None, @@ -427,7 +431,7 @@ impl HttpServer { config: self.config.clone(), query_engine: self.query_engine, query_router: self.query_router, - store: self.store, + sketch_index: self.sketch_index, adapter: adapter.clone(), fallback: self.config.adapter_config.fallback.clone(), hot_reload_config: self.hot_reload_config.clone(), @@ -518,7 +522,7 @@ impl HttpServer { config: self.config.clone(), query_engine: self.query_engine.clone(), query_router: self.query_router.clone(), - store: self.store.clone(), + sketch_index: self.sketch_index.clone(), adapter: adapter.clone(), fallback: self.config.adapter_config.fallback.clone(), hot_reload_config: self.hot_reload_config.clone(), @@ -1541,7 +1545,7 @@ async fn handle_runtime_info( // Delegate to adapter for protocol-specific handling state .adapter - .handle_runtime_info_with_headers(state.store.clone(), forwarding_headers) + .handle_runtime_info_with_headers(state.sketch_index.clone(), forwarding_headers) .await } @@ -1798,7 +1802,7 @@ mod tests { 15000, )); - let mut server = HttpServer::new(config, query_engine, store); + let mut server = { let _store = store; let idx = Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); HttpServer::new(config, query_engine, idx) }; if let Some(handle) = hot_reload { server = server.with_hot_reload_config(handle); } @@ -2045,7 +2049,7 @@ aggregations: streaming_config.clone(), 15000, )); - let server = HttpServer::new(config, query_engine, store) + let server = { let _store = store; let idx = Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); HttpServer::new(config, query_engine, idx) } .with_hot_reload_config(hot_reload) .with_schemas(schemas); server @@ -2550,7 +2554,7 @@ aggregations: let sc = StreamingConfig::new(map); Arc::new(crate::stores::sketch_db::SchemaRegistry::from_streaming_config(&sc)) }; - let server = HttpServer::new(config, query_engine, store) + let server = { let _store = store; let idx = Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); HttpServer::new(config, query_engine, idx) } .with_backfill_registry(registry) .with_schemas(schemas); server @@ -2901,7 +2905,7 @@ aggregations: 15000, )); let mut server = - HttpServer::new(config, query_engine, store).with_hot_reload_config(hot_reload); + { let _store = store; let idx = Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); HttpServer::new(config, query_engine, idx) }.with_hot_reload_config(hot_reload); for engine in extra_engines { server = server.with_query_engine(engine); } @@ -2946,7 +2950,7 @@ aggregations: streaming_arc, 15000, )); - let mut server = HttpServer::new(config, query_engine, store) + let mut server = { let _store = store; let idx = Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); HttpServer::new(config, query_engine, idx) } .with_hot_reload_config(hot_reload) .with_backend_storage_routing(Arc::new(routing)); for engine in extra_engines { @@ -3694,7 +3698,7 @@ aggregations: 15000, )); let routing_handle = HotReloadBackendStorageRouting::empty(); - let server = HttpServer::new(config, query_engine, store) + let server = { let _store = store; let idx = Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); HttpServer::new(config, query_engine, idx) } .with_hot_reload_config(hot_reload) .with_hot_reload_backend_storage_routing(routing_handle.clone()); let port = server.start_test_server().await.expect("start ok"); @@ -4152,7 +4156,7 @@ aggregations: 15000, )); let mut server = - HttpServer::new(config, query_engine, store).with_hot_reload_config(hot_reload); + { let _store = store; let idx = Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); HttpServer::new(config, query_engine, idx) }.with_hot_reload_config(hot_reload); for engine in engines { server = server.with_query_engine(engine); } @@ -4200,7 +4204,7 @@ aggregations: 15000, )); let cache = Arc::new(crate::query_engines::routing::FreshnessProbeCache::new()); - let server = HttpServer::new(config, query_engine, store).with_probe_cache(cache.clone()); + let server = { let _store = store; let idx = Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); HttpServer::new(config, query_engine, idx) }.with_probe_cache(cache.clone()); let port = server .start_test_server() .await @@ -4557,21 +4561,14 @@ async fn handle_store_metrics(State(state): State) -> axum::response:: use axum::http::StatusCode; use axum::response::IntoResponse; - match state.store.get_earliest_timestamp_per_aggregation_id() { - Ok(timestamps) => { - let body = serde_json::json!({ - "status": "success", - "aggregation_count": timestamps.len(), - "earliest_timestamps": timestamps}); - (StatusCode::OK, axum::Json(body)).into_response() - } - Err(e) => { - let body = serde_json::json!({ - "status": "error", - "error": format!("{}", e)}); - (StatusCode::INTERNAL_SERVER_ERROR, axum::Json(body)).into_response() - } - } + // M2.3.6g — earliest timestamps come from SketchIndex's per-sid + // `first_seen_unix_ms` metadata. Always succeeds (no I/O). + let timestamps = state.sketch_index.earliest_timestamps_per_sid(); + let body = serde_json::json!({ + "status": "success", + "sid_count": timestamps.len(), + "earliest_timestamps_per_sid": timestamps}); + (StatusCode::OK, axum::Json(body)).into_response() } // ─── StreamingConfig hot-reload (PR E) ─────────────────────────────────── diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index fb0d3f21..deacc826 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -610,7 +610,7 @@ async fn main() -> Result<()> { // design, §6). When precompute isn't enabled, the registry is // absent and the swap handler no-ops on schema reconciliation // (legacy per-batch reconcile in ingest still works). - let mut server = HttpServer::new(http_config, engine, store.clone()) + let mut server = HttpServer::new(http_config, engine, sketch_index.clone()) .with_hot_reload_config(hot_reload_config.clone()) .with_probe_cache(probe_cache.clone()); diff --git a/data_plane/src/stores/sketch_db/index/mod.rs b/data_plane/src/stores/sketch_db/index/mod.rs index 225185e2..872b8e85 100644 --- a/data_plane/src/stores/sketch_db/index/mod.rs +++ b/data_plane/src/stores/sketch_db/index/mod.rs @@ -903,6 +903,19 @@ impl SketchIndex { } impl SketchIndex { + /// Phase 5 M2.3.6g — runtime-info / diagnostic helper. Returns the + /// per-sid `first_seen_unix_ms` for every registered sid. The + /// legacy `Store::get_earliest_timestamp_per_aggregation_id` returned + /// an analogous `agg_id → ts` map; this is the SketchIndex + /// equivalent. HTTP server's `/api/v1/status/runtimeinfo` adapter + /// surfaces it under the JSON field `earliest_timestamp_per_sid`. + pub fn earliest_timestamps_per_sid(&self) -> std::collections::HashMap { + let g = self.instances.read().unwrap(); + g.iter() + .map(|(sid, m)| (*sid, m.first_seen_unix_ms.max(0) as u64)) + .collect() + } + /// Phase 5 M2.3.6e — write-side helper. Given an /// `AggregationConfig` and one `(PrecomputedOutput, AggregateCore)` /// pair (the shape both the live worker AND the backfill processor diff --git a/data_plane/src/tests/capability_miss_http_e2e_tests.rs b/data_plane/src/tests/capability_miss_http_e2e_tests.rs index c6098a30..6bef1182 100644 --- a/data_plane/src/tests/capability_miss_http_e2e_tests.rs +++ b/data_plane/src/tests/capability_miss_http_e2e_tests.rs @@ -175,7 +175,9 @@ async fn start_backend(controller_url: String, hot_reload: HotReloadStreamingCon port: 0, handle_http_requests: true, adapter_config}; - let server = HttpServer::new(config, engine, store).with_hot_reload_config(hot_reload.clone()); + let _store = store; + let idx = std::sync::Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); + let server = HttpServer::new(config, engine, idx).with_hot_reload_config(hot_reload.clone()); server .start_test_server() .await diff --git a/data_plane/src/tests/prometheus_forwarding_tests.rs b/data_plane/src/tests/prometheus_forwarding_tests.rs index 8a24a8d9..8e52d469 100644 --- a/data_plane/src/tests/prometheus_forwarding_tests.rs +++ b/data_plane/src/tests/prometheus_forwarding_tests.rs @@ -82,7 +82,9 @@ async fn setup_test_server(prometheus_port: u16) -> (HttpServer, u16) { 15000, // 15s scrape interval )); - let server = HttpServer::new(config, query_engine, store); + let _store = store; + let idx = std::sync::Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); + let server = HttpServer::new(config, query_engine, idx); let actual_port = server .start_test_server() .await @@ -170,7 +172,9 @@ async fn test_forwarding_disabled() { 15000, // 15s scrape interval )); - let server = HttpServer::new(config, query_engine, store); + let _store = store; + let idx = std::sync::Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); + let server = HttpServer::new(config, query_engine, idx); let server_port = server .start_test_server() .await @@ -220,7 +224,9 @@ async fn test_prometheus_server_unreachable() { 15000, // 15s scrape interval )); - let server = HttpServer::new(config, query_engine, store); + let _store = store; + let idx = std::sync::Arc::new(crate::stores::sketch_db::index::SketchIndex::new()); + let server = HttpServer::new(config, query_engine, idx); let server_port = server .start_test_server() .await