use std::any::Any;
use std::fs::File;
use std::io::BufReader;
use std::path::Path;
use std::sync::{Arc, Mutex};
use lora_store::{GraphStorage, GraphStorageMut, InMemoryGraph, MutationRecorder};
use lora_wal::{replay_dir, Lsn, Wal, WalConfig, WalMirror, WalRecorder};
use crate::database::Database;
use crate::error::{LoraError, LoraErrorCode};
use crate::live_store::LiveStore;
use crate::named::{DatabaseName, DatabaseOpenOptions};
use crate::plan_cache::PlanCache;
use crate::snapshot::{ManagedSnapshotStore, SnapshotConfig};
use crate::wal::archive::WalArchive;
use super::replay::replay_into;
impl Database<InMemoryGraph> {
pub fn in_memory() -> Self {
Self::from_graph(InMemoryGraph::new())
}
pub fn open_with_wal(wal_config: WalConfig) -> Result<Self, LoraError> {
match wal_config {
WalConfig::Disabled => Ok(Self::in_memory()),
WalConfig::Enabled {
dir,
sync_mode,
segment_target_bytes,
} => {
let mut graph = InMemoryGraph::new();
let (wal, events) = Wal::open(dir, sync_mode, segment_target_bytes, Lsn::ZERO)?;
replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
let recorder = Arc::new(WalRecorder::new(wal));
Ok(Self::from_graph_with_wal(graph, recorder, None, None))
}
}
}
pub fn open_with_wal_snapshots(
wal_config: WalConfig,
snapshot_config: SnapshotConfig,
) -> Result<Self, LoraError> {
let snapshot_store =
Arc::new(ManagedSnapshotStore::open(snapshot_config).map_err(LoraError::from_anyhow)?);
let mut graph = InMemoryGraph::new();
match wal_config {
WalConfig::Disabled => Err(LoraError::new(
LoraErrorCode::Config,
"managed snapshots require WAL enabled",
)),
WalConfig::Enabled {
dir,
sync_mode,
segment_target_bytes,
} => {
let snapshot_lsn = snapshot_store
.load_latest(&mut graph)
.map_err(LoraError::from_anyhow)?;
let (wal, events) = Wal::open(dir, sync_mode, segment_target_bytes, snapshot_lsn)?;
replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
let recorder = Arc::new(WalRecorder::new(wal));
Ok(Self::from_graph_with_wal(
graph,
recorder,
Some(snapshot_store),
None,
))
}
}
}
pub fn open_named(
database_name: impl AsRef<str>,
options: DatabaseOpenOptions,
) -> Result<Self, LoraError> {
let name = DatabaseName::parse(database_name.as_ref())?;
let archive = Arc::new(WalArchive::open(
options.database_path_for(&name),
options.max_database_bytes,
)?);
let mut graph = InMemoryGraph::new();
let snapshot_lsn = if let Some(bytes) = archive.snapshot_bytes()? {
let (payload, info) = crate::snapshot::decode_snapshot_bytes(&bytes, None)?;
graph.load_snapshot_payload(payload)?;
info.wal_lsn.map(Lsn::new).unwrap_or(Lsn::ZERO)
} else {
Lsn::ZERO
};
let (wal, events) = Wal::open(
archive.work_dir(),
options.sync_mode,
options.segment_target_bytes,
snapshot_lsn,
)?;
replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
let mirror: Arc<dyn WalMirror> = archive.clone();
let recorder = Arc::new(WalRecorder::new_with_mirror(wal, Some(mirror)));
recorder.flush()?;
Ok(Self::from_graph_with_wal(
graph,
recorder,
None,
Some(archive),
))
}
pub fn recover(
snapshot_path: impl AsRef<Path>,
wal_config: WalConfig,
) -> Result<Self, LoraError> {
let snapshot_path = snapshot_path.as_ref();
let mut graph = InMemoryGraph::new();
let snapshot_lsn = match File::open(snapshot_path) {
Ok(f) => {
let reader = BufReader::new(f);
let (payload, info) = crate::snapshot::read_snapshot_from(reader, None)?;
graph.load_snapshot_payload(payload)?;
info.wal_lsn.map(Lsn::new).unwrap_or(Lsn::ZERO)
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Lsn::ZERO,
Err(e) => return Err(e.into()),
};
match wal_config {
WalConfig::Disabled => Ok(Self::from_graph(graph)),
WalConfig::Enabled {
dir,
sync_mode,
segment_target_bytes,
} => {
if dir.exists() {
if let Ok(outcome) = replay_dir(&dir, Lsn::ZERO) {
if let Some(marker) = outcome.checkpoint_lsn_observed {
if marker > snapshot_lsn {
eprintln!(
"lora-wal: snapshot at LSN {} is older than the newest \
checkpoint marker on disk (LSN {}). Replaying every WAL \
record above LSN {}; consider passing the more recent \
snapshot to --restore-from.",
snapshot_lsn.raw(),
marker.raw(),
snapshot_lsn.raw()
);
}
}
}
}
let (wal, events) = Wal::open(dir, sync_mode, segment_target_bytes, snapshot_lsn)?;
replay_into(&mut graph, events).map_err(LoraError::from_anyhow)?;
let recorder = Arc::new(WalRecorder::new(wal));
Ok(Self::from_graph_with_wal(graph, recorder, None, None))
}
}
}
fn from_graph_with_wal(
mut graph: InMemoryGraph,
recorder: Arc<WalRecorder>,
snapshots: Option<Arc<ManagedSnapshotStore>>,
named_archive: Option<Arc<WalArchive>>,
) -> Self {
graph.set_mutation_recorder(Some(recorder.clone() as Arc<dyn MutationRecorder>));
Self {
store: Arc::new(LiveStore::new(Arc::new(graph))),
writer: Arc::new(Mutex::new(())),
lock_table: Arc::new(lora_store::LockTable::new()),
wal: Some(recorder),
snapshots,
named_archive,
plan_cache: Arc::new(PlanCache::new()),
changes: Arc::new(crate::changes::ChangeHub::default()),
}
}
}
impl<S> Database<S>
where
S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
{
pub(crate) fn new(store: Arc<LiveStore<S>>) -> Self {
Self {
store,
writer: Arc::new(Mutex::new(())),
lock_table: Arc::new(lora_store::LockTable::new()),
wal: None,
snapshots: None,
named_archive: None,
plan_cache: Arc::new(PlanCache::new()),
changes: Arc::new(crate::changes::ChangeHub::default()),
}
}
pub fn from_graph(graph: S) -> Self {
Self::new(Arc::new(LiveStore::new(Arc::new(graph))))
}
}