knut-thund 0.1.7

Þ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
//! Streaming semantics — the part of the IR that Spark Declarative Pipelines
//! cannot fully express yet and that the native Arrow/DataFusion backend
//! implements first-class.
//!
//! These types are deliberately modelled on the *consensus* of the research
//! landscape (Arroyo, Flink, DataFusion's `Boundedness`×`EmissionType`,
//! RisingWave): an event-time stream is described by **where rows come from**
//! ([`SourceSpec`]), **how event-time progress is tracked**
//! ([`Watermark`] — min-across-inputs with allowed lateness, the knob Arroyo
//! notably lacks), **when results are emitted** ([`Trigger`]), and **what is
//! written** ([`OutputMode`]). [`WindowSpec`] captures the windowed
//! aggregation shape (tumbling / hopping / session).

use serde::{Deserialize, Serialize};

/// Where a source flow reads its rows from, and whether that input is bounded
/// (a finite batch) or unbounded (a never-ending stream).
///
/// Mirrors the DataFusion `Boundedness` distinction: a [`SourceSpec`] tagged
/// `unbounded` forces the planner down the streaming path (incremental
/// operators, state, checkpoints); a bounded one is a plain batch scan.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SourceSpec {
    /// A finite scan over an object-store path / table (batch).
    Batch {
        /// URI or catalog name (`s3://…/*.parquet`, `iceberg.db.tbl`).
        uri: String,
        /// Storage/format hint (`parquet`, `iceberg`, `delta`, `csv`).
        format: String,
    },
    /// An unbounded Kafka-style topic stream.
    Kafka {
        /// Bootstrap servers / connection URI.
        bootstrap: String,
        /// Topic name.
        topic: String,
        /// Value format (`json`, `avro`, `arrow`).
        format: String,
    },
    /// Change-data-capture from an upstream table (Debezium-style changelog).
    Cdc {
        /// Source table identifier.
        source: String,
    },
    /// A file-drop directory watched for new files (the knut-bifrost pattern).
    FileDrop {
        /// Directory URI to watch.
        dir: String,
        /// File format.
        format: String,
    },
}

/// Event-time watermark policy for a streaming source.
///
/// The watermark is `max(event_time) - allowed_lateness`, combined across all
/// inputs as the **minimum** (so the slowest input gates progress). Rows
/// arriving more than `allowed_lateness` behind the watermark are dropped —
/// but, unlike Arroyo, the lateness is an explicit knob here.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Watermark {
    /// The event-time column the watermark is derived from.
    pub event_time_column: String,
    /// Allowed lateness, in milliseconds, before a row is considered late.
    pub allowed_lateness_ms: u64,
    /// How long (ms) a partition may be idle before it stops gating the
    /// combined watermark (Arroyo's `idle_micros` analogue). `0` = never idle.
    #[serde(default)]
    pub idle_timeout_ms: u64,
}

/// When a streaming query emits results.
///
/// Mirrors Spark Structured Streaming's trigger taxonomy, which SDP streaming
/// tables build on.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Trigger {
    /// Emit as soon as a micro-batch's worth of data is buffered (default).
    Continuous,
    /// Fixed micro-batch cadence.
    ProcessingTime {
        /// Interval in milliseconds.
        interval_ms: u64,
    },
    /// Process all currently-available data, then stop (one bounded pass over
    /// an otherwise unbounded source — the SDP "once" / `availableNow` flow).
    AvailableNow,
}

impl Default for Trigger {
    fn default() -> Self {
        Trigger::Continuous
    }
}

/// What a streaming flow writes to its target on each trigger.
///
/// The classic Structured Streaming output modes. `Append` is the only mode
/// SDP streaming tables natively support; `Update`/`Complete` map to the
/// native backend's updating-table semantics (retract-then-insert).
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OutputMode {
    /// Only new rows since the last trigger (append-only / streaming table).
    Append,
    /// Rows whose aggregate value changed (retract-then-insert; updating MV).
    Update,
    /// The full result table every trigger (small aggregates only).
    Complete,
}

impl Default for OutputMode {
    fn default() -> Self {
        OutputMode::Append
    }
}

/// A windowed-aggregation specification for a streaming flow.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WindowSpec {
    /// Fixed, non-overlapping windows.
    Tumbling {
        /// Window size in milliseconds.
        size_ms: u64,
    },
    /// Overlapping fixed windows that advance by `slide_ms`.
    Hopping {
        /// Window size in milliseconds.
        size_ms: u64,
        /// Slide/advance in milliseconds.
        slide_ms: u64,
    },
    /// Activity-bounded windows that close after a gap of inactivity.
    Session {
        /// Inactivity gap that closes a session, in milliseconds.
        gap_ms: u64,
    },
}