allsource_core/
store_refresh.rs1use 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#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)]
20pub struct RefreshReport {
21 pub new_events: usize,
23 pub new_parquet_files: usize,
25 pub wal_changed: bool,
27 pub newest_event_at: Option<DateTime<Utc>>,
29}
30
31#[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 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 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 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 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}