diff --git a/crates/temporalstore-rust/src/meta.rs b/crates/temporalstore-rust/src/meta.rs index 665d1c1ec..3256c5ca0 100644 --- a/crates/temporalstore-rust/src/meta.rs +++ b/crates/temporalstore-rust/src/meta.rs @@ -6100,6 +6100,95 @@ mod tests { assert_eq!(reopened.list_namespaces().namespaces.len(), 1); } + #[test] + fn the_version_report_counts_each_state_of_each_resource() { + // Nine numbers, and nothing checked any of them. They were nine + // separate walks of three collections; they are now one walk each, and + // a walk that tallies has to put every resource in the same place the + // separate filters did. + let meta = SingleNodeMeta::default(); + for (i, want) in ["normal", "frozen", "dropped"].iter().enumerate() { + let addr = format!("node-{i}"); + assert!(meta + .register_server(RegisterServerRequest { + numa_nodes: Vec::new(), + server_addr: addr.clone(), + node_id: i as u64 + 1, + location: "rack-1".to_string(), + binary_version: "v1".to_string(), + }) + .status + .ok); + assert!(meta + .register_proxy(RegisterProxyRequest { + proxy_addr: format!("proxy-{i}"), + namespace: String::new(), + location: "rack-1".to_string(), + config_version: 0, + binary_version: "v1".to_string(), + }) + .status + .ok); + let change = |endpoint: String| StateChangeRequest { + endpoint, + freeze_cooldown_ms: 0, + reason: FreezeReason::Unspecified, + }; + if *want == "frozen" { + assert!(meta.freeze_server(change(addr.clone())).status.ok); + assert!(meta.freeze_proxy(change(format!("proxy-{i}"))).status.ok); + } + if *want == "dropped" { + assert!(meta.drop_server(change(addr)).status.ok); + assert!(meta.drop_proxy(change(format!("proxy-{i}"))).status.ok); + } + } + assert!(meta + .add_namespace(AddNamespaceRequest { + namespace: "ns".to_string() + }) + .status + .ok); + for (i, want) in ["normal", "frozen", "dropped"].iter().enumerate() { + let name = format!("t{i}"); + assert!(meta + .add_table(AddTableRequest { + namespace: "ns".to_string(), + table_name: name.clone(), + first_shard_id: (i as u64 + 1) * 100, + shard_count: 1, + replica_count: 1, + partition_version: 0, + serving_options: TableServingOptions::default(), + }) + .status + .ok); + let request = DeleteTableRequest { + namespace: "ns".to_string(), + table_name: name, + }; + if *want == "frozen" { + assert!(meta.freeze_table(request).status.ok); + } else if *want == "dropped" { + assert!(meta.delete_table(request).status.ok); + } + } + + let report = meta.topology_version_report(TopologyVersionRequest::default()); + assert_eq!(report.normal_servers, 1, "servers: normal"); + assert_eq!(report.frozen_servers, 1, "servers: frozen"); + assert_eq!(report.dropped_servers, 1, "servers: dropped"); + assert_eq!(report.normal_proxies, 1, "proxies: normal"); + assert_eq!(report.frozen_proxies, 1, "proxies: frozen"); + assert_eq!(report.dropped_proxies, 1, "proxies: dropped"); + assert_eq!(report.normal_tables, 1, "tables: normal"); + assert_eq!(report.frozen_tables, 1, "tables: frozen"); + assert_eq!(report.dropped_tables, 1, "tables: dropped"); + assert_eq!(report.server_count, 3); + assert_eq!(report.proxy_count, 3); + assert_eq!(report.table_count, 3); + } + #[test] fn the_learned_shard_index_answers_what_the_scan_answered() { // The index exists because asking every table which shard it owns, once diff --git a/crates/temporalstore-rust/src/meta/topology_helpers.rs b/crates/temporalstore-rust/src/meta/topology_helpers.rs index b7537d2ce..d6dd41f6d 100644 --- a/crates/temporalstore-rust/src/meta/topology_helpers.rs +++ b/crates/temporalstore-rust/src/meta/topology_helpers.rs @@ -216,81 +216,90 @@ pub(super) fn stats_from_state(state: &MetaState, counters: &MetaCounters) -> Me pub(super) const TOPOLOGY_EVENT_HISTORY_LIMIT: usize = 256; +/// How many of a collection are normal, frozen and dropped, in one walk. +/// +/// The report wants all three, and asking for them separately walked the +/// collection three times. They are disjoint, so one walk answers all of them. +fn tally_states(states: impl Iterator) -> (usize, usize, usize) { + let mut normal = 0; + let mut frozen = 0; + let mut dropped = 0; + for state in states { + match state { + MetaEntityState::Normal => normal += 1, + MetaEntityState::Frozen => frozen += 1, + MetaEntityState::Dropped => dropped += 1, + // Any other state is counted by none of the three, exactly as the + // separate filters counted it by none of them. + _ => {} + } + } + (normal, frozen, dropped) +} + pub(super) fn topology_version_report_from_state( state: &MetaState, old_topology_version: u64, ) -> TopologyVersionReport { - let changed_tables = state - .tables - .values() - .filter(|table| table.info.topology_version > old_topology_version) - .map(|table| table.info.clone()) - .collect::>(); - let events = state - .topology_events - .iter() - .filter(|event| event.topology_version > old_topology_version) - .cloned() - .collect::>(); + // Nothing can be newer than a caller who is already current, so there is + // nothing to look for. Every table and every event carries a version that + // `record_topology_event` returned, which is the global version as it stood + // when it was bumped -- so none of them can be above it. + // + // This is the answer almost every check gets: callers ask whether the + // topology moved far more often than it moves, and producing this answer + // used to walk every table and every recorded event to conclude that none + // of them qualified. + let unchanged = old_topology_version >= state.topology_version; + let changed_tables = if unchanged { + Vec::new() + } else { + state + .tables + .values() + .filter(|table| table.info.topology_version > old_topology_version) + .map(|table| table.info.clone()) + .collect::>() + }; + let events = if unchanged { + Vec::new() + } else { + state + .topology_events + .iter() + .filter(|event| event.topology_version > old_topology_version) + .cloned() + .collect::>() + }; let event_history_truncated = old_topology_version < state.topology_version && state .topology_events .front() .is_some_and(|event| old_topology_version < event.topology_version.saturating_sub(1)); + let (normal_servers, frozen_servers, dropped_servers) = + tally_states(state.servers.values().map(|server| server.state)); + let (normal_proxies, frozen_proxies, dropped_proxies) = + tally_states(state.proxies.values().map(|proxy| proxy.state)); + let (normal_tables, frozen_tables, dropped_tables) = + tally_states(state.tables.values().map(|table| table.info.state)); TopologyVersionReport { status: Status::ok(), current_topology_version: state.topology_version, old_topology_version, - unchanged: old_topology_version >= state.topology_version, + unchanged, server_count: state.servers.len(), proxy_count: state.proxies.len(), table_count: state.tables.len(), shard_route_count: state.shards.len(), - normal_servers: state - .servers - .values() - .filter(|server| server.state == MetaEntityState::Normal) - .count(), - frozen_servers: state - .servers - .values() - .filter(|server| server.state == MetaEntityState::Frozen) - .count(), - dropped_servers: state - .servers - .values() - .filter(|server| server.state == MetaEntityState::Dropped) - .count(), - normal_proxies: state - .proxies - .values() - .filter(|proxy| proxy.state == MetaEntityState::Normal) - .count(), - frozen_proxies: state - .proxies - .values() - .filter(|proxy| proxy.state == MetaEntityState::Frozen) - .count(), - dropped_proxies: state - .proxies - .values() - .filter(|proxy| proxy.state == MetaEntityState::Dropped) - .count(), - normal_tables: state - .tables - .values() - .filter(|table| table.info.state == MetaEntityState::Normal) - .count(), - frozen_tables: state - .tables - .values() - .filter(|table| table.info.state == MetaEntityState::Frozen) - .count(), - dropped_tables: state - .tables - .values() - .filter(|table| table.info.state == MetaEntityState::Dropped) - .count(), + normal_servers, + frozen_servers, + dropped_servers, + normal_proxies, + frozen_proxies, + dropped_proxies, + normal_tables, + frozen_tables, + dropped_tables, changed_tables, events, event_history_truncated,