#![allow(clippy::print_stdout, clippy::print_stderr)]
use serde::Deserialize;
use spate::json::{JsonDeserializerBuilder, NdjsonFramer};
use spate::prelude::*;
use spate::s3::S3Source;
use spate_test::{TestEncoder, capture_sink};
use std::io::Write as _;
use std::time::Duration;
#[derive(Debug, Deserialize)]
struct Reading {
sensor: String,
value: f64,
}
fn config_yaml(root: &std::path::Path) -> String {
format!(
r#"
pipeline: {{ name: s3-backfill-demo, threads: 1 }}
checkpoint: {{ interval: 200ms }}
metrics: {{ exporter: none }}
source:
s3:
url: "file://{data}/"
deserializer:
json: {{}}
sink: {{ capture: {{}} }}
"#,
data = root.join("data").display(),
)
}
fn run_once(yaml: &str) -> Result<Vec<String>, Box<dyn std::error::Error>> {
let pipeline = Pipeline::from_config(PipelineConfig::from_str(yaml)?)?;
let source = S3Source::from_component_config(&pipeline.config().source, pipeline.io_handle())?
.with_framer(|| Box::new(NdjsonFramer::new(64 << 20)));
let deser_section = pipeline
.config()
.deserializer
.as_ref()
.ok_or("this pipeline requires a `deserializer` section")?
.clone();
let (sink, script) = capture_sink(1, 1);
let sink = sink.with_pool_config({
let mut cfg = SinkPoolConfig::default();
cfg.batch.linger = Duration::from_millis(50);
cfg
});
let report = pipeline
.sink(sink)?
.chains(move |ctx| {
let chunk_cfg = ctx.chunk();
let deser = JsonDeserializerBuilder::from_component(&deser_section)
.and_then(|b| b.for_source_framing(ctx.source_framing))
.expect("deserializer config")
.with_metrics(ctx.pipeline.clone(), "main")
.build_serde::<Reading>();
chain_owned::<Reading, _>(deser)
.with_metrics(ctx.pipeline, "main")
.filter(|r: &Reading| r.value.is_finite())
.map(|r: Reading| format!("{}={}", r.sensor, r.value).into_bytes())
.sink(
TestEncoder,
KeyHashRouter,
chunk_cfg,
ctx.queues,
ctx.budget,
)
.build()
})
.runtime_options(RuntimeOptions {
handle_signals: false,
..RuntimeOptions::default()
})
.run(source)?;
println!("pipeline exit: {:?}", report.state);
assert_eq!(report.state, ExitState::Completed);
Ok(script
.writes()
.iter()
.flat_map(|w| spate_test::decode_rows(&w.payload))
.map(|r| String::from_utf8_lossy(&r).into_owned())
.collect())
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
spate::telemetry::init(spate::telemetry::LogFormat::Pretty, "info");
let root = tempfile::tempdir()?;
std::fs::create_dir_all(root.path().join("data"))?;
std::fs::write(
root.path().join("data/2026-07-13.ndjson"),
concat!(
"{\"sensor\":\"kitchen\",\"value\":21.5}\n",
"{\"sensor\":\"attic\",\"value\":31.0}\n",
),
)?;
let mut gz = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
gz.write_all(
concat!(
"{\"sensor\":\"cellar\",\"value\":12.25}\n",
"{\"sensor\":\"hall\",\"value\":19.75}\n",
)
.as_bytes(),
)?;
std::fs::write(root.path().join("data/2026-07-14.ndjson.gz"), gz.finish()?)?;
let yaml = config_yaml(root.path());
println!("── run 1: fresh backfill ──");
let mut rows = run_once(&yaml)?;
rows.sort();
println!("rows written ({}): {rows:?}", rows.len());
assert_eq!(rows.len(), 4);
println!("\n── run 2: rerun on an ephemeral (solo) store ──");
let mut rows = run_once(&yaml)?;
rows.sort();
println!(
"rows written ({}): solo progress died with run 1, so the rerun replayed \
the prefix — inject a coordinator over a durable backend for real resume",
rows.len()
);
assert_eq!(rows.len(), 4);
Ok(())
}