Skip to content

Commit 8661b52

Browse files
zzylolclaude
andcommitted
refactor: switch SimpleMapStore from legacy to current store implementations
SimpleMapStore now wraps SimpleMapStoreGlobal/SimpleMapStorePerKey (the epoch-based implementations) instead of the legacy stores. Adds diagnostic_info() to both current stores. Legacy stores remain available for benchmark comparisons only. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent c0c739c commit 8661b52

3 files changed

Lines changed: 109 additions & 8 deletions

File tree

‎asap-query-engine/src/stores/simple_map_store/global.rs‎

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -154,6 +154,49 @@ impl SimpleMapStoreGlobal {
154154
cleanup_policy,
155155
}
156156
}
157+
158+
/// Collect diagnostic info about store contents.
159+
pub fn diagnostic_info(&self) -> super::StoreDiagnostics {
160+
use super::{AggregationDiagnostic, StoreDiagnostics};
161+
162+
let data = self.lock.lock().unwrap();
163+
let mut per_aggregation = Vec::new();
164+
let mut total_time_map_entries: usize = 0;
165+
let total_sketch_bytes: usize = 0;
166+
167+
for (&agg_id, per_key) in &data.stores {
168+
let time_map_len = per_key.current_epoch.window_count()
169+
+ per_key
170+
.sealed_epochs
171+
.values()
172+
.map(|e| e.distinct_window_count())
173+
.sum::<usize>();
174+
let read_counts_len = data
175+
.read_counts
176+
.get(&agg_id)
177+
.map(|rc| rc.len())
178+
.unwrap_or(0);
179+
total_time_map_entries += time_map_len;
180+
181+
let num_aggregate_objects = per_key.current_epoch.len()
182+
+ per_key.sealed_epochs.values().map(|e| e.entries.len()).sum::<usize>();
183+
184+
per_aggregation.push(AggregationDiagnostic {
185+
aggregation_id: agg_id,
186+
time_map_len,
187+
read_counts_len,
188+
num_aggregate_objects,
189+
sketch_bytes: 0, // skip serialization for diagnostics
190+
});
191+
}
192+
193+
StoreDiagnostics {
194+
num_aggregations: data.stores.len(),
195+
total_time_map_entries,
196+
total_sketch_bytes,
197+
per_aggregation,
198+
}
199+
}
157200
}
158201

159202
/// Extracted config fields needed inside the locked batch loop.

‎asap-query-engine/src/stores/simple_map_store/mod.rs‎

Lines changed: 23 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -7,17 +7,32 @@ use crate::data_model::{
77
AggregateCore, CleanupPolicy, LockStrategy, PrecomputedOutput, StreamingConfig,
88
};
99
use crate::stores::{Store, StoreResult, TimestampedBucketsMap};
10+
use global::SimpleMapStoreGlobal;
11+
use per_key::SimpleMapStorePerKey;
1012
use std::collections::HashMap;
1113
use std::sync::Arc;
1214

13-
pub use legacy::LegacySimpleMapStoreGlobal;
14-
pub use legacy::LegacySimpleMapStorePerKey;
15-
pub use legacy::{AggregationDiagnostic, StoreDiagnostics};
15+
/// Diagnostic snapshot from a single aggregation ID in the store.
16+
pub struct AggregationDiagnostic {
17+
pub aggregation_id: u64,
18+
pub time_map_len: usize,
19+
pub read_counts_len: usize,
20+
pub num_aggregate_objects: usize,
21+
pub sketch_bytes: usize,
22+
}
23+
24+
/// Diagnostic snapshot of the entire store.
25+
pub struct StoreDiagnostics {
26+
pub num_aggregations: usize,
27+
pub total_time_map_entries: usize,
28+
pub total_sketch_bytes: usize,
29+
pub per_aggregation: Vec<AggregationDiagnostic>,
30+
}
1631

1732
/// Enum wrapper that dispatches to either global or per-key lock implementation
1833
pub enum SimpleMapStore {
19-
Global(LegacySimpleMapStoreGlobal),
20-
PerKey(LegacySimpleMapStorePerKey),
34+
Global(SimpleMapStoreGlobal),
35+
PerKey(SimpleMapStorePerKey),
2136
}
2237

2338
impl SimpleMapStore {
@@ -26,7 +41,7 @@ impl SimpleMapStore {
2641
Self::new_with_strategy(streaming_config, cleanup_policy, LockStrategy::PerKey)
2742
}
2843

29-
/// Collect diagnostic info for memory leak investigation.
44+
/// Collect diagnostic info for memory investigation.
3045
pub fn diagnostic_info(&self) -> StoreDiagnostics {
3146
match self {
3247
SimpleMapStore::Global(store) => store.diagnostic_info(),
@@ -41,11 +56,11 @@ impl SimpleMapStore {
4156
lock_strategy: LockStrategy,
4257
) -> Self {
4358
match lock_strategy {
44-
LockStrategy::Global => SimpleMapStore::Global(LegacySimpleMapStoreGlobal::new(
59+
LockStrategy::Global => SimpleMapStore::Global(SimpleMapStoreGlobal::new(
4560
streaming_config,
4661
cleanup_policy,
4762
)),
48-
LockStrategy::PerKey => SimpleMapStore::PerKey(LegacySimpleMapStorePerKey::new(
63+
LockStrategy::PerKey => SimpleMapStore::PerKey(SimpleMapStorePerKey::new(
4964
streaming_config,
5065
cleanup_policy,
5166
)),

‎asap-query-engine/src/stores/simple_map_store/per_key.rs‎

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -183,6 +183,49 @@ impl SimpleMapStorePerKey {
183183
}
184184
}
185185

186+
/// Collect diagnostic info about store contents.
187+
pub fn diagnostic_info(&self) -> super::StoreDiagnostics {
188+
use super::{AggregationDiagnostic, StoreDiagnostics};
189+
190+
let mut per_aggregation = Vec::new();
191+
let mut total_time_map_entries: usize = 0;
192+
let total_sketch_bytes: usize = 0;
193+
194+
for entry in self.store.iter() {
195+
let agg_id = *entry.key();
196+
let data = match entry.value().read() {
197+
Ok(d) => d,
198+
Err(_) => continue,
199+
};
200+
let time_map_len = data.current_epoch.window_count()
201+
+ data
202+
.sealed_epochs
203+
.values()
204+
.map(|e| e.distinct_window_count())
205+
.sum::<usize>();
206+
let read_counts_len = data.read_counts.lock().map(|rc| rc.len()).unwrap_or(0);
207+
total_time_map_entries += time_map_len;
208+
209+
let num_aggregate_objects = data.current_epoch.len()
210+
+ data.sealed_epochs.values().map(|e| e.entries.len()).sum::<usize>();
211+
212+
per_aggregation.push(AggregationDiagnostic {
213+
aggregation_id: agg_id,
214+
time_map_len,
215+
read_counts_len,
216+
num_aggregate_objects,
217+
sketch_bytes: 0, // skip serialization for diagnostics
218+
});
219+
}
220+
221+
StoreDiagnostics {
222+
num_aggregations: self.store.len(),
223+
total_time_map_entries,
224+
total_sketch_bytes,
225+
per_aggregation,
226+
}
227+
}
228+
186229
fn cleanup_old_aggregates(
187230
&self,
188231
data: &mut StoreKeyData,

0 commit comments

Comments
 (0)