use std::collections::BTreeMap;
use std::path::PathBuf;
use std::sync::Arc;
use crate::archive::{Archive, Span};
use crate::error::Result;
use crate::segment::SegmentEncoder;
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct StreamCatalog {
pub name: String,
pub segments: u64,
pub sealed: Span,
pub live: Span,
pub bytes: u64,
}
impl StreamCatalog {
pub fn rows(&self) -> u64 {
self.sealed.rows + self.live.rows
}
pub fn span(&self) -> Option<(i64, i64)> {
let first = [self.sealed.first_ts, self.live.first_ts]
.into_iter()
.flatten()
.min()?;
let last = [self.sealed.last_ts, self.live.last_ts]
.into_iter()
.flatten()
.max()?;
Some((first, last))
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct SourceCatalog {
pub id: i64,
pub uuid: Option<String>,
pub labels: BTreeMap<String, String>,
pub metadata: BTreeMap<String, String>,
pub complete: bool,
pub clock_anchor_wall_ns: i64,
pub streams: Vec<StreamCatalog>,
}
impl SourceCatalog {
pub fn span(&self) -> Option<(i64, i64)> {
let spans: Vec<(i64, i64)> = self.streams.iter().filter_map(|s| s.span()).collect();
Some((
spans.iter().map(|s| s.0).min()?,
spans.iter().map(|s| s.1).max()?,
))
}
}
pub fn catalog(db: &Archive) -> Result<Vec<SourceCatalog>> {
db.read_snapshot(catalog_snapshotted)
}
fn catalog_snapshotted(db: &Archive) -> Result<Vec<SourceCatalog>> {
let mut out = Vec::new();
for src in db.read_sources()? {
let mut streams = Vec::new();
for name in db.all_streams(src.id)? {
let (segments, sealed) = db.segment_span(src.id, &name)?;
let live = db.live_wal_span(src.id, &name)?;
let bytes = db.stream_bytes(src.id, &name)?;
streams.push(StreamCatalog {
name,
segments,
sealed,
live,
bytes,
});
}
out.push(SourceCatalog {
id: src.id,
uuid: src.uuid,
labels: src.meta.labels,
metadata: src.meta.metadata,
complete: src.complete,
clock_anchor_wall_ns: src.meta.clock_anchor_wall_ns,
streams,
});
}
Ok(out)
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct Overview {
pub bytes: u64,
pub pages: crate::archive::PageStats,
pub sources: Vec<SourceCatalog>,
}
impl Overview {
pub fn segments(&self) -> u64 {
self.sources
.iter()
.flat_map(|s| &s.streams)
.map(|s| s.segments)
.sum()
}
pub fn rows(&self) -> u64 {
self.sources
.iter()
.flat_map(|s| &s.streams)
.map(|s| s.rows())
.sum()
}
pub fn span(&self) -> Option<(i64, i64)> {
let spans: Vec<(i64, i64)> = self.sources.iter().filter_map(|s| s.span()).collect();
Some((
spans.iter().map(|s| s.0).min()?,
spans.iter().map(|s| s.1).max()?,
))
}
pub fn free_bytes(&self) -> u64 {
self.pages.free as u64 * self.pages.page_size as u64
}
}
pub fn describe(db: &Archive) -> Result<Overview> {
db.read_snapshot(|db| {
Ok(Overview {
bytes: db.archive_bytes()?,
pages: db.page_stats()?,
sources: catalog_snapshotted(db)?,
})
})
}
pub fn probe(
db: &Archive,
source_id: i64,
stream: &str,
encoder: &dyn SegmentEncoder,
) -> Result<Option<Vec<u8>>> {
db.read_snapshot(|db| {
check_encoder_of(db, source_id, encoder)?;
if let Some((seq, _)) = db.read_segment_meta(source_id, stream)?.first() {
return db.read_segment_bytes(source_id, stream, *seq);
}
let live = db.live_wal(source_id, stream)?;
Ok(crate::segment::materialize(encoder, stream, &live)?.map(|t| t.bytes))
})
}
pub fn stream_range(
db: &Archive,
source_id: i64,
stream: &str,
start: i64,
end: i64,
encoder: &dyn SegmentEncoder,
) -> Result<Vec<Vec<u8>>> {
db.read_snapshot(|db| {
check_encoder_of(db, source_id, encoder)?;
let mut segments: Vec<Vec<u8>> = db
.segments_overlapping(source_id, stream, start, end)?
.into_iter()
.map(|s| s.bytes)
.collect();
let live: Vec<_> = db
.live_wal(source_id, stream)?
.into_iter()
.filter(|r| r.ts >= start && r.ts <= end)
.collect();
if let Some(tail) = crate::segment::materialize(encoder, stream, &live)? {
segments.push(tail.bytes);
}
Ok(segments)
})
}
#[derive(Debug)]
#[non_exhaustive]
pub struct SourceSegments {
pub labels: BTreeMap<String, String>,
pub metadata: BTreeMap<String, String>,
pub complete: bool,
pub streams: Vec<(String, Vec<Vec<u8>>)>,
}
pub fn read_archive(db: &Archive, encoder: &dyn SegmentEncoder) -> Result<Vec<SourceSegments>> {
db.read_snapshot(|db| read_archive_snapshotted(db, encoder))
}
fn read_archive_snapshotted(
db: &Archive,
encoder: &dyn SegmentEncoder,
) -> Result<Vec<SourceSegments>> {
let mut out = Vec::new();
for src in db.read_sources()? {
crate::segment::check_encoder(src.id, &src.meta.metadata, encoder)?;
let mut streams = Vec::new();
for stream in db.all_streams(src.id)? {
let segments = stream_segments_snapshotted(db, src.id, &stream, encoder)?;
if segments.is_empty() {
continue;
}
streams.push((stream, segments));
}
out.push(SourceSegments {
labels: src.meta.labels,
metadata: src.meta.metadata,
complete: src.complete,
streams,
});
}
Ok(out)
}
pub fn stream_segments(
db: &Archive,
source_id: i64,
stream: &str,
encoder: &dyn SegmentEncoder,
) -> Result<Vec<Vec<u8>>> {
db.read_snapshot(|db| {
check_encoder_of(db, source_id, encoder)?;
stream_segments_snapshotted(db, source_id, stream, encoder)
})
}
fn check_encoder_of(db: &Archive, source_id: i64, encoder: &dyn SegmentEncoder) -> Result<()> {
if encoder.version().is_none() {
return Ok(());
}
crate::segment::check_encoder(source_id, &db.source_metadata(source_id)?, encoder)
}
fn stream_segments_snapshotted(
db: &Archive,
source_id: i64,
stream: &str,
encoder: &dyn SegmentEncoder,
) -> Result<Vec<Vec<u8>>> {
let mut segments: Vec<Vec<u8>> = db
.read_segments(source_id, stream)?
.into_iter()
.map(|s| s.bytes)
.collect();
let live = db.live_wal(source_id, stream)?;
if let Some(tail) = crate::segment::materialize(encoder, stream, &live)? {
segments.push(tail.bytes);
}
Ok(segments)
}
pub fn stream_indexes(
db: &Archive,
source_id: i64,
stream: &str,
encoder: &dyn SegmentEncoder,
) -> Result<Vec<Option<Vec<u8>>>> {
db.read_snapshot(|db| {
check_encoder_of(db, source_id, encoder)?;
let mut out: Vec<Option<Vec<u8>>> = db
.read_segment_indexes(source_id, stream)?
.into_iter()
.map(|(_, index)| index)
.collect();
let live = db.live_wal(source_id, stream)?;
if let Some(tail) = crate::segment::materialize(encoder, stream, &live)? {
out.push(tail.index);
}
Ok(out)
})
}
pub enum SegmentBytes {
Bytes(Vec<Vec<u8>>),
AtPath {
path: PathBuf,
source_id: i64,
stream: String,
},
Shared {
db: Arc<std::sync::Mutex<Archive>>,
source_id: i64,
stream: String,
},
}
impl SegmentBytes {
pub fn at_path(path: PathBuf, source_id: i64, stream: String) -> Self {
SegmentBytes::AtPath {
path,
source_id,
stream,
}
}
pub fn shared(db: Arc<std::sync::Mutex<Archive>>, source_id: i64, stream: String) -> Self {
SegmentBytes::Shared {
db,
source_id,
stream,
}
}
pub fn all(&self, encoder: &dyn SegmentEncoder) -> Result<Vec<Vec<u8>>> {
match self {
SegmentBytes::Bytes(b) => Ok(b.clone()),
SegmentBytes::AtPath {
path,
source_id,
stream,
} => {
let db = Archive::open(path)?;
stream_segments(&db, *source_id, stream, encoder)
}
SegmentBytes::Shared {
db,
source_id,
stream,
} => {
let db = db.lock().unwrap_or_else(|e| e.into_inner());
stream_segments(&db, *source_id, stream, encoder)
}
}
}
}