knut-thund 0.1.1

Þ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
//! Pluggable execution backends.
//!
//! The [`ExecBackend`] trait is the seam that makes Þund's central promise
//! work: a pipeline authored once (as a [`crate::ir::Pipeline`]) can run on
//! *either* the native Arrow/DataFusion engine — Arroyo-style, RecordBatch
//! streams the whole way — *or* be lowered to Spark Declarative Pipelines and
//! executed on a cluster over Spark Connect, with no change to the IR.
//!
//! # The contract
//!
//! A backend is, at minimum, a function from an IR pipeline to a *run* that
//! emits a stream of [`RunEvent`]s. It also reports its [`Capabilities`] so
//! the planner can decide — before running — whether the IR is expressible on
//! that backend (e.g. SDP cannot express arbitrary output modes or, in OSS
//! Spark 4.1, expectations), rather than failing mid-run. This mirrors
//! DataFusion's `SanityCheckPlan` rule: reject unrunnable plans up front.
//!
//! Two backends ship:
//! - [`native`] — the native Arrow/DataFusion engine (feature `native`).
//! - [`spark`] — the Spark-Connect-SDP lowering backend (feature `spark`),
//!   reusing the `knut-pipelines` client.

pub mod native;
pub mod spark;

use crate::error::Result;
use crate::ir::Pipeline;
use serde::{Deserialize, Serialize};

/// What a backend can express/execute. Reported by [`ExecBackend::capabilities`]
/// so the planner can validate an IR against a backend before running.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Capabilities {
    /// Stable backend identifier (`"native"`, `"spark-connect-sdp"`).
    pub name: String,
    /// Can run bounded/batch flows.
    pub batch: bool,
    /// Can run unbounded/streaming flows.
    pub streaming: bool,
    /// Supports explicit event-time watermarks + windowing.
    pub event_time: bool,
    /// Supports CDC / apply-changes flows.
    pub cdc: bool,
    /// Supports data-quality expectations.
    pub expectations: bool,
    /// Output modes the backend can write.
    pub output_modes: Vec<String>,
}

/// One observable event from a running pipeline — the atom of the live
/// stream. Deliberately shaped to fold straight onto the `knut-pipelines`
/// [`knut_pipelines::PipelineRunEvent`] so the existing facett viewer renders
/// a native-backend run and an SDP run identically.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RunEvent {
    /// RFC3339 timestamp, if known.
    pub timestamp: Option<String>,
    /// The dataset/flow this event concerns, if attributable.
    pub element: Option<String>,
    /// Human-readable message.
    pub message: String,
    /// Lifecycle phase, if this event marks one.
    pub phase: Option<RunPhase>,
}

/// Coarse lifecycle phase of a flow/run (mirrors the SDP event lifecycle
/// QUEUED → PLANNING → RUNNING → COMPLETED/FAILED).
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RunPhase {
    /// Queued, not yet planned.
    Queued,
    /// Planning the execution.
    Planning,
    /// Actively running.
    Running,
    /// Finished successfully.
    Completed,
    /// Finished with an error.
    Failed,
}

/// A handle to a started run. For batch this completes; for streaming it runs
/// until cancelled. Backends return their own concrete type; the trait keeps
/// it minimal for the scaffold.
pub trait RunHandle {
    /// Drain the next batch of events, blocking until at least one is
    /// available or the run ends. Returns an empty vec when the run is over.
    fn poll_events(&mut self) -> Result<Vec<RunEvent>>;

    /// Request cancellation of a streaming run.
    fn cancel(&mut self) -> Result<()>;
}

/// A pluggable execution backend: lowers and runs a [`Pipeline`].
pub trait ExecBackend {
    /// The concrete run handle this backend returns.
    type Run: RunHandle;

    /// What this backend can express/run.
    fn capabilities(&self) -> Capabilities;

    /// Validate that `pipeline` is expressible on this backend, given its
    /// [`Capabilities`]. The default implementation checks the streaming /
    /// CDC / expectation / output-mode caps against the IR. Backends rarely
    /// need to override this.
    fn check(&self, pipeline: &Pipeline) -> Result<()> {
        let caps = self.capabilities();
        if pipeline.is_streaming() && !caps.streaming {
            return Err(crate::error::ThundError::Unsupported {
                backend: leak(caps.name.clone()),
                what: "streaming flows".into(),
            });
        }
        for f in &pipeline.flows {
            if matches!(f.kind, crate::ir::FlowKind::Cdc { .. }) && !caps.cdc {
                return Err(crate::error::ThundError::Unsupported {
                    backend: leak(caps.name.clone()),
                    what: format!("CDC flow `{}`", f.name),
                });
            }
            if !f.expectations.is_empty() && !caps.expectations {
                return Err(crate::error::ThundError::Unsupported {
                    backend: leak(caps.name.clone()),
                    what: format!("expectations on flow `{}`", f.name),
                });
            }
        }
        Ok(())
    }

    /// Lower + start the pipeline, returning a run handle that streams events.
    fn run(&self, pipeline: &Pipeline) -> Result<Self::Run>;
}

/// Helper: leak a backend name into a `'static` str for the error type.
/// Backend names come from a fixed small set, so the leak is bounded.
fn leak(s: String) -> &'static str {
    Box::leak(s.into_boxed_str())
}