use std::collections::BTreeMap;
use crate::archive::Archive;
use crate::error::{Error, Result};
use super::frame::{Frame, IndexKind, IndexState, NO_INDEX_STATE};
fn fold_index_state(state: IndexState, ts: i64, blob: &[u8]) -> IndexState {
const PRIME: u64 = 0x0000_0100_0000_01b3;
const BASIS_A: u64 = 0xcbf2_9ce4_8422_2325;
const BASIS_B: u64 = 0x9e37_79b9_7f4a_7c15;
let mut a = if state == NO_INDEX_STATE {
BASIS_A
} else {
state.0
};
let mut b = if state == NO_INDEX_STATE {
BASIS_B
} else {
state.1
};
for byte in ts.to_le_bytes().iter().chain(blob) {
a = (a ^ u64::from(*byte)).wrapping_mul(PRIME);
b = (b ^ u64::from(*byte)).wrapping_mul(PRIME.rotate_left(17));
}
if (a, b) == NO_INDEX_STATE {
(1, 0)
} else {
(a, b)
}
}
#[derive(Default)]
struct StreamCursor {
wal_after: Option<i64>,
index_after: Option<(i64, usize)>,
}
struct SourceCursor {
ordinal: u32,
source_id: i64,
streams: BTreeMap<String, StreamCursor>,
clock_after: Option<i64>,
seq: u64,
index_state: IndexState,
sent_full: bool,
}
pub struct ArchivePublisher {
sources: Vec<SourceCursor>,
}
impl std::fmt::Debug for ArchivePublisher {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ArchivePublisher")
.field("sources", &self.sources.len())
.finish_non_exhaustive()
}
}
impl ArchivePublisher {
pub fn tailing(db: &Archive) -> Result<(Self, Vec<Frame>)> {
db.read_snapshot(|db| Self::open(db, None))
}
pub fn catching_up(db: &Archive, since: i64) -> Result<(Self, Vec<Frame>)> {
db.read_snapshot(|db| Self::open(db, Some(since)))
}
pub fn next(&mut self, db: &Archive) -> Result<Vec<Frame>> {
db.read_snapshot(|db| {
let mut out = Vec::new();
for cursor in &mut self.sources {
cursor.poll(db, &mut out)?;
}
Ok(out)
})
}
pub fn ordinals(&self) -> Vec<(u32, i64)> {
self.sources
.iter()
.map(|c| (c.ordinal, c.source_id))
.collect()
}
fn open(db: &Archive, since: Option<i64>) -> Result<(Self, Vec<Frame>)> {
let mut sources = Vec::new();
let mut frames = Vec::new();
let watermarks = db.sealed_watermarks()?;
for (ordinal, rec) in db.read_sources()?.into_iter().enumerate() {
let ordinal = u32::try_from(ordinal).map_err(|_| {
Error::Message("an archive with more than u32::MAX sources".to_string())
})?;
frames.push(Frame::Handshake {
source: ordinal,
uuid: rec.uuid.clone(),
labels: rec.meta.labels.clone(),
metadata: rec.meta.metadata.clone(),
clock_anchor_wall_ns: rec.meta.clock_anchor_wall_ns,
complete: rec.complete,
});
let mut cursor = SourceCursor {
ordinal,
source_id: rec.id,
streams: BTreeMap::new(),
clock_after: None,
seq: 0,
index_state: NO_INDEX_STATE,
sent_full: false,
};
for name in db.caller_row_streams(rec.id)? {
cursor.streams.entry(name.clone()).or_default();
cursor.emit_index(db, &name, since.unwrap_or(i64::MIN), &mut frames)?;
}
for stream in db.all_streams(rec.id)? {
let entry = cursor.streams.entry(stream.clone()).or_default();
let watermark = watermarks
.get(&rec.id)
.and_then(|m| m.get(&stream))
.copied();
match since {
Some(since) => {
for segment in db.segments_overlapping(rec.id, &stream, since, i64::MAX)? {
frames.push(Frame::Segment {
source: ordinal,
stream: stream.clone(),
meta: segment.meta,
bytes: segment.bytes,
caller_index: segment.caller_index,
});
}
entry.wal_after = watermark;
}
None => {
let live = db.live_wal_span(rec.id, &stream)?;
cursor
.streams
.get_mut(&stream)
.expect("just inserted")
.wal_after = live.last_ts.or(watermark);
}
}
}
for (ts, offset) in db.read_clock_offsets(rec.id)? {
match since {
Some(since) if ts < since => continue,
_ => {}
}
if since.is_none() {
cursor.clock_after = Some(ts);
continue;
}
frames.push(Frame::ClockOffset {
source: ordinal,
ts,
offset_ns: offset,
});
cursor.clock_after = Some(ts);
}
sources.push(cursor);
}
Ok((ArchivePublisher { sources }, frames))
}
}
impl SourceCursor {
fn emit_index(
&mut self,
db: &Archive,
stream: &str,
from: i64,
out: &mut Vec<Frame>,
) -> Result<()> {
let entry = self.streams.entry(stream.to_string()).or_default();
let (start, mut skip) = match entry.index_after {
Some((ts, count)) => (ts, count),
None => (from, 0),
};
let rows = db.read_caller_rows(self.source_id, stream, start, i64::MAX)?;
let mut at_last: usize = 0;
let mut last_ts: Option<i64> = entry.index_after.map(|(ts, _)| ts);
for row in rows {
if skip > 0 && Some(row.ts) == last_ts {
skip -= 1;
at_last += 1;
continue;
}
if Some(row.ts) == last_ts {
at_last += 1;
} else {
last_ts = Some(row.ts);
at_last = 1;
}
self.index_state = fold_index_state(self.index_state, row.ts, &row.blob);
let kind = if self.sent_full {
IndexKind::Delta
} else {
IndexKind::Full
};
self.sent_full = true;
out.push(Frame::Index {
source: self.ordinal,
stream: stream.to_string(),
ts: row.ts,
kind,
state: self.index_state,
blob: row.blob,
});
}
if let Some(ts) = last_ts {
self.streams
.get_mut(stream)
.expect("inserted above")
.index_after = Some((ts, at_last));
}
Ok(())
}
fn poll(&mut self, db: &Archive, out: &mut Vec<Frame>) -> Result<()> {
for name in db.caller_row_streams(self.source_id)? {
self.emit_index(db, &name, i64::MIN, out)?;
}
for (ts, offset) in db.read_clock_offsets(self.source_id)? {
if self.clock_after.is_some_and(|after| ts <= after) {
continue;
}
out.push(Frame::ClockOffset {
source: self.ordinal,
ts,
offset_ns: offset,
});
self.clock_after = Some(ts);
}
let watermarks = db.sealed_watermarks()?;
let mut rows = Vec::new();
for stream in db.all_streams(self.source_id)? {
let entry = self.streams.entry(stream.clone()).or_default();
let watermark = watermarks
.get(&self.source_id)
.and_then(|m| m.get(&stream))
.copied();
if let (Some(w), Some(after)) = (watermark, entry.wal_after) {
if w > after {
return Err(Error::Message(format!(
"source {}: stream `{stream}` sealed past this publisher's cursor \
(watermark {w}, last row shipped {after}); those rows are no longer \
in the live tail. Reconnect with `catching_up` from {after}",
self.source_id
)));
}
}
let after = entry.wal_after;
let mut newest = after;
for row in db.live_wal(self.source_id, &stream)? {
if after.is_some_and(|a| row.ts <= a) {
continue;
}
newest = Some(newest.map_or(row.ts, |n| n.max(row.ts)));
rows.push(row);
}
entry.wal_after = newest;
}
rows.sort_by_key(|r| (r.ts, r.stream.clone()));
out.push(Frame::Rows {
source: self.ordinal,
seq: self.seq,
index_state: self.index_state,
rows,
});
self.seq += 1;
Ok(())
}
}