knut-thund 0.1.8

Þund — a Rust-native, Arrow-centric streaming dataflow engine (batch + streaming) with a pluggable execution backend: native Arrow/DataFusion or lower-to-Spark-Declarative-Pipelines via Spark Connect. The 'Airflow killer' authoring+runtime for knut.
Documentation
//! The **barrier/epoch checkpoint state store** (design
//! `thund-streaming-state-design.md`, phases 3-5) — the resurrection of
//! [`super::native::NativeBackend::state_root`].
//!
//! One **epoch = one trigger**. At each trigger boundary a barrier is injected:
//! the stateful operator snapshots its state to Parquet, the source's consumed
//! offset and the combined watermark are gathered, and a JSON **manifest** is
//! written **last** — the manifest write is the atomic **commit point**. Only
//! after it lands is `_thund_latest.json` flipped to the new epoch. A crash
//! before the manifest leaves a partial `epoch=N/` dir that recovery ignores
//! (latest still points at `N-1`): write-ahead, commit-by-manifest.
//!
//! The substrate is the one knut already writes Parquet through — parquet-rs
//! `ArrowWriter` over `object_store` (`file://…`, `s3://…`), exactly graphar's
//! `ChunkStore` — so workers stay stateless (no RocksDB; the Arroyo/RisingWave
//! disaggregated-state-on-object-storage position). Everything here is version-
//! matched to datafusion's Arrow via its `datafusion::object_store` /
//! `datafusion::parquet` re-exports, so the state store shares one Arrow tree.
#![cfg(feature = "native")]

use crate::error::{Result, ThundError};
use datafusion::arrow::array::RecordBatch;
use datafusion::arrow::datatypes::Schema;
use datafusion::object_store::{
    Error as OsError, ObjectStore, ObjectStoreExt, PutPayload, parse_url, path::Path as StorePath,
};
use datafusion::parquet::arrow::ArrowWriter;
use datafusion::parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use datafusion::parquet::file::reader::ChunkReader;
use serde::{Deserialize, Serialize};
use std::sync::Arc;

fn cp_err(ctx: &str, e: impl std::fmt::Display) -> ThundError {
    ThundError::Backend(format!("checkpoint {ctx}: {e}"))
}

/// The status a committed manifest carries. Recovery only trusts a manifest
/// with exactly this value, so a half-written epoch can never be resumed from.
pub(crate) const COMMITTED: &str = "COMMITTED";

/// The operator-state Parquet file name inside an `epoch=N/` directory.
const OPERATOR_FILE: &str = "operator-state.parquet";

/// The per-epoch manifest — the **commit point**. Written last inside its
/// `epoch=N/` directory, then pointed at by `_thund_latest.json`.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct CheckpointManifest {
    /// The epoch number (== the 0-based trigger index this epoch checkpoints).
    pub epoch: u64,
    /// `COMMITTED` once durable. Recovery ignores anything else.
    pub status: String,
    /// The combined event-time watermark `W` at this epoch (`None` before the
    /// first advance).
    pub watermark: Option<i64>,
    /// Per-input max observed event time at this epoch — restores the
    /// [`super::native`] watermark tracker exactly on recovery.
    pub input_maxima: Vec<(String, i64)>,
    /// Micro-batches (triggers) consumed **through** this epoch — the source
    /// offset the run resumes from (skip the already-consumed prefix).
    pub offset: usize,
    /// Total rows dropped as late across the stream so far (carried so the
    /// resumed run's counters continue, not reset).
    pub late_rows_dropped: u64,
    /// Windows closed-and-emitted so far (same reason).
    pub windows_closed: usize,
    /// The `(start, close_end)` tags of every window already emitted — restored
    /// so a resumed run never re-emits a window that closed pre-crash.
    pub emitted_windows: Vec<(i64, i64)>,
    /// Operator-state files under `epoch=N/` (empty when the operator held no
    /// state at this epoch — e.g. everything evicted).
    pub operator_files: Vec<String>,
}

impl CheckpointManifest {
    /// The manifest of an empty/cold operator at epoch 0 offset 0 — the value a
    /// first run starts from (kept internal; a cold store returns `None`).
    pub(crate) fn is_committed(&self) -> bool {
        self.status == COMMITTED
    }
}

/// The pointer file naming the newest committed epoch (flipped last, after the
/// epoch's manifest is durable).
#[derive(Debug, Clone, Serialize, Deserialize)]
struct LatestPointer {
    epoch: u64,
}

/// Parquet-backed barrier/epoch state store over `object_store`, rooted at
/// `{state_root}/{pipeline}/{flow_target}/`.
pub(crate) struct ParquetStateStore {
    store: Arc<dyn ObjectStore>,
    /// `{prefix}/{pipeline}/{flow_target}` — the per-flow checkpoint root.
    base: StorePath,
}

impl ParquetStateStore {
    /// Open (or lazily create) the store under `state_root` for one flow.
    ///
    /// `state_root` is an object-store URI (`file:///…`, `s3://bucket/pfx`);
    /// the store shares datafusion's `object_store`, so it reaches every backend
    /// graphar's `ChunkStore` does. No IO happens here — directories are created
    /// implicitly on the first `put`.
    pub(crate) fn open(state_root: &str, pipeline: &str, flow_target: &str) -> Result<Self> {
        let url = url::Url::parse(state_root)
            .map_err(|e| cp_err(&format!("parse state_root `{state_root}`"), e))?;
        let (store, prefix) = parse_url(&url).map_err(|e| cp_err("open object store", e))?;
        let base = prefix.join(pipeline).join(flow_target);
        Ok(Self {
            store: Arc::from(store),
            base,
        })
    }

    /// The `epoch=NNNNNN/` directory path for `epoch`.
    fn epoch_dir(&self, epoch: u64) -> StorePath {
        self.base.clone().join(format!("epoch={epoch:06}"))
    }

    /// Commit one epoch: write the operator state Parquet (if any), then the
    /// manifest (the atomic commit point), then flip `_thund_latest.json`.
    ///
    /// The manifest's `operator_files` is filled here from whether `state`
    /// carried rows, so the caller need not track file names. Write order is
    /// **state → manifest → latest**, so a crash between any two steps leaves an
    /// epoch that recovery skips (latest still points at the prior epoch).
    pub(crate) async fn commit(
        &self,
        mut manifest: CheckpointManifest,
        state: &[RecordBatch],
    ) -> Result<()> {
        let dir = self.epoch_dir(manifest.epoch);

        // 1. Operator state → Parquet (skip when the operator holds nothing).
        let has_state = state.iter().any(|b| b.num_rows() > 0);
        manifest.operator_files = if has_state {
            let schema = state
                .iter()
                .find(|b| b.num_rows() > 0)
                .map(|b| b.schema())
                .unwrap_or_else(|| Arc::new(Schema::empty()));
            let bytes = batches_to_parquet(state, schema)?;
            let path = dir.clone().join(OPERATOR_FILE);
            self.store
                .put(&path, PutPayload::from(bytes))
                .await
                .map_err(|e| cp_err("put operator state", e))?;
            vec![OPERATOR_FILE.to_string()]
        } else {
            Vec::new()
        };

        // 2. Manifest = the atomic commit point (status COMMITTED), written LAST
        //    inside the epoch dir.
        manifest.status = COMMITTED.to_string();
        let mbytes = serde_json::to_vec(&manifest).map_err(|e| cp_err("encode manifest", e))?;
        self.store
            .put(
                &dir.clone().join("_thund_checkpoint.json"),
                PutPayload::from(mbytes),
            )
            .await
            .map_err(|e| cp_err("put manifest", e))?;

        // 3. Flip the latest pointer — only now is epoch N the recovery target.
        let lbytes = serde_json::to_vec(&LatestPointer {
            epoch: manifest.epoch,
        })
        .map_err(|e| cp_err("encode latest", e))?;
        self.store
            .put(
                &self.base.clone().join("_thund_latest.json"),
                PutPayload::from(lbytes),
            )
            .await
            .map_err(|e| cp_err("put latest", e))?;
        Ok(())
    }

    /// The newest **committed** manifest, or `None` for a cold store (first run)
    /// or one whose pointer/epoch is missing or not `COMMITTED`.
    pub(crate) async fn latest(&self) -> Result<Option<CheckpointManifest>> {
        let latest = match self
            .get_json::<LatestPointer>(&self.base.clone().join("_thund_latest.json"))
            .await?
        {
            Some(l) => l,
            None => return Ok(None),
        };
        let manifest_path = self.epoch_dir(latest.epoch).join("_thund_checkpoint.json");
        match self.get_json::<CheckpointManifest>(&manifest_path).await? {
            Some(m) if m.is_committed() => Ok(Some(m)),
            _ => Ok(None), // partial/half-written epoch → ignore, cold start
        }
    }

    /// Restore the operator state `RecordBatch`es a manifest points at (empty
    /// when the operator held nothing at that epoch).
    pub(crate) async fn load_operator_state(
        &self,
        manifest: &CheckpointManifest,
    ) -> Result<Vec<RecordBatch>> {
        let mut out = Vec::new();
        let dir = self.epoch_dir(manifest.epoch);
        for f in &manifest.operator_files {
            let data = self
                .store
                .get(&dir.clone().join(f.as_str()))
                .await
                .map_err(|e| cp_err("get operator state", e))?
                .bytes()
                .await
                .map_err(|e| cp_err("read operator state", e))?;
            out.extend(parquet_to_batches(data)?);
        }
        Ok(out)
    }

    /// Fetch + JSON-decode an object, mapping a genuine miss to `None` (any
    /// other store error is surfaced).
    async fn get_json<T: for<'de> Deserialize<'de>>(&self, path: &StorePath) -> Result<Option<T>> {
        match self.store.get(path).await {
            Ok(r) => {
                let data = r.bytes().await.map_err(|e| cp_err("read json", e))?;
                let v = serde_json::from_slice::<T>(&data).map_err(|e| cp_err("decode json", e))?;
                Ok(Some(v))
            }
            Err(OsError::NotFound { .. }) => Ok(None),
            Err(e) => Err(cp_err("get json", e)),
        }
    }
}

/// Encode `batches` (all sharing `schema`) to in-memory Parquet bytes.
fn batches_to_parquet(batches: &[RecordBatch], schema: Arc<Schema>) -> Result<Vec<u8>> {
    let mut buf: Vec<u8> = Vec::new();
    let mut w = ArrowWriter::try_new(&mut buf, schema, None)
        .map_err(|e| cp_err("open parquet writer", e))?;
    for b in batches {
        if b.num_rows() > 0 {
            w.write(b).map_err(|e| cp_err("write parquet", e))?;
        }
    }
    w.close().map_err(|e| cp_err("close parquet", e))?;
    Ok(buf)
}

/// Decode Parquet `data` (any `ChunkReader` — the object-store `Bytes` here)
/// back into `RecordBatch`es.
fn parquet_to_batches<R: ChunkReader + 'static>(data: R) -> Result<Vec<RecordBatch>> {
    let reader = ParquetRecordBatchReaderBuilder::try_new(data)
        .map_err(|e| cp_err("open parquet reader", e))?
        .build()
        .map_err(|e| cp_err("build parquet reader", e))?;
    let mut out = Vec::new();
    for b in reader {
        out.push(b.map_err(|e| cp_err("read parquet batch", e))?);
    }
    Ok(out)
}

#[cfg(test)]
mod tests {
    use super::*;
    use datafusion::arrow::array::Int64Array;
    use datafusion::arrow::datatypes::{DataType, Field};

    fn state_batch(ids: &[i64]) -> RecordBatch {
        let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]));
        RecordBatch::try_new(schema, vec![Arc::new(Int64Array::from(ids.to_vec()))]).unwrap()
    }

    fn rt() -> tokio::runtime::Runtime {
        tokio::runtime::Builder::new_current_thread()
            .enable_all()
            .build()
            .unwrap()
    }

    /// A cold store (nothing committed) reports no latest checkpoint.
    #[test]
    fn cold_store_has_no_latest() {
        let dir = tempfile::tempdir().unwrap();
        let root = format!("file://{}", dir.path().display());
        let store = ParquetStateStore::open(&root, "p", "t").unwrap();
        rt().block_on(async {
            assert!(
                store.latest().await.unwrap().is_none(),
                "cold store: no latest"
            );
        });
    }

    /// Commit an epoch, then read the latest manifest + restore its operator
    /// state Parquet — the round-trip the recovery path relies on.
    #[test]
    fn commit_then_latest_and_restore_roundtrips() {
        let dir = tempfile::tempdir().unwrap();
        let root = format!("file://{}", dir.path().display());
        let store = ParquetStateStore::open(&root, "pipe", "flow").unwrap();
        let manifest = CheckpointManifest {
            epoch: 7,
            status: String::new(), // commit() stamps COMMITTED
            watermark: Some(250),
            input_maxima: vec![("events".into(), 300)],
            offset: 8,
            late_rows_dropped: 3,
            windows_closed: 2,
            emitted_windows: vec![(0, 100), (100, 200)],
            operator_files: Vec::new(),
        };
        let state = vec![state_batch(&[1, 2, 3])];
        rt().block_on(async {
            store.commit(manifest.clone(), &state).await.unwrap();

            let got = store.latest().await.unwrap().expect("a committed latest");
            assert_eq!(got.epoch, 7);
            assert!(got.is_committed());
            assert_eq!(got.watermark, Some(250));
            assert_eq!(got.offset, 8);
            assert_eq!(got.late_rows_dropped, 3);
            assert_eq!(got.windows_closed, 2);
            assert_eq!(got.emitted_windows, vec![(0, 100), (100, 200)]);
            assert_eq!(
                got.operator_files,
                vec!["operator-state.parquet".to_string()]
            );

            let restored = store.load_operator_state(&got).await.unwrap();
            let rows: usize = restored.iter().map(|b| b.num_rows()).sum();
            assert_eq!(rows, 3, "operator state restored row-for-row");
        });
    }

    /// A later committed epoch supersedes an earlier one (latest points at the
    /// newest), and an epoch that stored no rows restores as empty state.
    #[test]
    fn latest_tracks_newest_and_empty_state_is_ok() {
        let dir = tempfile::tempdir().unwrap();
        let root = format!("file://{}", dir.path().display());
        let store = ParquetStateStore::open(&root, "p", "t").unwrap();
        let base = CheckpointManifest {
            epoch: 0,
            status: String::new(),
            watermark: None,
            input_maxima: Vec::new(),
            offset: 0,
            late_rows_dropped: 0,
            windows_closed: 0,
            emitted_windows: Vec::new(),
            operator_files: Vec::new(),
        };
        rt().block_on(async {
            store
                .commit(
                    CheckpointManifest {
                        epoch: 0,
                        offset: 1,
                        ..base.clone()
                    },
                    &[state_batch(&[1])],
                )
                .await
                .unwrap();
            // Epoch 1 evicted everything → no operator files.
            store
                .commit(
                    CheckpointManifest {
                        epoch: 1,
                        offset: 2,
                        ..base.clone()
                    },
                    &[],
                )
                .await
                .unwrap();

            let got = store.latest().await.unwrap().expect("latest");
            assert_eq!(got.epoch, 1, "latest is the newest committed epoch");
            assert!(
                got.operator_files.is_empty(),
                "an evicted epoch stores no state"
            );
            assert!(store.load_operator_state(&got).await.unwrap().is_empty());
        });
    }
}