use std::collections::VecDeque;
use std::pin::Pin;
use std::sync::Arc;
use serde::{Deserialize, Serialize};
use tokio_stream::Stream;
use crate::channels::{QueryResult, ResultDiff};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct KeyedSnapshotRow {
pub k: u64,
pub v: serde_json::Value,
}
pub const DEFAULT_OUTBOX_CAPACITY: usize = 1000;
#[derive(Debug, Clone)]
pub struct QueryOutputState {
results: im::HashMap<u64, serde_json::Value>,
as_of_sequence: u64,
outbox: VecDeque<Arc<QueryResult>>,
outbox_capacity: usize,
initialized: bool,
generation: u64,
}
impl QueryOutputState {
const MAX_OUTBOX_CAPACITY: usize = 1_000_000;
pub fn new(outbox_capacity: usize) -> Self {
let effective_capacity = outbox_capacity.clamp(1, Self::MAX_OUTBOX_CAPACITY);
Self {
results: im::HashMap::new(),
as_of_sequence: 0,
outbox: VecDeque::with_capacity(effective_capacity.min(1024)),
outbox_capacity: effective_capacity,
initialized: false,
generation: 0,
}
}
pub fn apply_diffs(&mut self, diffs: &[ResultDiff]) {
for diff in diffs {
match diff {
ResultDiff::Add {
data,
row_signature,
} => {
self.results.insert(*row_signature, data.clone());
}
ResultDiff::Delete { row_signature, .. } => {
self.results.remove(row_signature);
}
ResultDiff::Update {
after,
row_signature,
..
} => {
self.results.insert(*row_signature, after.clone());
}
ResultDiff::Aggregation {
after,
row_signature,
..
} => {
self.results.insert(*row_signature, after.clone());
}
ResultDiff::Noop => {}
}
}
}
pub fn advance_sequence_and_push(&mut self, mut result: QueryResult) -> Arc<QueryResult> {
self.as_of_sequence = self.as_of_sequence.saturating_add(1);
result.sequence = self.as_of_sequence;
let arc_result = Arc::new(result);
self.push_outbox(arc_result.clone());
arc_result
}
pub fn apply_committed_sequence(
&mut self,
sequence: u64,
diffs: &[ResultDiff],
mut result: QueryResult,
) -> Arc<QueryResult> {
let expected = self.as_of_sequence.saturating_add(1);
if sequence != expected {
log::error!(
"committed output sequence {sequence} != next in-memory sequence {expected}"
);
}
self.apply_diffs(diffs);
self.as_of_sequence = sequence;
result.sequence = sequence;
let arc_result = Arc::new(result);
self.push_outbox(arc_result.clone());
arc_result
}
fn push_outbox(&mut self, arc_result: Arc<QueryResult>) {
if self.outbox.len() >= self.outbox_capacity {
self.outbox.pop_front();
}
self.outbox.push_back(arc_result);
}
pub fn get_results_as_vec(&self) -> Vec<serde_json::Value> {
self.results.values().cloned().collect()
}
pub fn outbox_capacity(&self) -> usize {
self.outbox_capacity
}
pub fn as_of_sequence(&self) -> u64 {
self.as_of_sequence
}
pub fn outbox_len(&self) -> usize {
self.outbox.len()
}
pub fn outbox_earliest_seq(&self) -> Option<u64> {
self.outbox.front().map(|r| r.sequence)
}
pub fn results_len(&self) -> usize {
self.results.len()
}
pub fn get_result(&self, row_signature: &u64) -> Option<&serde_json::Value> {
self.results.get(row_signature)
}
pub fn clone_results(&self) -> im::HashMap<u64, serde_json::Value> {
self.results.clone()
}
pub fn fetch_outbox_after(
&self,
after_sequence: u64,
) -> Result<Vec<Arc<QueryResult>>, OutboxGap> {
if after_sequence >= self.as_of_sequence {
return Ok(Vec::new());
}
let earliest = self
.outbox
.front()
.map(|r| r.sequence)
.unwrap_or(self.as_of_sequence + 1);
if after_sequence + 1 < earliest {
return Err(OutboxGap {
requested: after_sequence,
earliest_available: earliest,
latest_sequence: self.as_of_sequence,
config_hash: 0, });
}
let entries: Vec<Arc<QueryResult>> = self
.outbox
.iter()
.filter(|r| r.sequence > after_sequence)
.cloned()
.collect();
Ok(entries)
}
pub fn hydrate(
&mut self,
results: im::HashMap<u64, serde_json::Value>,
mut outbox: Vec<Arc<QueryResult>>,
as_of_sequence: u64,
generation: u64,
) {
debug_assert!(
!self.initialized,
"hydrate must run once from uninitialized QueryOutputState"
);
outbox.sort_by_key(|result| result.sequence);
if outbox.len() > self.outbox_capacity {
let skip = outbox.len() - self.outbox_capacity;
outbox = outbox.split_off(skip);
}
self.results = results;
self.outbox = VecDeque::from(outbox);
self.as_of_sequence = self.as_of_sequence.max(as_of_sequence);
self.generation = generation;
self.initialized = true;
}
pub fn reset(&mut self) {
self.reset_from_generation(self.generation);
}
pub fn reset_from_generation(&mut self, persisted: u64) {
self.results.clear();
self.outbox.clear();
self.as_of_sequence = 0;
self.generation = next_output_generation(persisted, self.generation);
self.initialized = true;
}
pub fn mark_initialized(&mut self) {
self.initialized = true;
}
pub fn initialized(&self) -> bool {
self.initialized
}
pub fn generation(&self) -> u64 {
self.generation
}
}
pub fn output_epoch_hash(config_hash: u64, generation: u64) -> u64 {
config_hash.wrapping_add(generation.wrapping_mul(0x9E3779B97F4A7C15))
}
pub(crate) fn next_output_generation(persisted: u64, ram: u64) -> u64 {
persisted.max(ram).saturating_add(1)
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum DurableOutputInconsistency {
#[error("outbox high-water {outbox_hwm} is ahead of stored result sequence {stored_sequence}")]
OutboxAheadOfSequence {
stored_sequence: u64,
outbox_hwm: u64,
},
#[error("durable outbox has a gap; retained sequences: {retained:?}")]
GappedOutbox { retained: Vec<u64> },
#[error("stored result sequence {stored_sequence} has no readable live rows")]
MissingLiveRows { stored_sequence: u64 },
#[error(
"sequence 0 against non-empty durable output (live_rows={live_rows}, outbox_hwm={outbox_hwm:?})"
)]
SequenceZeroAgainstDurableOutput {
live_rows: usize,
outbox_hwm: Option<u64>,
},
#[error("failed to deserialize durable outbox entry at sequence {sequence}: {message}")]
CorruptOutbox { sequence: u64, message: String },
#[error("failed to deserialize durable live row {row_signature}: {message}")]
CorruptLiveRow { row_signature: u64, message: String },
#[error("failed to read durable query output: {message}")]
ReadFailed { message: String },
}
impl DurableOutputInconsistency {
pub(crate) fn is_transient_read(&self) -> bool {
matches!(self, Self::ReadFailed { .. })
}
}
pub(crate) fn reconcile_durable_output(
stored_sequence: Option<u64>,
outbox_sequences: &[u64],
live_row_count: usize,
live_rows_readable: bool,
) -> Result<u64, DurableOutputInconsistency> {
let stored = stored_sequence.unwrap_or(0);
if !live_rows_readable && stored > 0 {
return Err(DurableOutputInconsistency::MissingLiveRows {
stored_sequence: stored,
});
}
let mut retained: Vec<u64> = outbox_sequences.to_vec();
retained.sort_unstable();
if !retained.is_empty() {
let unique_count = {
let mut deduped = retained.clone();
deduped.dedup();
deduped.len()
};
if unique_count != retained.len()
|| retained.windows(2).any(|window| window[1] != window[0] + 1)
{
return Err(DurableOutputInconsistency::GappedOutbox { retained });
}
}
let outbox_hwm = retained.last().copied();
if stored == 0 && (live_row_count > 0 || outbox_hwm.is_some()) {
return Err(
DurableOutputInconsistency::SequenceZeroAgainstDurableOutput {
live_rows: live_row_count,
outbox_hwm,
},
);
}
if let Some(hwm) = outbox_hwm {
if hwm > stored {
return Err(DurableOutputInconsistency::OutboxAheadOfSequence {
stored_sequence: stored,
outbox_hwm: hwm,
});
}
}
Ok(stored)
}
#[derive(Debug, Clone, PartialEq, thiserror::Error)]
#[error("Outbox gap: requested after seq {requested}, but earliest available is {earliest_available} (latest: {latest_sequence})")]
pub struct OutboxGap {
pub requested: u64,
pub earliest_available: u64,
pub latest_sequence: u64,
pub config_hash: u64,
}
#[derive(Debug, Clone, PartialEq, thiserror::Error)]
pub enum FetchError {
#[error("Query is not running (status: {status:?})")]
NotRunning {
status: crate::channels::ComponentStatus,
},
#[error("Timed out waiting for query to finish bootstrapping")]
TimedOut,
#[error(transparent)]
OutboxGap(#[from] OutboxGap),
}
#[derive(Debug, Clone)]
pub struct SnapshotResponse {
results: im::HashMap<u64, serde_json::Value>,
pub as_of_sequence: u64,
pub config_hash: u64,
pub output_generation: u64,
}
impl SnapshotResponse {
pub fn new(
results: im::HashMap<u64, serde_json::Value>,
as_of_sequence: u64,
config_hash: u64,
) -> Self {
Self {
results,
as_of_sequence,
config_hash,
output_generation: 0,
}
}
pub fn with_output_generation(mut self, generation: u64) -> Self {
self.output_generation = generation;
self
}
pub fn stream(self) -> impl Stream<Item = serde_json::Value> + Send {
tokio_stream::iter(self.results.into_iter().map(|(_, v)| v))
}
pub fn stream_keyed(self) -> impl Stream<Item = (u64, serde_json::Value)> + Send {
tokio_stream::iter(self.results)
}
pub fn to_vec(&self) -> Vec<serde_json::Value> {
self.results.values().cloned().collect()
}
pub fn len(&self) -> usize {
self.results.len()
}
pub fn is_empty(&self) -> bool {
self.results.is_empty()
}
}
#[derive(Debug, Clone)]
pub struct OutboxResponse {
pub results: Vec<Arc<QueryResult>>,
pub latest_sequence: u64,
pub config_hash: u64,
pub output_generation: u64,
}
pub struct SnapshotStream {
inner: Pin<Box<dyn Stream<Item = (u64, serde_json::Value)> + Send>>,
pub as_of_sequence: u64,
pub config_hash: u64,
}
impl SnapshotStream {
pub fn from_snapshot(snapshot: SnapshotResponse) -> Self {
let as_of_sequence = snapshot.as_of_sequence;
let config_hash = snapshot.config_hash;
Self {
inner: Box::pin(snapshot.stream_keyed()),
as_of_sequence,
config_hash,
}
}
pub fn from_stream(
stream: impl Stream<Item = serde_json::Value> + Send + 'static,
as_of_sequence: u64,
config_hash: u64,
) -> Self {
use tokio_stream::StreamExt;
Self {
inner: Box::pin(stream.map(|v| (0u64, v))),
as_of_sequence,
config_hash,
}
}
pub fn from_keyed_stream(
stream: impl Stream<Item = (u64, serde_json::Value)> + Send + 'static,
as_of_sequence: u64,
config_hash: u64,
) -> Self {
Self {
inner: Box::pin(stream),
as_of_sequence,
config_hash,
}
}
pub async fn collect_vec(self) -> Vec<serde_json::Value> {
use tokio_stream::StreamExt;
self.inner.map(|(_, v)| v).collect().await
}
pub async fn collect_keyed_vec(self) -> Vec<(u64, serde_json::Value)> {
use tokio_stream::StreamExt;
self.inner.collect().await
}
pub async fn collect_keyed_vec_capped(self, limit: usize) -> Vec<(u64, serde_json::Value)> {
use tokio_stream::StreamExt;
self.inner.take(limit).collect().await
}
pub async fn next_keyed(&mut self) -> Option<(u64, serde_json::Value)> {
use tokio_stream::StreamExt;
self.inner.next().await
}
}
impl Stream for SnapshotStream {
type Item = serde_json::Value;
fn poll_next(
mut self: Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Self::Item>> {
self.inner
.as_mut()
.poll_next(cx)
.map(|opt| opt.map(|(_, v)| v))
}
}
pub struct OutboxStream {
inner: Pin<Box<dyn Stream<Item = Arc<QueryResult>> + Send>>,
pub latest_sequence: u64,
pub config_hash: u64,
}
impl OutboxStream {
pub fn from_outbox(outbox: OutboxResponse) -> Self {
let latest_sequence = outbox.latest_sequence;
let config_hash = outbox.config_hash;
Self {
inner: Box::pin(tokio_stream::iter(outbox.results)),
latest_sequence,
config_hash,
}
}
pub fn from_stream(
stream: impl Stream<Item = Arc<QueryResult>> + Send + 'static,
latest_sequence: u64,
config_hash: u64,
) -> Self {
Self {
inner: Box::pin(stream),
latest_sequence,
config_hash,
}
}
pub async fn collect_vec(self) -> Vec<Arc<QueryResult>> {
use tokio_stream::StreamExt;
self.inner.collect().await
}
}
impl Stream for OutboxStream {
type Item = Arc<QueryResult>;
fn poll_next(
mut self: Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Self::Item>> {
self.inner.as_mut().poll_next(cx)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
fn make_query_result(query_id: &str, diffs: Vec<ResultDiff>) -> QueryResult {
QueryResult::new(
query_id.to_string(),
0, chrono::Utc::now(),
diffs,
HashMap::new(),
)
}
#[test]
fn test_apply_diffs_add() {
let mut state = QueryOutputState::new(10);
let diffs = vec![ResultDiff::Add {
data: serde_json::json!({"name": "Alice"}),
row_signature: 100,
}];
state.apply_diffs(&diffs);
assert_eq!(state.results.len(), 1);
assert_eq!(
state.results.get(&100),
Some(&serde_json::json!({"name": "Alice"}))
);
}
#[test]
fn test_apply_diffs_delete() {
let mut state = QueryOutputState::new(10);
state
.results
.insert(100, serde_json::json!({"name": "Alice"}));
let diffs = vec![ResultDiff::Delete {
data: serde_json::json!({"name": "Alice"}),
row_signature: 100,
}];
state.apply_diffs(&diffs);
assert_eq!(state.results.len(), 0);
}
#[test]
fn test_apply_diffs_update() {
let mut state = QueryOutputState::new(10);
state
.results
.insert(100, serde_json::json!({"name": "Alice"}));
let diffs = vec![ResultDiff::Update {
data: serde_json::json!({"name": "Bob"}),
before: serde_json::json!({"name": "Alice"}),
after: serde_json::json!({"name": "Bob"}),
grouping_keys: None,
row_signature: 100,
}];
state.apply_diffs(&diffs);
assert_eq!(state.results.len(), 1);
assert_eq!(
state.results.get(&100),
Some(&serde_json::json!({"name": "Bob"}))
);
}
#[test]
fn test_apply_diffs_aggregation() {
let mut state = QueryOutputState::new(10);
let diffs = vec![ResultDiff::Aggregation {
before: None,
after: serde_json::json!({"count": 5}),
row_signature: 200,
}];
state.apply_diffs(&diffs);
assert_eq!(state.results.len(), 1);
assert_eq!(
state.results.get(&200),
Some(&serde_json::json!({"count": 5}))
);
let diffs = vec![ResultDiff::Aggregation {
before: Some(serde_json::json!({"count": 5})),
after: serde_json::json!({"count": 10}),
row_signature: 200,
}];
state.apply_diffs(&diffs);
assert_eq!(state.results.len(), 1);
assert_eq!(
state.results.get(&200),
Some(&serde_json::json!({"count": 10}))
);
}
#[test]
fn test_apply_diffs_noop() {
let mut state = QueryOutputState::new(10);
state
.results
.insert(100, serde_json::json!({"name": "Alice"}));
let diffs = vec![ResultDiff::Noop];
state.apply_diffs(&diffs);
assert_eq!(state.results.len(), 1);
}
#[test]
fn test_advance_sequence_and_push() {
let mut state = QueryOutputState::new(3);
let result = make_query_result("q1", vec![]);
let arc = state.advance_sequence_and_push(result);
assert_eq!(arc.sequence, 1);
assert_eq!(state.as_of_sequence, 1);
assert_eq!(state.outbox.len(), 1);
assert_eq!(state.outbox.back().unwrap().sequence, 1);
let result = make_query_result("q1", vec![]);
let arc = state.advance_sequence_and_push(result);
assert_eq!(arc.sequence, 2);
assert_eq!(state.outbox.len(), 2);
}
#[test]
fn test_apply_committed_sequence_uses_durable_seq() {
let mut state = QueryOutputState::new(3);
let diffs = vec![ResultDiff::Add {
data: serde_json::json!({"name": "Alice"}),
row_signature: 1,
}];
let result = make_query_result("q1", diffs.clone());
let arc = state.apply_committed_sequence(1, &diffs, result);
assert_eq!(arc.sequence, 1);
assert_eq!(state.as_of_sequence(), 1);
assert_eq!(state.results_len(), 1);
assert_eq!(state.outbox_len(), 1);
}
#[test]
fn test_outbox_capacity_eviction() {
let mut state = QueryOutputState::new(3);
for _ in 0..5 {
let result = make_query_result("q1", vec![]);
state.advance_sequence_and_push(result);
}
assert_eq!(state.outbox.len(), 3);
assert_eq!(state.as_of_sequence, 5);
assert_eq!(state.outbox.front().unwrap().sequence, 3);
assert_eq!(state.outbox.back().unwrap().sequence, 5);
}
#[test]
fn test_fetch_outbox_after_caught_up() {
let mut state = QueryOutputState::new(10);
let result = make_query_result("q1", vec![]);
state.advance_sequence_and_push(result);
let entries = state.fetch_outbox_after(1).unwrap();
assert!(entries.is_empty());
let entries = state.fetch_outbox_after(100).unwrap();
assert!(entries.is_empty());
}
#[test]
fn test_fetch_outbox_after_returns_entries() {
let mut state = QueryOutputState::new(10);
for _ in 0..5 {
let result = make_query_result("q1", vec![]);
state.advance_sequence_and_push(result);
}
let entries = state.fetch_outbox_after(2).unwrap();
assert_eq!(entries.len(), 3);
assert_eq!(entries[0].sequence, 3);
assert_eq!(entries[1].sequence, 4);
assert_eq!(entries[2].sequence, 5);
}
#[test]
fn test_fetch_outbox_after_gap_error() {
let mut state = QueryOutputState::new(3);
for _ in 0..5 {
let result = make_query_result("q1", vec![]);
state.advance_sequence_and_push(result);
}
let err = state.fetch_outbox_after(0).unwrap_err();
assert_eq!(err.requested, 0);
assert_eq!(err.earliest_available, 3);
assert_eq!(err.latest_sequence, 5);
assert_eq!(err.config_hash, 0); }
#[test]
fn test_get_results_as_vec() {
let mut state = QueryOutputState::new(10);
state.results.insert(1, serde_json::json!({"a": 1}));
state.results.insert(2, serde_json::json!({"b": 2}));
let vec = state.get_results_as_vec();
assert_eq!(vec.len(), 2);
assert!(vec.contains(&serde_json::json!({"a": 1})));
assert!(vec.contains(&serde_json::json!({"b": 2})));
}
#[test]
fn test_snapshot_clone_is_independent() {
let mut state = QueryOutputState::new(10);
state
.results
.insert(1, serde_json::json!({"name": "Alice"}));
let snapshot = state.results.clone();
state.results.insert(1, serde_json::json!({"name": "Bob"}));
assert_eq!(
snapshot.get(&1),
Some(&serde_json::json!({"name": "Alice"}))
);
assert_eq!(
state.results.get(&1),
Some(&serde_json::json!({"name": "Bob"}))
);
}
#[test]
fn test_outbox_capacity_zero_clamped_to_one() {
let mut state = QueryOutputState::new(0);
assert_eq!(state.outbox_capacity, 1);
let result = make_query_result("q1", vec![]);
state.advance_sequence_and_push(result);
assert_eq!(state.outbox.len(), 1);
let result = make_query_result("q1", vec![]);
state.advance_sequence_and_push(result);
assert_eq!(state.outbox.len(), 1);
assert_eq!(state.outbox.front().unwrap().sequence, 2);
}
#[tokio::test]
async fn snapshot_stream_yields_all_values() {
use tokio_stream::StreamExt;
let mut map = im::HashMap::new();
map.insert(1, serde_json::json!({"id": 1}));
map.insert(2, serde_json::json!({"id": 2}));
map.insert(3, serde_json::json!({"id": 3}));
let snap = SnapshotResponse::new(map, 10, 42);
assert_eq!(snap.len(), 3);
assert!(!snap.is_empty());
let mut collected: Vec<serde_json::Value> = snap.stream().collect().await;
collected.sort_by_key(|v| v["id"].as_u64().unwrap());
assert_eq!(collected.len(), 3);
assert_eq!(collected[0]["id"], 1);
assert_eq!(collected[1]["id"], 2);
assert_eq!(collected[2]["id"], 3);
}
#[tokio::test]
async fn test_snapshot_stream_preserves_row_signatures() {
use tokio_stream::StreamExt;
let mut map = im::HashMap::new();
map.insert(11u64, serde_json::json!({"id": 1}));
map.insert(22u64, serde_json::json!({"id": 2}));
let snap = SnapshotResponse::new(map, 5, 7);
let mut keyed: Vec<(u64, serde_json::Value)> = snap.clone().stream_keyed().collect().await;
keyed.sort_by_key(|(sig, _)| *sig);
assert_eq!(keyed.len(), 2);
assert_eq!(keyed[0].0, 11);
assert_eq!(keyed[0].1["id"], 1);
assert_eq!(keyed[1].0, 22);
let stream = SnapshotStream::from_snapshot(snap);
assert_eq!(stream.as_of_sequence, 5);
let mut via_stream = stream.collect_keyed_vec().await;
via_stream.sort_by_key(|(sig, _)| *sig);
assert_eq!(via_stream.len(), 2);
assert_eq!(via_stream[0].0, 11);
assert_eq!(via_stream[1].0, 22);
}
#[tokio::test]
async fn test_snapshot_stream_from_bare_values_uses_zero_signature() {
let stream = SnapshotStream::from_stream(
tokio_stream::iter(vec![serde_json::json!({"id": 1})]),
0,
0,
);
let keyed = stream.collect_keyed_vec().await;
assert_eq!(keyed.len(), 1);
assert_eq!(
keyed[0].0, 0,
"bare-value stream rows have unknown signature 0"
);
}
#[test]
fn reconcile_matching_sequence_and_outbox_hwm() {
let seq = reconcile_durable_output(Some(4), &[1, 2, 3, 4], 3, true).unwrap();
assert_eq!(seq, 4);
}
#[test]
fn reconcile_does_not_lower_stored_sequence_below_outbox_hwm() {
let seq = reconcile_durable_output(Some(5), &[3, 4], 1, true).unwrap();
assert_eq!(seq, 5);
}
#[test]
fn reconcile_outbox_ahead_of_sequence_is_inconsistent() {
let err = reconcile_durable_output(Some(4), &[1, 2, 3, 4, 5], 3, true).unwrap_err();
assert_eq!(
err,
DurableOutputInconsistency::OutboxAheadOfSequence {
stored_sequence: 4,
outbox_hwm: 5,
}
);
}
#[test]
fn reconcile_gapped_outbox_is_inconsistent() {
let err = reconcile_durable_output(Some(5), &[1, 2, 4, 5], 2, true).unwrap_err();
assert!(matches!(
err,
DurableOutputInconsistency::GappedOutbox { retained } if retained == vec![1, 2, 4, 5]
));
}
#[test]
fn reconcile_duplicate_outbox_sequence_is_inconsistent() {
let err = reconcile_durable_output(Some(3), &[1, 2, 2, 3], 2, true).unwrap_err();
assert!(matches!(
err,
DurableOutputInconsistency::GappedOutbox { retained } if retained == vec![1, 2, 2, 3]
));
}
#[test]
fn reconcile_missing_live_rows_for_stored_sequence() {
let err = reconcile_durable_output(Some(3), &[1, 2, 3], 0, false).unwrap_err();
assert_eq!(
err,
DurableOutputInconsistency::MissingLiveRows { stored_sequence: 3 }
);
}
#[test]
fn reconcile_sequence_zero_against_live_rows_is_inconsistent() {
let err = reconcile_durable_output(None, &[], 2, true).unwrap_err();
assert_eq!(
err,
DurableOutputInconsistency::SequenceZeroAgainstDurableOutput {
live_rows: 2,
outbox_hwm: None,
}
);
}
#[test]
fn reconcile_sequence_zero_against_outbox_is_inconsistent() {
let err = reconcile_durable_output(None, &[1, 2, 3], 0, true).unwrap_err();
assert_eq!(
err,
DurableOutputInconsistency::SequenceZeroAgainstDurableOutput {
live_rows: 0,
outbox_hwm: Some(3),
}
);
}
#[test]
fn reconcile_empty_durable_output_is_sequence_zero() {
assert_eq!(reconcile_durable_output(None, &[], 0, true).unwrap(), 0);
assert_eq!(reconcile_durable_output(Some(0), &[], 0, true).unwrap(), 0);
}
#[test]
fn hydrate_installs_results_outbox_and_sequence() {
let mut state = QueryOutputState::new(10);
let mut results = im::HashMap::new();
results.insert(1, serde_json::json!({"id": "p1"}));
let outbox = vec![
Arc::new(make_query_result("q1", vec![])),
Arc::new(make_query_result("q1", vec![])),
];
let outbox: Vec<Arc<QueryResult>> = outbox
.into_iter()
.enumerate()
.map(|(i, result)| {
let mut owned = (*result).clone();
owned.sequence = (i as u64) + 1;
Arc::new(owned)
})
.collect();
state.hydrate(results.clone(), outbox.clone(), 2, 0);
assert_eq!(state.as_of_sequence(), 2);
assert_eq!(state.results_len(), 1);
assert_eq!(state.outbox_len(), 2);
assert_eq!(state.outbox_earliest_seq(), Some(1));
let fetched = state.fetch_outbox_after(0).unwrap();
assert_eq!(fetched.len(), 2);
assert_eq!(fetched[0].sequence, 1);
assert_eq!(fetched[1].sequence, 2);
}
#[test]
fn hydrate_from_empty_installs_sequence() {
let mut state = QueryOutputState::new(10);
assert_eq!(state.as_of_sequence(), 0);
state.hydrate(im::HashMap::new(), Vec::new(), 4, 0);
assert_eq!(state.as_of_sequence(), 4);
}
#[test]
fn hydrate_trims_outbox_to_capacity_keeping_newest() {
let mut state = QueryOutputState::new(2);
let outbox: Vec<Arc<QueryResult>> = (1..=4)
.map(|seq| {
let mut result = make_query_result("q1", vec![]);
result.sequence = seq;
Arc::new(result)
})
.collect();
state.hydrate(im::HashMap::new(), outbox, 4, 0);
assert_eq!(state.outbox_len(), 2);
assert_eq!(state.outbox_earliest_seq(), Some(3));
assert_eq!(state.as_of_sequence(), 4);
let fetched = state.fetch_outbox_after(2).unwrap();
assert_eq!(
fetched.iter().map(|r| r.sequence).collect::<Vec<_>>(),
vec![3, 4]
);
}
#[test]
fn reset_clears_hydrated_state() {
let mut state = QueryOutputState::new(10);
let mut results = im::HashMap::new();
results.insert(1, serde_json::json!({"id": "p1"}));
let mut result = make_query_result("q1", vec![]);
result.sequence = 1;
state.hydrate(results, vec![Arc::new(result)], 1, 0);
state.reset();
assert_eq!(state.as_of_sequence(), 0);
assert_eq!(state.results_len(), 0);
assert_eq!(state.outbox_len(), 0);
assert_eq!(state.generation(), 1);
assert!(state.initialized());
}
#[test]
fn reset_from_persisted_generation_does_not_restart_at_one() {
let mut state = QueryOutputState::new(10);
state.reset_from_generation(4);
assert_eq!(state.generation(), 5);
assert!(state.initialized());
assert_eq!(state.as_of_sequence(), 0);
}
#[test]
fn next_output_generation_always_bumps_from_disk() {
assert_eq!(next_output_generation(0, 0), 1);
assert_eq!(next_output_generation(1, 0), 2);
assert_eq!(next_output_generation(1, 1), 2);
assert_eq!(next_output_generation(4, 1), 5);
}
#[test]
fn reset_from_generation_bumps_existing_ram_generation() {
let mut state = QueryOutputState::new(10);
state.reset_from_generation(1);
assert_eq!(state.generation(), 2);
state.reset_from_generation(1);
assert_eq!(state.generation(), 3);
}
#[test]
fn output_epoch_hash_generation_zero_is_identity() {
assert_eq!(output_epoch_hash(42, 0), 42);
assert_ne!(output_epoch_hash(42, 1), 42);
}
}