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
//! The **Spark-Connect-SDP lowering** execution backend (feature `spark`).
//!
//! This backend does *not* run a query engine of its own. It lowers the Þund
//! IR to a Spark Declarative Pipelines (SDP) dataflow graph
//! ([`crate::ir::Pipeline::to_sdp`]) and drives a real Spark Connect server
//! through the existing `knut-pipelines` client — i.e. it reuses, rather than
//! reimplements, knut's Spark Connect connector and SDP proto.
//!
//! # Why reuse `knut-pipelines`
//!
//! `knut-pipelines` (feature `grpc`) already vendors the Spark Connect
//! `pipelines.proto` and implements the four-command define-and-run dance over
//! `SparkConnectService.ExecutePlan`:
//!
//! 1. `CreateDataflowGraph` → server `dataflow_graph_id`
//! 2. `DefineOutput` per dataset (table / MV / view / sink)
//! 3. `DefineFlow` per runnable flow (the flow's `Relation`/SQL)
//! 4. `StartRun` → opens the streamed event lifecycle
//!
//! So `SparkBackend` is a thin adaptor: lower the IR
//! ([`crate::ir::Pipeline::to_sdp`]), then — under the `spark` feature — call
//! [`knut_pipelines::PipelinesClient`] to `define_graph` + `start_run`, folding
//! the streamed `PipelineRunEvent`s into [`RunEvent`]s on the [`SparkRun`]. The
//! `define_graph` + `start_run` path is fully wired here; it needs a reachable
//! Spark Connect server to complete (there is no way to fake the cluster), so
//! without one the run returns an honest transport error rather than a stub.
//!
//! # Version requirements (from the SDP research)
//!
//! SDP **GA'd in Apache Spark 4.1.0** (2025-12-16), *not* 4.0. The
//! `pipelines.proto` lives at
//! `sql/connect/common/src/main/protobuf/spark/connect/pipelines.proto` and is
//! absent from branch-4.0. **Spark 4.2+** adds external **sinks**
//! (`create_sink`) and **persistent views**; **expectations** are not yet
//! fully upstreamed into OSS 4.1 — hence [`SparkBackend::capabilities`]
//! reports `expectations: false` by default and the planner refuses to lower
//! a pipeline that relies on them onto Spark (the IR keeps them; the native
//! backend runs them). Always generate the tonic stubs from the proto at the
//! *exact target Spark tag* — `DefineOutput`/`DefineFlow` field names drift
//! between 4.1.0, 4.1.2 and 4.2 master.

use super::{Capabilities, ExecBackend, RunEvent, RunHandle};
use crate::error::{Result, ThundError};
use crate::ir::Pipeline;

/// The Spark-Connect-SDP backend. Holds the Spark Connect URL it lowers onto.
#[derive(Debug, Clone)]
pub struct SparkBackend {
    /// Spark Connect endpoint (`sc://host:15002`).
    pub connect_url: String,
    /// Whether the target server is Spark 4.2+ (enables sinks + persistent
    /// views in the reported capabilities).
    pub spark_4_2_plus: bool,
}

impl SparkBackend {
    /// A backend lowering onto the Spark Connect server at `connect_url`
    /// (assumed Spark 4.1.x by default).
    pub fn new(connect_url: impl Into<String>) -> Self {
        SparkBackend {
            connect_url: connect_url.into(),
            spark_4_2_plus: false,
        }
    }

    /// Declare the target server to be Spark 4.2+, returning `self`.
    pub fn spark_4_2(mut self) -> Self {
        self.spark_4_2_plus = true;
        self
    }
}

impl ExecBackend for SparkBackend {
    type Run = SparkRun;

    fn capabilities(&self) -> Capabilities {
        Capabilities {
            name: "spark-connect-sdp".into(),
            batch: true,
            streaming: true,
            // SDP has streaming tables but no first-class explicit watermark /
            // window knobs over Connect; those degrade to Spark defaults.
            event_time: false,
            // AutoCDC (`AutoCdcFlowDetails`) is on master / 4.2-era; gate it.
            cdc: self.spark_4_2_plus,
            // Not fully upstreamed in OSS 4.1.
            expectations: false,
            // SDP streaming tables are append-only.
            output_modes: vec!["append".into()],
        }
    }

    fn run(&self, pipeline: &Pipeline) -> Result<Self::Run> {
        self.check(pipeline)?;
        pipeline.validate()?;

        // The faithful part: lower the IR graph shape to SDP. This works on
        // the default build (no `spark`) too — proving the lowering without a
        // live Spark.
        let sdp = pipeline.to_sdp();
        let lowered = sdp.validate().map_err(|e| ThundError::Backend(e.to_string()));
        // HONEST STATUS: the IR→SDP lowering + validation genuinely works, so it
        // emits a green row.
        crate::functional_status(
            "knut-thund/backend_spark",
            "lower_to_sdp",
            lowered.is_ok(),
            &self.connect_url,
        );
        lowered?;

        tracing::info!(
            target: "knut_thund::spark",
            pipeline = %pipeline.name,
            connect_url = %self.connect_url,
            datasets = sdp.datasets.len(),
            flows = sdp.flows.len(),
            "spark backend: lowered IR to SDP dataflow graph"
        );

        #[cfg(feature = "spark")]
        {
            exec::run_on_spark(self, pipeline, sdp)
        }
        #[cfg(not(feature = "spark"))]
        {
            let _ = &sdp;
            // Without the `spark` feature the `knut-pipelines` gRPC client is
            // not compiled in. The lowering above is proven; the live run needs
            // `--features spark`. Honest capability gap (not a scaffold TODO).
            crate::functional_status(
                "knut-thund/backend_spark",
                "live_run",
                false,
                "spark feature disabled: Spark Connect client not compiled in",
            );
            Err(ThundError::Backend(
                "spark backend live run requires the `spark` feature (Spark Connect client not \
                 compiled in); IR lowered to SDP graph successfully"
                    .into(),
            ))
        }
    }
}

/// A handle to a Spark-Connect SDP run.
///
/// The `knut-pipelines` client streams the run's `PipelineRunEvent`s to
/// completion; each is folded into a [`RunEvent`] and queued here, drained
/// once by [`RunHandle::poll_events`].
#[derive(Debug, Default)]
pub struct SparkRun {
    /// Events folded from the Spark Connect run stream, in arrival order.
    events: std::collections::VecDeque<RunEvent>,
}

impl RunHandle for SparkRun {
    fn poll_events(&mut self) -> Result<Vec<RunEvent>> {
        Ok(self.events.drain(..).collect())
    }

    fn cancel(&mut self) -> Result<()> {
        Ok(())
    }
}

/// The live Spark Connect driver (feature `spark`).
#[cfg(feature = "spark")]
mod exec {
    use super::{SparkBackend, SparkRun};
    use crate::backend::{RunEvent, RunPhase};
    use crate::error::{Result, ThundError};
    use crate::ir::Pipeline;
    use knut_pipelines::{
        DataflowGraph, PipelineEventSink, PipelineRun, PipelineRunEvent, PipelinesClient,
    };
    use std::collections::VecDeque;

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

    /// A sink that folds each `knut-pipelines` run event into a Þund
    /// [`RunEvent`], mapping the SDP severity onto a coarse phase where the
    /// message makes it obvious.
    struct Collector<'a> {
        events: &'a mut VecDeque<RunEvent>,
    }

    impl PipelineEventSink for Collector<'_> {
        fn on_event(&mut self, ev: &PipelineRunEvent) -> knut_pipelines::Result<()> {
            self.events.push_back(RunEvent {
                timestamp: ev.timestamp.clone(),
                element: ev.element.clone(),
                message: ev.message.clone(),
                phase: phase_of(&ev.message),
            });
            Ok(())
        }
    }

    /// Best-effort phase inference from an SDP event message (the proto carries
    /// only timestamp + message; the viewer derives the rest the same way).
    fn phase_of(message: &str) -> Option<RunPhase> {
        let m = message.to_ascii_lowercase();
        if m.contains("fail") || m.contains("error") {
            Some(RunPhase::Failed)
        } else if m.contains("complet") {
            Some(RunPhase::Completed)
        } else if m.contains("running") {
            Some(RunPhase::Running)
        } else if m.contains("queued") {
            Some(RunPhase::Queued)
        } else {
            None
        }
    }

    /// Connect, register the graph, and drive the run to completion, collecting
    /// its events. Synchronous entry point (the [`crate::backend::ExecBackend`]
    /// trait is sync); spins a current-thread Tokio runtime for the gRPC I/O.
    pub(crate) fn run_on_spark(
        backend: &SparkBackend,
        pipeline: &Pipeline,
        sdp: DataflowGraph,
    ) -> Result<SparkRun> {
        let rt = tokio::runtime::Builder::new_current_thread()
            .enable_all()
            .build()
            .map_err(|e| be("build tokio runtime", e))?;

        let result = rt.block_on(drive(backend, pipeline, &sdp));

        match &result {
            Ok(run) => crate::functional_status(
                "knut-thund/backend_spark",
                "live_run",
                true,
                &format!(
                    "{}: {} event(s) from {}",
                    pipeline.name,
                    run.events.len(),
                    backend.connect_url
                ),
            ),
            Err(e) => crate::functional_status(
                "knut-thund/backend_spark",
                "live_run",
                false,
                &format!("{}: {e}", backend.connect_url),
            ),
        }
        result
    }

    async fn drive(
        backend: &SparkBackend,
        pipeline: &Pipeline,
        sdp: &DataflowGraph,
    ) -> Result<SparkRun> {
        let mut client = PipelinesClient::connect(&backend.connect_url)
            .await
            .map_err(|e| be("connect", e))?;

        let graph_id = client
            .define_graph(sdp)
            .await
            .map_err(|e| be("define_graph", e))?;

        // Spark requires a storage root to actually materialise outputs; with
        // none we do a `dry` validation run rather than provoke a server-side
        // failure — an honest degrade, not a fake success.
        let dry = pipeline.storage.is_none();
        let run_spec = PipelineRun {
            graph_id: Some(graph_id),
            full_refresh_all: true,
            refresh_selection: Vec::new(),
            storage: pipeline.storage.clone(),
            dry,
        };

        let mut events = VecDeque::new();
        if dry {
            events.push_back(RunEvent {
                timestamp: Some(chrono::Utc::now().to_rfc3339()),
                element: None,
                message: "no storage root set — running a dry (validate-only) SDP run".into(),
                phase: Some(RunPhase::Planning),
            });
        }
        {
            let mut sink = Collector { events: &mut events };
            client
                .start_run(&run_spec, &mut sink)
                .await
                .map_err(|e| be("start_run", e))?;
        }
        Ok(SparkRun { events })
    }
}