use super::storage::MemoryStorage;
use crate::model::Quad;
use crate::parser::RdfFormat;
use crate::serializer::Serializer;
use crate::{OxirsError, Result};
use std::fs::{File, OpenOptions};
use std::io::{BufWriter, Read, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::Mutex;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SyncPolicy {
EveryN(u64),
OnFlush,
}
impl Default for SyncPolicy {
fn default() -> Self {
SyncPolicy::EveryN(1000)
}
}
#[derive(Debug)]
pub struct PersistentState {
path: PathBuf,
data_file: PathBuf,
append: Mutex<Option<BufWriter<File>>>,
dirty: AtomicBool,
ops_since_sync: AtomicU64,
sync_policy: SyncPolicy,
load_had_errors: AtomicBool,
}
impl PersistentState {
pub fn open(path: PathBuf, sync_policy: SyncPolicy, load_had_errors: bool) -> Result<Self> {
std::fs::create_dir_all(&path).map_err(|e| {
OxirsError::Io(format!(
"Failed to create directory {}: {e}",
path.display()
))
})?;
let data_file = path.join("data.nq");
let writer = Self::open_append_writer(&data_file)?;
Ok(Self {
path,
data_file,
append: Mutex::new(Some(writer)),
dirty: AtomicBool::new(false),
ops_since_sync: AtomicU64::new(0),
sync_policy,
load_had_errors: AtomicBool::new(load_had_errors),
})
}
pub fn path(&self) -> &Path {
&self.path
}
fn open_append_writer(data_file: &Path) -> Result<BufWriter<File>> {
let file = OpenOptions::new()
.create(true)
.append(true)
.open(data_file)
.map_err(|e| {
OxirsError::Io(format!(
"Failed to open append log {}: {e}",
data_file.display()
))
})?;
Ok(BufWriter::new(file))
}
pub fn mark_dirty(&self) {
self.dirty.store(true, Ordering::SeqCst);
}
pub fn is_dirty(&self) -> bool {
self.dirty.load(Ordering::SeqCst)
}
pub fn append_line(&self, line: &str) -> Result<()> {
let mut guard = self
.append
.lock()
.map_err(|e| OxirsError::Store(format!("append lock poisoned: {e}")))?;
if guard.is_none() {
*guard = Some(Self::open_append_writer(&self.data_file)?);
}
let writer = guard
.as_mut()
.ok_or_else(|| OxirsError::Store("append writer unavailable".to_string()))?;
writeln!(writer, "{line}")
.map_err(|e| OxirsError::Io(format!("Failed to append quad: {e}")))?;
let ops = self.ops_since_sync.fetch_add(1, Ordering::SeqCst) + 1;
if let SyncPolicy::EveryN(n) = self.sync_policy {
if n > 0 && ops % n == 0 {
Self::flush_writer(writer)?;
self.ops_since_sync.store(0, Ordering::SeqCst);
}
}
Ok(())
}
pub fn append_lines(&self, lines: &[String]) -> Result<()> {
if lines.is_empty() {
return Ok(());
}
let mut guard = self
.append
.lock()
.map_err(|e| OxirsError::Store(format!("append lock poisoned: {e}")))?;
if guard.is_none() {
*guard = Some(Self::open_append_writer(&self.data_file)?);
}
let writer = guard
.as_mut()
.ok_or_else(|| OxirsError::Store("append writer unavailable".to_string()))?;
for line in lines {
writeln!(writer, "{line}")
.map_err(|e| OxirsError::Io(format!("Failed to append quad: {e}")))?;
}
Self::flush_writer(writer)?;
self.ops_since_sync.store(0, Ordering::SeqCst);
Ok(())
}
fn flush_writer(writer: &mut BufWriter<File>) -> Result<()> {
writer
.flush()
.map_err(|e| OxirsError::Io(format!("Failed to flush append log: {e}")))?;
writer
.get_ref()
.sync_all()
.map_err(|e| OxirsError::Io(format!("Failed to fsync append log: {e}")))?;
Ok(())
}
pub fn flush(&self, storage: &MemoryStorage) -> Result<()> {
if self.is_dirty() {
self.compact(storage)
} else {
let mut guard = self
.append
.lock()
.map_err(|e| OxirsError::Store(format!("append lock poisoned: {e}")))?;
if let Some(writer) = guard.as_mut() {
Self::flush_writer(writer)?;
}
self.ops_since_sync.store(0, Ordering::SeqCst);
Ok(())
}
}
pub fn compact(&self, storage: &MemoryStorage) -> Result<()> {
if self.load_had_errors.load(Ordering::SeqCst) {
return Err(OxirsError::Store(format!(
"Refusing to compact {}: the initial load reported parse/open errors, so \
rewriting the file could discard data that failed to load. Fix or remove the \
corrupt data.nq and reopen the store.",
self.data_file.display()
)));
}
let tmp_file = self.data_file.with_extension("nq.tmp");
{
let file = File::create(&tmp_file).map_err(|e| {
OxirsError::Io(format!("Failed to create {}: {e}", tmp_file.display()))
})?;
let mut writer = BufWriter::new(file);
let serializer = Serializer::new(RdfFormat::NQuads);
for quad in storage.iter_quads() {
let line = serializer.serialize_quad_to_nquads(&quad)?;
writeln!(writer, "{line}").map_err(|e| {
OxirsError::Io(format!("Failed to write {}: {e}", tmp_file.display()))
})?;
}
writer.flush().map_err(|e| {
OxirsError::Io(format!("Failed to flush {}: {e}", tmp_file.display()))
})?;
writer.get_ref().sync_all().map_err(|e| {
OxirsError::Io(format!("Failed to fsync {}: {e}", tmp_file.display()))
})?;
}
let mut guard = self
.append
.lock()
.map_err(|e| OxirsError::Store(format!("append lock poisoned: {e}")))?;
*guard = None;
std::fs::rename(&tmp_file, &self.data_file).map_err(|e| {
if let Ok(w) = Self::open_append_writer(&self.data_file) {
*guard = Some(w);
}
OxirsError::Io(format!(
"Failed to atomically replace {}: {e}",
self.data_file.display()
))
})?;
*guard = Some(Self::open_append_writer(&self.data_file)?);
self.dirty.store(false, Ordering::SeqCst);
self.ops_since_sync.store(0, Ordering::SeqCst);
Ok(())
}
}
pub fn load_from_disk(data_file: &Path) -> Result<(MemoryStorage, bool)> {
let mut storage = MemoryStorage::new();
let mut file = File::open(data_file)
.map_err(|e| OxirsError::Io(format!("Failed to open {}: {e}", data_file.display())))?;
let mut bytes = Vec::new();
file.read_to_end(&mut bytes)
.map_err(|e| OxirsError::Io(format!("Failed to read {}: {e}", data_file.display())))?;
drop(file);
if bytes.is_empty() {
return Ok((storage, false));
}
let ends_with_newline = bytes.last() == Some(&b'\n');
let segments: Vec<&[u8]> = bytes.split(|&b| b == b'\n').collect();
let seg_count = segments.len();
let mut had_errors = false;
let mut torn_trailing: Option<usize> = None;
for (idx, seg) in segments.iter().enumerate() {
let is_last_segment = idx + 1 == seg_count;
if seg.iter().all(|b| b.is_ascii_whitespace()) {
continue;
}
match parse_nquads_line(seg) {
Ok(quads) => {
for quad in quads {
storage.insert_quad(quad);
}
}
Err(msg) => {
if is_last_segment && !ends_with_newline {
tracing::warn!(
"Dropping torn trailing line ({} bytes) in {}: {msg}",
seg.len(),
data_file.display()
);
torn_trailing = Some(seg.len());
} else {
had_errors = true;
tracing::warn!(
"Skipping unparseable line {} in {}: {msg}",
idx + 1,
data_file.display()
);
}
}
}
}
if let Some(torn_len) = torn_trailing {
let new_len = bytes.len().saturating_sub(torn_len) as u64;
match OpenOptions::new().write(true).open(data_file) {
Ok(f) => {
if let Err(e) = f.set_len(new_len) {
tracing::error!(
"Failed to truncate torn trailing line in {}: {e}",
data_file.display()
);
}
}
Err(e) => tracing::error!(
"Failed to open {} to truncate torn trailing line: {e}",
data_file.display()
),
}
} else if !ends_with_newline {
match OpenOptions::new().append(true).open(data_file) {
Ok(mut f) => {
if let Err(e) = f.write_all(b"\n") {
tracing::error!(
"Failed to normalize trailing newline in {}: {e}",
data_file.display()
);
}
}
Err(e) => tracing::error!(
"Failed to open {} to normalize trailing newline: {e}",
data_file.display()
),
}
}
Ok((storage, had_errors))
}
fn parse_nquads_line(line: &[u8]) -> std::result::Result<Vec<Quad>, String> {
use crate::format::{RdfFormat as FormatRdfFormat, RdfParser};
match line.iter().position(|&b| !b.is_ascii_whitespace()) {
None => return Ok(Vec::new()),
Some(i) if line[i] == b'#' => return Ok(Vec::new()),
_ => {}
}
let parser = RdfParser::new(FormatRdfFormat::NQuads);
let mut quads = Vec::new();
for item in parser.for_slice(line) {
match item {
Ok(quad) => quads.push(quad),
Err(e) => return Err(e.to_string()),
}
}
if quads.is_empty() {
return Err("line produced no quads".to_string());
}
Ok(quads)
}