use crate::error::FaucetError;
use crate::idempotency::DeliveryMode;
use crate::write_mode::WriteMode;
use futures_core::Stream;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::pin::Pin;
#[non_exhaustive]
#[derive(Clone, Copy, PartialEq, Eq, Debug, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum NativeFormat {
NdJson,
Csv,
Parquet,
ArrowIpc,
}
impl NativeFormat {
pub fn as_str(self) -> &'static str {
match self {
NativeFormat::NdJson => "ndjson",
NativeFormat::Csv => "csv",
NativeFormat::Parquet => "parquet",
NativeFormat::ArrowIpc => "arrow_ipc",
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct CsvDialect {
pub has_header: bool,
pub delimiter: u8,
}
impl Default for CsvDialect {
fn default() -> Self {
Self {
has_header: true,
delimiter: b',',
}
}
}
pub enum NativePayload {
Bytes(Vec<u8>),
Stream(Pin<Box<dyn Stream<Item = Result<Vec<u8>, FaucetError>> + Send>>),
}
impl std::fmt::Debug for NativePayload {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
NativePayload::Bytes(b) => f.debug_tuple("Bytes").field(&b.len()).finish(),
NativePayload::Stream(_) => f.write_str("Stream(<byte stream>)"),
}
}
}
#[derive(Debug)]
pub struct NativeBatch {
pub format: NativeFormat,
pub payload: NativePayload,
pub csv: CsvDialect,
pub records: Option<u64>,
pub bookmark: Option<Value>,
}
impl NativeBatch {
pub fn bytes(format: NativeFormat, payload: Vec<u8>) -> Self {
Self {
format,
payload: NativePayload::Bytes(payload),
csv: CsvDialect::default(),
records: None,
bookmark: None,
}
}
pub fn with_bookmark(mut self, bookmark: Option<Value>) -> Self {
self.bookmark = bookmark;
self
}
pub fn with_records(mut self, records: Option<u64>) -> Self {
self.records = records;
self
}
pub fn with_csv(mut self, csv: CsvDialect) -> Self {
self.csv = csv;
self
}
}
#[derive(Clone, Debug)]
pub struct NativeLoadCapability {
pub format: NativeFormat,
pub mechanism: &'static str,
pub write_modes: &'static [WriteMode],
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct NativeLoadContext {
pub write_mode: WriteMode,
pub first_batch: bool,
}
#[derive(Clone, Copy, Debug)]
pub struct NativePlanInputs<'a> {
pub source_formats: &'a [NativeFormat],
pub sink_caps: &'a [NativeLoadCapability],
pub has_transforms: bool,
pub has_governance: bool,
pub delivery: DeliveryMode,
pub write_mode: WriteMode,
pub has_dlq: bool,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct NativePlan {
pub format: NativeFormat,
pub mechanism: &'static str,
}
pub fn plan_native_transfer(inputs: &NativePlanInputs<'_>) -> Option<NativePlan> {
if inputs.has_transforms || inputs.has_governance || inputs.has_dlq {
return None;
}
if inputs.delivery != DeliveryMode::AtLeastOnce {
return None;
}
for &format in inputs.source_formats {
if let Some(cap) = inputs
.sink_caps
.iter()
.find(|cap| cap.format == format && cap.write_modes.contains(&inputs.write_mode))
{
return Some(NativePlan {
format: cap.format,
mechanism: cap.mechanism,
});
}
}
None
}
#[cfg(test)]
mod tests {
use super::*;
fn bq_caps() -> Vec<NativeLoadCapability> {
vec![
NativeLoadCapability {
format: NativeFormat::NdJson,
mechanism: "bigquery-load-job",
write_modes: &[WriteMode::Append, WriteMode::Overwrite],
},
NativeLoadCapability {
format: NativeFormat::Csv,
mechanism: "bigquery-load-job",
write_modes: &[WriteMode::Append, WriteMode::Overwrite],
},
]
}
fn base<'a>(
source_formats: &'a [NativeFormat],
sink_caps: &'a [NativeLoadCapability],
) -> NativePlanInputs<'a> {
NativePlanInputs {
source_formats,
sink_caps,
has_transforms: false,
has_governance: false,
delivery: DeliveryMode::AtLeastOnce,
write_mode: WriteMode::Append,
has_dlq: false,
}
}
#[test]
fn matches_when_format_and_prereqs_hold() {
let caps = bq_caps();
let src = [NativeFormat::Csv];
let plan = plan_native_transfer(&base(&src, &caps));
assert_eq!(
plan,
Some(NativePlan {
format: NativeFormat::Csv,
mechanism: "bigquery-load-job",
})
);
}
#[test]
fn source_preference_order_breaks_ties() {
let caps = bq_caps();
let src = [NativeFormat::NdJson, NativeFormat::Csv];
assert_eq!(
plan_native_transfer(&base(&src, &caps)).unwrap().format,
NativeFormat::NdJson
);
let src = [NativeFormat::Csv, NativeFormat::NdJson];
assert_eq!(
plan_native_transfer(&base(&src, &caps)).unwrap().format,
NativeFormat::Csv
);
}
#[test]
fn no_match_when_formats_disjoint() {
let caps = bq_caps();
let src = [NativeFormat::Parquet, NativeFormat::ArrowIpc];
assert_eq!(plan_native_transfer(&base(&src, &caps)), None);
}
#[test]
fn empty_source_or_sink_yields_none() {
let caps = bq_caps();
assert_eq!(plan_native_transfer(&base(&[], &caps)), None);
let src = [NativeFormat::Csv];
assert_eq!(plan_native_transfer(&base(&src, &[])), None);
}
#[test]
fn transforms_and_governance_gates_block_unconditionally() {
let caps = bq_caps();
let src = [NativeFormat::Csv];
let mut inp = base(&src, &caps);
inp.has_transforms = true;
assert_eq!(plan_native_transfer(&inp), None, "transforms must block");
let mut inp = base(&src, &caps);
inp.has_governance = true;
assert_eq!(plan_native_transfer(&inp), None, "governance must block");
}
#[test]
fn governance_and_dlq_gates_are_not_capability_waivable() {
let caps = vec![NativeLoadCapability {
format: NativeFormat::Csv,
mechanism: "tolerant",
write_modes: &[
WriteMode::Append,
WriteMode::Overwrite,
WriteMode::Upsert,
WriteMode::Delete,
],
}];
let src = [NativeFormat::Csv];
let mut inp = base(&src, &caps);
inp.has_transforms = true;
assert_eq!(plan_native_transfer(&inp), None, "transforms must block");
let mut inp = base(&src, &caps);
inp.has_governance = true;
assert_eq!(plan_native_transfer(&inp), None, "governance must block");
let mut inp = base(&src, &caps);
inp.has_dlq = true;
assert_eq!(plan_native_transfer(&inp), None, "DLQ must block");
}
#[test]
fn exactly_once_is_rejected_structurally() {
let caps = bq_caps();
let src = [NativeFormat::Csv];
let mut inp = base(&src, &caps);
inp.delivery = DeliveryMode::ExactlyOnce;
assert_eq!(plan_native_transfer(&inp), None);
}
#[test]
fn write_mode_gate_blocks_unlisted_modes() {
let caps = bq_caps();
let src = [NativeFormat::Csv];
let mut inp = base(&src, &caps);
inp.write_mode = WriteMode::Upsert;
assert_eq!(plan_native_transfer(&inp), None);
let mut inp = base(&src, &caps);
inp.write_mode = WriteMode::Overwrite;
assert!(plan_native_transfer(&inp).is_some());
}
#[test]
fn dlq_gate_blocks_unconditionally() {
let caps = bq_caps();
let src = [NativeFormat::Csv];
let mut inp = base(&src, &caps);
inp.has_dlq = true;
assert_eq!(plan_native_transfer(&inp), None);
}
#[test]
fn format_helpers() {
assert_eq!(NativeFormat::NdJson.as_str(), "ndjson");
assert_eq!(NativeFormat::Csv.as_str(), "csv");
assert_eq!(NativeFormat::Parquet.as_str(), "parquet");
assert_eq!(NativeFormat::ArrowIpc.as_str(), "arrow_ipc");
assert_eq!(CsvDialect::default().delimiter, b',');
assert!(CsvDialect::default().has_header);
}
#[test]
fn native_batch_builders_and_debug() {
let b = NativeBatch::bytes(NativeFormat::NdJson, b"{}\n".to_vec())
.with_bookmark(Some(serde_json::json!({"o": 1})))
.with_records(Some(1))
.with_csv(CsvDialect {
has_header: false,
delimiter: b'\t',
});
assert_eq!(b.format, NativeFormat::NdJson);
assert_eq!(b.records, Some(1));
assert!(b.bookmark.is_some());
assert!(!b.csv.has_header);
let dbg = format!("{:?}", b.payload);
assert!(dbg.contains("Bytes"), "{dbg}");
}
#[test]
fn native_payload_stream_debug() {
let s = NativePayload::Stream(Box::pin(futures::stream::empty()));
assert!(format!("{s:?}").contains("byte stream"));
}
#[test]
fn load_context_carries_first_batch_and_mode() {
let ctx = NativeLoadContext {
write_mode: WriteMode::Overwrite,
first_batch: true,
};
assert!(ctx.first_batch);
assert_eq!(ctx.write_mode, WriteMode::Overwrite);
}
}