use crate::buffer::{BufferMode, SegmentWriter};
use crate::shared_state::SharedState;
pub fn drain_thread_local(shared: &SharedState) {
crate::encoder::drain_to_collector(&shared.collector);
}
pub fn drain_encoded_batches(shared: &SharedState) -> Vec<Vec<u8>> {
crate::encoder::drain_to_collector(&shared.collector);
let mut out = Vec::new();
while let Some(batch) = shared.collector.next() {
out.push(batch.into_encoded_bytes());
}
out
}
pub fn drain_into<M: BufferMode>(
shared: &SharedState,
writer: &mut SegmentWriter<M>,
) -> std::io::Result<()> {
crate::encoder::drain_to_collector(&shared.collector);
while let Some(batch) = shared.collector.next() {
writer.write_encoded_batch(&batch)?;
}
Ok(())
}
pub fn write_event<M: BufferMode>(
writer: &mut SegmentWriter<M>,
event: &dyn crate::encoder::Encodable,
) -> std::io::Result<()> {
let bytes = crate::encoder::encode_single(event);
writer.write_encoded_batch(&crate::collector::Batch::new(bytes, 1))
}
#[cfg(feature = "pipeline")]
pub use pipeline_helpers::*;
#[cfg(feature = "pipeline")]
mod pipeline_helpers {
use crate::dump::{DumpId, DumpTrigger};
use crate::fs::Fs;
use crate::pipeline::SegmentProcessor;
use crate::worker::WorkerLoop;
use std::io::Write as _;
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
fn dev_null_sink() -> metrique::writer::BoxEntrySink {
metrique::writer::sink::DevNullSink::boxed()
}
fn segment_with_epoch(epoch_secs: u64) -> Vec<u8> {
use dial9_trace_format::encoder::Encoder;
let mut enc = Encoder::new_to(Vec::new()).unwrap();
enc.write_infallible(&crate::format::ClockSyncEvent {
timestamp_ns: 1,
realtime_ns: epoch_secs * 1_000_000_000,
});
enc.into_inner()
}
pub fn new_dump_id() -> DumpId {
DumpId::new()
}
pub fn disk_segment(
path: impl Into<std::path::PathBuf>,
index: u32,
) -> crate::sealed::SegmentRef {
crate::sealed::SegmentRef::Disk(crate::sealed::SealedSegment::new_for_test(
path.into(),
index,
))
}
pub fn new_dump_completion(
dump_id: DumpId,
triggered_at: std::time::SystemTime,
time_range: (std::time::SystemTime, std::time::SystemTime),
segments_processed: usize,
metadata: Vec<(String, String)>,
failed: bool,
) -> crate::dump::DumpCompletion {
crate::dump::DumpCompletion {
dump_id,
triggered_at,
time_range,
segments_processed,
metadata,
failed,
}
}
pub async fn run_pipeline_continuous(
segments: Vec<Vec<u8>>,
processors: Vec<Box<dyn SegmentProcessor>>,
poll_interval: Duration,
) -> std::io::Result<()> {
let fs = Fs::new_in_memory(64 * 1024 * 1024, 1024)?;
for (index, bytes) in segments.iter().enumerate() {
let mut handle = fs.create_segment(Path::new("x"))?;
handle.write_all(bytes)?;
fs.seal(handle, Path::new("x"), index as u32)?;
}
fs.mark_writer_done();
let stop = CancellationToken::new();
let mut worker =
WorkerLoop::new(fs, poll_interval, processors, stop, dev_null_sink(), None).await?;
worker.run().await;
Ok(())
}
pub struct TriggeredPipeline {
pub trigger: DumpTrigger,
fs: Arc<Fs>,
stop: CancellationToken,
join: JoinHandle<()>,
}
impl std::fmt::Debug for TriggeredPipeline {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TriggeredPipeline").finish_non_exhaustive()
}
}
impl TriggeredPipeline {
pub fn seal(&self, index: u32, epoch_secs: u64) {
let mut handle = self.fs.create_segment(Path::new("x")).unwrap();
handle.write_all(&segment_with_epoch(epoch_secs)).unwrap();
self.fs.seal(handle, Path::new("x"), index).unwrap();
}
pub async fn shutdown(self) {
self.stop.cancel();
let _ = self.join.await;
}
}
pub fn spawn_triggered_pipeline(
processors: Vec<Box<dyn SegmentProcessor>>,
) -> TriggeredPipeline {
let fs = Fs::new_in_memory(64 * 1024, 1024).unwrap();
let (trigger, rx) = crate::dump::channel();
let stop = CancellationToken::new();
let worker_fs = Arc::clone(&fs);
let worker_stop = stop.clone();
let join = tokio::spawn(async move {
let mut worker = WorkerLoop::new(
worker_fs,
Duration::from_millis(10),
processors,
worker_stop,
dev_null_sink(),
Some(rx),
)
.await
.expect("initialize worker");
worker.run().await;
});
TriggeredPipeline {
trigger,
fs,
stop,
join,
}
}
}