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}