use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use serde::{Deserialize, Serialize};
use crate::events::{
BrokerClock, DurableEventJournal, JournalConfig, JournalError, LifecycleStatus, SystemClock,
TraverseEvent,
};
use super::{PrivateTraceEntry, PublicTraceEntry};
const TRACE_RECORD_EVENT_TYPE: &str = "dev.traverse.trace.recorded";
const TRACE_RECORD_OWNER: &str = "traverse-runtime";
const TRACE_RECORD_VERSION: &str = "1.0.0";
const RECOVERY_REPLAY_PAGE_SIZE: usize = 256;
#[derive(Debug, PartialEq, Eq)]
pub enum TraceJournalError {
Journal(JournalError),
Serialize(String),
Deserialize(String),
LockPoisoned,
}
impl std::fmt::Display for TraceJournalError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Journal(error) => write!(f, "trace journal failure: {error}"),
Self::Serialize(message) => write!(f, "trace record serialization failed: {message}"),
Self::Deserialize(message) => {
write!(f, "durable trace record deserialization failed: {message}")
}
Self::LockPoisoned => write!(f, "durable trace journal lock is poisoned"),
}
}
}
impl std::error::Error for TraceJournalError {}
impl From<JournalError> for TraceJournalError {
fn from(error: JournalError) -> Self {
Self::Journal(error)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
struct DurableTraceRecord {
public: PublicTraceEntry,
#[serde(default, skip_serializing_if = "Option::is_none")]
private: Option<PrivateTraceEntry>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct TraceRecoveryReport {
pub recovered_trace_ids: Vec<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct TracePrunedEvidence {
pub workspace_root: PathBuf,
pub deleted_segment_paths: Vec<PathBuf>,
}
pub struct DurableTraceJournal {
root: PathBuf,
journal: DurableEventJournal,
recovery: TraceRecoveryReport,
}
impl DurableTraceJournal {
pub fn open(root: &Path, config: JournalConfig) -> Result<Self, TraceJournalError> {
Self::open_with_clock(root, config, Arc::new(SystemClock))
}
pub fn open_with_clock(
root: &Path,
config: JournalConfig,
clock: Arc<dyn BrokerClock>,
) -> Result<Self, TraceJournalError> {
let journal = DurableEventJournal::open(root, config, clock)?;
let recovery = Self::compute_recovery_report(&journal)?;
Ok(Self {
root: root.to_path_buf(),
journal,
recovery,
})
}
fn compute_recovery_report(
journal: &DurableEventJournal,
) -> Result<TraceRecoveryReport, TraceJournalError> {
let mut recovered_trace_ids = Vec::new();
let mut cursor = "0".to_string();
loop {
let page = journal.replay_from(&cursor, RECOVERY_REPLAY_PAGE_SIZE)?;
if page.is_empty() {
break;
}
for (next_cursor, event) in &page {
cursor.clone_from(next_cursor);
let record: DurableTraceRecord = serde_json::from_value(event.data.clone())
.map_err(|error| TraceJournalError::Deserialize(error.to_string()))?;
recovered_trace_ids.push(record.public.id);
}
}
Ok(TraceRecoveryReport {
recovered_trace_ids,
})
}
#[must_use]
pub fn recovery_report(&self) -> &TraceRecoveryReport {
&self.recovery
}
pub fn record(
&mut self,
public: &PublicTraceEntry,
private: Option<&PrivateTraceEntry>,
) -> Result<String, TraceJournalError> {
let record = DurableTraceRecord {
public: public.clone(),
private: private.cloned(),
};
let data = serde_json::to_value(&record)
.map_err(|error| TraceJournalError::Serialize(error.to_string()))?;
let event = TraverseEvent {
id: public.id.clone(),
source: public.source.clone(),
event_type: TRACE_RECORD_EVENT_TYPE.to_string(),
datacontenttype: "application/json".to_string(),
time: public.time.clone(),
data,
owner: TRACE_RECORD_OWNER.to_string(),
version: TRACE_RECORD_VERSION.to_string(),
lifecycle_status: LifecycleStatus::Active,
deduplication_id: None,
ordering_scope: None,
correlation_id: None,
causation_id: None,
subject_id: None,
actor_id: None,
};
Ok(self.journal.append(&event)?)
}
pub fn prune(&mut self) -> Result<TracePrunedEvidence, TraceJournalError> {
let deleted_segment_paths = self.journal.prune()?;
Ok(TracePrunedEvidence {
workspace_root: self.root.clone(),
deleted_segment_paths,
})
}
}
pub trait TraceDurabilitySink: Send + Sync {
fn record(
&self,
public: &PublicTraceEntry,
private: Option<&PrivateTraceEntry>,
) -> Result<String, TraceJournalError>;
}
impl TraceDurabilitySink for Mutex<DurableTraceJournal> {
fn record(
&self,
public: &PublicTraceEntry,
private: Option<&PrivateTraceEntry>,
) -> Result<String, TraceJournalError> {
let mut journal = self.lock().map_err(|_| TraceJournalError::LockPoisoned)?;
journal.record(public, private)
}
}
pub struct DurableTraceConfig {
pub sink: Arc<dyn TraceDurabilitySink>,
pub fail_closed: bool,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn trace_journal_error_display_covers_every_variant() {
let cases: Vec<TraceJournalError> = vec![
TraceJournalError::Journal(JournalError::InvalidCursor("bad cursor".to_string())),
TraceJournalError::Serialize("boom".to_string()),
TraceJournalError::Deserialize("boom".to_string()),
TraceJournalError::LockPoisoned,
];
for error in &cases {
assert!(
!error.to_string().is_empty(),
"Display must produce a non-empty string for {error:?}"
);
}
}
}