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