1use std::any::Any;
19use std::fs::File;
20use std::io::BufReader;
21use std::path::Path;
22use std::sync::{Arc, Mutex};
23
24use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph, MutationRecorder};
25use lora_wal::{replay_dir, Lsn, Wal, WalConfig, WalMirror, WalRecorder};
26
27use crate::database::Database;
28use crate::error::{LoraError, LoraErrorCode};
29use crate::live_store::LiveStore;
30use crate::named::{DatabaseName, DatabaseOpenOptions};
31use crate::plan_cache::PlanCache;
32use crate::snapshot::{ManagedSnapshotStore, SnapshotConfig};
33use crate::wal::archive::WalArchive;
34
35use super::replay::replay_into;
36
37impl Database<InMemoryGraph> {
38 pub fn in_memory() -> Self {
40 Self::from_graph(InMemoryGraph::new())
41 }
42
43 pub fn open_with_wal(wal_config: WalConfig) -> Result<Self, LoraError> {
54 match wal_config {
55 WalConfig::Disabled => Ok(Self::in_memory()),
56 WalConfig::Enabled {
57 dir,
58 sync_mode,
59 segment_target_bytes,
60 } => {
61 let mut graph = InMemoryGraph::new();
62 let (wal, events) = Wal::open(dir, sync_mode, segment_target_bytes, Lsn::ZERO)?;
63 replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
64 let recorder = Arc::new(WalRecorder::new(wal));
65 Ok(Self::from_graph_with_wal(graph, recorder, None, None))
66 }
67 }
68 }
69
70 pub fn open_with_wal_snapshots(
77 wal_config: WalConfig,
78 snapshot_config: SnapshotConfig,
79 ) -> Result<Self, LoraError> {
80 let snapshot_store =
81 Arc::new(ManagedSnapshotStore::open(snapshot_config).map_err(LoraError::from_anyhow)?);
82 let mut graph = InMemoryGraph::new();
83
84 match wal_config {
85 WalConfig::Disabled => Err(LoraError::new(
86 LoraErrorCode::Config,
87 "managed snapshots require WAL enabled",
88 )),
89 WalConfig::Enabled {
90 dir,
91 sync_mode,
92 segment_target_bytes,
93 } => {
94 let snapshot_lsn = snapshot_store
95 .load_latest(&mut graph)
96 .map_err(LoraError::from_anyhow)?;
97 let (wal, events) = Wal::open(dir, sync_mode, segment_target_bytes, snapshot_lsn)?;
98 replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
99 let recorder = Arc::new(WalRecorder::new(wal));
100 Ok(Self::from_graph_with_wal(
101 graph,
102 recorder,
103 Some(snapshot_store),
104 None,
105 ))
106 }
107 }
108 }
109
110 pub fn open_named(
117 database_name: impl AsRef<str>,
118 options: DatabaseOpenOptions,
119 ) -> Result<Self, LoraError> {
120 let name = DatabaseName::parse(database_name.as_ref())?;
121 let archive = Arc::new(WalArchive::open(
122 options.database_path_for(&name),
123 options.max_database_bytes,
124 )?);
125 let mut graph = InMemoryGraph::new();
126 let snapshot_lsn = if let Some(bytes) = archive.snapshot_bytes()? {
127 let (payload, info) = crate::snapshot::decode_snapshot_bytes(&bytes, None)?;
128 graph.load_snapshot_payload(payload)?;
129 info.wal_lsn.map(Lsn::new).unwrap_or(Lsn::ZERO)
130 } else {
131 Lsn::ZERO
132 };
133 let (wal, events) = Wal::open(
134 archive.work_dir(),
135 options.sync_mode,
136 options.segment_target_bytes,
137 snapshot_lsn,
138 )?;
139 replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
140 let mirror: Arc<dyn WalMirror> = archive.clone();
141 let recorder = Arc::new(WalRecorder::new_with_mirror(wal, Some(mirror)));
142 recorder.flush()?;
146 Ok(Self::from_graph_with_wal(
147 graph,
148 recorder,
149 None,
150 Some(archive),
151 ))
152 }
153
154 pub fn recover(
173 snapshot_path: impl AsRef<Path>,
174 wal_config: WalConfig,
175 ) -> Result<Self, LoraError> {
176 let snapshot_path = snapshot_path.as_ref();
177 let mut graph = InMemoryGraph::new();
178 let snapshot_lsn = match File::open(snapshot_path) {
179 Ok(f) => {
180 let reader = BufReader::new(f);
181 let (payload, info) = crate::snapshot::read_snapshot_from(reader, None)?;
182 graph.load_snapshot_payload(payload)?;
183 info.wal_lsn.map(Lsn::new).unwrap_or(Lsn::ZERO)
184 }
185 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Lsn::ZERO,
186 Err(e) => return Err(e.into()),
187 };
188
189 match wal_config {
190 WalConfig::Disabled => Ok(Self::from_graph(graph)),
191 WalConfig::Enabled {
192 dir,
193 sync_mode,
194 segment_target_bytes,
195 } => {
196 if dir.exists() {
202 if let Ok(outcome) = replay_dir(&dir, Lsn::ZERO) {
203 if let Some(marker) = outcome.checkpoint_lsn_observed {
204 if marker > snapshot_lsn {
205 eprintln!(
206 "lora-wal: snapshot at LSN {} is older than the newest \
207 checkpoint marker on disk (LSN {}). Replaying every WAL \
208 record above LSN {}; consider passing the more recent \
209 snapshot to --restore-from.",
210 snapshot_lsn.raw(),
211 marker.raw(),
212 snapshot_lsn.raw()
213 );
214 }
215 }
216 }
217 }
218
219 let (wal, events) = Wal::open(dir, sync_mode, segment_target_bytes, snapshot_lsn)?;
220 replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
221 let recorder = Arc::new(WalRecorder::new(wal));
222 Ok(Self::from_graph_with_wal(graph, recorder, None, None))
223 }
224 }
225 }
226
227 fn from_graph_with_wal(
230 mut graph: InMemoryGraph,
231 recorder: Arc<WalRecorder>,
232 snapshots: Option<Arc<ManagedSnapshotStore>>,
233 named_archive: Option<Arc<WalArchive>>,
234 ) -> Self {
235 graph.set_mutation_recorder(Some(recorder.clone() as Arc<dyn MutationRecorder>));
236 Self {
237 store: Arc::new(LiveStore::new(Arc::new(graph))),
238 writer: Arc::new(Mutex::new(())),
239 lock_table: Arc::new(lora_store::LockTable::new()),
240 wal: Some(recorder),
241 snapshots,
242 named_archive,
243 plan_cache: Arc::new(PlanCache::new()),
244 }
245 }
246}
247
248impl<S> Database<S>
249where
250 S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
251{
252 pub(crate) fn new(store: Arc<LiveStore<S>>) -> Self {
254 Self {
255 store,
256 writer: Arc::new(Mutex::new(())),
257 lock_table: Arc::new(lora_store::LockTable::new()),
258 wal: None,
259 snapshots: None,
260 named_archive: None,
261 plan_cache: Arc::new(PlanCache::new()),
262 }
263 }
264
265 pub fn from_graph(graph: S) -> Self {
267 Self::new(Arc::new(LiveStore::new(Arc::new(graph))))
268 }
269}