1use std::collections::HashMap;
89
90use anyhow::{Context, Result};
91use crate::snapshot::BackupTarget;
92use crate::stream::{
93 for_each_frame_in_generation, list_and_parse_generation_manifests, validate_generation_chain,
94 CoreWalSeam, WalGeneration, WalInsertSeam, Watermark,
95};
96
97pub trait WarmApplier: WalInsertSeam {
102 fn trim_page_cache(&self, target_kb: i64) -> Result<()>;
108}
109
110impl WarmApplier for CoreWalSeam {
111 fn trim_page_cache(&self, target_kb: i64) -> Result<()> {
112 self.trim_page_cache_kb(target_kb)
113 }
114}
115
116#[derive(Debug, Clone, Copy, PartialEq, Eq)]
120pub struct WalPullerConfig {
121 pub max_appliers: usize,
125 pub max_disk_bytes: u64,
129 pub applier_cache_kb: i64,
131}
132
133impl Default for WalPullerConfig {
134 fn default() -> Self {
135 Self {
136 max_appliers: 1_000,
137 max_disk_bytes: 500 * 1024 * 1024 * 1024,
138 applier_cache_kb: 64,
139 }
140 }
141}
142
143#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
145pub struct PullReport {
146 pub frames_pulled: u64,
149 pub checkpoint_seq: u32,
152 pub appliers_fanned: usize,
154}
155
156struct AttachedApplier {
157 seam: Box<dyn WarmApplier>,
158 dest_size_bytes: u64,
159}
160
161pub struct WalPuller<'a> {
166 target: &'a BackupTarget,
167 page_size: usize,
168 cfg: WalPullerConfig,
169 appliers: HashMap<String, AttachedApplier>,
170 disk_used_bytes: u64,
171 pulled: Option<Watermark>,
175 pulled_generation: Option<WalGeneration>,
180}
181
182impl<'a> WalPuller<'a> {
183 pub fn new(target: &'a BackupTarget, page_size: usize, cfg: WalPullerConfig) -> Self {
184 Self {
185 target,
186 page_size,
187 cfg,
188 appliers: HashMap::new(),
189 disk_used_bytes: 0,
190 pulled: None,
191 pulled_generation: None,
192 }
193 }
194
195 pub fn page_size(&self) -> usize {
196 self.page_size
197 }
198
199 pub fn attached_count(&self) -> usize {
200 self.appliers.len()
201 }
202
203 pub fn disk_used_bytes(&self) -> u64 {
204 self.disk_used_bytes
205 }
206
207 pub fn has_pulled(&self) -> bool {
211 self.pulled.is_some()
212 }
213
214 pub fn attach(
228 &mut self,
229 name: impl Into<String>,
230 seam: Box<dyn WarmApplier>,
231 dest_size_bytes: u64,
232 ) -> Result<()> {
233 anyhow::ensure!(
234 self.pulled.is_none(),
235 "WalPuller v1 does not support attaching an applier mid-stream — materialize it via \
236 stream::restore_latest_stream and hand it to a fresh WalPuller, or attach before the \
237 first pull_once()"
238 );
239 let name = name.into();
240 anyhow::ensure!(
241 !self.appliers.contains_key(&name),
242 "applier {name:?} is already attached"
243 );
244 anyhow::ensure!(
245 self.appliers.len() < self.cfg.max_appliers,
246 "FD budget exhausted: {} appliers already attached (max_appliers={})",
247 self.appliers.len(),
248 self.cfg.max_appliers,
249 );
250 let projected = self.disk_used_bytes.saturating_add(dest_size_bytes);
251 anyhow::ensure!(
252 projected <= self.cfg.max_disk_bytes,
253 "disk budget exhausted: attaching {name:?} ({dest_size_bytes} bytes) would use \
254 {projected} bytes total, max_disk_bytes={}",
255 self.cfg.max_disk_bytes,
256 );
257
258 seam.wal_insert_begin()
259 .with_context(|| format!("wal_insert_begin for applier {name:?}"))?;
260 seam.trim_page_cache(self.cfg.applier_cache_kb)
261 .with_context(|| format!("trimming page cache for applier {name:?}"))?;
262
263 self.disk_used_bytes = projected;
264 self.appliers
265 .insert(name, AttachedApplier { seam, dest_size_bytes });
266 Ok(())
267 }
268
269 pub fn detach(&mut self, name: &str) -> Result<Box<dyn WarmApplier>> {
274 let applier = self
275 .appliers
276 .remove(name)
277 .with_context(|| format!("no applier attached under {name:?}"))?;
278 self.disk_used_bytes = self.disk_used_bytes.saturating_sub(applier.dest_size_bytes);
279 applier
280 .seam
281 .wal_insert_end(false)
282 .with_context(|| format!("wal_insert_end for applier {name:?}"))?;
283 Ok(applier.seam)
284 }
285
286 pub fn promote(&mut self, name: &str) -> Result<Box<dyn WarmApplier>> {
292 self.detach(name)
293 }
294
295 pub async fn pull_once(&mut self) -> Result<PullReport> {
307 let manifests = list_and_parse_generation_manifests(self.target).await?;
308 if manifests.is_empty() {
309 return Ok(PullReport {
310 frames_pulled: 0,
311 checkpoint_seq: self.pulled.map(|p| p.checkpoint_seq).unwrap_or(0),
312 appliers_fanned: self.appliers.len(),
313 });
314 }
315 let chain = validate_generation_chain(&manifests)?;
316 anyhow::ensure!(
317 chain.page_size == self.page_size,
318 "generation chain page_size {} does not match this WalPuller's page_size {}",
319 chain.page_size,
320 self.page_size,
321 );
322 if let Some(prior) = self.pulled_generation {
330 anyhow::ensure!(
331 chain.generation.is_provably_same_as(&prior),
332 "WAL restart since the last pull ({} -> {}) — attached appliers \
333 hold frames from the old WAL generation and can't be trusted to compose with the \
334 new one; re-materialize them via a fresh restore_latest_stream and start a new \
335 WalPuller",
336 prior.describe(),
337 chain.generation.describe(),
338 );
339 }
340 let start_frame = self.pulled.map(|p| p.last_frame + 1).unwrap_or(1);
341 if chain.total_frames < start_frame {
342 return Ok(PullReport {
343 frames_pulled: 0,
344 checkpoint_seq: chain.generation.checkpoint_seq,
345 appliers_fanned: self.appliers.len(),
346 });
347 }
348
349 let mut frames_pulled = 0u64;
350 let mut last_frame_no = start_frame - 1;
351 let appliers = &self.appliers;
352 for m in &manifests {
353 if m.last_frame < start_frame {
354 continue; }
356 frames_pulled += for_each_frame_in_generation(
363 self.target,
364 m,
365 start_frame,
366 |frame_no, bytes| {
367 for (applier_name, applier) in appliers.iter() {
368 applier.seam.wal_insert_frame(frame_no, bytes).with_context(|| {
369 format!("fanning frame {frame_no} to applier {applier_name:?}")
370 })?;
371 }
372 last_frame_no = frame_no;
373 Ok(())
374 },
375 )
376 .await?;
377 }
378 self.pulled = Some(Watermark {
379 checkpoint_seq: chain.generation.checkpoint_seq,
380 last_frame: last_frame_no,
381 });
382 self.pulled_generation = Some(chain.generation);
383 Ok(PullReport {
384 frames_pulled,
385 checkpoint_seq: chain.generation.checkpoint_seq,
386 appliers_fanned: self.appliers.len(),
387 })
388 }
389}
390
391#[cfg(test)]
392mod tests {
393 use super::*;
394 use crate::backpressure::BackpressureConfig;
395 use crate::snapshot::BackupTarget;
396 use crate::stream::{tail_frames, FrameInfo, StreamConfig, WalSeam, WAL_FRAME_HEADER_SIZE};
397 use object_store::memory::InMemory;
398 use object_store::path::Path as ObjPath;
399 use object_store::{
400 CopyOptions, GetOptions, GetResult, ListResult, MultipartUpload, ObjectMeta, ObjectStore,
401 ObjectStoreExt, PutMultipartOptions, PutOptions, PutPayload, PutResult, Result as OsResult,
402 };
403 use std::cell::RefCell;
404 use std::fmt;
405 use std::sync::atomic::{AtomicU64, Ordering};
406 use std::sync::Arc;
407
408 fn fresh_target() -> BackupTarget {
411 BackupTarget {
412 store: Arc::new(InMemory::new()),
413 prefix: "backups".into(),
414 }
415 }
416
417 fn stream_cfg() -> StreamConfig<'static> {
418 StreamConfig {
419 base_snapshot_key: "backups/snapshots/snapshot-00000000000000000001.db",
420 page_size: 4096,
421 backpressure: BackpressureConfig::default(),
422 rpo_target: None,
423 epoch: 0,
424 owner: None,
425 pointer_generation: 0,
426 }
427 }
428
429 fn puller_cfg(max_appliers: usize, max_disk_bytes: u64) -> WalPullerConfig {
430 WalPullerConfig {
431 max_appliers,
432 max_disk_bytes,
433 applier_cache_kb: 64,
434 }
435 }
436
437 struct MockWal {
442 state: RefCell<MockState>,
443 }
444 struct MockState {
445 checkpoint_seq: u32,
446 frames: Vec<FrameInfo>,
447 }
448 impl MockWal {
449 fn new() -> Self {
450 Self { state: RefCell::new(MockState { checkpoint_seq: 0, frames: Vec::new() }) }
451 }
452 fn append(&self, page_no: u32, db_size: u32) {
453 self.state.borrow_mut().frames.push(FrameInfo { page_no, db_size });
454 }
455 fn restart(&self) {
456 let mut s = self.state.borrow_mut();
457 s.checkpoint_seq += 1;
458 s.frames.clear();
459 }
460 }
461 impl WalSeam for MockWal {
462 fn wal_state(&self) -> Result<Watermark> {
463 let s = self.state.borrow();
464 Ok(Watermark { checkpoint_seq: s.checkpoint_seq, last_frame: s.frames.len() as u64 })
465 }
466 fn wal_get_frame(&self, frame_no: u64, buf: &mut [u8]) -> Result<FrameInfo> {
467 let s = self.state.borrow();
468 let f = s.frames[frame_no as usize - 1];
469 buf[0..4].copy_from_slice(&f.page_no.to_be_bytes());
470 buf[4..8].copy_from_slice(&f.db_size.to_be_bytes());
471 buf[8..WAL_FRAME_HEADER_SIZE].fill(0);
472 buf[WAL_FRAME_HEADER_SIZE..].fill(frame_no as u8);
473 Ok(f)
474 }
475 fn wal_auto_actions_disable(&self) {}
476 }
477
478 #[derive(Clone)]
485 struct MockApplier(std::rc::Rc<RefCell<Vec<MockEvent>>>);
486 #[derive(Debug, Clone, PartialEq, Eq)]
487 enum MockEvent {
488 Begin,
489 TrimCache { target_kb: i64 },
490 Frame { frame_no: u64, fill: u8 },
491 End { force_commit: bool },
492 }
493 impl MockApplier {
494 fn new() -> Self {
495 Self(std::rc::Rc::new(RefCell::new(Vec::new())))
496 }
497 fn events(&self) -> Vec<MockEvent> {
498 self.0.borrow().clone()
499 }
500 }
501 impl WalInsertSeam for MockApplier {
502 fn wal_insert_begin(&self) -> Result<()> {
503 self.0.borrow_mut().push(MockEvent::Begin);
504 Ok(())
505 }
506 fn wal_insert_frame(&self, frame_no: u64, frame: &[u8]) -> Result<()> {
507 self.0
508 .borrow_mut()
509 .push(MockEvent::Frame { frame_no, fill: frame[WAL_FRAME_HEADER_SIZE] });
510 Ok(())
511 }
512 fn wal_insert_end(&self, force_commit: bool) -> Result<()> {
513 self.0.borrow_mut().push(MockEvent::End { force_commit });
514 Ok(())
515 }
516 }
517 impl WarmApplier for MockApplier {
518 fn trim_page_cache(&self, target_kb: i64) -> Result<()> {
519 self.0.borrow_mut().push(MockEvent::TrimCache { target_kb });
520 Ok(())
521 }
522 }
523
524 struct FailingTrimApplier;
527 impl WalInsertSeam for FailingTrimApplier {
528 fn wal_insert_begin(&self) -> Result<()> {
529 Ok(())
530 }
531 fn wal_insert_frame(&self, _: u64, _: &[u8]) -> Result<()> {
532 Ok(())
533 }
534 fn wal_insert_end(&self, _: bool) -> Result<()> {
535 Ok(())
536 }
537 }
538 impl WarmApplier for FailingTrimApplier {
539 fn trim_page_cache(&self, _: i64) -> Result<()> {
540 anyhow::bail!("simulated trim failure")
541 }
542 }
543
544 struct CountingStore {
549 inner: Arc<dyn ObjectStore>,
550 gets: AtomicU64,
551 }
552 impl CountingStore {
553 fn new(inner: Arc<dyn ObjectStore>) -> Self {
554 Self { inner, gets: AtomicU64::new(0) }
555 }
556 fn get_count(&self) -> u64 {
557 self.gets.load(Ordering::SeqCst)
558 }
559 }
560 impl fmt::Display for CountingStore {
561 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
562 write!(f, "CountingStore({})", self.inner)
563 }
564 }
565 impl fmt::Debug for CountingStore {
566 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
567 write!(f, "CountingStore({:?})", self.inner)
568 }
569 }
570 #[async_trait::async_trait]
571 impl ObjectStore for CountingStore {
572 async fn put_opts(
573 &self,
574 location: &ObjPath,
575 payload: PutPayload,
576 opts: PutOptions,
577 ) -> OsResult<PutResult> {
578 self.inner.put_opts(location, payload, opts).await
579 }
580 async fn put_multipart_opts(
581 &self,
582 location: &ObjPath,
583 opts: PutMultipartOptions,
584 ) -> OsResult<Box<dyn MultipartUpload>> {
585 self.inner.put_multipart_opts(location, opts).await
586 }
587 async fn get_opts(&self, location: &ObjPath, options: GetOptions) -> OsResult<GetResult> {
588 self.gets.fetch_add(1, Ordering::SeqCst);
589 self.inner.get_opts(location, options).await
590 }
591 fn delete_stream(
592 &self,
593 locations: futures_util::stream::BoxStream<'static, OsResult<ObjPath>>,
594 ) -> futures_util::stream::BoxStream<'static, OsResult<ObjPath>> {
595 self.inner.delete_stream(locations)
596 }
597 fn list(
598 &self,
599 prefix: Option<&ObjPath>,
600 ) -> futures_util::stream::BoxStream<'static, OsResult<ObjectMeta>> {
601 self.inner.list(prefix)
602 }
603 async fn list_with_delimiter(&self, prefix: Option<&ObjPath>) -> OsResult<ListResult> {
604 self.inner.list_with_delimiter(prefix).await
605 }
606 async fn copy_opts(
607 &self,
608 from: &ObjPath,
609 to: &ObjPath,
610 options: CopyOptions,
611 ) -> OsResult<()> {
612 self.inner.copy_opts(from, to, options).await
613 }
614 }
615
616 #[test]
619 fn attach_enforces_fd_budget() {
620 let target = fresh_target();
621 let mut puller = WalPuller::new(&target, 4096, puller_cfg(1, u64::MAX));
622 puller.attach("a", Box::new(MockApplier::new()), 0).unwrap();
623 let err = puller.attach("b", Box::new(MockApplier::new()), 0).unwrap_err();
624 assert!(format!("{err}").contains("FD budget"), "err was {err}");
625 assert_eq!(puller.attached_count(), 1);
626 }
627
628 #[test]
629 fn attach_enforces_disk_budget() {
630 let target = fresh_target();
631 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, 100));
632 puller.attach("a", Box::new(MockApplier::new()), 60).unwrap();
633 let err = puller.attach("b", Box::new(MockApplier::new()), 60).unwrap_err();
634 assert!(format!("{err}").contains("disk budget"), "err was {err}");
635 assert_eq!(puller.disk_used_bytes(), 60);
636 }
637
638 #[test]
639 fn attach_rejects_duplicate_name() {
640 let target = fresh_target();
641 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
642 puller.attach("a", Box::new(MockApplier::new()), 0).unwrap();
643 let err = puller.attach("a", Box::new(MockApplier::new()), 0).unwrap_err();
644 assert!(format!("{err}").contains("already attached"), "err was {err}");
645 }
646
647 #[test]
650 fn attach_begins_session_and_trims_cache() {
651 let target = fresh_target();
652 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
653 let a = MockApplier::new();
654 puller.attach("a", Box::new(a.clone()), 0).unwrap();
655 assert_eq!(a.events(), vec![MockEvent::Begin, MockEvent::TrimCache { target_kb: 64 }]);
656 }
657
658 #[test]
659 fn attach_propagates_trim_failure_without_registering_applier() {
660 let target = fresh_target();
661 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
662 let err = puller.attach("a", Box::new(FailingTrimApplier), 0).unwrap_err();
663 assert!(format!("{err}").contains("trimming page cache"), "err was {err}");
664 assert_eq!(puller.attached_count(), 0, "a failed attach must not register");
665 }
666
667 #[tokio::test]
668 async fn attach_refuses_after_first_pull() {
669 let seam = MockWal::new();
670 seam.append(1, 1); let target = fresh_target();
672 let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
673
674 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
675 assert!(!puller.has_pulled());
676 let report = puller.pull_once().await.unwrap();
677 assert_eq!(report.frames_pulled, 1, "a real frame must actually be fanned");
678 assert!(puller.has_pulled());
679 let err = puller
680 .attach("late", Box::new(MockApplier::new()), 0)
681 .unwrap_err();
682 assert!(format!("{err}").contains("mid-stream"), "err was {err}");
683 }
684
685 #[tokio::test]
690 async fn attach_still_allowed_after_an_empty_pull() {
691 let target = fresh_target();
692 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
693 let report = puller.pull_once().await.unwrap();
694 assert_eq!(report.frames_pulled, 0);
695 assert!(!puller.has_pulled());
696 puller.attach("a", Box::new(MockApplier::new()), 0).unwrap();
697 assert_eq!(puller.attached_count(), 1);
698 }
699
700 #[tokio::test]
707 async fn pull_once_downloads_each_frame_once_regardless_of_applier_count() {
708 let seam = MockWal::new();
709 seam.append(1, 0);
710 seam.append(2, 0);
711 seam.append(3, 3); let raw_target = fresh_target();
713 let _ = tail_frames(&seam, &raw_target, &stream_cfg()).await.unwrap();
714
715 let counting = Arc::new(CountingStore::new(raw_target.store.clone()));
716 let counted_target = BackupTarget { store: counting.clone(), prefix: "backups".into() };
717
718 let mut puller = WalPuller::new(&counted_target, 4096, puller_cfg(10, u64::MAX));
719 for name in ["a", "b", "c"] {
720 puller.attach(name, Box::new(MockApplier::new()), 0).unwrap();
721 }
722 let report = puller.pull_once().await.unwrap();
723 assert_eq!(report.frames_pulled, 3);
724 assert_eq!(report.appliers_fanned, 3);
725 assert_eq!(
726 counting.get_count(),
727 2,
728 "1 batch object holding all 3 frames + 1 generation manifest, once each"
729 );
730 }
731
732 #[tokio::test]
737 async fn pull_once_fans_identical_frames_to_every_applier() {
738 let seam = MockWal::new();
739 seam.append(1, 0);
740 seam.append(2, 2); let target = fresh_target();
742 let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
743
744 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
745 let a = MockApplier::new();
746 let b = MockApplier::new();
747 puller.attach("a", Box::new(a.clone()), 0).unwrap();
748 puller.attach("b", Box::new(b.clone()), 0).unwrap();
749 puller.pull_once().await.unwrap();
750
751 let want = vec![
752 MockEvent::Begin,
753 MockEvent::TrimCache { target_kb: 64 },
754 MockEvent::Frame { frame_no: 1, fill: 1 },
755 MockEvent::Frame { frame_no: 2, fill: 2 },
756 ];
757 assert_eq!(a.events(), want);
758 assert_eq!(b.events(), want);
759 }
760
761 #[tokio::test]
764 async fn pull_once_is_incremental_across_calls() {
765 let seam = MockWal::new();
766 seam.append(1, 1); let target = fresh_target();
768 let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
769
770 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
771 puller.attach("a", Box::new(MockApplier::new()), 0).unwrap();
772 let first = puller.pull_once().await.unwrap();
773 assert_eq!(first.frames_pulled, 1);
774
775 let second = puller.pull_once().await.unwrap();
777 assert_eq!(second.frames_pulled, 0);
778
779 seam.append(2, 2); let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
781 let third = puller.pull_once().await.unwrap();
782 assert_eq!(third.frames_pulled, 1, "only the new frame, not a re-fan of frame 1");
783 }
784
785 #[tokio::test]
787 async fn pull_once_with_no_manifests_is_a_noop() {
788 let target = fresh_target();
789 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
790 let report = puller.pull_once().await.unwrap();
791 assert_eq!(report, PullReport { frames_pulled: 0, checkpoint_seq: 0, appliers_fanned: 0 });
792 }
793
794 #[tokio::test]
797 async fn pull_once_rejects_page_size_mismatch() {
798 let seam = MockWal::new();
799 seam.append(1, 1);
800 let target = fresh_target();
801 let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
802
803 let mut puller = WalPuller::new(&target, 8192, puller_cfg(10, u64::MAX));
804 let err = puller.pull_once().await.unwrap_err();
805 assert!(format!("{err}").contains("page_size"), "err was {err}");
806 }
807
808 #[tokio::test]
811 async fn pull_once_refuses_across_a_restart() {
812 let seam = MockWal::new();
813 seam.append(1, 1); let target = fresh_target();
815 let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
816
817 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
818 puller.attach("a", Box::new(MockApplier::new()), 0).unwrap();
819 puller.pull_once().await.unwrap();
820
821 seam.restart();
822 seam.append(1, 1);
823 let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
824
825 let err = puller.pull_once().await.unwrap_err();
826 assert!(format!("{err}").contains("WAL restart"), "err was {err}");
827 }
828
829 #[test]
832 fn detach_frees_disk_budget_and_ends_session() {
833 let target = fresh_target();
834 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, 100));
835 let a = MockApplier::new();
836 puller.attach("a", Box::new(a.clone()), 60).unwrap();
837 assert_eq!(puller.disk_used_bytes(), 60);
838
839 puller.detach("a").unwrap();
840 assert_eq!(puller.disk_used_bytes(), 0);
841 assert_eq!(puller.attached_count(), 0);
842 assert_eq!(a.events().last(), Some(&MockEvent::End { force_commit: false }));
843
844 puller.attach("b", Box::new(MockApplier::new()), 60).unwrap();
846 assert_eq!(puller.disk_used_bytes(), 60);
847 }
848
849 #[test]
850 fn detach_unknown_name_errors() {
851 let target = fresh_target();
852 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
853 assert!(puller.detach("ghost").is_err());
854 }
855
856 #[test]
857 fn promote_ends_session_with_no_forced_commit() {
858 let target = fresh_target();
859 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
860 let a = MockApplier::new();
861 puller.attach("a", Box::new(a.clone()), 0).unwrap();
862 puller.promote("a").unwrap();
863 assert_eq!(puller.attached_count(), 0);
864 assert_eq!(a.events().last(), Some(&MockEvent::End { force_commit: false }));
865 }
866
867 struct TempDb(std::path::PathBuf);
870 impl TempDb {
871 fn new(tag: &str) -> Self {
872 TempDb(std::env::temp_dir().join(format!(
873 "turso-backup-puller-{tag}-{}-{}.db",
874 std::process::id(),
875 std::time::SystemTime::now()
876 .duration_since(std::time::UNIX_EPOCH)
877 .unwrap()
878 .as_nanos(),
879 )))
880 }
881 fn path(&self) -> &str {
882 self.0.to_str().unwrap()
883 }
884 }
885 impl Drop for TempDb {
886 fn drop(&mut self) {
887 for sfx in ["", "-wal", "-shm"] {
888 let _ = std::fs::remove_file(format!("{}{sfx}", self.0.display()));
889 }
890 }
891 }
892
893 async fn seed_rows(path: &str, start: i64, count: i64) {
894 let db = turso::Builder::new_local(path).build().await.unwrap();
895 let conn = db.connect().unwrap();
896 conn.execute("CREATE TABLE IF NOT EXISTS t (id INTEGER PRIMARY KEY, v TEXT)", ())
897 .await
898 .unwrap();
899 conn.execute("BEGIN", ()).await.unwrap();
900 for i in start..start + count {
901 conn.execute("INSERT INTO t (id, v) VALUES (?, ?)", (i, format!("v{i}")))
902 .await
903 .unwrap();
904 }
905 conn.execute("COMMIT", ()).await.unwrap();
906 }
907
908 async fn count_rows(path: &str) -> i64 {
909 let db = turso::Builder::new_local(path).build().await.unwrap();
910 let conn = db.connect().unwrap();
911 let mut r = conn.query("SELECT COUNT(*) FROM t", ()).await.unwrap();
912 let row = r.next().await.unwrap().unwrap();
913 row.get::<i64>(0).unwrap()
914 }
915
916 #[tokio::test]
922 async fn live_pull_once_replays_onto_a_real_applier() {
923 let src = TempDb::new("live-src");
924 let dest = TempDb::new("live-dest");
925 seed_rows(src.path(), 0, 20).await;
926
927 let target = fresh_target();
928 let base_key = match crate::snapshot::snapshot_and_upload(src.path(), &target)
929 .await
930 .unwrap()
931 {
932 crate::snapshot::SnapshotOutcome::Uploaded { key, .. } => key,
933 other => panic!("expected a fresh Uploaded snapshot, got {other:?}"),
934 };
935
936 seed_rows(src.path(), 1000, 5).await;
938 {
939 let seam = crate::stream::CoreWalSeam::open(src.path()).unwrap();
940 let cfg = StreamConfig {
941 base_snapshot_key: &base_key,
942 page_size: 4096,
943 backpressure: BackpressureConfig::default(),
944 rpo_target: None,
945 epoch: 0,
946 owner: None,
947 pointer_generation: 0,
948 };
949 tail_frames(&seam, &target, &cfg).await.unwrap();
950 }
951
952 let base_bytes = target
955 .store
956 .get(&ObjPath::from(base_key.clone()))
957 .await
958 .unwrap()
959 .bytes()
960 .await
961 .unwrap();
962 std::fs::write(dest.path(), &base_bytes).unwrap();
963
964 let applier = crate::stream::CoreWalSeam::open(dest.path()).unwrap();
965 let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
966 puller.attach("dest", Box::new(applier), 0).unwrap();
967 let report = puller.pull_once().await.unwrap();
968 assert!(report.frames_pulled > 0);
969
970 puller.promote("dest").unwrap();
971 assert_eq!(count_rows(dest.path()).await, 25, "20 base + 5 streamed via the puller");
972 }
973}