Skip to content

Commit f43ee29

Browse files
zzylolclaude
andcommitted
Add CapturingOutputSink and worker correctness unit tests
Adds a CapturingOutputSink to output_sink.rs that stores all emitted (PrecomputedOutput, Box<dyn AggregateCore>) pairs in a Mutex<Vec> for inspection in tests, alongside drain() and len() helpers. Adds 6 unit tests in worker.rs covering: - Raw mode: each sample forwarded as SumAccumulator with sum == value - Tumbling window: correct boundary and aggregated sum on window close - Sliding window pane sharing: single sample emits in both overlapping windows with the correct sum via snapshot/take pane merge - GROUP BY separate emits: two series on the same worker produce independent per-series accumulators (no ingest-time cross-series merge) - Late data Drop: sample behind watermark - allowed_lateness not emitted - Late data ForwardToStore: late sample for evicted pane emitted as mini-accumulator with correct window bounds and sum Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
1 parent 369836c commit f43ee29

3 files changed

Lines changed: 439 additions & 1 deletion

File tree

‎asap-query-engine/Cargo.toml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,7 @@ path = "src/bin/precompute_engine.rs"
6363
name = "test_e2e_precompute"
6464
path = "src/bin/test_e2e_precompute.rs"
6565

66+
6667
[dev-dependencies]
6768
tempfile = "3.20.0"
6869

‎asap-query-engine/src/precompute_engine/output_sink.rs‎

Lines changed: 38 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
use crate::data_model::{AggregateCore, PrecomputedOutput};
22
use crate::stores::Store;
3-
use std::sync::Arc;
3+
use std::sync::{Arc, Mutex};
44
use tracing::debug_span;
55

66
/// Trait for emitting completed window outputs.
@@ -61,6 +61,43 @@ impl OutputSink for RawPassthroughSink {
6161
}
6262
}
6363

64+
/// A capturing sink for testing that stores all emitted outputs.
65+
pub struct CapturingOutputSink {
66+
pub captured: Mutex<Vec<(PrecomputedOutput, Box<dyn AggregateCore>)>>,
67+
}
68+
69+
impl CapturingOutputSink {
70+
pub fn new() -> Self {
71+
Self {
72+
captured: Mutex::new(Vec::new()),
73+
}
74+
}
75+
76+
pub fn drain(&self) -> Vec<(PrecomputedOutput, Box<dyn AggregateCore>)> {
77+
self.captured.lock().unwrap().drain(..).collect()
78+
}
79+
80+
pub fn len(&self) -> usize {
81+
self.captured.lock().unwrap().len()
82+
}
83+
}
84+
85+
impl Default for CapturingOutputSink {
86+
fn default() -> Self {
87+
Self::new()
88+
}
89+
}
90+
91+
impl OutputSink for CapturingOutputSink {
92+
fn emit_batch(
93+
&self,
94+
outputs: Vec<(PrecomputedOutput, Box<dyn AggregateCore>)>,
95+
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
96+
self.captured.lock().unwrap().extend(outputs);
97+
Ok(())
98+
}
99+
}
100+
64101
/// A no-op sink for testing that just counts emitted batches.
65102
pub struct NoopOutputSink {
66103
pub emit_count: std::sync::atomic::AtomicU64,

0 commit comments

Comments
 (0)