use crate::event::context::StageType;
use crate::journal::reader::JournalReader;
use crate::journal::{ArchiveStatus, StatusDerivation};
use crate::{ChainEvent, StageId};
use async_trait::async_trait;
use std::path::{Path, PathBuf};
use thiserror::Error;
#[derive(Debug, Error)]
pub enum ReplayError {
#[error("Replay archive path is not a directory: {path}")]
ArchivePathNotDirectory { path: PathBuf },
#[error("Replay archive is missing run_manifest.json at {path}")]
MissingManifest { path: PathBuf },
#[error("unsupported journal schema version: {journal_schema_version} (supported: {supported}); re-record the run with this build of ObzenFlow")]
UnsupportedJournalSchemaVersion {
journal_schema_version: String,
supported: &'static str,
},
#[error("Replay archive capability '{capability}' has unsupported version {found:?} (supported: {supported}); re-record the run with this build of ObzenFlow")]
UnsupportedArchiveCapability {
capability: &'static str,
found: Option<u64>,
supported: u32,
},
#[error("Replay archive system.log missing at {path}")]
MissingSystemLog { path: PathBuf },
#[error("Replay archive status is '{status:?}' and replay requires a completed or cancelled archive; re-run with --allow-incomplete-archive to override")]
IncompleteArchive { status: ArchiveStatus },
#[error("Stage '{stage_key}' not found in run manifest")]
StageNotInManifest { stage_key: String },
#[error("Stage '{stage_key}' is not a source in archive (archived: {archived_type:?}, expected: {expected_type:?})")]
StageTypeMismatch {
stage_key: String,
archived_type: StageType,
expected_type: StageType,
},
#[error("Replay archive journal missing at {path}")]
MissingJournal { path: PathBuf },
#[error("Replay archive journal appears corrupted at position {record_position} in {path}: {message}")]
CorruptedArchive {
path: PathBuf,
record_position: u64,
message: String,
},
#[error("Replay archive I/O error: {message}")]
Io {
message: String,
#[source]
source: std::io::Error,
},
#[error("Replay archive parse error: {message}")]
Parse { message: String },
}
#[async_trait]
pub trait ReplayArchive: Send + Sync {
async fn open_source_reader(
&self,
stage_key: &str,
expected_type: StageType,
) -> Result<Box<dyn JournalReader<ChainEvent>>, ReplayError>;
async fn open_effect_history(
&self,
stage_key: &str,
) -> Result<Box<dyn JournalReader<ChainEvent>>, ReplayError>;
fn source_data_journal_path(&self, stage_key: &str) -> Result<PathBuf, ReplayError>;
fn archive_flow_id(&self) -> &str;
fn archived_stage_id(&self, stage_key: &str) -> Result<StageId, ReplayError>;
fn archive_status(&self) -> ArchiveStatus;
fn status_derivation(&self) -> StatusDerivation;
fn allow_incomplete_archive(&self) -> bool;
fn source_stage_keys(&self) -> Vec<String>;
fn archive_path(&self) -> &Path;
fn manifest_capability(&self, _name: &str) -> Option<u32> {
None
}
fn bounded_direct_fact_admission(
&self,
) -> &[crate::journal::archive::manifest::RunManifestDirectFactAdmission] {
&[]
}
fn max_recorded_generation(&self) -> crate::ReaderGeneration {
crate::ReaderGeneration(0)
}
fn max_recorded_admission_seq(&self) -> crate::AdmissionSeq {
crate::AdmissionSeq(0)
}
}