use crate::metrics::observations::ObservationRegistry;
use obzenflow_core::event::observability::ObservationRecorder;
use obzenflow_core::{FlowId, WriterId};
use std::collections::{HashMap, HashSet};
use std::sync::{Arc, Mutex};
use obzenflow_core::journal::ArchiveStatus;
use obzenflow_core::{MiddlewareExecutionScope, ReaderGeneration, StageId};
use crate::messaging::upstream_subscription::StageInputPosition;
use obzenflow_core::journal::archive::ReplayArchive;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RuntimeMode {
Live,
Replay,
Resume,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct ExecutionPosition {
pub stage_id: StageId,
pub position: StageInputPosition,
pub generation: Option<ReaderGeneration>,
}
#[derive(Debug, Clone, Copy)]
pub enum ExecutionPositionSource {
Data {
stage_id: StageId,
position: StageInputPosition,
generation: Option<ReaderGeneration>,
},
FlowControl {
stage_id: StageId,
cause: Option<ExecutionPosition>,
},
StageLifecycle { stage_id: StageId },
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SourceExecutionPhase {
Live,
Replaying,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SourceReplayExhaustion {
Terminate,
ContinueLive,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum HeartbeatExecutionPolicy {
Active,
Suppressed,
DormantUntilLive,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum EffectPortRegistrationPolicy {
Required,
OptionalStrictReplay,
}
pub trait ExecutionStrategy: std::fmt::Debug + Send + Sync {
fn scope_at(&self, at: ExecutionPosition) -> MiddlewareExecutionScope;
fn stage_scope(&self, stage: StageId) -> MiddlewareExecutionScope;
fn scope_for_generation(
&self,
stage: StageId,
_generation: ReaderGeneration,
) -> MiddlewareExecutionScope {
self.stage_scope(stage)
}
fn missing_outcome_is_corruption(&self, at: ExecutionPosition) -> bool;
fn in_doubt_effect_is_fatal(&self) -> bool;
fn effect_port_registration_policy(&self) -> EffectPortRegistrationPolicy;
fn source_phase_for(&self, stage: StageId) -> SourceExecutionPhase;
fn source_replay_exhausted(&self, stage: StageId) -> SourceReplayExhaustion;
fn heartbeat_policy_for(&self, stage: StageId) -> HeartbeatExecutionPolicy;
}
#[derive(Debug)]
struct Live;
impl ExecutionStrategy for Live {
fn scope_at(&self, _at: ExecutionPosition) -> MiddlewareExecutionScope {
MiddlewareExecutionScope::LiveHandler
}
fn stage_scope(&self, _stage: StageId) -> MiddlewareExecutionScope {
MiddlewareExecutionScope::LiveHandler
}
fn missing_outcome_is_corruption(&self, _at: ExecutionPosition) -> bool {
false
}
fn in_doubt_effect_is_fatal(&self) -> bool {
false
}
fn effect_port_registration_policy(&self) -> EffectPortRegistrationPolicy {
EffectPortRegistrationPolicy::Required
}
fn source_phase_for(&self, _stage: StageId) -> SourceExecutionPhase {
SourceExecutionPhase::Live
}
fn source_replay_exhausted(&self, _stage: StageId) -> SourceReplayExhaustion {
SourceReplayExhaustion::Terminate
}
fn heartbeat_policy_for(&self, _stage: StageId) -> HeartbeatExecutionPolicy {
HeartbeatExecutionPolicy::Active
}
}
#[derive(Debug)]
struct Replay {
incomplete: bool,
}
impl Replay {
fn reconstruction_scope(&self) -> MiddlewareExecutionScope {
if self.incomplete {
MiddlewareExecutionScope::ResumeHandler
} else {
MiddlewareExecutionScope::StrictReplayHandler
}
}
}
impl ExecutionStrategy for Replay {
fn scope_at(&self, _at: ExecutionPosition) -> MiddlewareExecutionScope {
self.reconstruction_scope()
}
fn stage_scope(&self, _stage: StageId) -> MiddlewareExecutionScope {
self.reconstruction_scope()
}
fn missing_outcome_is_corruption(&self, _at: ExecutionPosition) -> bool {
!self.incomplete
}
fn in_doubt_effect_is_fatal(&self) -> bool {
!self.incomplete
}
fn effect_port_registration_policy(&self) -> EffectPortRegistrationPolicy {
if self.incomplete {
EffectPortRegistrationPolicy::Required
} else {
EffectPortRegistrationPolicy::OptionalStrictReplay
}
}
fn source_phase_for(&self, _stage: StageId) -> SourceExecutionPhase {
SourceExecutionPhase::Replaying
}
fn source_replay_exhausted(&self, _stage: StageId) -> SourceReplayExhaustion {
SourceReplayExhaustion::Terminate
}
fn heartbeat_policy_for(&self, _stage: StageId) -> HeartbeatExecutionPolicy {
HeartbeatExecutionPolicy::Suppressed
}
}
#[derive(Debug)]
pub struct ResumeState {
resume_generation: ReaderGeneration,
frontier_by_stage: Mutex<HashMap<StageId, ReaderGeneration>>,
infinite_sources: Mutex<HashSet<StageId>>,
recorded_effect_seq_max: Mutex<HashMap<StageId, StageInputPosition>>,
recorded_delivered_high_water: Mutex<HashMap<StageId, u64>>,
}
impl ResumeState {
fn new(resume_generation: ReaderGeneration) -> Self {
Self {
resume_generation,
frontier_by_stage: Mutex::new(HashMap::new()),
infinite_sources: Mutex::new(HashSet::new()),
recorded_effect_seq_max: Mutex::new(HashMap::new()),
recorded_delivered_high_water: Mutex::new(HashMap::new()),
}
}
fn recorded_effect_seq_max(&self, stage: StageId) -> Option<StageInputPosition> {
self.recorded_effect_seq_max
.lock()
.expect("resume effect high-water lock poisoned")
.get(&stage)
.copied()
}
fn frontier(&self, stage: StageId) -> ReaderGeneration {
self.frontier_by_stage
.lock()
.expect("resume frontier lock poisoned")
.get(&stage)
.copied()
.unwrap_or_default()
}
fn stage_is_live(&self, stage: StageId) -> bool {
self.frontier(stage) >= self.resume_generation
}
fn is_infinite_source(&self, stage: StageId) -> bool {
self.infinite_sources
.lock()
.expect("resume source registry lock poisoned")
.contains(&stage)
}
}
#[derive(Debug, Clone)]
pub struct ResumeControl(Arc<ResumeState>);
impl ResumeControl {
pub fn resume_generation(&self) -> ReaderGeneration {
self.0.resume_generation
}
pub fn record_generation_boundary(&self, stage: StageId, generation: ReaderGeneration) {
let mut frontier = self
.0
.frontier_by_stage
.lock()
.expect("resume frontier lock poisoned");
let entry = frontier.entry(stage).or_default();
if generation > *entry {
*entry = generation;
}
}
pub fn register_infinite_source(&self, stage: StageId) {
self.0
.infinite_sources
.lock()
.expect("resume source registry lock poisoned")
.insert(stage);
}
pub fn record_effect_high_water(&self, stage: StageId, max: StageInputPosition) {
self.0
.recorded_effect_seq_max
.lock()
.expect("resume effect high-water lock poisoned")
.insert(stage, max);
}
pub fn record_delivered_high_water(&self, stage: StageId, count: u64) {
self.0
.recorded_delivered_high_water
.lock()
.expect("resume delivered high-water lock poisoned")
.insert(stage, count);
}
pub fn recorded_delivered_high_water(&self, stage: StageId) -> Option<u64> {
self.0
.recorded_delivered_high_water
.lock()
.expect("resume delivered high-water lock poisoned")
.get(&stage)
.copied()
}
}
#[derive(Debug)]
struct Resume {
state: Arc<ResumeState>,
}
impl ExecutionStrategy for Resume {
fn scope_at(&self, at: ExecutionPosition) -> MiddlewareExecutionScope {
match at.generation {
Some(generation) if generation >= self.state.resume_generation => {
MiddlewareExecutionScope::LiveHandler
}
Some(_) => MiddlewareExecutionScope::ResumeHandler,
None => self.stage_scope(at.stage_id),
}
}
fn stage_scope(&self, stage: StageId) -> MiddlewareExecutionScope {
if self.state.stage_is_live(stage) {
MiddlewareExecutionScope::LiveHandler
} else {
MiddlewareExecutionScope::ResumeHandler
}
}
fn scope_for_generation(
&self,
_stage: StageId,
generation: ReaderGeneration,
) -> MiddlewareExecutionScope {
if generation >= self.state.resume_generation {
MiddlewareExecutionScope::LiveHandler
} else {
MiddlewareExecutionScope::ResumeHandler
}
}
fn missing_outcome_is_corruption(&self, at: ExecutionPosition) -> bool {
match self.state.recorded_effect_seq_max(at.stage_id) {
Some(max) => at.position <= max,
None => false,
}
}
fn in_doubt_effect_is_fatal(&self) -> bool {
false
}
fn effect_port_registration_policy(&self) -> EffectPortRegistrationPolicy {
EffectPortRegistrationPolicy::Required
}
fn source_phase_for(&self, stage: StageId) -> SourceExecutionPhase {
if self.state.stage_is_live(stage) {
SourceExecutionPhase::Live
} else {
SourceExecutionPhase::Replaying
}
}
fn source_replay_exhausted(&self, stage: StageId) -> SourceReplayExhaustion {
if self.state.is_infinite_source(stage) {
SourceReplayExhaustion::ContinueLive
} else {
SourceReplayExhaustion::Terminate
}
}
fn heartbeat_policy_for(&self, stage: StageId) -> HeartbeatExecutionPolicy {
if self.state.stage_is_live(stage) {
HeartbeatExecutionPolicy::Active
} else {
HeartbeatExecutionPolicy::DormantUntilLive
}
}
}
fn archive_is_sealed(status: ArchiveStatus) -> bool {
matches!(status, ArchiveStatus::Completed | ArchiveStatus::Cancelled)
}
fn replay_incomplete(status: ArchiveStatus, allow_incomplete: bool) -> bool {
!archive_is_sealed(status) && allow_incomplete
}
#[derive(Clone)]
pub struct RuntimeExecution {
strategy: Arc<dyn ExecutionStrategy>,
archive: Option<Arc<dyn ReplayArchive>>,
effect_cursors: Arc<crate::effects::EffectCursorCoordinator>,
observations: Arc<ObservationRegistry>,
resume: Option<ResumeControl>,
}
impl std::fmt::Debug for RuntimeExecution {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("RuntimeExecution")
.field("strategy", &self.strategy)
.field(
"archive",
&self.archive.as_ref().map(|a| a.archive_flow_id()),
)
.finish()
}
}
impl RuntimeExecution {
pub fn observations(&self) -> &Arc<ObservationRegistry> {
&self.observations
}
pub fn new(mode: RuntimeMode, archive: Option<Arc<dyn ReplayArchive>>) -> Self {
let mut resume = None;
let strategy: Arc<dyn ExecutionStrategy> = match mode {
RuntimeMode::Live => Arc::new(Live),
RuntimeMode::Replay => {
let incomplete = archive
.as_deref()
.map(|a| replay_incomplete(a.archive_status(), a.allow_incomplete_archive()))
.unwrap_or(false);
Arc::new(Replay { incomplete })
}
RuntimeMode::Resume => {
let archive_ref = archive
.as_deref()
.expect("--resume-from requires an archive; config layer must reject earlier");
let resume_generation =
ReaderGeneration(archive_ref.max_recorded_generation().0 + 1);
let state = Arc::new(ResumeState::new(resume_generation));
resume = Some(ResumeControl(Arc::clone(&state)));
Arc::new(Resume { state })
}
};
Self {
strategy,
archive,
effect_cursors: Arc::new(crate::effects::EffectCursorCoordinator::default()),
observations: Arc::new(ObservationRegistry::default()),
resume,
}
}
pub fn host_observations_allowed(&self) -> bool {
self.strategy.effect_port_registration_policy()
!= EffectPortRegistrationPolicy::OptionalStrictReplay
}
pub fn observation_recorder(
&self,
flow_id: FlowId,
observer: WriterId,
) -> Arc<dyn ObservationRecorder> {
let scope = crate::metrics::observations::scope(self, flow_id);
self.observations.activate_scope(scope);
Arc::new(
self.observations
.capture_owner(scope, observer, self.clone()),
)
}
pub fn resume_control(&self) -> Option<&ResumeControl> {
self.resume.as_ref()
}
#[cfg(test)]
pub(crate) fn from_effect_runtime_mode(
mode: crate::effects::EffectRuntimeMode,
archive: Option<Arc<dyn ReplayArchive>>,
) -> Self {
use crate::effects::EffectRuntimeMode;
let strategy: Arc<dyn ExecutionStrategy> = match mode {
EffectRuntimeMode::Live => Arc::new(Live),
EffectRuntimeMode::ReplayStrict => Arc::new(Replay { incomplete: false }),
EffectRuntimeMode::ResumeIncomplete => Arc::new(Replay { incomplete: true }),
};
Self {
strategy,
archive,
effect_cursors: Arc::new(crate::effects::EffectCursorCoordinator::default()),
observations: Arc::new(ObservationRegistry::default()),
resume: None,
}
}
pub fn archive_for_io(&self) -> Option<&Arc<dyn ReplayArchive>> {
self.archive.as_ref()
}
pub fn handler_scope_for(&self, source: ExecutionPositionSource) -> MiddlewareExecutionScope {
match source {
ExecutionPositionSource::Data {
stage_id,
position,
generation,
} => self.strategy.scope_at(ExecutionPosition {
stage_id,
position,
generation,
}),
ExecutionPositionSource::FlowControl { stage_id, cause } => match cause {
Some(at) => self.strategy.scope_at(at),
None => self.strategy.stage_scope(stage_id),
},
ExecutionPositionSource::StageLifecycle { stage_id } => {
self.strategy.stage_scope(stage_id)
}
}
}
pub fn dispatch_scope(
&self,
stage_id: StageId,
position: Option<StageInputPosition>,
generation: Option<ReaderGeneration>,
) -> MiddlewareExecutionScope {
match (position, generation) {
(Some(p), generation) => self.scope_at(ExecutionPosition {
stage_id,
position: p,
generation,
}),
(None, Some(generation)) => self.strategy.scope_for_generation(stage_id, generation),
(None, None) => self.stage_scope(stage_id),
}
}
pub fn is_reconstructing(&self, at: ExecutionPosition) -> bool {
self.strategy.scope_at(at).is_deterministic_replay()
}
pub fn scope_at(&self, at: ExecutionPosition) -> MiddlewareExecutionScope {
self.strategy.scope_at(at)
}
pub fn stage_scope(&self, stage: StageId) -> MiddlewareExecutionScope {
self.strategy.stage_scope(stage)
}
pub fn missing_outcome_is_corruption(&self, at: ExecutionPosition) -> bool {
self.strategy.missing_outcome_is_corruption(at)
}
pub fn in_doubt_effect_is_fatal(&self) -> bool {
self.strategy.in_doubt_effect_is_fatal()
}
pub fn effect_port_registration_policy(&self) -> EffectPortRegistrationPolicy {
self.strategy.effect_port_registration_policy()
}
pub(crate) fn effect_cursor_coordinator(
&self,
) -> &Arc<crate::effects::EffectCursorCoordinator> {
&self.effect_cursors
}
pub fn source_phase_for(&self, stage: StageId) -> SourceExecutionPhase {
self.strategy.source_phase_for(stage)
}
pub fn source_replay_exhausted(&self, stage: StageId) -> SourceReplayExhaustion {
self.strategy.source_replay_exhausted(stage)
}
pub fn heartbeat_policy_for(&self, stage: StageId) -> HeartbeatExecutionPolicy {
self.strategy.heartbeat_policy_for(stage)
}
}
#[cfg(test)]
mod tests {
use super::*;
use obzenflow_core::MiddlewareExecutionScope as Scope;
fn at() -> ExecutionPosition {
ExecutionPosition {
stage_id: StageId::new(),
position: StageInputPosition(0),
generation: None,
}
}
#[test]
fn live_strategy_answers() {
let s = Live;
assert_eq!(s.scope_at(at()), Scope::LiveHandler);
assert_eq!(s.stage_scope(StageId::new()), Scope::LiveHandler);
assert!(!s.missing_outcome_is_corruption(at()));
assert_eq!(
s.source_phase_for(StageId::new()),
SourceExecutionPhase::Live
);
assert_eq!(
s.heartbeat_policy_for(StageId::new()),
HeartbeatExecutionPolicy::Active
);
assert!(!s.scope_at(at()).is_deterministic_replay());
}
#[test]
fn replay_sealed_strategy_answers() {
let s = Replay { incomplete: false };
assert_eq!(s.scope_at(at()), Scope::StrictReplayHandler);
assert_eq!(s.stage_scope(StageId::new()), Scope::StrictReplayHandler);
assert!(s.missing_outcome_is_corruption(at()));
assert_eq!(
s.source_phase_for(StageId::new()),
SourceExecutionPhase::Replaying
);
assert_eq!(
s.source_replay_exhausted(StageId::new()),
SourceReplayExhaustion::Terminate
);
assert_eq!(
s.heartbeat_policy_for(StageId::new()),
HeartbeatExecutionPolicy::Suppressed
);
assert!(s.scope_at(at()).is_deterministic_replay());
}
#[test]
fn replay_incomplete_strategy_answers() {
let s = Replay { incomplete: true };
assert_eq!(s.scope_at(at()), Scope::ResumeHandler);
assert!(!s.missing_outcome_is_corruption(at()));
assert!(s.scope_at(at()).is_deterministic_replay());
assert_eq!(
s.source_replay_exhausted(StageId::new()),
SourceReplayExhaustion::Terminate
);
}
#[test]
fn resume_strategy_answers_track_the_frontier() {
let state = Arc::new(ResumeState::new(ReaderGeneration(1)));
let control = ResumeControl(Arc::clone(&state));
let s = Resume {
state: Arc::clone(&state),
};
let stage = StageId::new();
let positioned = ExecutionPosition {
stage_id: stage,
position: StageInputPosition(3),
generation: None,
};
assert_eq!(s.scope_at(positioned), Scope::ResumeHandler);
assert_eq!(s.stage_scope(stage), Scope::ResumeHandler);
assert!(s.scope_at(positioned).is_deterministic_replay());
assert!(!s.missing_outcome_is_corruption(positioned));
assert_eq!(s.source_phase_for(stage), SourceExecutionPhase::Replaying);
assert_eq!(
s.heartbeat_policy_for(stage),
HeartbeatExecutionPolicy::DormantUntilLive
);
assert_eq!(
s.source_replay_exhausted(stage),
SourceReplayExhaustion::Terminate
);
control.register_infinite_source(stage);
assert_eq!(
s.source_replay_exhausted(stage),
SourceReplayExhaustion::ContinueLive
);
control.record_generation_boundary(stage, ReaderGeneration(1));
assert_eq!(s.scope_at(positioned), Scope::LiveHandler);
assert_eq!(s.stage_scope(stage), Scope::LiveHandler);
assert_eq!(s.source_phase_for(stage), SourceExecutionPhase::Live);
assert_eq!(
s.heartbeat_policy_for(stage),
HeartbeatExecutionPolicy::Active
);
control.record_generation_boundary(stage, ReaderGeneration(0));
assert_eq!(s.stage_scope(stage), Scope::LiveHandler);
let other = StageId::new();
assert_eq!(s.stage_scope(other), Scope::ResumeHandler);
}
#[test]
fn resume_scope_at_reads_the_event_generation() {
let state = Arc::new(ResumeState::new(ReaderGeneration(1)));
let control = ResumeControl(Arc::clone(&state));
let s = Resume {
state: Arc::clone(&state),
};
let stage = StageId::new();
let at = |generation: Option<ReaderGeneration>| ExecutionPosition {
stage_id: stage,
position: StageInputPosition(3),
generation,
};
assert_eq!(
s.scope_at(at(Some(ReaderGeneration(0)))),
Scope::ResumeHandler
);
assert_eq!(
s.scope_at(at(Some(ReaderGeneration(1)))),
Scope::LiveHandler
);
assert_eq!(s.scope_at(at(None)), Scope::ResumeHandler);
assert_eq!(
s.scope_for_generation(stage, ReaderGeneration(0)),
Scope::ResumeHandler
);
assert_eq!(
s.scope_for_generation(stage, ReaderGeneration(1)),
Scope::LiveHandler
);
control.record_generation_boundary(stage, ReaderGeneration(1));
assert_eq!(s.scope_at(at(None)), Scope::LiveHandler);
assert_eq!(
s.scope_at(at(Some(ReaderGeneration(0)))),
Scope::ResumeHandler
);
assert_eq!(
s.scope_for_generation(stage, ReaderGeneration(0)),
Scope::ResumeHandler
);
}
#[test]
fn resume_effect_miss_is_positional_not_generational() {
let state = Arc::new(ResumeState::new(ReaderGeneration(1)));
let control = ResumeControl(Arc::clone(&state));
let s = Resume {
state: Arc::clone(&state),
};
let stage = StageId::new();
let at = |p: u64| ExecutionPosition {
stage_id: stage,
position: StageInputPosition(p),
generation: None,
};
assert!(!s.missing_outcome_is_corruption(at(1)));
control.record_effect_high_water(stage, StageInputPosition(5));
assert!(s.missing_outcome_is_corruption(at(1)));
assert!(s.missing_outcome_is_corruption(at(5)));
assert!(!s.missing_outcome_is_corruption(at(6)));
control.record_generation_boundary(stage, ReaderGeneration(1));
assert!(s.missing_outcome_is_corruption(at(5)));
assert!(!s.missing_outcome_is_corruption(at(6)));
}
#[test]
fn replay_incomplete_derivation_matches_old_policy() {
assert!(!replay_incomplete(ArchiveStatus::Completed, true));
assert!(!replay_incomplete(ArchiveStatus::Cancelled, true));
assert!(replay_incomplete(ArchiveStatus::Failed, true));
assert!(replay_incomplete(ArchiveStatus::Unknown, true));
assert!(!replay_incomplete(ArchiveStatus::Failed, false));
assert!(!replay_incomplete(ArchiveStatus::Unknown, false));
}
#[test]
fn live_runtime_execution_is_never_reconstructing() {
let exec = RuntimeExecution::new(RuntimeMode::Live, None);
assert!(exec.archive_for_io().is_none());
assert!(!exec.is_reconstructing(at()));
assert_eq!(
exec.heartbeat_policy_for(StageId::new()),
HeartbeatExecutionPolicy::Active
);
}
#[test]
fn handler_scope_for_resolves_every_source() {
let exec = RuntimeExecution::new(RuntimeMode::Live, None);
let sid = StageId::new();
let p = StageInputPosition(7);
assert_eq!(
exec.handler_scope_for(ExecutionPositionSource::Data {
stage_id: sid,
position: p,
generation: None
}),
Scope::LiveHandler
);
assert_eq!(
exec.handler_scope_for(ExecutionPositionSource::FlowControl {
stage_id: sid,
cause: Some(ExecutionPosition {
stage_id: sid,
position: p,
generation: None
})
}),
Scope::LiveHandler
);
assert_eq!(
exec.handler_scope_for(ExecutionPositionSource::FlowControl {
stage_id: sid,
cause: None
}),
Scope::LiveHandler
);
assert_eq!(
exec.handler_scope_for(ExecutionPositionSource::StageLifecycle { stage_id: sid }),
Scope::LiveHandler
);
}
}