knut-thund 0.2.0

Þ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.
//! 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;

/// The barrier/epoch checkpoint state store (feature `native`) the incremental
/// streaming path snapshots + recovers through. Compiled only with the
/// DataFusion tree, since it shares datafusion's Arrow/Parquet/object_store.
#[cfg(feature = "native")]
pub(crate) mod checkpoint;

use crate::error::Result;
use crate::ir::{Dataset, OutputType, 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,
    /// Advertises **incremental streaming state** — an event-time watermark
    /// tracker + windowed/keyed operator state + barrier checkpoints under a
    /// `state_root` (the opt-in [`native::StreamingMode::Incremental`] path).
    /// A *capability advertised*, never a `check` gate: a pipeline never fails
    /// [`ExecBackend::check`] for lacking it, it just runs on recompute.
    #[serde(default)]
    pub incremental_state: bool,
    /// Supports **non-identity partition transforms** in
    /// [`crate::ir::Dataset::partition_cols`] — `months(event_ts)`,
    /// `bucket(16, id)`, `truncate(10, name)` and friends, as opposed to bare
    /// column names.
    ///
    /// This is a `check` gate, and deliberately so. Spark's SDP rejects them
    /// server-side (`PartitionHelper#applyPartitioning` accepts only
    /// `IdentityTransform`, and `CLUSTER BY` dies in the same place), so
    /// without this the failure surfaces as an opaque
    /// `PIPELINE_SQL_GRAPH_ELEMENT_REGISTRATION_ERROR` from the JVM *after* a
    /// graph has been created — or worse, is masked by the storage-less
    /// degrade-to-dry-run path. Failing in `check` keeps it local and legible.
    ///
    /// `#[serde(default)]` = `false`, the conservative reading for any
    /// capability record written before this field existed.
    #[serde(default)]
    pub partition_transforms: bool,
    /// Supports executing **graph-sink** datasets — a `falkordb-node` /
    /// `falkordb-edge` [`OutputType::Sink`] carrying the rel2graph
    /// `knut.node.*` / `knut.edge.*` properties, i.e. an idempotent Cypher
    /// `UNWIND … MERGE` graph load.
    ///
    /// This is a `check` gate, and deliberately so. Only the `NativeBackend`
    /// (built `native + rel2graph`) can execute a graph load — it drives the
    /// mutation through the existing `knut_bifrost::sink::FalkorDbSink`. The
    /// SDP `SparkBackend` structurally **cannot**: an SDP flow is a SQL
    /// `Relation` and an SDP output is a `TableDetails`, with nowhere to attach
    /// a Cypher closure / `foreachPartition` / FalkorDB client. Handing it a
    /// graph-sink pipeline would `DefineOutput` an unresolvable `falkordb-node`
    /// format and fail obscurely server-side. Failing in `check` keeps it local
    /// and legible — the operator is told to use `NativeBackend` or the PySpark
    /// projection instead.
    ///
    /// `#[serde(default)]` = `false`, the conservative reading for any
    /// capability record written before this field existed.
    #[serde(default)]
    pub graph_sinks: bool,
    /// Output modes the backend can write.
    pub output_modes: Vec<String>,
}

/// `true` when a `partition_cols` entry is a transform rather than a bare
/// column name.
///
/// Deliberately syntactic — a parenthesis is what separates `months(event_ts)`
/// from `event_month`. The IR stores these verbatim as strings (the
/// L2-additive choice that keeps existing RON specs valid), so this is the
/// classification available without a second typed model.
pub(crate) fn is_partition_transform(col: &str) -> bool {
    col.contains('(')
}

/// `true` when `ds` is a rel2graph **graph-sink** dataset — an
/// [`OutputType::Sink`] whose `format` is `falkordb-node` / `falkordb-edge`
/// AND which carries the matching rel2graph target property
/// (`knut.node.label` / `knut.edge.rel`).
///
/// Feature-independent by design: the backend-trait `check` gate compiles on
/// every build (the `rel2graph` module, and its `prop::*` constants, are behind
/// `#[cfg(feature = "rel2graph")]`), so the keys are matched as literals — the
/// same pair the native `graph_sink::graph_elem` classifier uses. Both the
/// format and the property are required so a hand-authored `Sink` with an
/// unrelated format is never mistaken for a graph load.
pub(crate) fn is_graph_sink(ds: &Dataset) -> bool {
    if ds.output_type != OutputType::Sink {
        return false;
    }
    match ds.format.as_deref() {
        Some("falkordb-node") => ds.properties.contains_key("knut.node.label"),
        Some("falkordb-edge") => ds.properties.contains_key("knut.edge.rel"),
        _ => false,
    }
}

/// Describe a graph-sink dataset's rel2graph target for a diagnostic — the node
/// label it MERGEs (`knut.node.label`) or the edge relationship it writes
/// (`knut.edge.rel`), so the `check` refusal names *which* graph element the
/// pipeline could not lower, not just the dataset name.
pub(crate) fn graph_sink_target(ds: &Dataset) -> String {
    if let Some(label) = ds.properties.get("knut.node.label") {
        format!("node label `{label}`")
    } else if let Some(rel) = ds.properties.get("knut.edge.rel") {
        format!("edge relationship `{rel}`")
    } else {
        "graph element".to_string()
    }
}

/// One flow's honest record of what the SDP lowering (`Pipeline::to_sdp`)
/// could **not** carry — the streaming/quality knobs that exist in the Þund IR
/// but have no home in a `knut_pipelines::Flow`/`DataflowGraph` and are silently
/// dropped when a pipeline is lowered onto Spark Declarative Pipelines.
///
/// This is the concrete surface behind the design's promise that "the streaming
/// knobs and expectations SDP cannot express are dropped here and surfaced by
/// the backend capability report" — made per-pipeline and machine-readable,
/// rather than only the static, whole-backend [`Capabilities`]. It is the
/// *honest* half of the Iceberg↔Spark streaming gap: the native backend RUNS
/// these knobs; the Spark path drops them and says so, right down to which flow
/// and which knob.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SdpFlowDrop {
    /// The flow whose lowering is lossy.
    pub flow: String,
    /// Human-readable descriptions of each knob SDP could not carry, in a
    /// stable order (watermark, window, trigger, output-mode, CDC, expectations).
    pub dropped: Vec<String>,
}

/// The per-pipeline record of everything [`crate::ir::Pipeline::to_sdp`] cannot
/// express for a given Spark target — the "capability-report the drop" surface.
///
/// A pipeline is **lossless** onto SDP iff this report is empty. When it is not,
/// each [`SdpFlowDrop`] names a flow and the knobs dropped on that flow, so a
/// caller (korp's Spark view, a CLI `--explain`, a test) can show the operator
/// *exactly* what running on Spark loses versus running on the native engine —
/// no silent degradation.
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct SdpLoweringReport {
    /// The backend name this report is computed for (`"spark-connect-sdp"`).
    pub backend: String,
    /// Whether the report assumed a Spark 4.2+ target (enables CDC / sinks /
    /// persistent views; a 4.1 target drops CDC entirely).
    pub spark_4_2_plus: bool,
    /// One entry per flow whose lowering is lossy; empty ⇒ lossless.
    pub flow_drops: Vec<SdpFlowDrop>,
}

impl SdpLoweringReport {
    /// `true` when the pipeline lowers to SDP with no dropped knobs.
    pub fn is_lossless(&self) -> bool {
        self.flow_drops.is_empty()
    }

    /// The lossy flows' drop record for `flow`, if any.
    pub fn flow(&self, flow: &str) -> Option<&SdpFlowDrop> {
        self.flow_drops.iter().find(|d| d.flow == flow)
    }

    /// Every dropped-knob description across all flows, flattened — for a quick
    /// `contains(...)` assertion or a one-line summary.
    pub fn all_dropped(&self) -> Vec<String> {
        self.flow_drops
            .iter()
            .flat_map(|d| d.dropped.iter().cloned())
            .collect()
    }
}

/// 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(),
            });
        }
        if !caps.partition_transforms {
            for d in &pipeline.datasets {
                if let Some(t) = d.partition_cols.iter().find(|c| is_partition_transform(c)) {
                    return Err(crate::error::ThundError::Unsupported {
                        backend: leak(caps.name.clone()),
                        what: format!(
                            "partition transform `{t}` on dataset `{}` \
                             (this backend takes identity partition columns only — \
                             materialize the value as a column and partition by that)",
                            d.name
                        ),
                    });
                }
            }
        }
        if !caps.graph_sinks {
            let graph_sinks: Vec<&Dataset> = pipeline
                .datasets
                .iter()
                .filter(|d| is_graph_sink(d))
                .collect();
            if let Some(d) = graph_sinks.first() {
                return Err(crate::error::ThundError::Unsupported {
                    backend: leak(caps.name.clone()),
                    what: format!(
                        "graph-sink pipelines — SparkBackend cannot execute the {n} graph-sink \
                         dataset(s) in this pipeline (a rel2graph Cypher graph load). \
                         First offender: dataset `{name}` (format `{fmt}`, {target}). \
                         WHY: SDP over Spark Connect has no shape to carry a graph write — an SDP \
                         flow lowers to a SQL relation and an SDP output to a table, neither of \
                         which can hold the idempotent Cypher `UNWIND … MERGE` a graph load is; \
                         `DefineOutput`-ing the unresolvable `{fmt}` format would only fail \
                         obscurely server-side after a graph was already created. \
                         WHAT TO DO: run this pipeline on the NativeBackend (build \
                         `--features native,rel2graph`), which executes the load through \
                         `knut_bifrost::sink::FalkorDbSink`; OR render the PySpark projection \
                         (`GraphPlan::to_pyspark`) and run that graph load as a Spark job.",
                        n = graph_sinks.len(),
                        name = d.name,
                        fmt = d.format.as_deref().unwrap_or(""),
                        target = graph_sink_target(d),
                    ),
                });
            }
        }
        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())
}