1use crate::destination::gcs::GcsStore;
27use crate::manifest::{MANIFEST_FILENAME, ManifestStatus, RunManifest};
28use crate::pipeline::validate_manifest::MANIFEST_MAX_BYTES;
29use anyhow::{Context, Result, bail};
30
31#[derive(Debug, Clone, PartialEq, Eq)]
35pub struct LoadIntegrity {
36 pub source_rows: Option<u64>,
41 pub file_rows: u64,
44 pub manifests: usize,
46}
47
48impl LoadIntegrity {
49 pub fn chain_prefix(&self) -> String {
53 let src = self
54 .source_rows
55 .map_or_else(|| "?".to_string(), |n| n.to_string());
56 format!("source {src} → files {}", self.file_rows)
57 }
58}
59
60#[allow(private_interfaces)]
77pub fn fetch_manifests_keyed(
78 store: &GcsStore,
79 gcs_prefix: &str,
80) -> Result<Vec<(String, RunManifest)>> {
81 let (_, base) = crate::load::split_gs_uri(gcs_prefix)?;
82 let keys = list_manifest_keys(store, base)?;
83 keys.into_iter()
84 .map(|key| {
85 let sz = store.stat_size(&key)?;
90 if sz > MANIFEST_MAX_BYTES {
91 bail!(
92 "manifest {key} is {sz} bytes, over the {MANIFEST_MAX_BYTES}-byte cap — \
93 refusing to read a possibly-hostile manifest into memory (CWE-400)"
94 );
95 }
96 let bytes = store.read(&key)?;
97 let m = serde_json::from_slice::<RunManifest>(&bytes)
98 .with_context(|| format!("parsing manifest {key}"))?;
99 Ok((key, m))
100 })
101 .collect()
102}
103
104#[allow(private_interfaces)]
108pub fn select_load_uris(
109 store: &GcsStore,
110 gcs_prefix: &str,
111 new: &[(String, RunManifest)],
112) -> Result<Vec<String>> {
113 let (bucket, base) = crate::load::split_gs_uri(gcs_prefix)?;
114 let all_parquet: Vec<String> = store
115 .list_files(base)?
116 .into_iter()
117 .filter(|k| k.ends_with(".parquet"))
118 .collect();
119 Ok(select_load_keys(new, &all_parquet)
120 .into_iter()
121 .map(|k| format!("gs://{bucket}/{k}"))
122 .collect())
123}
124
125fn resolve_parts<'a>(
131 manifest_key: &'a str,
132 m: &'a RunManifest,
133) -> impl Iterator<Item = String> + 'a {
134 let dir = manifest_key.rsplit_once('/').map(|(d, _)| d).unwrap_or("");
135 m.parts.iter().map(move |p| {
136 if dir.is_empty() {
137 p.path.clone()
138 } else {
139 format!("{dir}/{}", p.path)
140 }
141 })
142}
143
144pub fn select_load_keys(new: &[(String, RunManifest)], all_parquet: &[String]) -> Vec<String> {
155 use std::collections::BTreeSet;
156 let present: std::collections::HashSet<&str> = all_parquet.iter().map(String::as_str).collect();
157 let mut selected: BTreeSet<String> = BTreeSet::new();
158 for (key, m) in new {
159 let mut resolved_any = false;
160 for full in resolve_parts(key, m) {
161 if present.contains(full.as_str()) {
162 selected.insert(full);
163 resolved_any = true;
164 }
165 }
166 if !resolved_any {
167 return all_parquet.to_vec();
169 }
170 }
171 selected.into_iter().collect()
172}
173
174#[allow(private_interfaces)]
199pub fn gc_orphans(
200 store: &GcsStore,
201 gcs_prefix: &str,
202 keyed: &[(String, RunManifest)],
203 active: bool,
204) -> Result<(usize, u64)> {
205 let (_bucket, base) = crate::load::split_gs_uri(gcs_prefix)?;
206 let mut keep: std::collections::HashSet<String> = std::collections::HashSet::new();
210 let mut terminal: std::collections::HashSet<String> = std::collections::HashSet::new();
211 for (key, m) in keyed {
212 match m.status {
213 ManifestStatus::Success => keep.extend(resolve_parts(key, m)),
214 ManifestStatus::Running => {}
219 ManifestStatus::Failed | ManifestStatus::Interrupted => {
220 terminal.extend(resolve_parts(key, m))
221 }
222 }
223 }
224 let mut removed = 0usize;
225 let mut removed_bytes = 0u64;
226 for key in store.list_files(base)? {
227 if !key.ends_with(".parquet") || keep.contains(&key) {
228 continue;
229 }
230 if active && !terminal.contains(&key) {
234 log::warn!(
235 "gc_orphans: sparing unmanifested `{key}` — a run is active on this prefix \
236 (run-status ledger); it is GC'd once no run is active or it gets a manifest"
237 );
238 continue;
239 }
240 removed_bytes += store.stat_size(&key).unwrap_or(0);
241 store.remove(&key)?;
242 removed += 1;
243 }
244 for (key, m) in keyed {
252 if m.status == ManifestStatus::Running && is_superseded(m, keyed) {
253 removed_bytes += store.stat_size(key).unwrap_or(0);
254 store.remove(key)?;
255 removed += 1;
256 }
257 }
258 Ok((removed, removed_bytes))
259}
260
261pub fn has_active_running_manifest(keyed: &[(String, RunManifest)]) -> bool {
271 keyed
272 .iter()
273 .any(|(_, m)| m.status == ManifestStatus::Running && !is_superseded(m, keyed))
274}
275
276fn is_superseded(m: &RunManifest, keyed: &[(String, RunManifest)]) -> bool {
283 keyed
284 .iter()
285 .any(|(_, o)| o.export_name == m.export_name && o.started_at > m.started_at)
286}
287
288pub fn latest_full(keyed: Vec<(String, RunManifest)>) -> Vec<(String, RunManifest)> {
301 keyed
302 .into_iter()
303 .max_by(|a, b| {
304 match (
305 chrono::DateTime::parse_from_rfc3339(&a.1.finished_at).ok(),
306 chrono::DateTime::parse_from_rfc3339(&b.1.finished_at).ok(),
307 ) {
308 (Some(x), Some(y)) => x.cmp(&y),
309 _ => a.1.finished_at.cmp(&b.1.finished_at),
310 }
311 })
312 .into_iter()
313 .collect()
314}
315
316pub fn select_runs(
326 keyed: Vec<(String, RunManifest)>,
327 loaded: &std::collections::HashSet<String>,
328 mode: crate::load::plan::LoadMode,
329) -> Vec<(String, RunManifest)> {
330 let keyed: Vec<(String, RunManifest)> = keyed
335 .into_iter()
336 .filter(|(_, m)| m.status != ManifestStatus::Running)
337 .collect();
338 match mode {
339 crate::load::plan::LoadMode::Full => latest_full(keyed),
340 _ => keyed
341 .into_iter()
342 .filter(|(_, m)| !loaded.contains(&m.run_id))
343 .collect(),
344 }
345}
346
347fn list_manifest_keys(store: &GcsStore, base: &str) -> Result<Vec<String>> {
349 let all: Vec<String> = store
350 .list_files(base)?
351 .into_iter()
352 .filter(|k| is_manifest_key(k))
353 .collect();
354 let run_unique: Vec<String> = all
362 .iter()
363 .filter(|k| is_run_unique_manifest(k.rsplit('/').next().unwrap_or("")))
364 .cloned()
365 .collect();
366 Ok(if run_unique.is_empty() {
367 all
368 } else {
369 run_unique
370 })
371}
372
373fn is_manifest_key(key: &str) -> bool {
378 let base = key.rsplit('/').next().unwrap_or("");
379 base == MANIFEST_FILENAME || is_run_unique_manifest(base)
380}
381
382fn is_run_unique_manifest(base: &str) -> bool {
385 crate::manifest::is_run_unique_manifest_name(base)
387}
388
389pub fn ensure_single_export(keyed: &[(String, RunManifest)]) -> Result<()> {
406 let mut names: std::collections::BTreeSet<&str> = keyed
407 .iter()
408 .map(|(_, m)| crate::manifest::snapshot_family(m.export_name.as_str()))
409 .filter(|n| !n.is_empty())
410 .collect();
411 if names.len() > 1 {
412 let listed: Vec<&str> = std::mem::take(&mut names).into_iter().collect();
413 bail!(
414 "the load prefix holds manifests from {} distinct exports ({}) — the load sums every \
415 manifest under the prefix and cleanup wipes it recursively, so loading a shared \
416 prefix would cross-contaminate the row count and could delete a sibling export's \
417 parts. Give each export a DISTINCT destination prefix (or scope the load prefix to a \
418 single export) and re-run.",
419 listed.len(),
420 listed.join(", ")
421 );
422 }
423 Ok(())
424}
425
426pub fn reconcile(manifests: &[RunManifest], allow_source_drift: bool) -> Result<LoadIntegrity> {
441 if manifests.is_empty() {
442 bail!(
443 "no `{MANIFEST_FILENAME}` found under the export prefix — refusing to load \
444 unverified files. A rivet export writes a manifest on success; its absence \
445 means the run never completed (or points at the wrong prefix)."
446 );
447 }
448
449 let mut file_rows: u64 = 0;
450 let mut source_rows: u64 = 0;
451 let mut any_source = false;
452
453 for m in manifests {
454 if m.status != ManifestStatus::Success {
456 bail!(
457 "manifest for run `{}` (export `{}`) is {:?}, not Success — refusing to load a \
458 partial export",
459 m.run_id,
460 m.export_name,
461 m.status
462 );
463 }
464 m.validate_self_consistency().map_err(|e| {
467 anyhow::anyhow!(
468 "manifest for run `{}` (export `{}`) is internally inconsistent: {e} — refusing \
469 to load",
470 m.run_id,
471 m.export_name
472 )
473 })?;
474
475 let rows = u64::try_from(m.row_count).with_context(|| {
476 format!(
477 "manifest for run `{}` has a negative row_count ({})",
478 m.run_id, m.row_count
479 )
480 })?;
481 file_rows += rows;
482
483 if let Some(src) = m
486 .source
487 .extraction
488 .as_ref()
489 .and_then(|x| x.source_row_count)
490 {
491 let src = u64::try_from(src).with_context(|| {
492 format!(
493 "manifest for run `{}` has a negative source_row_count ({src})",
494 m.run_id
495 )
496 })?;
497 any_source = true;
498 source_rows += src;
499 if src != rows {
500 if allow_source_drift {
501 eprintln!(
502 "warning: source→file drift for run `{}` (export `{}`): source had {src} \
503 rows, extracted {rows} (--allow-source-drift)",
504 m.run_id, m.export_name
505 );
506 } else {
507 bail!(
508 "source→file mismatch for run `{}` (export `{}`): source had {src} rows \
509 but {rows} were extracted — the extract dropped {} row(s). Investigate \
510 before loading, or pass --allow-source-drift to override.",
511 m.run_id,
512 m.export_name,
513 src.abs_diff(rows)
514 );
515 }
516 }
517 }
518 }
519
520 Ok(LoadIntegrity {
521 source_rows: any_source.then_some(source_rows),
522 file_rows,
523 manifests: manifests.len(),
524 })
525}
526
527#[cfg(test)]
528mod tests {
529 use super::*;
530 use crate::manifest::{
531 ExtractionMetadata, ManifestDestination, ManifestPart, ManifestSource, PartStatus,
532 };
533
534 fn manifest(run: &str, rows: i64, source: Option<i64>) -> RunManifest {
537 RunManifest {
538 row_hash: None,
539 manifest_version: crate::manifest::MANIFEST_VERSION,
540 run_id: run.into(),
541 export_name: "orders".into(),
542 mode: "batch".into(),
543 started_at: "t".into(),
544 finished_at: "t".into(),
545 status: ManifestStatus::Success,
546 source: ManifestSource {
547 engine: "pg".into(),
548 schema: None,
549 table: None,
550 extraction: source.map(|n| ExtractionMetadata {
551 strategy: "full".into(),
552 cursor_column: None,
553 cursor_type: None,
554 cursor_low: None,
555 cursor_high: None,
556 source_row_count: Some(n),
557 }),
558 },
559 destination: ManifestDestination {
560 kind: "gcs".into(),
561 uri: "gs://b/p".into(),
562 },
563 format: "parquet".into(),
564 compression: "zstd".into(),
565 schema_fingerprint: "xxh3:0".into(),
566 row_count: rows,
567 part_count: 1,
568 parts: vec![ManifestPart {
569 part_id: 0,
570 path: "part-000000.parquet".into(),
571 rows,
572 size_bytes: 1,
573 content_fingerprint: "xxh3:0".into(),
574 content_md5: String::new(),
575 status: PartStatus::Committed,
576 }],
577 column_checksums: None,
578 checksum_key_column: None,
579 }
580 }
581
582 #[test]
583 fn sums_file_and_source_rows_across_manifests() {
584 let ms = vec![manifest("r1", 100, Some(100)), manifest("r2", 40, Some(40))];
585 let got = reconcile(&ms, false).unwrap();
586 assert_eq!(got.file_rows, 140);
587 assert_eq!(got.source_rows, Some(140));
588 assert_eq!(got.manifests, 2);
589 }
590
591 #[test]
592 fn ensure_single_export_refuses_a_prefix_shared_by_two_exports() {
593 let keyed = |m: RunManifest| ("gs://b/p/manifest-x.json".to_string(), m);
594 let orders = keyed(manifest("r1", 10, None)); let mut cust_m = manifest("r2", 10, None);
596 cust_m.export_name = "customers".into();
597 assert!(ensure_single_export(&[orders.clone(), keyed(cust_m)]).is_err());
600 assert!(ensure_single_export(&[orders.clone(), keyed(manifest("r3", 5, None))]).is_ok());
602 let mut legacy = manifest("r0", 3, None);
605 legacy.export_name = String::new();
606 assert!(ensure_single_export(&[orders, keyed(legacy)]).is_ok());
607 }
608
609 #[test]
615 fn ensure_single_export_admits_the_cdc_snapshot_leg_of_the_same_export() {
616 let keyed = |m: RunManifest| ("gs://b/p/manifest-x.json".to_string(), m);
617 let drain = keyed(manifest("r1", 10, None)); let mut snap = manifest("r2", 10, None);
619 snap.export_name = "orders__snapshot_orders".into();
620 assert!(ensure_single_export(&[drain.clone(), keyed(snap)]).is_ok());
622 let mut foreign = manifest("r3", 10, None);
624 foreign.export_name = "customers__snapshot_customers".into();
625 assert!(ensure_single_export(&[drain, keyed(foreign)]).is_err());
626 }
627
628 #[test]
629 fn source_rows_is_none_when_no_manifest_probed_the_source() {
630 let ms = vec![manifest("r1", 100, None), manifest("r2", 40, None)];
631 let got = reconcile(&ms, false).unwrap();
632 assert_eq!(got.file_rows, 140);
633 assert_eq!(
634 got.source_rows, None,
635 "unprobed source is unknown, not zero"
636 );
637 }
638
639 #[test]
640 fn source_rows_present_even_if_only_some_manifests_probed() {
641 let ms = vec![manifest("r1", 100, Some(100)), manifest("r2", 40, None)];
642 let got = reconcile(&ms, false).unwrap();
643 assert_eq!(got.source_rows, Some(100));
644 }
645
646 #[test]
647 fn empty_manifests_refuses_to_load() {
648 let err = reconcile(&[], false).unwrap_err().to_string();
649 assert!(err.contains("refusing to load"), "{err}");
650 }
651
652 #[test]
653 fn non_success_manifest_refuses_to_load() {
654 let mut m = manifest("r1", 100, Some(100));
655 m.status = ManifestStatus::Interrupted;
656 let err = reconcile(&[m], false).unwrap_err().to_string();
657 assert!(err.contains("not Success"), "{err}");
658 }
659
660 #[test]
661 fn self_inconsistent_manifest_refuses_to_load() {
662 let mut m = manifest("r1", 100, Some(100));
665 m.row_count = 999; let err = reconcile(&[m], false).unwrap_err().to_string();
667 assert!(err.contains("inconsistent"), "{err}");
668 }
669
670 #[test]
671 fn source_file_mismatch_hard_fails_by_default() {
672 let m = manifest("r1", 100, Some(120));
674 let err = reconcile(&[m], false).unwrap_err().to_string();
675 assert!(err.contains("source→file mismatch"), "{err}");
676 assert!(err.contains("dropped 20"), "{err}");
677 }
678
679 #[test]
680 fn source_file_mismatch_is_allowed_under_the_override() {
681 let m = manifest("r1", 100, Some(120));
682 let got = reconcile(&[m], true).expect("--allow-source-drift proceeds");
683 assert_eq!(got.file_rows, 100);
684 assert_eq!(
685 got.source_rows,
686 Some(120),
687 "the probed source count is still surfaced"
688 );
689 }
690
691 fn keyed(key: &str, run: &str, part: &str) -> (String, RunManifest) {
693 let mut m = manifest(run, 10, Some(10));
694 m.parts[0].path = part.into();
695 (key.to_string(), m)
696 }
697
698 #[test]
699 fn select_load_keys_picks_only_the_new_runs_parts() {
700 let all = vec![
702 "base/r1-000.parquet".to_string(),
703 "base/r2-000.parquet".to_string(),
704 ];
705 let new = vec![keyed("base/manifest-r2.json", "r2", "r2-000.parquet")];
706 assert_eq!(
707 select_load_keys(&new, &all),
708 vec!["base/r2-000.parquet".to_string()],
709 "loads r2's part only — not r1's already-loaded file"
710 );
711 }
712
713 #[test]
714 fn select_load_uris_lists_and_re_prefixes_selected_keys_with_the_source_bucket() {
715 let (store, _g) = fs_store(&[
721 ("base/r1-000.parquet", b"a".to_vec()),
722 ("base/manifest-r1.json", b"{}".to_vec()), ]);
724 let new = vec![keyed("base/manifest-r1.json", "r1", "r1-000.parquet")];
725 assert_eq!(
726 select_load_uris(&store, "gs://my-bucket/base", &new).unwrap(),
727 vec!["gs://my-bucket/base/r1-000.parquet".to_string()],
728 "the selected key, re-prefixed with the source bucket — never the manifest"
729 );
730 }
731
732 #[test]
733 fn select_load_keys_resolves_a_snapshot_subprefix_manifest() {
734 let all = vec!["base/snapshot/snap-000.parquet".to_string()];
737 let new = vec![keyed(
738 "base/snapshot/manifest-r1.json",
739 "r1",
740 "snap-000.parquet",
741 )];
742 assert_eq!(
743 select_load_keys(&new, &all),
744 vec!["base/snapshot/snap-000.parquet".to_string()]
745 );
746 }
747
748 #[test]
749 fn select_load_keys_falls_back_to_full_listing_when_a_manifest_has_no_present_part() {
750 let all = vec!["base/a.parquet".to_string(), "base/b.parquet".to_string()];
753 let new = vec![keyed("base/manifest-r1.json", "r1", "missing.parquet")];
754 assert_eq!(
755 select_load_keys(&new, &all),
756 all,
757 "unresolvable part → blanket fallback"
758 );
759 }
760
761 #[test]
762 fn select_load_keys_empty_new_set_selects_nothing_never_the_full_listing() {
763 let all = vec![
768 "base/r1-000.parquet".to_string(),
769 "base/r2-000.parquet".to_string(),
770 ];
771 assert!(
772 select_load_keys(&[], &all).is_empty(),
773 "no new runs ⇒ load nothing, not everything"
774 );
775 }
776
777 #[test]
778 fn gc_orphans_removes_unmanifested_parquet_only() {
779 let (store, _g) = fs_store(&[
780 ("base/r1-000.parquet", b"aa".to_vec()), ("base/orphan.parquet", b"junk".to_vec()), ("base/manifest.json", b"{}".to_vec()), ("base/_SUCCESS", b"".to_vec()), ]);
785 let keyed = vec![keyed("base/manifest-r1.json", "r1", "r1-000.parquet")];
786 let (removed, bytes) = gc_orphans(&store, "gs://b/base", &keyed, false).unwrap();
788 assert_eq!(removed, 1, "only the unmanifested part is removed");
789 assert_eq!(bytes, 4, "'junk' is 4 bytes");
790 let mut left = store.list_files("base").unwrap();
791 left.sort();
792 assert_eq!(
793 left,
794 vec![
795 "base/_SUCCESS".to_string(),
796 "base/manifest.json".to_string(),
797 "base/r1-000.parquet".to_string(),
798 ],
799 "the manifested part, the manifest, and _SUCCESS all survive"
800 );
801 }
802
803 #[test]
804 fn gc_orphans_of_an_all_manifested_prefix_removes_nothing() {
805 let (store, _g) = fs_store(&[("base/r1-000.parquet", b"a".to_vec())]);
806 let keyed = vec![keyed("base/manifest-r1.json", "r1", "r1-000.parquet")];
807 assert_eq!(
808 gc_orphans(&store, "gs://b/base", &keyed, false).unwrap().0,
809 0
810 );
811 }
812
813 #[test]
814 fn gc_orphans_keeps_a_snapshot_subprefix_manifested_part() {
815 let (store, _g) = fs_store(&[
816 ("base/snapshot/s-000.parquet", b"a".to_vec()), ("base/orphan.parquet", b"x".to_vec()), ]);
819 let keyed = vec![keyed("base/snapshot/manifest-s.json", "s", "s-000.parquet")];
820 let (removed, _) = gc_orphans(&store, "gs://b/base", &keyed, false).unwrap();
821 assert_eq!(removed, 1, "the top-level orphan goes");
822 assert_eq!(
823 store.list_files("base/snapshot").unwrap(),
824 vec!["base/snapshot/s-000.parquet".to_string()],
825 "the snapshot-subprefix manifested part is kept"
826 );
827 }
828
829 #[test]
830 fn gc_orphans_deletes_a_terminal_runs_parts_even_while_a_run_is_active() {
831 let (store, _g) = fs_store(&[("base/f-000.parquet", b"x".to_vec())]);
835 let mut kv = keyed("base/manifest-f.json", "f", "f-000.parquet");
836 kv.1.status = ManifestStatus::Failed;
837 assert_eq!(gc_orphans(&store, "gs://b/base", &[kv], true).unwrap().0, 1);
838 }
839
840 #[test]
841 fn gc_orphans_spares_an_unmanifested_part_while_a_run_is_active() {
842 let (store, _g) = fs_store(&[
848 ("base/r1-000.parquet", b"aa".to_vec()), ("base/inflight.parquet", b"x".to_vec()), ]);
851 let keyed = vec![keyed("base/manifest-r1.json", "r1", "r1-000.parquet")];
852 let (removed, _) = gc_orphans(&store, "gs://b/base", &keyed, true).unwrap();
853 assert_eq!(
854 removed, 0,
855 "an unmanifested part is spared while a run is active"
856 );
857 assert!(
858 store
859 .list_files("base")
860 .unwrap()
861 .iter()
862 .any(|k| k.ends_with("inflight.parquet")),
863 "the live run's in-flight part must survive gc_orphans"
864 );
865 }
866
867 #[test]
868 fn gc_orphans_collects_an_unmanifested_part_when_no_run_is_active() {
869 let (store, _g) = fs_store(&[
872 ("base/r1-000.parquet", b"aa".to_vec()),
873 ("base/dead-orphan.parquet", b"x".to_vec()),
874 ]);
875 let keyed = vec![keyed("base/manifest-r1.json", "r1", "r1-000.parquet")];
876 assert_eq!(
877 gc_orphans(&store, "gs://b/base", &keyed, false).unwrap().0,
878 1,
879 "with no active run, an unmanifested orphan is collected"
880 );
881 }
882
883 fn running(run: &str, started_at: &str) -> (String, RunManifest) {
885 let mut m = manifest(run, 0, None);
886 m.status = ManifestStatus::Running;
887 m.started_at = started_at.into();
888 m.finished_at = String::new();
889 m.parts.clear();
890 (format!("base/manifest-{run}.json"), m)
891 }
892
893 #[test]
894 fn has_active_running_manifest_true_for_a_lone_running_marker() {
895 assert!(has_active_running_manifest(&[running(
896 "r1",
897 "2026-01-01T00:00:00Z"
898 )]));
899 }
900
901 #[test]
902 fn has_active_running_manifest_false_when_none_is_running() {
903 assert!(!has_active_running_manifest(&[keyed(
905 "k",
906 "r1",
907 "p.parquet"
908 )]));
909 }
910
911 #[test]
912 fn has_active_running_manifest_false_for_a_superseded_running_marker() {
913 let r1 = running("r1", "2026-01-01T00:00:00Z");
917 let mut r2 = manifest("r2", 10, Some(10)); r2.started_at = "2026-01-02T00:00:00Z".into();
919 assert!(!has_active_running_manifest(&[
920 r1,
921 ("base/manifest-r2.json".into(), r2)
922 ]));
923 }
924
925 #[test]
926 fn select_runs_drops_a_running_manifest() {
927 let run_marker = running("r_running", "2026-01-02T00:00:00Z");
931 let ok = keyed("base/manifest-r_ok.json", "r_ok", "r_ok-000.parquet"); let sel = select_runs(
933 vec![run_marker, ok],
934 &std::collections::HashSet::new(),
935 crate::load::plan::LoadMode::Incremental,
936 );
937 assert_eq!(sel.len(), 1, "only the Success run is selected");
938 assert_eq!(sel[0].1.run_id, "r_ok");
939 assert!(sel.iter().all(|(_, m)| m.status != ManifestStatus::Running));
940 }
941
942 #[test]
943 fn gc_orphans_removes_a_superseded_running_marker_but_spares_a_live_one() {
944 let (store, _g) = fs_store(&[
949 ("base/manifest-r1.json", b"{}".to_vec()), ("base/manifest-r2.json", b"{}".to_vec()), ]);
952 let r1 = running("r1", "2026-01-01T00:00:01Z");
953 let r2 = running("r2", "2026-01-01T00:00:02Z");
954 let (removed, _) = gc_orphans(&store, "gs://b/base", &[r1, r2], true).unwrap();
955 assert_eq!(removed, 1, "only the superseded running marker is removed");
956 let left = store.list_files("base").unwrap();
957 assert!(
958 !left.iter().any(|k| k.ends_with("manifest-r1.json")),
959 "the superseded (dead) running marker is deleted"
960 );
961 assert!(
962 left.iter().any(|k| k.ends_with("manifest-r2.json")),
963 "the live (non-superseded) running marker survives — it is the active signal"
964 );
965 }
966
967 fn keyed_at(run: &str, finished_at: &str) -> (String, RunManifest) {
968 let mut m = manifest(run, 100, None);
969 m.finished_at = finished_at.into();
970 (format!("base/manifest-{run}.json"), m)
971 }
972
973 #[test]
974 fn latest_full_picks_the_newest_snapshot_not_all() {
975 let keyed = vec![
978 keyed_at("r1", "2026-01-01T00:00:00Z"),
979 keyed_at("r3", "2026-01-03T00:00:00Z"),
980 keyed_at("r2", "2026-01-02T00:00:00Z"),
981 ];
982 let sel = latest_full(keyed);
983 assert_eq!(sel.len(), 1, "exactly one snapshot, never all");
984 assert_eq!(sel[0].1.run_id, "r3", "the newest by finished_at");
985 }
986
987 #[test]
988 fn latest_full_re_materializes_even_when_the_latest_is_already_loaded() {
989 let keyed = vec![
994 keyed_at("r1", "2026-01-01T00:00:00Z"),
995 keyed_at("r2", "2026-01-02T00:00:00Z"),
996 ];
997 let sel = latest_full(keyed);
999 assert_eq!(sel.len(), 1);
1000 assert_eq!(sel[0].1.run_id, "r2", "always the latest, loaded or not");
1001 }
1002
1003 #[test]
1004 fn latest_full_of_no_staged_runs_is_empty_so_the_caller_no_ops_without_truncating() {
1005 assert!(latest_full(Vec::new()).is_empty());
1008 }
1009
1010 #[test]
1011 fn latest_full_orders_by_parsed_instant_not_lexical_bytes() {
1012 let keyed = vec![
1015 keyed_at("older", "2026-01-01T00:00:00Z"),
1016 keyed_at("newer", "2026-01-01T00:00:00.500Z"),
1017 ];
1018 let sel = latest_full(keyed);
1019 assert_eq!(sel.len(), 1);
1020 assert_eq!(
1021 sel[0].1.run_id, "newer",
1022 "the fractional-second run is the newer instant"
1023 );
1024 }
1025
1026 #[test]
1027 fn select_runs_full_picks_the_latest_even_when_loaded_and_even_when_stateless() {
1028 use crate::load::plan::LoadMode;
1029 use std::collections::HashSet;
1030 let keyed = vec![
1031 keyed_at("r1", "2026-01-01T00:00:00Z"),
1032 keyed_at("r2", "2026-01-02T00:00:00Z"),
1033 ];
1034 let loaded = HashSet::from(["r2".to_string()]);
1036 let sel = select_runs(keyed.clone(), &loaded, LoadMode::Full);
1037 assert_eq!(sel.len(), 1);
1038 assert_eq!(sel[0].1.run_id, "r2", "Full picks latest, loaded or not");
1039 let sel = select_runs(keyed, &HashSet::new(), LoadMode::Full);
1042 assert_eq!(sel.len(), 1, "stateless Full is not a blanket load");
1043 assert_eq!(sel[0].1.run_id, "r2");
1044 }
1045
1046 #[test]
1047 fn select_runs_append_modes_filter_loaded_and_load_all_when_stateless() {
1048 use crate::load::plan::LoadMode;
1049 use std::collections::HashSet;
1050 let keyed = vec![
1051 keyed_at("r1", "2026-01-01T00:00:00Z"),
1052 keyed_at("r2", "2026-01-02T00:00:00Z"),
1053 ];
1054 let loaded = HashSet::from(["r1".to_string()]);
1055 for mode in [LoadMode::Incremental, LoadMode::Cdc] {
1056 let sel = select_runs(keyed.clone(), &loaded, mode);
1057 assert_eq!(sel.len(), 1, "{mode:?}: only the unloaded run");
1058 assert_eq!(sel[0].1.run_id, "r2");
1059 assert_eq!(
1061 select_runs(keyed.clone(), &HashSet::new(), mode).len(),
1062 2,
1063 "{mode:?}: stateless loads every run"
1064 );
1065 }
1066 }
1067
1068 #[test]
1069 fn is_run_unique_manifest_needs_both_prefix_and_json() {
1070 assert!(is_run_unique_manifest("manifest-20260101T000000.json"));
1071 assert!(!is_run_unique_manifest("manifest.json")); assert!(!is_run_unique_manifest("manifest-abc.txt")); assert!(!is_run_unique_manifest("data.json")); }
1075
1076 #[test]
1077 fn chain_prefix_renders_source_and_files() {
1078 let known = LoadIntegrity {
1079 source_rows: Some(100),
1080 file_rows: 100,
1081 manifests: 1,
1082 };
1083 assert_eq!(known.chain_prefix(), "source 100 → files 100");
1084 let unknown = LoadIntegrity {
1085 source_rows: None,
1086 file_rows: 40,
1087 manifests: 1,
1088 };
1089 assert_eq!(unknown.chain_prefix(), "source ? → files 40");
1090 }
1091
1092 #[test]
1093 fn is_manifest_key_matches_only_the_final_segment() {
1094 assert!(is_manifest_key("gs://b/p/manifest.json"));
1095 assert!(is_manifest_key("manifest.json"));
1096 assert!(!is_manifest_key("gs://b/p/part-0.parquet"));
1097 assert!(!is_manifest_key("gs://b/p/x_manifest.json"));
1098 }
1099
1100 fn fs_store(files: &[(&str, Vec<u8>)]) -> (GcsStore, tempfile::TempDir) {
1106 let dir = tempfile::tempdir().unwrap();
1107 for (rel, bytes) in files {
1108 let p = dir.path().join(rel);
1109 std::fs::create_dir_all(p.parent().unwrap()).unwrap();
1110 std::fs::write(p, bytes).unwrap();
1111 }
1112 let store = GcsStore::open_fs(dir.path().to_str().unwrap()).unwrap();
1113 (store, dir)
1114 }
1115
1116 fn manifest_bytes(run: &str, rows: i64, source: Option<i64>) -> Vec<u8> {
1117 serde_json::to_vec(&manifest(run, rows, source)).unwrap()
1118 }
1119
1120 #[test]
1121 fn list_manifest_keys_prefers_run_unique_copies_over_the_canonical_pointer() {
1122 let (store, _g) = fs_store(&[
1126 ("base/manifest.json", b"{}".to_vec()),
1127 ("base/manifest-r1.json", b"{}".to_vec()),
1128 ("base/manifest-r2.json", b"{}".to_vec()),
1129 ("base/part-0.parquet", b"x".to_vec()), ]);
1131 let mut keys = list_manifest_keys(&store, "base").unwrap();
1132 keys.sort();
1133 assert_eq!(
1134 keys,
1135 vec![
1136 "base/manifest-r1.json".to_string(),
1137 "base/manifest-r2.json".to_string(),
1138 ]
1139 );
1140 }
1141
1142 #[test]
1143 fn list_manifest_keys_falls_back_to_the_canonical_name_for_a_single_run() {
1144 let (store, _g) = fs_store(&[("base/manifest.json", b"{}".to_vec())]);
1147 assert_eq!(
1148 list_manifest_keys(&store, "base").unwrap(),
1149 vec!["base/manifest.json".to_string()]
1150 );
1151 }
1152
1153 #[test]
1154 fn fetch_manifests_keyed_reads_and_parses_every_run_copy_under_the_prefix() {
1155 let (store, _g) = fs_store(&[
1159 ("base/manifest.json", manifest_bytes("r2", 40, Some(40))),
1160 (
1161 "base/manifest-r1.json",
1162 manifest_bytes("r1", 100, Some(100)),
1163 ),
1164 ("base/manifest-r2.json", manifest_bytes("r2", 40, Some(40))),
1165 ]);
1166 let manifests: Vec<_> = fetch_manifests_keyed(&store, "gs://my-bucket/base")
1167 .unwrap()
1168 .into_iter()
1169 .map(|(_, m)| m)
1170 .collect();
1171 assert_eq!(manifests.len(), 2);
1172 let integrity = reconcile(&manifests, false).unwrap();
1173 assert_eq!(integrity.file_rows, 140);
1174 assert_eq!(integrity.manifests, 2);
1175 }
1176
1177 #[test]
1178 fn fetch_manifests_keyed_names_the_key_when_a_manifest_is_unparseable() {
1179 let (store, _g) = fs_store(&[("base/manifest.json", b"{ not json".to_vec())]);
1180 let err = fetch_manifests_keyed(&store, "gs://my-bucket/base")
1181 .unwrap_err()
1182 .to_string();
1183 assert!(
1184 err.contains("parsing manifest") && err.contains("base/manifest.json"),
1185 "error should name the offending key: {err}"
1186 );
1187 }
1188
1189 const FAKE_GCS_ENDPOINT: &str = "http://127.0.0.1:4443";
1200 const FAKE_GCS_BUCKET: &str = "rivet-load-emulator";
1201
1202 fn fake_gcs_store() -> GcsStore {
1207 let created = std::process::Command::new("curl")
1208 .args([
1209 "-s",
1210 "-X",
1211 "POST",
1212 &format!("{FAKE_GCS_ENDPOINT}/storage/v1/b?project=rivet-test"),
1213 "-H",
1214 "Content-Type: application/json",
1215 "-d",
1216 &format!("{{\"name\":\"{FAKE_GCS_BUCKET}\"}}"),
1217 ])
1218 .output();
1219 assert!(
1220 created.is_ok_and(|o| o.status.success()),
1221 "could not reach fake-gcs to create the bucket — is `docker compose up -d fake-gcs` running on :4443?"
1222 );
1223 let cfg = crate::config::DestinationConfig {
1224 destination_type: crate::config::DestinationType::Gcs,
1225 bucket: Some(FAKE_GCS_BUCKET.into()),
1226 endpoint: Some(FAKE_GCS_ENDPOINT.into()),
1227 allow_anonymous: true,
1228 ..Default::default()
1229 };
1230 GcsStore::new(&cfg).expect("build GcsStore against fake-gcs")
1231 }
1232
1233 fn drain(store: &GcsStore, prefix: &str) {
1241 for key in store.list_files(prefix).unwrap() {
1242 store.remove(&key).unwrap();
1243 }
1244 }
1245
1246 #[test]
1247 #[ignore = "emulator: needs `docker compose up -d fake-gcs` (fsouza/fake-gcs-server :4443)"]
1248 fn storage_contract_over_fake_gcs() {
1249 let store = fake_gcs_store();
1250 let prefix = "load-contract/orders";
1253 drain(&store, prefix);
1254 let gs = format!("gs://{FAKE_GCS_BUCKET}/{prefix}");
1255
1256 store
1259 .put(
1260 &format!("{prefix}/manifest-r1.json"),
1261 &manifest_bytes("r1", 100, Some(100)),
1262 )
1263 .unwrap();
1264 store
1265 .put(&format!("{prefix}/part-000000.parquet"), b"rows-of-r1")
1266 .unwrap();
1267 store
1268 .put(&format!("{prefix}/orphan.parquet"), b"crash-leftover")
1269 .unwrap();
1270
1271 let keyed = fetch_manifests_keyed(&store, &gs).unwrap();
1273 assert_eq!(keyed.len(), 1, "the run's manifest, read back over GCS");
1274
1275 let manifests: Vec<_> = keyed.iter().map(|(_, m)| m.clone()).collect();
1277 assert_eq!(
1278 reconcile(&manifests, false).unwrap().file_rows,
1279 100,
1280 "file_rows drives the count-gate; a bad GCS read would corrupt it"
1281 );
1282
1283 assert_eq!(
1285 select_load_uris(&store, &gs, &keyed).unwrap(),
1286 vec![format!(
1287 "gs://{FAKE_GCS_BUCKET}/{prefix}/part-000000.parquet"
1288 )],
1289 "load pulls the manifested part, not the unmanifested crash orphan"
1290 );
1291
1292 let (removed, _bytes) = gc_orphans(&store, &gs, &keyed, false).unwrap();
1295 assert_eq!(
1296 removed, 1,
1297 "exactly the orphan parquet is GC'd over real GCS"
1298 );
1299 let mut left = store.list_files(prefix).unwrap();
1300 left.sort();
1301 assert_eq!(
1302 left,
1303 vec![
1304 format!("{prefix}/manifest-r1.json"),
1305 format!("{prefix}/part-000000.parquet"),
1306 ],
1307 "the manifested part + its manifest survive the orphan GC"
1308 );
1309
1310 drain(&store, prefix);
1314 assert!(
1315 store.list_files(prefix).unwrap().is_empty(),
1316 "teardown left the prefix clean over real GCS"
1317 );
1318 }
1319}