Skip to main content

allsource_core/
store_refresh.rs

1//! Catch a read-only replica up with what its writer has made durable since boot.
2//!
3//! A replica is never reopened to refresh: `ParquetStorage::new` deletes
4//! `*.parquet.tmp` files, which on a live data dir are the writer's in-flight
5//! flushes. Refresh only lists, stats and reads.
6
7use std::{
8    collections::HashSet,
9    path::{Path, PathBuf},
10    sync::Arc,
11};
12
13use chrono::{DateTime, Utc};
14
15use super::EventStore;
16use crate::{error::Result, infrastructure::persistence::wal::WalSegmentStamp};
17
18/// What one [`EventStore::refresh_from_disk`] call added.
19#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)]
20pub struct RefreshReport {
21    /// Events not in memory before this call.
22    pub new_events: usize,
23    /// Parquet files read for the first time.
24    pub new_parquet_files: usize,
25    /// Whether any WAL segment changed, which forces a WAL re-read.
26    pub wal_changed: bool,
27    /// Newest timestamp among `new_events`.
28    pub newest_event_at: Option<DateTime<Utc>>,
29}
30
31/// The disk a replica has already read, so a refresh reads only the difference.
32#[derive(Default)]
33pub(crate) struct RefreshState {
34    seen_parquet: HashSet<PathBuf>,
35    wal_stamps: Option<Vec<WalSegmentStamp>>,
36}
37
38impl RefreshState {
39    pub(crate) fn mark_parquet_seen(&mut self, files: impl IntoIterator<Item = PathBuf>) {
40        self.seen_parquet.extend(files);
41    }
42
43    pub(crate) fn set_wal_stamps(&mut self, stamps: Vec<WalSegmentStamp>) {
44        self.wal_stamps = Some(stamps);
45    }
46}
47
48impl EventStore {
49    /// Load what a live writer has made durable since this replica last looked.
50    ///
51    /// Reads Parquet files this replica has not read yet, and re-reads the WAL
52    /// when any segment's size or mtime changed. Every event goes through the
53    /// id dedup in `append_loaded_event`, so overlap between the WAL, a fresh
54    /// checkpoint and a compacted file adds nothing twice. Never writes, never
55    /// creates or deletes a file.
56    ///
57    /// A no-op on a writable store: a writer already holds everything it wrote,
58    /// and `WriteAheadLog::recover` would reset its sequence counter.
59    pub fn refresh_from_disk(&self) -> Result<RefreshReport> {
60        let mut report = RefreshReport::default();
61        if !self.read_only {
62            return Ok(report);
63        }
64
65        let mut state = self.refresh_state.lock();
66        let mut loaded = Vec::new();
67        let mut tenants = HashSet::new();
68        // Marked seen only once every fallible step has passed: a failure below
69        // drops `loaded`, and its files must be read again next time.
70        let mut newly_seen = Vec::new();
71
72        if let Some(storage) = self.storage.as_ref().map(Arc::clone) {
73            let storage = storage.read();
74            let files = storage.list_parquet_files()?;
75            let present: HashSet<&Path> = files.iter().map(PathBuf::as_path).collect();
76            // Compaction removes files whose events are already in memory.
77            state
78                .seen_parquet
79                .retain(|path| present.contains(path.as_path()));
80            for path in &files {
81                if state.seen_parquet.contains(path) {
82                    continue;
83                }
84                let tenant = storage.tenant_id_for_file(path);
85                match storage.load_events_from_file_path(path, &tenant) {
86                    Ok(events) => {
87                        loaded.extend(events);
88                        newly_seen.push(path.clone());
89                        tenants.insert(tenant);
90                        report.new_parquet_files += 1;
91                    }
92                    Err(error) => tracing::debug!(
93                        file = %path.display(),
94                        error = %error,
95                        "refresh: Parquet file not readable yet; retrying on the next refresh"
96                    ),
97                }
98            }
99        }
100
101        if let Some(wal) = self.wal.as_ref() {
102            // Stamped before reading, so an append that lands mid-read changes
103            // the next stamp and is picked up then.
104            let stamps = wal.segment_stamps()?;
105            if state.wal_stamps.as_ref() != Some(&stamps) {
106                loaded.extend(wal.recover()?);
107                state.wal_stamps = Some(stamps);
108                report.wal_changed = true;
109            }
110        }
111        state.seen_parquet.extend(newly_seen);
112
113        let _resident = self.cache_residency_gate.read();
114        for event in loaded {
115            let at = event.timestamp;
116            if self.append_loaded_event(event) {
117                report.new_events += 1;
118                report.newest_event_at = report.newest_event_at.max(Some(at));
119            }
120        }
121        for tenant in &tenants {
122            self.tenant_loader.mark_loaded(tenant);
123        }
124
125        Ok(report)
126    }
127}