use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use chrono::{NaiveDate, Utc};
use tokio::sync::mpsc;
use trusty_common::control_bus::HarnessEvent;
use super::config::{LogConfig, WRITE_CHANNEL_CAPACITY, day_file_name};
use super::error::LogError;
use super::format::encode_line;
use super::recovery::{earliest_seq, recover_next_seq};
use super::replay::{ReplayItem, replay_since};
use super::retention::enforce_retention;
#[derive(Debug, Clone, Copy)]
pub(crate) struct RecoveredState {
pub next_seq: u64,
}
pub(crate) struct DurableLog {
dir: PathBuf,
tx: mpsc::Sender<HarnessEvent>,
#[allow(dead_code)]
written: Arc<AtomicU64>,
earliest_retained_seq: Arc<AtomicU64>,
}
impl DurableLog {
pub(crate) async fn open(config: LogConfig) -> Result<(Self, RecoveredState), LogError> {
Self::open_with_capacity(config, WRITE_CHANNEL_CAPACITY).await
}
pub(crate) async fn open_with_capacity(
config: LogConfig,
channel_capacity: usize,
) -> Result<(Self, RecoveredState), LogError> {
trusty_common::uds::prepare_socket_dir(&config.dir).map_err(|source| {
LogError::PrepareDir {
path: config.dir.clone(),
source,
}
})?;
let today = Utc::now().date_naive();
let survivors = enforce_retention(&config.dir, today, config.retain_days).await?;
let next_seq = recover_next_seq(&survivors).await?;
let earliest = earliest_seq(&survivors).await?;
let (tx, rx) = mpsc::channel(channel_capacity.max(1));
let written = Arc::new(AtomicU64::new(0));
let earliest_retained_seq = Arc::new(AtomicU64::new(earliest.unwrap_or(0)));
tokio::spawn(run_writer(
config.clone(),
rx,
Arc::clone(&written),
Arc::clone(&earliest_retained_seq),
));
Ok((
Self {
dir: config.dir,
tx,
written,
earliest_retained_seq,
},
RecoveredState { next_seq },
))
}
pub(crate) fn enqueue(&self, event: HarnessEvent) -> bool {
self.tx.try_send(event).is_ok()
}
#[cfg(test)]
pub(crate) fn written(&self) -> u64 {
self.written.load(Ordering::Relaxed)
}
pub(crate) async fn replay_since(&self, since_seq: u64) -> Result<Vec<ReplayItem>, LogError> {
replay_since(&self.dir, since_seq, self.earliest_retained_seq()).await
}
fn earliest_retained_seq(&self) -> Option<u64> {
match self.earliest_retained_seq.load(Ordering::Relaxed) {
0 => None,
n => Some(n),
}
}
}
async fn run_writer(
config: LogConfig,
mut rx: mpsc::Receiver<HarnessEvent>,
written: Arc<AtomicU64>,
earliest_retained_seq: Arc<AtomicU64>,
) {
let mut open_day: Option<NaiveDate> = None;
let mut file: Option<tokio::fs::File> = None;
let mut current_path: Option<PathBuf> = None;
while let Some(event) = rx.recv().await {
let today = Utc::now().date_naive();
if open_day != Some(today) {
match rotate(&config, today, &earliest_retained_seq, file.as_mut()).await {
Ok((f, path)) => {
file = Some(f);
current_path = Some(path);
open_day = Some(today);
}
Err(e) => {
tracing::error!(
error = %e,
"event log: could not rotate to today's file; this event \
will not be persisted"
);
continue;
}
}
}
let (Some(f), Some(path)) = (file.as_mut(), current_path.as_deref()) else {
continue;
};
if let Err(e) = write_line(f, path, &event).await {
tracing::error!(
error = %e,
seq = event.seq,
"event log: write failed; this event will not be persisted"
);
continue;
}
written.fetch_add(1, Ordering::Relaxed);
}
if let Some(f) = file.as_mut() {
let _ = tokio::io::AsyncWriteExt::flush(f).await;
let _ = f.sync_data().await;
}
}
pub(super) async fn rotate(
config: &LogConfig,
today: NaiveDate,
earliest_retained_seq: &Arc<AtomicU64>,
previous: Option<&mut tokio::fs::File>,
) -> Result<(tokio::fs::File, PathBuf), LogError> {
if let Some(f) = previous {
let _ = tokio::io::AsyncWriteExt::flush(f).await;
if let Err(source) = f.sync_data().await {
tracing::error!(
error = %source,
"event log: fsync of the outgoing day file failed at rotation; \
its un-synced tail may not survive a crash"
);
}
}
let survivors = enforce_retention(&config.dir, today, config.retain_days).await?;
let earliest = earliest_seq(&survivors).await?;
earliest_retained_seq.store(earliest.unwrap_or(0), Ordering::Relaxed);
let path = config.dir.join(day_file_name(today));
let file = open_append_0600(&path).await?;
Ok((file, path))
}
#[cfg(unix)]
async fn open_append_0600(path: &std::path::Path) -> Result<tokio::fs::File, LogError> {
use std::os::unix::fs::PermissionsExt as _;
let file = tokio::fs::OpenOptions::new()
.create(true)
.append(true)
.mode(trusty_common::uds::SOCKET_MODE)
.open(path)
.await
.map_err(|source| LogError::Io {
op: "open",
path: path.to_path_buf(),
source,
})?;
if let Err(e) = file
.set_permissions(std::fs::Permissions::from_mode(
trusty_common::uds::SOCKET_MODE,
))
.await
{
tracing::warn!(
error = %e,
path = %path.display(),
"event log: could not re-assert 0600 permissions on an existing file"
);
}
Ok(file)
}
#[cfg(not(unix))]
async fn open_append_0600(path: &std::path::Path) -> Result<tokio::fs::File, LogError> {
tokio::fs::OpenOptions::new()
.create(true)
.append(true)
.open(path)
.await
.map_err(|source| LogError::Io {
op: "open",
path: path.to_path_buf(),
source,
})
}
pub(super) async fn write_line(
file: &mut tokio::fs::File,
path: &Path,
event: &HarnessEvent,
) -> Result<(), LogError> {
use tokio::io::AsyncWriteExt as _;
let line = encode_line(event);
file.write_all(&line).await.map_err(|source| LogError::Io {
op: "write",
path: path.to_path_buf(),
source,
})?;
file.flush().await.map_err(|source| LogError::Io {
op: "flush",
path: path.to_path_buf(),
source,
})
}