use crate::testing::FlowTestHarness;
use obzenflow_core::event::journal_record::JournalRecord;
use obzenflow_core::event::payloads::system_payload::{MetricsCoordinationEvent, SystemPayload};
use obzenflow_core::event::{SystemEvent, WriterId};
use obzenflow_core::journal::Journal;
use obzenflow_core::StageId;
use std::sync::Arc;
use std::time::Duration;
use thiserror::Error;
#[derive(Debug, Error)]
pub enum MetricsBarrierError {
#[error("metrics shut down without publishing the final buffer")]
ShutdownWithoutDrain,
#[error(
"flow handle has no system journal; \
cannot construct MetricsBarrier on a flow built without one"
)]
MissingSystemJournal,
#[error("unknown stage `{0}` for MetricsBarrier::try_on_stage")]
UnknownStage(String),
#[error("ambiguous stage `{0}`: multiple stages share this name")]
AmbiguousStage(String),
#[error("failed to read system journal: {0}")]
JournalRead(String),
#[error("flow handle has no topology; cannot resolve stage names")]
MissingTopology,
}
pub struct MetricsBarrier {
system_journal: Arc<dyn Journal<SystemEvent>>,
stage_writer_key: Option<String>,
baseline_offset: u64,
}
impl MetricsBarrier {
pub async fn try_on_stage(
handle: &FlowTestHarness,
stage_name: &str,
) -> Result<Self, MetricsBarrierError> {
let system_journal = handle
.system_journal()
.ok_or(MetricsBarrierError::MissingSystemJournal)?;
let stage_id = resolve_stage_id(handle, stage_name)?;
let stage_writer_key = Some(WriterId::from(stage_id).to_string());
let baseline_offset = current_journal_offset(&system_journal).await?;
Ok(Self {
system_journal,
stage_writer_key,
baseline_offset,
})
}
pub async fn try_on_flow(handle: &FlowTestHarness) -> Result<Self, MetricsBarrierError> {
let system_journal = handle
.system_journal()
.ok_or(MetricsBarrierError::MissingSystemJournal)?;
let baseline_offset = current_journal_offset(&system_journal).await?;
Ok(Self {
system_journal,
stage_writer_key: None,
baseline_offset,
})
}
pub async fn wait_for_stage_seq(&self, target_seq: u64) -> Result<(), MetricsBarrierError> {
let writer_key = self
.stage_writer_key
.as_ref()
.expect("wait_for_stage_seq called on a non-stage barrier");
let mut scan_from = self.baseline_offset;
loop {
let envelopes = read_journal_from(&self.system_journal, scan_from).await?;
let next_scan_from = scan_from + envelopes.len() as u64;
for env in envelopes {
if let SystemPayload::MetricsCoordination(MetricsCoordinationEvent::Exported {
watermark,
}) = &env.payload
{
if let Some(seq) = watermark.clocks.get(writer_key) {
if *seq >= target_seq {
return Ok(());
}
}
}
}
scan_from = next_scan_from;
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
pub async fn wait_for_drained(&self) -> Result<(), MetricsBarrierError> {
let mut scan_from = self.baseline_offset;
loop {
let envelopes = read_journal_from(&self.system_journal, scan_from).await?;
let next_scan_from = scan_from + envelopes.len() as u64;
for env in envelopes {
match &env.payload {
SystemPayload::MetricsCoordination(MetricsCoordinationEvent::Drained) => {
return Ok(())
}
SystemPayload::MetricsCoordination(MetricsCoordinationEvent::Shutdown) => {
return Err(MetricsBarrierError::ShutdownWithoutDrain)
}
_ => {}
}
}
scan_from = next_scan_from;
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
}
fn resolve_stage_id(
handle: &FlowTestHarness,
stage_name: &str,
) -> Result<StageId, MetricsBarrierError> {
use crate::id_conversions::StageIdExt;
let topology = handle
.topology()
.ok_or(MetricsBarrierError::MissingTopology)?;
let mut matches: Vec<StageId> = topology
.stages()
.filter(|s| s.name == stage_name)
.map(|s| StageId::from_topology_id(s.id))
.collect();
match matches.len() {
0 => Err(MetricsBarrierError::UnknownStage(stage_name.to_string())),
1 => Ok(matches.remove(0)),
_ => Err(MetricsBarrierError::AmbiguousStage(stage_name.to_string())),
}
}
async fn current_journal_offset(
journal: &Arc<dyn Journal<SystemEvent>>,
) -> Result<u64, MetricsBarrierError> {
let mut reader = journal
.reader()
.await
.map_err(|e| MetricsBarrierError::JournalRead(e.to_string()))?;
let mut count: u64 = 0;
loop {
match reader.next().await {
Ok(Some(_)) => count += 1,
Ok(None) => return Ok(count),
Err(e) => return Err(MetricsBarrierError::JournalRead(e.to_string())),
}
}
}
async fn read_journal_from(
journal: &Arc<dyn Journal<SystemEvent>>,
from: u64,
) -> Result<Vec<JournalRecord<SystemPayload>>, MetricsBarrierError> {
let mut reader = journal
.reader_from(from)
.await
.map_err(|e| MetricsBarrierError::JournalRead(e.to_string()))?;
let mut envelopes = Vec::new();
loop {
match reader.next().await {
Ok(Some(env)) => envelopes.push(env),
Ok(None) => return Ok(envelopes),
Err(e) => return Err(MetricsBarrierError::JournalRead(e.to_string())),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::id_conversions::StageIdExt;
use crate::metrics::observations::ObservationRegistry;
use crate::pipeline::fsm::PipelineFsmEvent;
use crate::pipeline::handle::FlowHandleExtras;
use crate::pipeline::{FlowHandle, PipelineState};
use crate::supervised_base::{ChannelBuilder, HandleBuilder, SupervisorTaskBuilder};
use obzenflow_core::event::journal_record::JournalRecord;
use obzenflow_core::event::observability::NoObservations;
use obzenflow_core::event::payloads::system_payload::MetricsCoordinationEvent;
use obzenflow_core::event::vector_clock::VectorClock;
use obzenflow_core::event::{
JournalEvent, JournalWriterId, SystemEvent, SystemPayload, WriterId,
};
use obzenflow_core::id::JournalId;
use obzenflow_core::journal::journal_error::JournalError;
use obzenflow_core::journal::journal_owner::JournalOwner;
use obzenflow_core::journal::reader::JournalReader;
use obzenflow_core::journal::Journal;
use obzenflow_core::StageId;
use obzenflow_topology::TopologyBuilder;
use std::sync::{Arc, Mutex};
use std::time::Duration;
struct MemoryJournal<T: JournalEvent> {
id: JournalId,
owner: Option<JournalOwner>,
events: Arc<Mutex<Vec<JournalRecord<T::Payload>>>>,
}
impl<T: JournalEvent> Default for MemoryJournal<T> {
fn default() -> Self {
Self {
id: JournalId::new(),
owner: None,
events: Arc::new(Mutex::new(Vec::new())),
}
}
}
struct MemoryJournalReader<T: JournalEvent> {
events: Arc<Mutex<Vec<JournalRecord<T::Payload>>>>,
pos: usize,
}
#[async_trait::async_trait]
impl<T> JournalReader<T> for MemoryJournalReader<T>
where
T: JournalEvent,
{
async fn next(&mut self) -> Result<Option<JournalRecord<T::Payload>>, JournalError> {
let guard = self
.events
.lock()
.expect("MemoryJournalReader: poisoned lock");
if self.pos >= guard.len() {
return Ok(None);
}
let envelope = guard[self.pos].clone();
drop(guard);
self.pos += 1;
Ok(Some(envelope))
}
fn position(&self) -> u64 {
self.pos as u64
}
}
#[async_trait::async_trait]
impl<T> Journal<T> for MemoryJournal<T>
where
T: JournalEvent + 'static,
{
fn id(&self) -> &JournalId {
&self.id
}
fn owner(&self) -> Option<&JournalOwner> {
self.owner.as_ref()
}
async fn append(
&self,
event: T,
mut options: obzenflow_core::journal::AppendOptions<'_, T>,
) -> Result<JournalRecord<T::Payload>, JournalError> {
let event = options.capture.prepare(0, event);
let envelope = JournalRecord::new(JournalWriterId::from(self.id), event);
let mut guard = self.events.lock().expect("MemoryJournal: poisoned lock");
guard.push(envelope.clone());
Ok(envelope)
}
async fn read_all_unordered(&self) -> Result<Vec<JournalRecord<T::Payload>>, JournalError> {
let guard = self.events.lock().expect("MemoryJournal: poisoned lock");
Ok(guard.clone())
}
async fn read_event(
&self,
event_id: &obzenflow_core::event::types::EventId,
) -> Result<Option<JournalRecord<T::Payload>>, JournalError> {
let guard = self.events.lock().expect("MemoryJournal: poisoned lock");
Ok(guard.iter().find(|e| e.id() == event_id).cloned())
}
async fn reader_from(
&self,
position: u64,
) -> Result<Box<dyn JournalReader<T>>, JournalError> {
Ok(Box::new(MemoryJournalReader {
events: Arc::clone(&self.events),
pos: position as usize,
}))
}
async fn read_last_n(
&self,
count: usize,
) -> Result<Vec<JournalRecord<T::Payload>>, JournalError> {
let guard = self.events.lock().expect("MemoryJournal: poisoned lock");
let len = guard.len();
let start = len.saturating_sub(count);
Ok(guard[start..].iter().rev().cloned().collect())
}
}
fn harness_with_system_journal(
system_journal: Arc<dyn Journal<SystemEvent>>,
topology: Option<Arc<obzenflow_topology::Topology>>,
) -> FlowTestHarness {
let (event_sender, _event_receiver, state_watcher) =
ChannelBuilder::<PipelineFsmEvent, PipelineState>::new().build(PipelineState::Created);
let supervisor_task = SupervisorTaskBuilder::<PipelineState>::new("dummy_pipeline")
.spawn_for_test(
|| async move { Ok::<(), Box<dyn std::error::Error + Send + Sync>>(()) },
);
let standard_handle = HandleBuilder::new()
.with_event_sender(event_sender)
.with_state_watcher(state_watcher)
.with_supervisor_task(supervisor_task)
.build_standard()
.expect("dummy handle should build");
let extras = FlowHandleExtras {
observations: Arc::new(ObservationRegistry::default()),
host_observations: Arc::new(NoObservations),
stage_cleanup: Vec::new(),
published_outcome: Default::default(),
metrics: Default::default(),
operational_failure: Default::default(),
topology,
flow_name: "dummy".to_string(),
contract_attachments: None,
system_journal: Some(system_journal),
pipeline_writer_id: obzenflow_core::event::WriterId::from(
obzenflow_core::id::SystemId::new(),
),
liveness_snapshots: None,
run_substrate: obzenflow_core::journal::factory::RunSubstrateState::Ephemeral,
flow_effective_config: None,
};
let handle = FlowHandle::new(standard_handle, extras);
FlowTestHarness::from_parts(handle, Vec::new()).expect("empty stage journals")
}
fn harness_without_system_journal(
topology: Option<Arc<obzenflow_topology::Topology>>,
) -> FlowTestHarness {
let (event_sender, _event_receiver, state_watcher) =
ChannelBuilder::<PipelineFsmEvent, PipelineState>::new().build(PipelineState::Created);
let supervisor_task = SupervisorTaskBuilder::<PipelineState>::new("dummy_pipeline")
.spawn_for_test(
|| async move { Ok::<(), Box<dyn std::error::Error + Send + Sync>>(()) },
);
let standard_handle = HandleBuilder::new()
.with_event_sender(event_sender)
.with_state_watcher(state_watcher)
.with_supervisor_task(supervisor_task)
.build_standard()
.expect("dummy handle should build");
let extras = FlowHandleExtras {
observations: Arc::new(ObservationRegistry::default()),
host_observations: Arc::new(NoObservations),
stage_cleanup: Vec::new(),
published_outcome: Default::default(),
metrics: Default::default(),
operational_failure: Default::default(),
topology,
flow_name: "dummy".to_string(),
contract_attachments: None,
system_journal: None,
pipeline_writer_id: obzenflow_core::event::WriterId::from(
obzenflow_core::id::SystemId::new(),
),
liveness_snapshots: None,
run_substrate: obzenflow_core::journal::factory::RunSubstrateState::Ephemeral,
flow_effective_config: None,
};
let handle = FlowHandle::new(standard_handle, extras);
FlowTestHarness::from_parts(handle, Vec::new()).expect("empty stage journals")
}
#[tokio::test]
async fn try_on_flow_errors_when_system_journal_is_missing() {
let harness = harness_without_system_journal(None);
let err = MetricsBarrier::try_on_flow(&harness)
.await
.err()
.expect("expected MissingSystemJournal");
assert!(
matches!(err, MetricsBarrierError::MissingSystemJournal),
"unexpected error: {err:?}"
);
}
#[tokio::test]
async fn try_on_stage_errors_when_topology_is_missing() {
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(MemoryJournal::default());
let harness = harness_with_system_journal(system_journal, None);
let err = MetricsBarrier::try_on_stage(&harness, "stage")
.await
.err()
.expect("expected MissingTopology");
assert!(
matches!(err, MetricsBarrierError::MissingTopology),
"unexpected error: {err:?}"
);
}
#[tokio::test]
async fn try_on_stage_errors_on_unknown_stage_name() {
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(MemoryJournal::default());
let mut topology_builder = TopologyBuilder::new();
topology_builder.add_stage(Some("present".to_string()));
topology_builder.add_stage(Some("sink".to_string()));
let topology = Arc::new(topology_builder.build_unchecked().expect("topology"));
let harness = harness_with_system_journal(system_journal, Some(topology));
let err = MetricsBarrier::try_on_stage(&harness, "missing")
.await
.err()
.expect("expected UnknownStage");
assert!(
matches!(err, MetricsBarrierError::UnknownStage(ref name) if name == "missing"),
"unexpected error: {err:?}"
);
}
#[tokio::test]
async fn try_on_stage_errors_on_ambiguous_stage_name() {
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(MemoryJournal::default());
let mut topology_builder = TopologyBuilder::new();
topology_builder.add_stage(Some("dup".to_string()));
topology_builder.add_stage(Some("dup".to_string()));
topology_builder.add_stage(Some("sink".to_string()));
let topology = Arc::new(topology_builder.build_unchecked().expect("topology"));
let harness = harness_with_system_journal(system_journal, Some(topology));
let err = MetricsBarrier::try_on_stage(&harness, "dup")
.await
.err()
.expect("expected AmbiguousStage");
assert!(
matches!(err, MetricsBarrierError::AmbiguousStage(ref name) if name == "dup"),
"unexpected error: {err:?}"
);
}
#[tokio::test]
async fn wait_for_drained_resolves_when_event_appended_before_wait_begins() {
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(MemoryJournal::default());
let harness = harness_with_system_journal(system_journal.clone(), None);
let barrier = MetricsBarrier::try_on_flow(&harness)
.await
.expect("construct barrier");
system_journal
.append(
SystemEvent::new(
WriterId::from(StageId::new()),
SystemPayload::MetricsCoordination(MetricsCoordinationEvent::Drained),
),
Default::default(),
)
.await
.expect("append drained");
tokio::time::timeout(Duration::from_secs(1), barrier.wait_for_drained())
.await
.expect("wait should resolve within timeout")
.expect("wait should succeed");
}
#[tokio::test]
async fn wait_for_drained_rejects_shutdown_without_final_publication() {
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(MemoryJournal::default());
let harness = harness_with_system_journal(system_journal.clone(), None);
let barrier = MetricsBarrier::try_on_flow(&harness)
.await
.expect("construct barrier");
system_journal
.append(
SystemEvent::new(
WriterId::from(StageId::new()),
SystemPayload::MetricsCoordination(MetricsCoordinationEvent::Shutdown),
),
Default::default(),
)
.await
.expect("append shutdown");
tokio::time::timeout(Duration::from_secs(1), barrier.wait_for_drained())
.await
.expect("wait should resolve within timeout")
.expect_err("shutdown is not a successful final publication");
}
#[tokio::test]
async fn wait_for_stage_seq_resolves_when_covering_export_is_already_present() {
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(MemoryJournal::default());
let mut topology_builder = TopologyBuilder::new();
let stage_topo_id = topology_builder.add_stage(Some("stage".to_string()));
topology_builder.add_stage(Some("sink".to_string()));
let topology = Arc::new(topology_builder.build_unchecked().expect("topology"));
let stage_id = StageId::from_topology_id(stage_topo_id);
let writer_key = WriterId::from(stage_id).to_string();
let harness = harness_with_system_journal(system_journal.clone(), Some(topology));
let barrier = MetricsBarrier::try_on_stage(&harness, "stage")
.await
.expect("construct stage barrier");
let mut watermark = VectorClock::new();
watermark.clocks.insert(writer_key, 5);
system_journal
.append(
SystemEvent::new(
WriterId::from(stage_id),
SystemPayload::MetricsCoordination(MetricsCoordinationEvent::Exported {
watermark,
}),
),
Default::default(),
)
.await
.expect("append exported");
tokio::time::timeout(Duration::from_secs(1), barrier.wait_for_stage_seq(3))
.await
.expect("wait should resolve within timeout")
.expect("wait should succeed");
}
#[tokio::test]
async fn wait_for_stage_seq_filters_unrelated_writer_progress() {
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(MemoryJournal::default());
let mut topology_builder = TopologyBuilder::new();
let stage_topo_id = topology_builder.add_stage(Some("stage".to_string()));
let other_topo_id = topology_builder.add_stage(Some("other".to_string()));
topology_builder.add_stage(Some("sink".to_string()));
let topology = Arc::new(topology_builder.build_unchecked().expect("topology"));
let stage_id = StageId::from_topology_id(stage_topo_id);
let stage_key = WriterId::from(stage_id).to_string();
let other_id = StageId::from_topology_id(other_topo_id);
let other_key = WriterId::from(other_id).to_string();
let harness = harness_with_system_journal(system_journal.clone(), Some(topology));
let barrier = MetricsBarrier::try_on_stage(&harness, "stage")
.await
.expect("construct stage barrier");
let mut wait = tokio::spawn(async move { barrier.wait_for_stage_seq(5).await });
let mut watermark = VectorClock::new();
watermark.clocks.insert(other_key, 5);
system_journal
.append(
SystemEvent::new(
WriterId::from(other_id),
SystemPayload::MetricsCoordination(MetricsCoordinationEvent::Exported {
watermark,
}),
),
Default::default(),
)
.await
.expect("append unrelated exported");
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(
!wait.is_finished(),
"wait should not resolve on unrelated writer watermark"
);
let mut watermark = VectorClock::new();
watermark.clocks.insert(stage_key, 5);
system_journal
.append(
SystemEvent::new(
WriterId::from(stage_id),
SystemPayload::MetricsCoordination(MetricsCoordinationEvent::Exported {
watermark,
}),
),
Default::default(),
)
.await
.expect("append covering exported");
tokio::time::timeout(Duration::from_secs(1), &mut wait)
.await
.expect("wait should resolve within timeout")
.expect("task should join")
.expect("wait should succeed");
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn wait_for_stage_seq_completes_under_paused_time_via_timer_polling() {
let system_journal: Arc<dyn Journal<SystemEvent>> = Arc::new(MemoryJournal::default());
let mut topology_builder = TopologyBuilder::new();
let stage_topo_id = topology_builder.add_stage(Some("stage".to_string()));
topology_builder.add_stage(Some("sink".to_string()));
let topology = Arc::new(topology_builder.build_unchecked().expect("topology"));
let stage_id = StageId::from_topology_id(stage_topo_id);
let stage_key = WriterId::from(stage_id).to_string();
let harness = harness_with_system_journal(system_journal.clone(), Some(topology));
let barrier = MetricsBarrier::try_on_stage(&harness, "stage")
.await
.expect("construct stage barrier");
let system_journal_for_writer = system_journal.clone();
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(100)).await;
let mut watermark = VectorClock::new();
watermark.clocks.insert(stage_key, 2);
system_journal_for_writer
.append(
SystemEvent::new(
WriterId::from(stage_id),
SystemPayload::MetricsCoordination(MetricsCoordinationEvent::Exported {
watermark,
}),
),
Default::default(),
)
.await
.expect("append exported");
});
let mut wait_task = tokio::spawn(async move { barrier.wait_for_stage_seq(2).await });
let (watchdog_tx, mut watchdog_rx) = tokio::sync::oneshot::channel::<()>();
std::thread::spawn(move || {
std::thread::sleep(std::time::Duration::from_secs(2));
let _ = watchdog_tx.send(());
});
tokio::select! {
res = &mut wait_task => {
res.expect("task should join").expect("wait should succeed");
}
_ = &mut watchdog_rx => {
wait_task.abort();
panic!("wait_for_stage_seq did not complete under paused time");
}
}
}
}