rivet-cli 0.21.1

Rivet: PostgreSQL/MySQL/SQL Server/MongoDB → Parquet/CSV (local, S3, GCS, Azure). Crate name rivet-cli; binary rivet.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
//! CDC export runner — the `mode: cdc` arm of `rivet run`.
//!
//! Runs a change stream configured from the YAML (not the `rivet cdc` CLI),
//! writes typed Parquet/CSV through the commit-seam sink, and records a
//! `RunSummary` + metric row so a CDC export appears in `rivet metrics` and the
//! run aggregate exactly like a batch export. Mirrors `run_export_job`'s contract
//! — `(Result<()>, RunSummary)`, metric recorded internally — so the orchestrator
//! treats a CDC export like any other.

use std::path::PathBuf;

use super::finalize::finalize_run_report;
use super::summary::RunSummary;
use crate::config::{Config, ExportConfig};
use crate::error::Result;
use crate::source::cdc::{CdcCapture, CdcConfig, CdcEngine, CdcEngineOpts, DrainMode, run_capture};
use crate::state::StateStore;

/// Run one `mode: cdc` export end to end, then record + report it like a batch
/// export. The metric row is written here (as `run_export_job` does); the
/// `RunSummary` is returned for the run aggregate.
/// The run-start warning for a CDC export that requested batch-only
/// `meta_columns` (`exported_at` / `row_hash`). CDC has its own sink and schema
/// (`__op`/`__pos`/`__seq` + the typed after-image), so those columns are NOT
/// injected — silently, before this. `None` when no meta column was requested.
/// Pure so it can be unit-tested without a live stream (mirrors
/// `detect::sparse_chunk_warning`).
fn cdc_ignored_meta_warning(export: &ExportConfig) -> Option<String> {
    export.meta_columns.any_enabled().then(|| {
        format!(
            "export '{}': mode: cdc ignores meta_columns (exported_at / row_hash) — the CDC \
             output carries its own __op/__pos/__seq columns plus the typed after-image, and \
             the batch meta columns are NOT added. Remove meta_columns from this export, or use \
             a batch mode if you need them.",
            export.name
        )
    })
}

pub(super) fn run_cdc_export(
    config_path: &str,
    config: &Config,
    export: &ExportConfig,
    state: &StateStore,
) -> (Result<()>, RunSummary) {
    // No-silent-config-drop: warn (don't hard-fail — a shared default may set
    // meta_columns for a mixed batch+CDC config) that the request has no effect.
    if let Some(msg) = cdc_ignored_meta_warning(export) {
        log::warn!("{msg}");
    }
    let started = std::time::Instant::now();
    let run_id = format!(
        "{}_{}",
        export.name,
        chrono::Utc::now().format("%Y%m%dT%H%M%S%3f")
    );

    let result = run_cdc_inner(config, export, &run_id);
    let duration_ms = started.elapsed().as_millis() as i64;

    let mut summary = match &result {
        // One manifest per captured table (a multi-table `tables:` stream
        // produces several); the export-level summary carries the totals.
        Ok(manifests) => {
            let bytes: u64 = manifests
                .iter()
                .flat_map(|m| &m.parts)
                .map(|p| p.size_bytes)
                .sum();
            cdc_summary(
                &run_id,
                export,
                "success",
                manifests.iter().map(|m| m.row_count).sum(),
                manifests.iter().map(|m| m.part_count as usize).sum(),
                bytes,
                duration_ms,
                None,
            )
        }
        Err(e) => cdc_summary(
            &run_id,
            export,
            "failed",
            0,
            0,
            0,
            duration_ms,
            Some(crate::redact::redact_error(e)),
        ),
    };

    // Record the run in the journal (so `rivet journal` shows a CDC run like a
    // batch one): one FileWritten per committed part, then the RunCompleted outcome.
    if let Ok(manifests) = &result {
        for (i, part) in manifests.iter().flat_map(|m| &m.parts).enumerate() {
            summary
                .journal
                .record(crate::journal::RunEvent::FileWritten {
                    file_name: part.path.clone(),
                    rows: part.rows,
                    bytes: part.size_bytes,
                    part_index: i,
                });
        }
    }
    summary
        .journal
        .record(crate::journal::RunEvent::RunCompleted {
            status: summary.status.clone(),
            error_message: summary.error_message.clone(),
            duration_ms,
        });
    if let Err(e) = state.store_journal(&summary.journal) {
        log::warn!(
            "cdc: journal persist failed for export '{}': {:#}",
            export.name,
            e
        );
    }

    record_metric(state, config, export, &summary);
    finalize_run_report(config_path, &summary, "cdc");
    (result.map(|_| ()), summary)
}

/// `cdc.initial: snapshot` — the anchor-then-snapshot half, run by the
/// orchestrator BEFORE the CDC drain. Returns the synthesized `mode: full`
/// exports still pending (their `snapshot/_SUCCESS` marker absent), after
/// ensuring the anchor exists. Anchor-BEFORE-snapshot is the whole point: a
/// change landing mid-snapshot is then also in the stream — an overlap the
/// PK+`__op` dedupe absorbs, never a gap.
pub(super) fn initial_snapshot_pending(
    config: &Config,
    export: &ExportConfig,
    state: &StateStore,
) -> Result<Vec<ExportConfig>> {
    let cdc = export.cdc.clone().unwrap_or_default();
    if cdc.initial != Some(crate::config::CdcInitialMode::Snapshot) {
        return Ok(Vec::new());
    }
    let url = config.source.resolve_url()?;
    let tls = config.source.tls.as_ref();

    let slot = cdc
        .slot
        .clone()
        .unwrap_or_else(|| crate::config::DEFAULT_PG_SLOT.to_string());

    // Each table's snapshot destination + whether its snapshot is already done.
    // Scanned BEFORE the anchor step, because a completed snapshot is resume
    // EVIDENCE: a missing server-side anchor after one must fail loud, never
    // silently re-anchor at "current" (which would skip every change since the
    // drop while reporting success — finding #28).
    let (tables, multi) = match (&export.tables, &export.table) {
        (Some(ts), _) => (ts.clone(), true),
        (None, Some(t)) => (vec![t.clone()], false),
        (None, None) => anyhow::bail!("export '{}': cdc mode requires `table:`", export.name),
    };
    let mut table_dests = Vec::with_capacity(tables.len());
    let mut done_flags = Vec::with_capacity(tables.len());
    for t in &tables {
        let table_dcfg = if multi {
            dest_for_table(&export.destination, t)
        } else {
            export.destination.clone()
        };
        let snap_dcfg = dest_for_table(&table_dcfg, "snapshot");
        let dest = crate::destination::create_destination(&snap_dcfg)?;
        // The state DB is authoritative (survives `cleanup_source` wiping the
        // bucket); the GCS `snapshot/_SUCCESS` marker stays a legacy co-signal so
        // pre-v14 runs and setups without state still skip correctly.
        let done = state.snapshot_done(&export.name, t)? || dest.head("_SUCCESS")?.is_some();
        table_dests.push((t.clone(), snap_dcfg));
        done_flags.push(done);
    }

    // The pure decision: which tables still need a snapshot, and whether prior
    // evidence forces the fail-loud anchor guard.
    let ckpt_resume = cdc
        .checkpoint
        .as_deref()
        .map(std::path::Path::new)
        .and_then(|p| crate::source::cdc::Position::load(p).ok().flatten())
        .is_some();
    let (pending_idx, resume_expected) = snapshot_plan(&done_flags, ckpt_resume);

    // The anchor — one entry point; the engine's AnchorModel decides the
    // mechanism (idempotent: a present anchor is never moved).
    CdcEngine::from_url(&url)?.ensure_anchor(
        &url,
        &slot,
        cdc.checkpoint.as_deref().map(std::path::Path::new),
        tls,
        resume_expected,
    )?;

    let mut pending = Vec::new();
    for idx in pending_idx {
        let (t, snap_dcfg) = &table_dests[idx];
        pending.push(synth_snapshot_export(export, t, snap_dcfg));
    }
    Ok(pending)
}

/// Synthesize the `mode: full` snapshot export for one CDC table — the batch
/// leg run BEFORE the CDC drain. Pure (no I/O) so the schema-consistency
/// invariants below are unit-testable.
fn synth_snapshot_export(
    export: &ExportConfig,
    table: &str,
    snap_dcfg: &crate::config::DestinationConfig,
) -> ExportConfig {
    let mut synth = export.clone();
    synth.name = format!("{}__snapshot_{table}", export.name);
    synth.mode = crate::config::ExportMode::Full;
    synth.table = Some(table.to_string());
    synth.tables = None;
    synth.cdc = None;
    synth.destination = snap_dcfg.clone();
    // NEVER inherit skip_empty: an EMPTY table with skip_empty=true would
    // write no snapshot/_SUCCESS, so the marker check re-snapshots on
    // every run forever. An empty snapshot must still complete (manifest +
    // _SUCCESS with 0 rows) for the handoff to converge.
    synth.skip_empty = false;
    // NEVER inherit batch meta_columns into the snapshot leg. The snapshot AND
    // the CDC stream BOTH load into the SAME `<table>__changes` log (the snapshot
    // rows with __op/__pos/__seq NULL — see load/cdc.rs::dedup_view_sql), and the
    // current-state view is `SELECT * EXCEPT(__op,__pos,__seq,__rn)` over it. The
    // CDC sink cannot carry exported_at/row_hash, so injecting them on the
    // snapshot ONLY gives the snapshot parquet extra columns the CDC parquet
    // lacks — appending both into one `__changes` table then mismatches, and any
    // meta column that survived would leak into the current-state view (populated
    // for backfill rows, NULL for every CDC-updated row). Clearing here keeps both
    // legs' columns identical; the run-start warn tells the operator the meta
    // columns are dropped for the whole CDC export.
    synth.meta_columns = Default::default();
    synth
}

/// The pure `initial: snapshot` decision, split out of the I/O in
/// [`initial_snapshot_pending`] so it can be unit-tested. Given, per table in
/// order, whether its snapshot is already `done` (state DB OR the legacy GCS
/// marker) and whether a checkpoint position survives (`ckpt_resume`), returns
/// the indices still PENDING a snapshot and whether the anchor step must treat a
/// missing server-side anchor as resume evidence.
///
/// A `done` snapshot is never re-run — the state DB remembers it even after
/// `cleanup_source` wiped the bucket marker. `resume_expected` is `true` when
/// ANY prior evidence exists — a live checkpoint OR any done snapshot — so a
/// lost server-side anchor fails LOUD instead of silently re-anchoring at
/// "current" (finding #28).
fn snapshot_plan(done_flags: &[bool], ckpt_resume: bool) -> (Vec<usize>, bool) {
    let pending = done_flags
        .iter()
        .enumerate()
        .filter_map(|(i, &done)| (!done).then_some(i))
        .collect();
    let resume_expected = ckpt_resume || done_flags.iter().any(|&d| d);
    (pending, resume_expected)
}

/// A multi-table stream lands each table under its own sub-prefix of the
/// export's destination (`<base>/<table>/`), so every table's prefix is
/// self-describing (its own parts + `manifest.json` + `_SUCCESS`), exactly like
/// N single-table exports — minus the N−1 extra slots/connections.
pub(crate) fn dest_for_table(
    base: &crate::config::DestinationConfig,
    table: &str,
) -> crate::config::DestinationConfig {
    let mut d = base.clone();
    match d.destination_type {
        crate::config::DestinationType::Local => {
            let p = d.path.take().unwrap_or_else(|| ".".into());
            d.path = Some(format!("{}/{}", p.trim_end_matches('/'), table));
        }
        crate::config::DestinationType::Stdout => {}
        // Cloud destinations: extend the key prefix. Cloud prefixes are
        // LITERAL key prefixes — the destination concatenates `prefix + key`
        // with no separator (the docs' `prefix: exports/` convention) — so the
        // sub-prefix must supply both its slashes itself, or every object
        // lands as a mangled flat key (`<prefix>/<table>cdc-….parquet`).
        _ => {
            let pfx = d.prefix.take().unwrap_or_default();
            let base = pfx.trim_end_matches('/');
            d.prefix = Some(if base.is_empty() {
                format!("{table}/")
            } else {
                format!("{base}/{table}/")
            });
        }
    }
    d
}

/// Build the capture from the config + export and drive it through the shared
/// [`crate::source::cdc::run_capture`] (the same assembler the `rivet cdc` CLI
/// uses). Returns the manifests (one per captured table) so the caller records
/// the metric + journal.
fn run_cdc_inner(
    config: &Config,
    export: &ExportConfig,
    run_id: &str,
) -> Result<Vec<crate::manifest::RunManifest>> {
    let url = config.source.resolve_url()?;
    let cdc = export.cdc.clone().unwrap_or_default();
    // `tables:` (multi-table, one stream) or the single `table:` — validation
    // guarantees exactly one of them is set.
    let (tables, multi) = match (&export.tables, &export.table) {
        (Some(ts), _) => (ts.clone(), true),
        (None, Some(t)) => (vec![t.clone()], false),
        (None, None) => {
            anyhow::bail!("export '{}': cdc mode requires `table:`", export.name)
        }
    };
    // Per-table destinations must outlive the borrowed CaptureOutputs.
    let mut wired: Vec<(String, Box<dyn crate::destination::Destination>, String)> =
        Vec::with_capacity(tables.len());
    for t in &tables {
        let dcfg = if multi {
            dest_for_table(&export.destination, t)
        } else {
            export.destination.clone()
        };
        let dest = crate::destination::create_destination(&dcfg)?;
        // Finding #44, early check: refuse BEFORE the first part lands if the
        // prefix belongs to the other pipeline shape (config error, fail the
        // run cleanly — the write-seam guard stays as the backstop).
        crate::manifest::guard_manifest_mode(dest.as_ref(), "cdc")?;
        let uri = dcfg
            .path
            .clone()
            .or_else(|| dcfg.prefix.clone())
            .unwrap_or_default();
        wired.push((t.clone(), dest, uri));
    }
    // `columns:` type overrides, narrowed per table: bare keys apply to every
    // captured table; qualified keys ("table.column") only to theirs, winning
    // over bare — so one table's override can never bleed into a same-named
    // column elsewhere.
    let all_overrides =
        crate::plan::build::parse_column_overrides_pub(&export.columns, &export.name)?;
    let outputs = wired
        .iter()
        .map(|(t, d, u)| crate::source::cdc::CaptureOutput {
            table: t.clone(),
            dest: d.as_ref(),
            dest_uri: u.clone(),
            overrides: crate::types::overrides_for_table(
                &all_overrides,
                t.rsplit('.').next().unwrap_or(t),
            ),
        })
        .collect();
    let now = chrono::Utc::now().to_rfc3339();

    // `until_current` defaults to `true` (bounded, scheduler-friendly). An explicit
    // `false` opts into a long-lived continuous stream — surface it so it is a
    // deliberate choice, never a silent never-terminating run.
    if !cdc.until_current {
        log::warn!(
            "cdc: `until_current: false` runs a CONTINUOUS stream (a long-lived daemon that never \
             exits on its own). For the scheduler model, omit it or set `until_current: true` (the \
             default) to drain to the log end and exit."
        );
    }

    run_capture(CdcCapture {
        cdc_cfg: CdcConfig {
            url,
            checkpoint: cdc.checkpoint.as_ref().map(PathBuf::from),
            drain: DrainMode::from_until_current(cdc.until_current),
            tls: config.source.tls.clone(),
            engine: match config.source.source_type {
                crate::config::SourceType::Mysql => CdcEngineOpts::Mysql {
                    server_id: cdc
                        .server_id
                        .unwrap_or(crate::config::DEFAULT_MYSQL_SERVER_ID),
                },
                crate::config::SourceType::Postgres => CdcEngineOpts::Postgres {
                    slot: cdc
                        .slot
                        .clone()
                        .unwrap_or_else(|| crate::config::DEFAULT_PG_SLOT.to_string()),
                },
                crate::config::SourceType::Mssql => CdcEngineOpts::Mssql {
                    capture_instance: cdc.capture_instance.clone(),
                },
                crate::config::SourceType::Mongo => {
                    CdcEngineOpts::Mongo {
                        canonical: config.source.mongo.as_ref().is_some_and(|m| {
                            matches!(m.json, crate::config::MongoJsonMode::Canonical)
                        }),
                    }
                }
            },
        },
        outputs,
        format: export.format,
        max_events: cdc.max_events,
        rollover: cdc.rollover.unwrap_or(100_000),
        rollover_memory_bytes: cdc.rollover_memory_mb.map(|mb| mb * 1024 * 1024),
        run_id: run_id.to_string(),
        started_at: now,
    })
}

/// Build the per-run summary (mirrors `synthetic_failed_summary`'s shape, for a
/// CDC run). CDC has no plan/tuning/cursor, so those fields default.
#[allow(clippy::too_many_arguments)]
fn cdc_summary(
    run_id: &str,
    export: &ExportConfig,
    status: &str,
    total_rows: i64,
    files: usize,
    bytes: u64,
    duration_ms: i64,
    error_message: Option<String>,
) -> RunSummary {
    // Only the fields a CDC run has; the batch-specific rest (cursor, quality,
    // chunk, reconcile, …) stay at RunSummary::default() (None / 0 / empty).
    RunSummary {
        cursor_column: None,
        cursor_low: None,
        cursor_high: None,
        run_id: run_id.to_string(),
        export_name: export.name.clone(),
        status: status.to_string(),
        total_rows,
        files_produced: files,
        bytes_written: bytes,
        files_committed: files,
        duration_ms,
        error_message,
        tuning_profile: "cdc".into(),
        format: export.format.label().to_string(),
        mode: "cdc".into(),
        destination_uri: export.destination.path.clone(),
        journal: crate::journal::RunJournal::new(run_id, &export.name),
        ..Default::default()
    }
}

/// Write the `export_metrics` row (built directly — no plan coupling, unlike the
/// batch `build_metric_row`) so a CDC run is queryable like a batch run.
fn record_metric(state: &StateStore, config: &Config, export: &ExportConfig, summary: &RunSummary) {
    let source_type = config
        .source
        .resolve_url()
        .ok()
        .and_then(|u| CdcEngine::from_url(&u).ok().map(CdcEngine::label))
        .map(|s| s.to_string());
    // Only the fields a CDC run actually has; the batch-specific rest (chunk_size,
    // cursor, quality, …) stay at their MetricRow::default() (None / 0).
    let row = crate::state::MetricRow {
        export_name: summary.export_name.clone(),
        run_id: summary.run_id.clone(),
        duration_ms: summary.duration_ms,
        total_rows: summary.total_rows,
        peak_rss_mb: Some(summary.peak_rss_mb),
        status: summary.status.clone(),
        error_message: summary.error_message.clone(),
        tuning_profile: Some("cdc".to_string()),
        format: Some(summary.format.clone()),
        mode: Some("cdc".to_string()),
        files_produced: summary.files_produced as i64,
        bytes_written: summary.bytes_written as i64,
        files_committed: summary.files_committed as i64,
        source_type,
        destination_type: Some(export.destination.destination_type.label().to_string()),
        rivet_version: Some(env!("CARGO_PKG_VERSION").to_string()),
        longest_chunk_ms: summary.journal.longest_chunk_ms(),
        ..Default::default()
    };
    if let Err(e) = state.record_metric_full(&row) {
        log::warn!(
            "cdc: failed to record metric for export '{}': {:#}",
            export.name,
            e
        );
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::config::{DestinationConfig, DestinationType};

    // A CDC export that requests batch-only meta_columns must WARN (the CDC sink
    // never injects them), not silently drop the request. Pure fn so this needs
    // no live stream.
    #[test]
    fn cdc_warns_when_meta_columns_are_requested_on_a_cdc_export() {
        let mut e = crate::config::sample_export("orders");
        // No meta columns → no warning.
        e.meta_columns = Default::default();
        assert!(
            cdc_ignored_meta_warning(&e).is_none(),
            "no warning without meta_columns"
        );
        // exported_at requested → warns, naming the export + the CDC-native cols.
        e.meta_columns.exported_at = true;
        let msg = cdc_ignored_meta_warning(&e).expect("must warn on exported_at");
        assert!(msg.contains("orders"), "names the export: {msg}");
        assert!(
            msg.contains("meta_columns"),
            "names the ignored knob: {msg}"
        );
        assert!(
            msg.contains("__op"),
            "points at the CDC-native columns: {msg}"
        );
        // row_hash alone also warns.
        e.meta_columns.exported_at = false;
        e.meta_columns.row_hash = true;
        assert!(
            cdc_ignored_meta_warning(&e).is_some(),
            "row_hash alone must warn too"
        );
    }

    // The snapshot (batch) leg and the CDC stream are two legs of ONE dataset the
    // load view merges by PK. batch-only meta_columns injected on the snapshot
    // ONLY would diverge the two legs' columns (base table has cols the changelog
    // lacks) and break the load — the exact mismatch a mixed snapshot+CDC load
    // hits. The synth snapshot must therefore NOT inherit meta_columns.
    #[test]
    fn snapshot_leg_does_not_inherit_batch_meta_columns() {
        let mut e = crate::config::sample_export("orders");
        e.meta_columns.exported_at = true;
        e.meta_columns.row_hash = true;
        let dcfg = DestinationConfig {
            destination_type: DestinationType::Local,
            path: Some("/tmp/snap".into()),
            ..Default::default()
        };
        let synth = synth_snapshot_export(&e, "orders", &dcfg);
        assert!(
            !synth.meta_columns.any_enabled(),
            "the snapshot leg must match the CDC leg's columns — no batch meta_columns"
        );
        // The other snapshot invariants stay intact.
        assert_eq!(synth.mode, crate::config::ExportMode::Full);
        assert!(
            synth.cdc.is_none(),
            "snapshot leg is a plain batch, not CDC"
        );
        assert!(!synth.skip_empty, "snapshot must complete even when empty");
        assert_eq!(synth.table.as_deref(), Some("orders"));
    }

    // RED test for the finding: cloud prefixes are LITERAL key prefixes —
    // `cloud.rs` concatenates `prefix + key` with NO separator (hence the
    // docs' `prefix: exports/` convention). The multi-table sub-prefix must
    // therefore supply both slashes itself; without the trailing one, every
    // object of a `tables:` export lands as `<prefix>/<table>cdc-….parquet`
    // (mangled flat keys) — observed live on a real GCS bucket.
    #[test]
    fn cloud_table_sub_prefix_carries_its_own_trailing_slash() {
        let gcs = DestinationConfig {
            destination_type: DestinationType::Gcs,
            bucket: Some("b".into()),
            prefix: Some("exports".into()),
            ..Default::default()
        };
        let d = dest_for_table(&gcs, "orders");
        let prefix = d.prefix.unwrap();
        // The exact join the cloud destination performs:
        let key = format!("{prefix}{}", "cdc-r-000000.parquet");
        assert_eq!(
            key, "exports/orders/cdc-r-000000.parquet",
            "cloud keys are prefix ++ name — the sub-prefix must end with '/'"
        );

        // A trailing slash on the configured prefix must not double up.
        let gcs2 = DestinationConfig {
            prefix: Some("exports/".into()),
            ..gcs.clone()
        };
        assert_eq!(
            dest_for_table(&gcs2, "orders").prefix.unwrap(),
            "exports/orders/"
        );

        // No configured prefix: the table becomes the whole prefix.
        let gcs3 = DestinationConfig {
            prefix: None,
            ..gcs
        };
        assert_eq!(dest_for_table(&gcs3, "orders").prefix.unwrap(), "orders/");
    }

    #[test]
    fn local_table_sub_path_is_a_plain_directory_join() {
        let local = DestinationConfig {
            destination_type: DestinationType::Local,
            path: Some("/data/cdc/".into()),
            ..Default::default()
        };
        assert_eq!(
            dest_for_table(&local, "orders").path.unwrap(),
            "/data/cdc/orders",
            "local paths go through the filesystem join — no trailing slash needed"
        );
    }

    // ── snapshot_plan: the pure `initial: snapshot` decision ─────────────────

    #[test]
    fn snapshot_plan_first_run_snapshots_all_with_no_resume_evidence() {
        // Nothing done, no checkpoint → snapshot every table, and this is a
        // genuine first anchor (resume_expected=false).
        assert_eq!(snapshot_plan(&[false, false], false), (vec![0, 1], false));
    }

    #[test]
    fn snapshot_plan_all_done_snapshots_nothing_but_keeps_resume_evidence() {
        // Every snapshot already done — the state DB remembers even after
        // `cleanup_source` wiped the bucket marker → re-snapshot NOTHING; and
        // that prior evidence forces the fail-loud anchor guard (#28).
        assert_eq!(
            snapshot_plan(&[true, true], false),
            (Vec::<usize>::new(), true)
        );
    }

    #[test]
    fn snapshot_plan_partial_snapshots_only_the_undone() {
        // One table done, one not → snapshot only the undone; a done sibling is
        // still resume evidence.
        assert_eq!(snapshot_plan(&[true, false], false), (vec![1], true));
    }

    #[test]
    fn snapshot_plan_checkpoint_alone_is_resume_evidence() {
        // No snapshot done but a live checkpoint survives → still snapshot (the
        // marker is gone), yet the checkpoint alone makes a lost anchor fail loud.
        assert_eq!(snapshot_plan(&[false], true), (vec![0], true));
    }

    #[test]
    fn snapshot_plan_no_evidence_is_a_legitimate_first_anchor() {
        // Nothing done, no checkpoint → snapshot, and no evidence means the
        // anchor is a legitimate first anchor, not a loud failure.
        assert_eq!(snapshot_plan(&[false], false), (vec![0], false));
    }
}