Skip to main content

rivet/config/
export.rs

1//! Per-export configuration: query/table/mode, chunking, format, destination link.
2//!
3//! `SchemaDriftPolicy` lives here because it is only ever read via
4//! [`ExportConfig::on_schema_drift`].
5
6use std::path::Path;
7
8use schemars::JsonSchema;
9use serde::{Deserialize, Serialize};
10
11use super::IncrementalCursorMode;
12use super::destination::DestinationConfig;
13use super::format::{CompressionProfile, CompressionType, FormatType, ParquetConfig};
14use super::resolve::{parse_file_size, resolve_vars};
15use crate::tuning::TuningConfig;
16
17/// What to do when structural schema drift is detected (column added, removed, or retyped).
18///
19/// ```yaml
20/// exports:
21///   - name: orders
22///     on_schema_drift: fail   # warn (default), continue, fail
23/// ```
24/// How deep `--validate` must verify each part's integrity.
25#[derive(Debug, Deserialize, Serialize, JsonSchema, Clone, Copy, PartialEq, Eq, Default)]
26#[serde(rename_all = "snake_case")]
27pub enum VerifyMode {
28    /// Accept size-only verification when no content checksum is available.
29    #[default]
30    Size,
31    /// Require every part's content to be MD5-verified against the store's
32    /// listing; fail validation for any part that is only size-verified.
33    Content,
34}
35
36impl VerifyMode {
37    /// Whether content (not just size) verification is required.
38    pub fn requires_content(self) -> bool {
39        matches!(self, VerifyMode::Content)
40    }
41}
42
43#[derive(Debug, Deserialize, Serialize, JsonSchema, Clone, Copy, PartialEq, Eq, Default)]
44#[serde(rename_all = "snake_case")]
45pub enum SchemaDriftPolicy {
46    /// Log a warning and continue. The new schema fingerprint is stored. (Default.)
47    #[default]
48    Warn,
49    /// Silently accept schema changes — store the new schema, no log output.
50    Continue,
51    /// Abort the run with a non-zero exit. The schema store is NOT updated so the
52    /// next run will detect the same change again.
53    Fail,
54}
55#[derive(Debug, Deserialize, JsonSchema, Clone)]
56#[serde(deny_unknown_fields)]
57pub struct ExportConfig {
58    pub name: String,
59    #[serde(default)]
60    pub query: Option<String>,
61    pub query_file: Option<String>,
62    /// Shortcut for `query: "SELECT * FROM <schema>.<table>"`.
63    ///
64    /// Accepts `table` or `schema.table` with ASCII-only identifiers
65    /// (`[A-Za-z_][A-Za-z0-9_]*`). Generates an unquoted single-table
66    /// query so the Postgres NUMERIC catalog-hint resolver recognises it
67    /// and auto-types `numeric(p,s)` columns without manual overrides.
68    ///
69    /// Mutually exclusive with `query` and `query_file`.
70    #[serde(default)]
71    pub table: Option<String>,
72    /// CDC only: capture **several** tables through ONE change stream (one
73    /// PostgreSQL slot / one MySQL binlog connection) instead of one export —
74    /// and one slot — per table. Each table's parts land under
75    /// `<destination>/<table>/` with their own `manifest.json` + `_SUCCESS`;
76    /// the checkpoint (stream position) is shared. Mutually exclusive with
77    /// `table:`. Not yet supported for SQL Server (capture instances are
78    /// per-table).
79    #[serde(default)]
80    pub tables: Option<Vec<String>>,
81    #[serde(default = "default_mode")]
82    pub mode: ExportMode,
83    /// Change-data-capture settings, required when `mode: cdc`. Reuses the
84    /// export's `table`, `destination`, and `format`; carries only the
85    /// CDC-specific knobs (resume checkpoint, per-engine stream params).
86    #[serde(default)]
87    pub cdc: Option<CdcExportConfig>,
88    pub cursor_column: Option<String>,
89    /// Secondary column for [`IncrementalCursorMode::Coalesce`] only (see ADR-0007).
90    #[serde(default)]
91    pub cursor_fallback_column: Option<String>,
92    /// How primary (and optional fallback) columns drive incremental progression.
93    #[serde(default)]
94    pub incremental_cursor_mode: IncrementalCursorMode,
95    pub chunk_column: Option<String>,
96    #[serde(default)]
97    pub chunk_dense: bool,
98    #[serde(default = "default_chunk_size")]
99    pub chunk_size: usize,
100    /// Target memory budget per chunk in MB. When set, `chunk_size` is derived
101    /// from this budget at plan-build time using a `pg_class` row-size estimate
102    /// (`pg_relation_size / reltuples`), clamped to `[10_000, 5_000_000]` rows.
103    ///
104    /// Mutually exclusive with an explicit non-default `chunk_size:`. Only
105    /// applies to `mode: chunked` on a Postgres source using the `table:`
106    /// shortcut (the row-size probe needs a known relation).
107    ///
108    /// ```yaml
109    /// exports:
110    ///   - name: page_views
111    ///     table: public.page_views
112    ///     mode: chunked
113    ///     chunk_size_memory_mb: 256
114    /// ```
115    #[serde(default)]
116    pub chunk_size_memory_mb: Option<u64>,
117    /// Divide the column range into exactly this many equal chunks.
118    /// Mutually exclusive with `chunk_dense` and `chunk_by_days`.
119    /// When set, `chunk_size` is computed dynamically from min/max.
120    pub chunk_count: Option<usize>,
121    pub chunk_by_days: Option<u32>,
122    /// Keyset (seek) pagination on this single index-backed unique key — the
123    /// source-safe shape for tables without a single-integer PK (OPT-4). The
124    /// column MUST be backed by a usable index (PK or unique); the planner
125    /// refuses a non-indexed key rather than emit a full-scan + filesort query.
126    pub chunk_by_key: Option<String>,
127    /// Concurrent chunk/page workers (default 1). On a RANGE chunk (`chunk_column`)
128    /// or KEYSET (`chunk_by_key`) export, `parallel: N` fans the table into N
129    /// ROW-percentile ranges that seek concurrently over separate connections — the
130    /// half-open intervals partition the key, so the union reads every row exactly
131    /// once (structural parity, all engines). Extraction is I/O-bound, so the win
132    /// plateaus early (~3x at N=4, little beyond). SWEET SPOT: indexed tables up to
133    /// ~10M rows at `parallel: 4`. `rivet init` scaffolds a row-scaled value
134    /// (<=500K -> 1, <5M -> 2, >=5M -> 4); a preflight warns past ~5M rows (peak RSS
135    /// ~= N x chunk_size). Beyond ~10M the KEYSET boundary sampler (an index OFFSET
136    /// skip) grows costly at setup — prefer a range `chunk_column` there.
137    #[serde(default = "default_parallel")]
138    pub parallel: usize,
139
140    /// Advisory execution wave (1 = highest priority, run first). Written by
141    /// `rivet plan` from the source-aware prioritization score (see ADR-0006)
142    /// and consumed by `rivet apply`, which runs exports wave-by-wave in
143    /// ascending order. `None` = unscheduled (apply treats it as the last wave).
144    /// Operators may hand-edit it; a later `rivet plan` refreshes it in place.
145    #[serde(default, skip_serializing_if = "Option::is_none")]
146    pub wave: Option<u32>,
147
148    /// Whether this export is cheap enough to run concurrently with its
149    /// wave-mates under `rivet apply --parallel-export-processes`. Written by
150    /// `rivet plan` (true when the source-aware cost class is `Low`, i.e.
151    /// < ~100K rows); a heavier table already chunk-parallelizes internally, so
152    /// two of them at once would overload the source. `None`/`false` → the
153    /// export runs alone within its wave. Operators may hand-edit it; a later
154    /// `rivet plan` refreshes it in place.
155    #[serde(default, skip_serializing_if = "Option::is_none")]
156    pub parallel_safe: Option<bool>,
157    pub time_column: Option<String>,
158    #[serde(default = "default_time_column_type")]
159    pub time_column_type: TimeColumnType,
160    pub days_window: Option<u32>,
161
162    /// Date/time output partitioning: split this export's rows into one
163    /// destination sub-prefix per calendar bucket of this **DATE or TIMESTAMP**
164    /// column, bucketed by `partition_granularity`
165    /// (`day` / `month` / `year`), in a Hive-style `col=value/` layout
166    /// (`created_at=2023-01-01/`, `created_at=2023-01/`, `created_at=2023/`).
167    /// Requires a `{partition}` token in `destination.path` /
168    /// `destination.prefix`.
169    ///
170    /// This is **not** arbitrary value partitioning: the column's min/max is
171    /// read and parsed as a date to generate contiguous calendar buckets, so a
172    /// non-temporal column (e.g. `partition_by: status`) fails at run time with
173    /// "could not parse partition min `<value>` from column `<col>` as a date".
174    /// To split by a categorical column, write one export per value with a
175    /// `WHERE` filter instead.
176    ///
177    /// Orthogonal to `mode`: each partition runs the export's own mode, so
178    /// `mode: chunked` chunks *within* a day. Rows whose partition column is
179    /// NULL land in `col=__HIVE_DEFAULT_PARTITION__/` (Hive default partition)
180    /// so no row is silently dropped. Not compatible with `mode: time_window`.
181    ///
182    /// ```yaml
183    /// exports:
184    ///   - name: events
185    ///     table: events
186    ///     partition_by: created_at        # must be a DATE or TIMESTAMP column
187    ///     partition_granularity: day
188    ///     destination:
189    ///       type: s3
190    ///       bucket: my-bucket
191    ///       prefix: "events/{partition}/"   # → events/created_at=2023-01-01/
192    /// ```
193    #[serde(default)]
194    pub partition_by: Option<String>,
195
196    /// Calendar bucket width for `partition_by`:
197    /// `day` (default), `month`, or `year`. Determines how the partition
198    /// column's date/timestamp range is split into contiguous Hive buckets
199    /// (`col=2023-01-01/` / `col=2023-01/` / `col=2023/`). Has no effect
200    /// unless `partition_by` is set.
201    #[serde(default)]
202    pub partition_granularity: PartitionGranularity,
203    pub format: FormatType,
204    #[serde(default)]
205    pub compression: CompressionType,
206    pub compression_level: Option<u32>,
207    pub compression_profile: Option<CompressionProfile>,
208    #[serde(default)]
209    pub skip_empty: bool,
210    pub destination: DestinationConfig,
211    /// Integrity depth required of `--validate` for this export's parts.
212    /// `size` (default) accepts size-only verification; `content` requires every
213    /// part's content MD5 to be checked against the store's listing (no
214    /// download) and **fails** validation for any part that could only be
215    /// size-verified — e.g. a part too large to upload as a single PUT (raise
216    /// `max_file_size` down so it fits), or a backend that exposes no checksum.
217    #[serde(default)]
218    pub verify: VerifyMode,
219    #[serde(default)]
220    pub meta_columns: MetaColumns,
221    #[serde(default)]
222    pub quality: Option<QualityConfig>,
223    /// Rotate to a new part when the current file reaches this size.
224    /// Accepts `B`/`KB`/`MB`/`GB` (case-insensitive) or a bare byte count;
225    /// a fractional value is allowed (`1.5GB`). Units are binary (IEC-style):
226    /// `KB` = 1024 bytes, `MB` = 1024 KB, `GB` = 1024 MB. Example: `256MB`.
227    pub max_file_size: Option<String>,
228    /// Persist per-chunk / per-page progress so a **crashed** run resumes from the
229    /// last durably committed point instead of re-reading from the start. This is
230    /// pure crash-recovery: a *clean* re-run (the prior run finished) still does a
231    /// full pass — it never silently skips already-exported rows. Safe to enable
232    /// on any table; `rivet init` defaults it on for chunked and keyset exports.
233    #[serde(default)]
234    pub chunk_checkpoint: bool,
235    /// Keyset only (`chunk_by_key`): on a **clean** re-run, continue from the last
236    /// exported key — pull ONLY rows with a key past the high-water mark. This is
237    /// incremental-by-key, correct ONLY for APPEND-ONLY tables (a mutable row whose
238    /// key already passed is silently never re-read). Opt-in and off by default;
239    /// crash-recovery does not need it (that is `chunk_checkpoint`). For a mutable
240    /// table use `mode: incremental` on a timestamp cursor instead.
241    #[serde(default)]
242    pub keyset_incremental: bool,
243    pub chunk_max_attempts: Option<u32>,
244    #[serde(default)]
245    pub tuning: Option<TuningConfig>,
246    /// Optional logical group for shared source capacity (replica, host). Advisory prioritization only.
247    #[serde(default)]
248    pub source_group: Option<String>,
249    /// Hint (Epic C / ADR-0006) that this export should always be treated as reconcile-heavy
250    /// by planning, independent of the `--reconcile` CLI flag. Advisory only.
251    #[serde(default)]
252    pub reconcile_required: bool,
253
254    /// Per-column type overrides (roadmap §8). Keys are column names; values
255    /// are short type strings such as `decimal(18,2)`, `timestamp_tz`, `json`.
256    ///
257    /// ```yaml
258    /// exports:
259    ///   - name: payments
260    ///     columns:
261    ///       amount: decimal(18,2)
262    ///       fee: decimal(18,6)
263    ///       created_at: timestamp_tz
264    /// ```
265    ///
266    /// Overrides take priority over autodetection and are validated at
267    /// plan time — an invalid type string fails before the export runs.
268    #[serde(default)]
269    pub columns: std::collections::HashMap<String, String>,
270
271    /// Downstream warehouse this export targets (`bigquery` / `bq`,
272    /// `duckdb`). When set, `rivet check --type-report` resolves each column
273    /// against it (native type, honest autoload type, recovery hint) without
274    /// needing `--target` on the CLI — the CLI flag still wins when both are
275    /// present. The Parquet interchange stays target-neutral (ADR-0014 T2);
276    /// `target:` only drives guidance and the future load-schema artifact.
277    ///
278    /// ```yaml
279    /// exports:
280    ///   - name: payments
281    ///     target: bigquery
282    /// ```
283    #[serde(default)]
284    pub target: Option<String>,
285
286    /// Per-export overrides for the top-level `load:` block (`pk`,
287    /// `cleanup_source`, `gc_orphans`, `cluster_by`, `allow_source_drift`); any
288    /// field omitted here inherits the top-level value. The warehouse `target`
289    /// is shared and stays in the top-level `load:` — it cannot be overridden
290    /// per export.
291    ///
292    /// ```yaml
293    /// load: { target: bigquery, project: p, dataset: d }   # shared default
294    /// exports:
295    ///   - name: orders
296    ///     table: orders
297    ///     mode: cdc
298    ///     load: { pk: [id] }                                # this table's pk
299    /// ```
300    ///
301    /// Raw JSON (parsed by the load module) so `config` carries no load types —
302    /// mirrors the top-level [`crate::config::Config::load`].
303    #[serde(default)]
304    pub load: Option<serde_json::Value>,
305
306    /// Policy applied when structural schema drift is detected (column added, removed, or retyped).
307    /// Defaults to `warn`: log a warning and continue.
308    #[serde(default)]
309    pub on_schema_drift: SchemaDriftPolicy,
310
311    /// Growth-factor threshold for data shape drift warnings (Epic 8).
312    /// When a string/binary column's max observed byte length in the current run
313    /// exceeds `stored_max * shape_drift_warn_factor`, Rivet logs a warning.
314    /// `None` uses the default of 2.0. Set to `0.0` to disable shape tracking.
315    ///
316    /// **Scope: single-batch exports only** (`mode: full` / `incremental`). Shape
317    /// tracking needs the run-wide per-column max byte length, which only the
318    /// single sink accumulates; the multi-part runners (chunked, keyset, and
319    /// parallel-Mongo) each write through several sinks and do not aggregate it,
320    /// so this factor is inert there. (Schema drift — added/dropped/retyped
321    /// columns — IS enforced on every path via `on_schema_drift`.)
322    #[serde(default)]
323    pub shape_drift_warn_factor: Option<f64>,
324
325    /// Parquet row group tuning. Only meaningful when `format: parquet`.
326    /// When absent, the parquet library default (1,048,576 rows/group) is used.
327    #[serde(default)]
328    pub parquet: Option<ParquetConfig>,
329}
330
331impl ExportConfig {
332    /// Resolve the effective `(CompressionType, level)` for this export.
333    /// `compression_profile` takes precedence over `compression` + `compression_level`.
334    ///
335    /// L24: when a profile is set *and* a conflicting explicit codec/level was
336    /// written, warn once that the profile wins rather than silently dropping the
337    /// explicit choice. An explicit codec is only detectable when it differs from
338    /// the `#[serde(default)]` (Zstd) — a literal `compression: zstd` alongside a
339    /// profile is indistinguishable from an omitted field and stays silent.
340    pub fn effective_compression(&self) -> (CompressionType, Option<u32>) {
341        if let Some(profile) = self.compression_profile {
342            let explicit_codec =
343                (self.compression != CompressionType::default()).then_some(self.compression);
344            if let Some(msg) = super::format::compression_profile_override_warning(
345                profile,
346                explicit_codec,
347                self.compression_level,
348            ) {
349                log::warn!("export '{}': {}", self.name, msg);
350            }
351            profile.to_codec()
352        } else {
353            (self.compression, self.compression_level)
354        }
355    }
356
357    pub fn max_file_size_bytes(&self) -> Option<u64> {
358        self.max_file_size
359            .as_ref()
360            .and_then(|s| parse_file_size(s).ok())
361    }
362
363    pub fn resolve_query(
364        &self,
365        config_dir: &Path,
366        params: Option<&std::collections::HashMap<String, String>>,
367    ) -> crate::error::Result<String> {
368        // table: shortcut takes precedence — already validated by
369        // `validate_business_rules` to be mutually exclusive with query/query_file.
370        if let Some(tbl) = &self.table {
371            validate_table_shortcut_ident(&self.name, tbl)?;
372            return Ok(format!("SELECT * FROM {tbl}"));
373        }
374        match (&self.query, &self.query_file) {
375            (Some(q), None) => {
376                if params.is_some() {
377                    resolve_vars(q, params)
378                } else {
379                    Ok(q.clone())
380                }
381            }
382            (None, Some(file)) => {
383                let file_path = std::path::Path::new(file);
384                // SecOps: block absolute paths and `..` traversal components.
385                if file_path.is_absolute() {
386                    anyhow::bail!(
387                        "export '{}': query_file must be a relative path: '{}'",
388                        self.name,
389                        file
390                    );
391                }
392                if file_path
393                    .components()
394                    .any(|c| c == std::path::Component::ParentDir)
395                {
396                    anyhow::bail!(
397                        "export '{}': query_file path must not contain '..': '{}'",
398                        self.name,
399                        file
400                    );
401                }
402                let joined = config_dir.join(file);
403                // Canonicalize-based check catches symlink-based evasion for files
404                // that already exist on disk.
405                if let Ok(canonical) = joined.canonicalize() {
406                    let base = config_dir
407                        .canonicalize()
408                        .unwrap_or_else(|_| config_dir.to_path_buf());
409                    if !canonical.starts_with(&base) {
410                        anyhow::bail!(
411                            "export '{}': query_file '{}' resolves outside the config directory",
412                            self.name,
413                            file
414                        );
415                    }
416                }
417                let raw = std::fs::read_to_string(&joined)?;
418                resolve_vars(&raw, params)
419            }
420            (Some(_), Some(_)) => {
421                anyhow::bail!(
422                    "export '{}': specify either 'query' or 'query_file', not both",
423                    self.name
424                )
425            }
426            (None, None) => {
427                anyhow::bail!(
428                    "export '{}': must specify exactly one of 'query', 'query_file', or 'table'",
429                    self.name
430                )
431            }
432        }
433    }
434}
435
436/// Validate the value of the `table:` YAML shortcut.
437///
438/// Accepts ASCII identifiers in the form `<table>` or `<schema>.<table>`. Each
439/// segment must match `[A-Za-z_][A-Za-z0-9_]*`. Anything else (quoted
440/// identifiers, exotic chars, three-part names, SQL injection attempts) is
441/// rejected — the user should fall back to `query:` for those cases.
442///
443/// The bound on identifier shape keeps generated SQL safe to interpolate
444/// without quoting and ensures the generated `SELECT * FROM <ident>` form is
445/// recognised by the PG catalog-hint parser ([src/source/postgres.rs]).
446fn validate_table_shortcut_ident(export_name: &str, raw: &str) -> crate::error::Result<()> {
447    let trimmed = raw.trim();
448    if trimmed.is_empty() {
449        anyhow::bail!("export '{export_name}': 'table' is empty");
450    }
451    let parts: Vec<&str> = trimmed.split('.').collect();
452    if parts.len() > 2 {
453        anyhow::bail!(
454            "export '{export_name}': 'table' must be `<name>` or `<schema>.<name>` (got '{raw}')"
455        );
456    }
457    for part in &parts {
458        if part.is_empty() {
459            anyhow::bail!("export '{export_name}': 'table' has an empty segment in '{raw}'");
460        }
461        let mut chars = part.chars();
462        let first = chars.next().unwrap();
463        if !(first.is_ascii_alphabetic() || first == '_') {
464            anyhow::bail!(
465                "export '{export_name}': 'table' segment '{part}' must start with a letter or underscore (use 'query:' for quoted identifiers)"
466            );
467        }
468        if !chars.all(|c| c.is_ascii_alphanumeric() || c == '_') {
469            anyhow::bail!(
470                "export '{export_name}': 'table' segment '{part}' contains non-identifier characters (use 'query:' for quoted identifiers)"
471            );
472        }
473    }
474    Ok(())
475}
476
477#[derive(Debug, Deserialize, Serialize, JsonSchema, Clone)]
478#[serde(deny_unknown_fields)]
479pub struct QualityConfig {
480    pub row_count_min: Option<usize>,
481    pub row_count_max: Option<usize>,
482    #[serde(default)]
483    pub null_ratio_max: std::collections::HashMap<String, f64>,
484    #[serde(default)]
485    pub unique_columns: Vec<String>,
486    /// Cap on the number of distinct values tracked per column during uniqueness checks.
487    /// When the limit is hit, a Warn issue is emitted and tracking stops for that column.
488    /// Prevents unbounded HashSet growth on high-cardinality columns.
489    pub unique_max_entries: Option<usize>,
490}
491
492#[derive(Debug, Deserialize, Serialize, JsonSchema, Clone, Default)]
493#[serde(deny_unknown_fields)]
494pub struct MetaColumns {
495    #[serde(default)]
496    pub exported_at: bool,
497    #[serde(default)]
498    pub row_hash: bool,
499}
500
501impl MetaColumns {
502    /// True iff the operator asked for at least one meta column. The batch
503    /// runners inject these at the shared `ExportSink` seam; the CDC path has
504    /// its OWN sink (`__op`/`__pos`/`__seq` + typed after-image) and does not,
505    /// so a CDC run uses this to warn that the request has no effect.
506    pub fn any_enabled(&self) -> bool {
507        self.exported_at || self.row_hash
508    }
509}
510
511fn default_mode() -> ExportMode {
512    ExportMode::Full
513}
514
515pub(crate) fn default_chunk_size() -> usize {
516    100_000
517}
518
519fn default_parallel() -> usize {
520    1
521}
522
523fn default_time_column_type() -> TimeColumnType {
524    TimeColumnType::Timestamp
525}
526
527/// `until_current` defaults to `true` — the OSS model is the BOUNDED, scheduler-
528/// driven drain ("read to the log end and exit"). Continuous streaming
529/// (`until_current: false`) is an explicit opt-in; making it the default silently
530/// put a hand-written CDC config onto the never-terminating streaming path.
531fn default_true() -> bool {
532    true
533}
534
535#[derive(Debug, Deserialize, JsonSchema, Clone, Copy, PartialEq, Eq)]
536#[serde(rename_all = "snake_case")]
537pub enum ExportMode {
538    Full,
539    Incremental,
540    Chunked,
541    TimeWindow,
542    /// Log-based change data capture (see [`CdcExportConfig`]): stream
543    /// INSERT/UPDATE/DELETE from the source's transaction log instead of querying
544    /// the table. Reuses the export's `table` / `destination` / `format`.
545    Cdc,
546}
547
548/// Default PostgreSQL logical slot when `cdc.slot` is omitted — shared by the
549/// runner ([`crate::pipeline`]'s cdc job) and config validation, so the
550/// same-slot conflict check sees the value that will actually be used.
551pub const DEFAULT_PG_SLOT: &str = "rivet_slot";
552/// Default MySQL replica `server_id` when `cdc.server_id` is omitted (see
553/// [`DEFAULT_PG_SLOT`] for why this is a shared const).
554pub const DEFAULT_MYSQL_SERVER_ID: u32 = 4271;
555
556/// What the FIRST CDC run does before draining changes.
557#[derive(Debug, Deserialize, Serialize, JsonSchema, Clone, Copy, PartialEq, Eq)]
558#[serde(rename_all = "lowercase")]
559pub enum CdcInitialMode {
560    /// Anchor-then-snapshot: create the resume anchor (PostgreSQL slot /
561    /// MySQL binlog checkpoint / SQL Server LSN checkpoint) FIRST, then run a
562    /// full batch snapshot of each table into `<destination>[/<table>]/snapshot/`,
563    /// then drain CDC. Because the anchor predates the snapshot read, anything
564    /// changed during the snapshot also appears in the change stream — an
565    /// overlap (dedupe by PK + `__op`), never a gap. The safe switch ordering,
566    /// enforced by construction instead of operator discipline.
567    Snapshot,
568}
569
570/// Per-export CDC settings, required when `mode: cdc`. The output `table`,
571/// `destination`, and `format` come from the export itself; this carries only the
572/// CDC-specific knobs (resume + per-engine stream params).
573#[derive(Debug, Deserialize, Serialize, JsonSchema, Clone)]
574pub struct CdcExportConfig {
575    /// First-run behaviour: `snapshot` = anchor → full snapshot → drain (see
576    /// [`CdcInitialMode`]). Omitted ⇒ capture changes only (the default; the
577    /// operator owns the initial load).
578    #[serde(default)]
579    pub initial: Option<CdcInitialMode>,
580    /// Persist/resume the source log position to this file. Omit to tail from the
581    /// current position without checkpointing.
582    pub checkpoint: Option<String>,
583    /// Catch up to the source's current end and exit (a bounded run), instead of
584    /// streaming indefinitely — ideal for a scheduler. For MySQL this is a
585    /// non-blocking binlog dump; PostgreSQL / SQL Server already drain-and-exit.
586    /// **Defaults to `true`** (bounded): the OSS model is scheduler-driven, and
587    /// omitting this must NOT silently start a never-terminating stream. Set it to
588    /// `false` to opt into continuous streaming explicitly.
589    #[serde(default = "default_true")]
590    pub until_current: bool,
591    /// Stop after N change events (default: until end of stream / interrupted).
592    pub max_events: Option<usize>,
593    /// Rows per output part file (default 100000). A part also rolls at a
594    /// transaction boundary, so it never splits a transaction. Larger ⇒ fewer,
595    /// bigger files but more drain memory — the PostgreSQL peek reads a part's
596    /// worth per batch, so drain RSS is O(rollover). Tune per workload: raise it
597    /// to cut file count, lower it to cap memory on a small extractor.
598    pub rollover: Option<usize>,
599    /// Roll a part once its buffered changes reach this many MB, whichever comes
600    /// first with `rollover`. Caps the in-memory buffer and the part file size by
601    /// bytes instead of a fixed row count — predictable for tables with wide
602    /// (large JSON / blob) rows, mirroring the batch path's `batch_size_memory_mb`.
603    pub rollover_memory_mb: Option<usize>,
604    /// MySQL replica server-id for the binlog connection (default 4271; must be
605    /// distinct from the source's and any other replica).
606    pub server_id: Option<u32>,
607    /// PostgreSQL logical replication slot name (default `rivet_slot`).
608    pub slot: Option<String>,
609    /// SQL Server CDC capture instance, e.g. `dbo_orders` — required for
610    /// `sqlserver://` sources.
611    pub capture_instance: Option<String>,
612}
613
614// Hand-written so the Rust `Default` MATCHES the serde default: `until_current`
615// must be `true` (bounded). The derived `Default` would use `bool::default()` =
616// `false`, and serde's `default = "default_true"` only affects Deserialize — so
617// `CdcExportConfig::default()` would silently mean `DrainMode::Continuous`. That
618// default is reached on the drain path (`cdc_job.rs`: `export.cdc.clone()
619// .unwrap_or_default()`) whenever a `mode: cdc` export omits the whole `cdc:`
620// block (valid for PG/MySQL), turning a minimal bounded drain into a
621// never-terminating daemon that persists no resume position — the exact footgun
622// `default_true` was added to kill. Mirrors `MongoConfig`'s hand-written Default.
623impl Default for CdcExportConfig {
624    fn default() -> Self {
625        Self {
626            initial: None,
627            checkpoint: None,
628            until_current: true,
629            max_events: None,
630            rollover: None,
631            rollover_memory_mb: None,
632            server_id: None,
633            slot: None,
634            capture_instance: None,
635        }
636    }
637}
638
639#[derive(Debug, Deserialize, Serialize, JsonSchema, Clone, Copy, PartialEq, Eq)]
640#[serde(rename_all = "lowercase")]
641pub enum TimeColumnType {
642    Timestamp,
643    Unix,
644}
645
646/// Calendar bucket width for date/timestamp output partitioning
647/// ([`ExportConfig::partition_by`]). The partition column must be a DATE or
648/// TIMESTAMP column; this picks how its range is split into contiguous Hive
649/// buckets. It is not a knob for partitioning by arbitrary column values.
650#[derive(Debug, Deserialize, Serialize, JsonSchema, Clone, Copy, PartialEq, Eq, Default)]
651#[serde(rename_all = "lowercase")]
652pub enum PartitionGranularity {
653    /// One bucket per calendar day (`col=2023-01-01/`). Default.
654    #[default]
655    Day,
656    /// One bucket per calendar month (`col=2023-01/`).
657    Month,
658    /// One bucket per calendar year (`col=2023/`).
659    Year,
660}
661
662/// Canonical fully-populated [`ExportConfig`] for tests across the crate.
663///
664/// One place lists every field, so adding a field is a single-site edit (the
665/// compiler still flags this literal if a field is missing). Test call sites
666/// take this baseline and override only the fields they exercise, rather than
667/// hand-writing the full struct — see `plan::build` and `preflight` tests.
668#[cfg(test)]
669pub(crate) fn sample_export(name: &str) -> ExportConfig {
670    ExportConfig {
671        name: name.into(),
672        target: None,
673        load: None,
674        verify: VerifyMode::Size,
675        query: Some("SELECT 1".into()),
676        query_file: None,
677        table: None,
678        tables: None,
679        mode: ExportMode::Full,
680        cdc: None,
681        cursor_column: None,
682        cursor_fallback_column: None,
683        incremental_cursor_mode: Default::default(),
684        chunk_column: None,
685        chunk_dense: false,
686        chunk_size: 100_000,
687        chunk_size_memory_mb: None,
688        chunk_count: None,
689        chunk_by_days: None,
690        chunk_by_key: None,
691        parallel: 1,
692        wave: None,
693        parallel_safe: None,
694        time_column: None,
695        time_column_type: TimeColumnType::Timestamp,
696        days_window: None,
697        partition_by: None,
698        partition_granularity: PartitionGranularity::Day,
699        format: FormatType::Parquet,
700        compression: CompressionType::None,
701        compression_level: None,
702        compression_profile: None,
703        skip_empty: false,
704        destination: crate::config::DestinationConfig {
705            destination_type: crate::config::DestinationType::Local,
706            path: Some("/tmp".into()),
707            ..Default::default()
708        },
709        meta_columns: MetaColumns::default(),
710        quality: None,
711        max_file_size: None,
712        chunk_checkpoint: false,
713        keyset_incremental: false,
714        chunk_max_attempts: None,
715        tuning: None,
716        source_group: None,
717        reconcile_required: false,
718        columns: Default::default(),
719        on_schema_drift: Default::default(),
720        shape_drift_warn_factor: None,
721        parquet: None,
722    }
723}
724
725#[cfg(test)]
726mod tests {
727    use super::*;
728
729    // ── ExportConfig::max_file_size_bytes ───────────────────────────────────
730
731    fn make_export_yaml(name: &str, extra: &str) -> ExportConfig {
732        let yaml = format!(
733            "name: {name}\nquery: \"SELECT 1\"\nformat: parquet\ndestination:\n  type: local\n  path: /tmp\n{extra}"
734        );
735        serde_yaml_ng::from_str(&yaml).expect("parse ExportConfig")
736    }
737
738    #[test]
739    fn max_file_size_bytes_none_when_unset() {
740        let exp = make_export_yaml("no_limit", "");
741        assert!(exp.max_file_size_bytes().is_none());
742    }
743
744    #[test]
745    fn max_file_size_bytes_parses_mb() {
746        let exp = make_export_yaml("sized", "max_file_size: \"128MB\"\n");
747        assert_eq!(exp.max_file_size_bytes(), Some(128 * 1024 * 1024));
748    }
749
750    #[test]
751    fn max_file_size_bytes_parses_gb() {
752        let exp = make_export_yaml("sized_gb", "max_file_size: \"2GB\"\n");
753        assert_eq!(exp.max_file_size_bytes(), Some(2 * 1024 * 1024 * 1024));
754    }
755
756    #[test]
757    fn max_file_size_bytes_returns_none_on_invalid() {
758        let exp = make_export_yaml("bad_size", "max_file_size: \"notanumber\"\n");
759        assert!(exp.max_file_size_bytes().is_none());
760    }
761
762    // ── ExportConfig::resolve_query ─────────────────────────────────────────
763
764    // Build a minimal ExportConfig directly, bypassing Config::from_yaml validation.
765    // This lets us test the four branches inside resolve_query itself, including
766    // the (both-set / neither-set) error paths that are normally prevented by the
767    // top-level validator.
768    fn make_export_direct(query: Option<&str>, query_file: Option<&str>) -> ExportConfig {
769        ExportConfig {
770            query: query.map(|s| s.to_string()),
771            query_file: query_file.map(|s| s.to_string()),
772            ..sample_export("test")
773        }
774    }
775
776    fn params(pairs: &[(&str, &str)]) -> std::collections::HashMap<String, String> {
777        pairs
778            .iter()
779            .map(|(k, v)| (k.to_string(), v.to_string()))
780            .collect()
781    }
782
783    #[test]
784    fn resolve_query_inline_no_params_returns_query_as_is() {
785        let exp = make_export_direct(Some("SELECT id FROM orders"), None);
786        let q = exp.resolve_query(Path::new("/tmp"), None).unwrap();
787        assert_eq!(q, "SELECT id FROM orders");
788    }
789
790    #[test]
791    fn resolve_query_inline_with_params_substitutes_vars() {
792        let exp = make_export_direct(Some("SELECT ${col} FROM ${table}"), None);
793        let p = params(&[("col", "id"), ("table", "orders")]);
794        let q = exp.resolve_query(Path::new("/tmp"), Some(&p)).unwrap();
795        assert_eq!(q, "SELECT id FROM orders");
796    }
797
798    #[test]
799    fn resolve_query_inline_params_empty_map_is_noop() {
800        let exp = make_export_direct(Some("SELECT 1"), None);
801        let p = params(&[]);
802        let q = exp.resolve_query(Path::new("/tmp"), Some(&p)).unwrap();
803        assert_eq!(q, "SELECT 1");
804    }
805
806    #[test]
807    fn resolve_query_inline_missing_var_returns_error() {
808        // SAFETY: test-only; this binary is single-threaded in the test runner context.
809        unsafe { std::env::remove_var("UNSET_RIVET_TEST_VAR") };
810        let exp = make_export_direct(Some("SELECT ${UNSET_RIVET_TEST_VAR}"), None);
811        let p = params(&[]);
812        let result = exp.resolve_query(Path::new("/tmp"), Some(&p));
813        assert!(result.is_err());
814        let msg = format!("{:#}", result.unwrap_err());
815        assert!(
816            msg.contains("UNSET_RIVET_TEST_VAR") || msg.contains("not set"),
817            "got: {msg}"
818        );
819    }
820
821    #[test]
822    fn resolve_query_file_reads_content() {
823        let dir = tempfile::TempDir::new().unwrap();
824        let sql_path = dir.path().join("query.sql");
825        std::fs::write(&sql_path, "SELECT * FROM customers").unwrap();
826        let exp = make_export_direct(None, Some("query.sql"));
827        let q = exp.resolve_query(dir.path(), None).unwrap();
828        assert_eq!(q, "SELECT * FROM customers");
829    }
830
831    #[test]
832    fn resolve_query_file_with_params_substitutes() {
833        let dir = tempfile::TempDir::new().unwrap();
834        let sql_path = dir.path().join("q.sql");
835        std::fs::write(&sql_path, "SELECT ${col} FROM ${tbl}").unwrap();
836        let exp = make_export_direct(None, Some("q.sql"));
837        let p = params(&[("col", "name"), ("tbl", "users")]);
838        let q = exp.resolve_query(dir.path(), Some(&p)).unwrap();
839        assert_eq!(q, "SELECT name FROM users");
840    }
841
842    // ── `table:` shortcut ───────────────────────────────────────────────────
843
844    #[test]
845    fn resolve_query_table_shortcut_qualified() {
846        let mut exp = make_export_direct(None, None);
847        exp.table = Some("public.users".into());
848        let q = exp.resolve_query(Path::new("/tmp"), None).unwrap();
849        assert_eq!(q, "SELECT * FROM public.users");
850    }
851
852    #[test]
853    fn resolve_query_table_shortcut_unqualified() {
854        let mut exp = make_export_direct(None, None);
855        exp.table = Some("orders".into());
856        let q = exp.resolve_query(Path::new("/tmp"), None).unwrap();
857        assert_eq!(q, "SELECT * FROM orders");
858    }
859
860    #[test]
861    fn resolve_query_table_shortcut_rejects_three_part_name() {
862        let mut exp = make_export_direct(None, None);
863        exp.table = Some("db.public.users".into());
864        let err = exp.resolve_query(Path::new("/tmp"), None).unwrap_err();
865        let msg = format!("{err:#}");
866        assert!(msg.contains("<schema>.<name>"), "got: {msg}");
867    }
868
869    #[test]
870    fn resolve_query_table_shortcut_rejects_sql_injection() {
871        for bad in [
872            "users; DROP TABLE x",
873            "users--",
874            "users'",
875            "users\"",
876            "public.\"My Table\"",
877            "0starts_with_digit",
878            "",
879            ".trailing",
880            "leading.",
881            "two..dots",
882        ] {
883            let mut exp = make_export_direct(None, None);
884            exp.table = Some(bad.into());
885            assert!(
886                exp.resolve_query(Path::new("/tmp"), None).is_err(),
887                "should reject `table:` value '{bad}'",
888            );
889        }
890    }
891
892    #[test]
893    fn resolve_query_table_shortcut_takes_precedence_over_query() {
894        let mut exp = make_export_direct(Some("SELECT id FROM x"), None);
895        exp.table = Some("public.y".into());
896        let q = exp.resolve_query(Path::new("/tmp"), None).unwrap();
897        assert_eq!(q, "SELECT * FROM public.y");
898    }
899
900    #[test]
901    fn resolve_query_file_missing_returns_error() {
902        let dir = tempfile::TempDir::new().unwrap();
903        let exp = make_export_direct(None, Some("nonexistent.sql"));
904        let result = exp.resolve_query(dir.path(), None);
905        assert!(result.is_err());
906        let msg = format!("{:#}", result.unwrap_err());
907        assert!(
908            msg.contains("nonexistent.sql") || msg.contains("No such file"),
909            "got: {msg}"
910        );
911    }
912
913    #[test]
914    fn resolve_query_both_set_returns_error() {
915        let mut exp = make_export_direct(Some("SELECT 1"), None);
916        exp.query_file = Some("file.sql".into());
917        let result = exp.resolve_query(Path::new("/tmp"), None);
918        assert!(result.is_err());
919        let msg = format!("{:#}", result.unwrap_err());
920        assert!(
921            msg.contains("not both") || msg.contains("query_file"),
922            "got: {msg}"
923        );
924    }
925
926    #[test]
927    fn resolve_query_neither_set_returns_error() {
928        let exp = make_export_direct(None, None);
929        let result = exp.resolve_query(Path::new("/tmp"), None);
930        assert!(result.is_err());
931        let msg = format!("{:#}", result.unwrap_err());
932        assert!(
933            msg.contains("query") || msg.contains("query_file"),
934            "got: {msg}"
935        );
936    }
937
938    // ── SecOps: query_file path traversal prevention ──────────────────────────
939
940    #[test]
941    fn resolve_query_file_dotdot_is_rejected() {
942        let dir = tempfile::TempDir::new().unwrap();
943        let exp = make_export_direct(None, Some("../secret.sql"));
944        let result = exp.resolve_query(dir.path(), None);
945        assert!(result.is_err());
946        let msg = format!("{:#}", result.unwrap_err());
947        assert!(
948            msg.contains("..") || msg.contains("traversal"),
949            "got: {msg}"
950        );
951    }
952
953    #[test]
954    fn resolve_query_file_nested_dotdot_is_rejected() {
955        let dir = tempfile::TempDir::new().unwrap();
956        let exp = make_export_direct(None, Some("subdir/../../etc/passwd"));
957        let result = exp.resolve_query(dir.path(), None);
958        assert!(result.is_err());
959        let msg = format!("{:#}", result.unwrap_err());
960        assert!(
961            msg.contains("..") || msg.contains("traversal"),
962            "got: {msg}"
963        );
964    }
965
966    #[test]
967    fn resolve_query_file_absolute_path_is_rejected() {
968        let dir = tempfile::TempDir::new().unwrap();
969        let exp = make_export_direct(None, Some("/etc/passwd"));
970        let result = exp.resolve_query(dir.path(), None);
971        assert!(result.is_err());
972        let msg = format!("{:#}", result.unwrap_err());
973        assert!(
974            msg.contains("relative") || msg.contains("absolute"),
975            "got: {msg}"
976        );
977    }
978
979    #[test]
980    fn resolve_query_file_in_subdir_is_allowed() {
981        let dir = tempfile::TempDir::new().unwrap();
982        let subdir = dir.path().join("queries");
983        std::fs::create_dir(&subdir).unwrap();
984        std::fs::write(subdir.join("orders.sql"), "SELECT * FROM orders").unwrap();
985        let exp = make_export_direct(None, Some("queries/orders.sql"));
986        let q = exp.resolve_query(dir.path(), None).unwrap();
987        assert_eq!(q, "SELECT * FROM orders");
988    }
989}