Skip to main content

faucet_core/
native.rs

1//! Native byte-passthrough fast-load (#633) — the third capability-negotiated
2//! transfer mechanism, alongside the Arrow **columnar** path (RFC 0002 / #375)
3//! and **staged bulk load** (#528).
4//!
5//! When a source can emit its records as raw bytes in a wire format
6//! ([`NativeFormat`]) and a sink can bulk-load that exact format directly, and
7//! no `Value`-shaped stage sits between them, the pipeline streams the bytes
8//! source → sink **without ever materializing [`serde_json::Value`] or an Arrow
9//! `RecordBatch`** — the minimum-memory, minimum-CPU path.
10//!
11//! This is the byte analogue of the columnar path: columnar is for typed
12//! columnar sources (parquet/delta) → typed sinks; byte-passthrough is for
13//! format-matched wire-byte pairs (a bulk-export CSV → BigQuery CSV load;
14//! an S3 `.jsonl` object → BigQuery NDJSON load) with no object-store hop, where
15//! the *destination* does the CSV/JSON → typed-column casting.
16//!
17//! Like the other two mechanisms it is **opt-in and additive**: a source
18//! advertises formats via [`Source::native_output_formats`](crate::Source::native_output_formats)
19//! and a sink advertises mechanisms via
20//! [`Sink::native_load_capabilities`](crate::Sink::native_load_capabilities). The
21//! pipeline uses the path only when [`plan_native_transfer`] finds a sink
22//! capability whose format the source offers and whose write modes cover the
23//! run — after the planner's own unconditional gates (no transforms, no
24//! quality/contract/masking pass, no DLQ, at-least-once only) have passed.
25//! The "no transformation between" gate is additionally enforced structurally —
26//! a [`TransformingSource`](crate::TransformingSource) does not override
27//! `native_output_formats`, so any attached transform makes the wrapped source
28//! advertise no formats and the fast path falls through to the `Value` path.
29//!
30//! The core types carry no heavy dependencies (byte payloads are `Vec<u8>` /
31//! boxed byte streams, matching [`crate::staging`]), so the negotiation is
32//! always compiled — it is a first-class part of the `Source`/`Sink` contract,
33//! not a Cargo feature. Individual connector implementations may still gate their
34//! `stream_native` / `load_native` bodies behind their own features.
35
36use crate::error::FaucetError;
37use crate::idempotency::DeliveryMode;
38use crate::write_mode::WriteMode;
39use futures_core::Stream;
40use serde::{Deserialize, Serialize};
41use serde_json::Value;
42use std::pin::Pin;
43
44/// A wire format that can move source → sink as raw bytes, bypassing per-record
45/// `serde_json::Value` (and Arrow) materialization. This is the negotiation key:
46/// a source advertises the formats it can emit and a sink the formats it can
47/// load, and [`plan_native_transfer`] matches them.
48#[non_exhaustive]
49#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash, Serialize, Deserialize)]
50#[serde(rename_all = "snake_case")]
51pub enum NativeFormat {
52    /// Newline-delimited JSON — one JSON object per line.
53    NdJson,
54    /// CSV / TSV with the dialect carried on the [`NativeBatch`].
55    Csv,
56    /// Apache Parquet file bytes.
57    Parquet,
58    /// Apache Arrow IPC stream bytes.
59    ArrowIpc,
60}
61
62impl NativeFormat {
63    /// A stable lowercase token for logs, metrics labels, and config.
64    pub fn as_str(self) -> &'static str {
65        match self {
66            NativeFormat::NdJson => "ndjson",
67            NativeFormat::Csv => "csv",
68            NativeFormat::Parquet => "parquet",
69            NativeFormat::ArrowIpc => "arrow_ipc",
70        }
71    }
72}
73
74/// CSV specifics carried on a [`NativeBatch`] (ignored for other formats). The
75/// sink must be able to honor the source's dialect, or refuse the match.
76#[derive(Clone, Copy, Debug, PartialEq, Eq)]
77pub struct CsvDialect {
78    /// Whether the first row is a header row the sink should skip on load.
79    pub has_header: bool,
80    /// The field delimiter byte (e.g. `b','` or `b'\t'`).
81    pub delimiter: u8,
82}
83
84impl Default for CsvDialect {
85    fn default() -> Self {
86        Self {
87            has_header: true,
88            delimiter: b',',
89        }
90    }
91}
92
93/// The bytes of one native batch — either fully buffered or a byte stream.
94///
95/// The [`Stream`](NativePayload::Stream) variant is what gives true O(1) memory:
96/// the sink consumes it (e.g. into a resumable upload) without ever buffering the
97/// whole batch. The [`Bytes`](NativePayload::Bytes) variant is the simple case,
98/// bounded by the source's page / batch size.
99pub enum NativePayload {
100    /// A fully-buffered byte batch.
101    Bytes(Vec<u8>),
102    /// A streamed byte batch consumed chunk-by-chunk.
103    Stream(Pin<Box<dyn Stream<Item = Result<Vec<u8>, FaucetError>> + Send>>),
104}
105
106impl std::fmt::Debug for NativePayload {
107    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
108        match self {
109            NativePayload::Bytes(b) => f.debug_tuple("Bytes").field(&b.len()).finish(),
110            NativePayload::Stream(_) => f.write_str("Stream(<byte stream>)"),
111        }
112    }
113}
114
115/// One native-format byte batch handed from a source's
116/// [`stream_native`](crate::Source::stream_native) to a sink's
117/// [`load_native`](crate::Sink::load_native).
118///
119/// `bookmark` carries the same checkpoint semantics as
120/// [`StreamPage`](crate::StreamPage): whenever it is `Some`, the pipeline flushes
121/// the sink and persists the bookmark before polling the next batch.
122#[derive(Debug)]
123pub struct NativeBatch {
124    /// The wire format of `payload`.
125    pub format: NativeFormat,
126    /// The batch bytes.
127    pub payload: NativePayload,
128    /// CSV dialect (only meaningful when `format == Csv`).
129    pub csv: CsvDialect,
130    /// Row count if the source knows it (used for metrics; `None` if unknown).
131    pub records: Option<u64>,
132    /// Checkpoint bookmark, or `None` for a mid-stream batch.
133    pub bookmark: Option<Value>,
134}
135
136impl NativeBatch {
137    /// A fully-buffered batch with no bookmark and no known count.
138    pub fn bytes(format: NativeFormat, payload: Vec<u8>) -> Self {
139        Self {
140            format,
141            payload: NativePayload::Bytes(payload),
142            csv: CsvDialect::default(),
143            records: None,
144            bookmark: None,
145        }
146    }
147
148    /// Set the checkpoint bookmark (builder-style).
149    pub fn with_bookmark(mut self, bookmark: Option<Value>) -> Self {
150        self.bookmark = bookmark;
151        self
152    }
153
154    /// Set the known row count (builder-style).
155    pub fn with_records(mut self, records: Option<u64>) -> Self {
156        self.records = records;
157        self
158    }
159
160    /// Set the CSV dialect (builder-style).
161    pub fn with_csv(mut self, csv: CsvDialect) -> Self {
162        self.csv = csv;
163        self
164    }
165}
166
167/// An efficient native-load mechanism a sink offers. A sink returns one of
168/// these per mechanism it supports from
169/// [`Sink::native_load_capabilities`](crate::Sink::native_load_capabilities).
170///
171/// A capability declares only what genuinely varies per mechanism: the format
172/// and the write modes it can honor. Everything a byte passthrough can never
173/// support — transforms, the quality/contract/masking governance passes,
174/// per-row DLQ routing, exactly-once delivery — is gated **unconditionally by
175/// the pipeline** ([`plan_native_transfer`] rejects those runs outright). No
176/// capability can waive those gates: they need `Value` access or a commit-token
177/// protocol the byte path structurally does not have, so a declarative opt-out
178/// could only ship silently-unenforced governance.
179#[derive(Clone, Debug)]
180pub struct NativeLoadCapability {
181    /// The wire format this mechanism consumes.
182    pub format: NativeFormat,
183    /// A short, stable label for logs / metrics (e.g. `"bigquery-load-job"`).
184    pub mechanism: &'static str,
185    /// Write modes this mechanism can honor (typically `Append` and/or
186    /// `Overwrite`; `Upsert`/`Delete` need per-row keys and so are never
187    /// passthrough-eligible — the pipeline rejects them before matching).
188    pub write_modes: &'static [WriteMode],
189}
190
191/// Per-call context the pipeline hands each
192/// [`load_native`](crate::Sink::load_native) invocation.
193#[derive(Clone, Copy, Debug, PartialEq, Eq)]
194pub struct NativeLoadContext {
195    /// The effective write mode for the run.
196    pub write_mode: WriteMode,
197    /// `true` for the first batch of the run — the signal an overwrite sink uses
198    /// to truncate-then-append (e.g. BigQuery `WRITE_TRUNCATE` on the first load,
199    /// `WRITE_APPEND` thereafter).
200    pub first_batch: bool,
201}
202
203/// Inputs to the pure [`plan_native_transfer`] negotiation.
204#[derive(Clone, Copy, Debug)]
205pub struct NativePlanInputs<'a> {
206    /// Formats the source can emit, in **preference order** (first = best).
207    pub source_formats: &'a [NativeFormat],
208    /// Mechanisms the sink offers.
209    pub sink_caps: &'a [NativeLoadCapability],
210    /// Whether any transform stage is attached to the pipeline.
211    pub has_transforms: bool,
212    /// Whether any quality / contract / masking / schema-drift pass is active.
213    pub has_governance: bool,
214    /// The run's delivery mode.
215    pub delivery: DeliveryMode,
216    /// The run's effective write mode.
217    pub write_mode: WriteMode,
218    /// Whether a DLQ is configured.
219    pub has_dlq: bool,
220}
221
222/// The negotiated native transfer: which format to move and which sink mechanism
223/// loads it.
224#[derive(Clone, Copy, Debug, PartialEq, Eq)]
225pub struct NativePlan {
226    /// The wire format both sides agreed on.
227    pub format: NativeFormat,
228    /// The sink mechanism's label.
229    pub mechanism: &'static str,
230}
231
232/// Decide whether a native byte-passthrough fast path applies, and which.
233///
234/// Deterministic and pure. First the **pipeline-owned gates** run — these are
235/// unconditional and no capability can waive them:
236///
237/// - transforms or any quality/contract/masking pass (`has_transforms` /
238///   `has_governance`): those need per-record `Value` access, which a byte
239///   passthrough never produces — running them would silently skip governance;
240/// - a configured DLQ (`has_dlq`): a byte batch cannot be split into per-row
241///   successes/failures, so the DLQ would be silently inert;
242/// - any delivery other than [`DeliveryMode::AtLeastOnce`]: the native runner
243///   implements no commit-token protocol, so exactly-once is rejected
244///   structurally rather than sold as a capability flag.
245///
246/// Then it walks the source's formats **in preference order** and returns the
247/// first one for which some sink capability matches on format and lists the
248/// run's write mode. Source preference order therefore breaks ties when the
249/// sink offers several matching formats. Returns `None` when no mechanism
250/// qualifies — the caller then falls through to the columnar or `Value` path
251/// unchanged.
252pub fn plan_native_transfer(inputs: &NativePlanInputs<'_>) -> Option<NativePlan> {
253    if inputs.has_transforms || inputs.has_governance || inputs.has_dlq {
254        return None;
255    }
256    if inputs.delivery != DeliveryMode::AtLeastOnce {
257        return None;
258    }
259    for &format in inputs.source_formats {
260        if let Some(cap) = inputs
261            .sink_caps
262            .iter()
263            .find(|cap| cap.format == format && cap.write_modes.contains(&inputs.write_mode))
264        {
265            return Some(NativePlan {
266                format: cap.format,
267                mechanism: cap.mechanism,
268            });
269        }
270    }
271    None
272}
273
274#[cfg(test)]
275mod tests {
276    use super::*;
277
278    /// A BigQuery-like capability: NDJSON + CSV, append or overwrite.
279    fn bq_caps() -> Vec<NativeLoadCapability> {
280        vec![
281            NativeLoadCapability {
282                format: NativeFormat::NdJson,
283                mechanism: "bigquery-load-job",
284                write_modes: &[WriteMode::Append, WriteMode::Overwrite],
285            },
286            NativeLoadCapability {
287                format: NativeFormat::Csv,
288                mechanism: "bigquery-load-job",
289                write_modes: &[WriteMode::Append, WriteMode::Overwrite],
290            },
291        ]
292    }
293
294    fn base<'a>(
295        source_formats: &'a [NativeFormat],
296        sink_caps: &'a [NativeLoadCapability],
297    ) -> NativePlanInputs<'a> {
298        NativePlanInputs {
299            source_formats,
300            sink_caps,
301            has_transforms: false,
302            has_governance: false,
303            delivery: DeliveryMode::AtLeastOnce,
304            write_mode: WriteMode::Append,
305            has_dlq: false,
306        }
307    }
308
309    #[test]
310    fn matches_when_format_and_prereqs_hold() {
311        let caps = bq_caps();
312        let src = [NativeFormat::Csv];
313        let plan = plan_native_transfer(&base(&src, &caps));
314        assert_eq!(
315            plan,
316            Some(NativePlan {
317                format: NativeFormat::Csv,
318                mechanism: "bigquery-load-job",
319            })
320        );
321    }
322
323    #[test]
324    fn source_preference_order_breaks_ties() {
325        let caps = bq_caps();
326        // Source prefers NdJson; both offered by the sink → NdJson wins.
327        let src = [NativeFormat::NdJson, NativeFormat::Csv];
328        assert_eq!(
329            plan_native_transfer(&base(&src, &caps)).unwrap().format,
330            NativeFormat::NdJson
331        );
332        // Reverse the preference → Csv wins.
333        let src = [NativeFormat::Csv, NativeFormat::NdJson];
334        assert_eq!(
335            plan_native_transfer(&base(&src, &caps)).unwrap().format,
336            NativeFormat::Csv
337        );
338    }
339
340    #[test]
341    fn no_match_when_formats_disjoint() {
342        let caps = bq_caps();
343        let src = [NativeFormat::Parquet, NativeFormat::ArrowIpc];
344        assert_eq!(plan_native_transfer(&base(&src, &caps)), None);
345    }
346
347    #[test]
348    fn empty_source_or_sink_yields_none() {
349        let caps = bq_caps();
350        assert_eq!(plan_native_transfer(&base(&[], &caps)), None);
351        let src = [NativeFormat::Csv];
352        assert_eq!(plan_native_transfer(&base(&src, &[])), None);
353    }
354
355    #[test]
356    fn transforms_and_governance_gates_block_unconditionally() {
357        let caps = bq_caps();
358        let src = [NativeFormat::Csv];
359        let mut inp = base(&src, &caps);
360        inp.has_transforms = true;
361        assert_eq!(plan_native_transfer(&inp), None, "transforms must block");
362        let mut inp = base(&src, &caps);
363        inp.has_governance = true;
364        assert_eq!(plan_native_transfer(&inp), None, "governance must block");
365    }
366
367    #[test]
368    fn governance_and_dlq_gates_are_not_capability_waivable() {
369        // Even the most permissive capability shape cannot opt out of the
370        // pipeline-owned governance/DLQ gates — those need `Value` access, so
371        // a waiver could only ship silently-unenforced governance.
372        let caps = vec![NativeLoadCapability {
373            format: NativeFormat::Csv,
374            mechanism: "tolerant",
375            write_modes: &[
376                WriteMode::Append,
377                WriteMode::Overwrite,
378                WriteMode::Upsert,
379                WriteMode::Delete,
380            ],
381        }];
382        let src = [NativeFormat::Csv];
383        let mut inp = base(&src, &caps);
384        inp.has_transforms = true;
385        assert_eq!(plan_native_transfer(&inp), None, "transforms must block");
386        let mut inp = base(&src, &caps);
387        inp.has_governance = true;
388        assert_eq!(plan_native_transfer(&inp), None, "governance must block");
389        let mut inp = base(&src, &caps);
390        inp.has_dlq = true;
391        assert_eq!(plan_native_transfer(&inp), None, "DLQ must block");
392    }
393
394    #[test]
395    fn exactly_once_is_rejected_structurally() {
396        // The native runner implements no commit-token protocol, so no
397        // capability can declare exactly-once support.
398        let caps = bq_caps();
399        let src = [NativeFormat::Csv];
400        let mut inp = base(&src, &caps);
401        inp.delivery = DeliveryMode::ExactlyOnce;
402        assert_eq!(plan_native_transfer(&inp), None);
403    }
404
405    #[test]
406    fn write_mode_gate_blocks_unlisted_modes() {
407        let caps = bq_caps();
408        let src = [NativeFormat::Csv];
409        let mut inp = base(&src, &caps);
410        inp.write_mode = WriteMode::Upsert;
411        assert_eq!(plan_native_transfer(&inp), None);
412        // Overwrite is allowed by the bq caps.
413        let mut inp = base(&src, &caps);
414        inp.write_mode = WriteMode::Overwrite;
415        assert!(plan_native_transfer(&inp).is_some());
416    }
417
418    #[test]
419    fn dlq_gate_blocks_unconditionally() {
420        let caps = bq_caps();
421        let src = [NativeFormat::Csv];
422        let mut inp = base(&src, &caps);
423        inp.has_dlq = true;
424        assert_eq!(plan_native_transfer(&inp), None);
425    }
426
427    #[test]
428    fn format_helpers() {
429        assert_eq!(NativeFormat::NdJson.as_str(), "ndjson");
430        assert_eq!(NativeFormat::Csv.as_str(), "csv");
431        assert_eq!(NativeFormat::Parquet.as_str(), "parquet");
432        assert_eq!(NativeFormat::ArrowIpc.as_str(), "arrow_ipc");
433        assert_eq!(CsvDialect::default().delimiter, b',');
434        assert!(CsvDialect::default().has_header);
435    }
436
437    #[test]
438    fn native_batch_builders_and_debug() {
439        let b = NativeBatch::bytes(NativeFormat::NdJson, b"{}\n".to_vec())
440            .with_bookmark(Some(serde_json::json!({"o": 1})))
441            .with_records(Some(1))
442            .with_csv(CsvDialect {
443                has_header: false,
444                delimiter: b'\t',
445            });
446        assert_eq!(b.format, NativeFormat::NdJson);
447        assert_eq!(b.records, Some(1));
448        assert!(b.bookmark.is_some());
449        assert!(!b.csv.has_header);
450        // Debug on the payload prints the length, not the bytes.
451        let dbg = format!("{:?}", b.payload);
452        assert!(dbg.contains("Bytes"), "{dbg}");
453    }
454
455    #[test]
456    fn native_payload_stream_debug() {
457        let s = NativePayload::Stream(Box::pin(futures::stream::empty()));
458        assert!(format!("{s:?}").contains("byte stream"));
459    }
460
461    #[test]
462    fn load_context_carries_first_batch_and_mode() {
463        let ctx = NativeLoadContext {
464            write_mode: WriteMode::Overwrite,
465            first_batch: true,
466        };
467        assert!(ctx.first_batch);
468        assert_eq!(ctx.write_mode, WriteMode::Overwrite);
469    }
470}