use crate::{diagnostics::Observation, ImModule};
use std::sync::Arc;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum SyncStage {
Recovery,
InventoryPage,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum SyncResult {
Started,
Success,
Failed,
Cancelled,
}
impl SyncResult {
pub const fn as_str(self) -> &'static str {
match self {
Self::Started => "started",
Self::Success => "success",
Self::Failed => "failed",
Self::Cancelled => "cancelled",
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct SyncRecord {
pub stage: SyncStage,
pub result: SyncResult,
pub started_at_ms: u64,
pub completed_at_ms: u64,
pub total_pages: Option<u64>,
}
pub trait SyncObserver: Send + Sync {
fn record(&self, record: SyncRecord);
}
impl<F: Fn(SyncRecord) + Send + Sync> SyncObserver for F {
fn record(&self, record: SyncRecord) {
self(record);
}
}
#[derive(Default)]
pub(crate) struct Timing {
observer: Option<Arc<dyn SyncObserver>>,
epoch: u64,
corr_floor: u64,
pending_operations: u64,
started: Option<u64>,
session: String,
page_started: Option<u64>,
pub page_index: u64,
pub total_pages: Option<u64>,
pub inventory_committed: bool,
}
impl ImModule {
pub fn with_sync_observer(mut self, observer: Arc<dyn SyncObserver>) -> Self {
self.sync_timing.observer = Some(observer);
self
}
pub(crate) fn observe_sync_tick(&mut self, effects: &[helix_core::Effect], corr_floor: u64) {
let epoch = self.state.recovery_session.session_epoch;
if epoch != 0 && epoch != self.sync_timing.epoch {
self.finish_sync_observation(SyncResult::Cancelled);
self.sync_timing.epoch = epoch;
self.sync_timing.corr_floor = corr_floor;
self.sync_timing.pending_operations = 0;
self.sync_timing.started = Some(self.diagnostics.now_ms);
self.sync_timing.session = self.state.connection_id.clone().unwrap_or_default();
self.sync_timing.page_index = 0;
self.sync_timing.total_pages = None;
self.sync_timing.inventory_committed = false;
self.emit_sync_observation(
SyncStage::Recovery,
SyncResult::Started,
self.diagnostics.now_ms,
);
}
if self.sync_timing.started.is_none() {
return;
}
for effect in effects {
let corr = match effect {
helix_core::Effect::Http { corr, .. }
| helix_core::Effect::Persist { corr, .. }
| helix_core::Effect::PersistAtomic { corr, .. } => *corr,
_ => continue,
};
if corr.raw() >= self.sync_timing.corr_floor
&& self
.state
.corr_map
.get(&corr)
.is_some_and(is_recovery_operation)
{
self.sync_timing.pending_operations += 1;
}
}
if matches!(
self.state.recovery_session.phase,
crate::sync_session::RecoveryPhase::Failed
| crate::sync_session::RecoveryPhase::Blocked
) {
self.finish_sync_observation(SyncResult::Failed);
} else if self.state.startup_channel_projection_ready
&& self.sync_timing.pending_operations == 0
&& self.sync_timing.inventory_committed
&& self.state.increment_pull.is_none()
&& self.state.channel_sync_persist_inflight == 0
&& !self.state.channel_sync_batch_pending
&& self.state.sync_scheduler.is_idle()
&& !self.state.recovery_session.has_pending_commits()
{
self.finish_sync_observation(SyncResult::Success);
}
}
pub(crate) fn observe_sync_reply(&mut self, tick: &helix_core::Tick) -> bool {
let helix_core::Tick::PortReply { corr, outcome } = tick else {
return false;
};
if self.sync_timing.started.is_none() || corr.raw() < self.sync_timing.corr_floor {
return false;
}
let Some(context) = self.state.corr_map.get(corr) else {
return false;
};
if !is_recovery_operation(context) {
return false;
}
self.sync_timing.pending_operations = self.sync_timing.pending_operations.saturating_sub(1);
if matches!(outcome, helix_core::tick::PortOutcome::Err(_)) {
self.finish_sync_observation(SyncResult::Failed);
} else if matches!(
context,
crate::state::CorrelationContext::IncrementBatchPersist { .. }
) {
self.sync_timing.inventory_committed = true;
}
true
}
pub(crate) fn start_page_observation(&mut self) {
if self.sync_timing.started.is_none() {
return;
}
self.sync_timing.page_index += 1;
self.sync_timing.page_started = Some(self.diagnostics.now_ms);
self.emit_sync_observation(
SyncStage::InventoryPage,
SyncResult::Started,
self.diagnostics.now_ms,
);
}
pub(crate) fn finish_page_observation(&mut self, result: SyncResult) {
if let Some(started) = self.sync_timing.page_started.take() {
self.emit_sync_observation(SyncStage::InventoryPage, result, started);
}
}
pub(crate) fn finish_sync_observation(&mut self, result: SyncResult) {
self.finish_page_observation(result);
if let Some(started) = self.sync_timing.started.take() {
self.emit_sync_observation(SyncStage::Recovery, result, started);
}
}
fn emit_sync_observation(&self, stage: SyncStage, result: SyncResult, started: u64) {
let now = self.diagnostics.now_ms;
let record = SyncRecord {
stage,
result,
started_at_ms: started,
completed_at_ms: now,
total_pages: self.sync_timing.total_pages,
};
if let Some(observer) = &self.sync_timing.observer {
observer.record(record);
}
self.diagnose(Observation {
event: match (stage, result) {
(SyncStage::Recovery, SyncResult::Started) => "sync_recovery_started",
(SyncStage::Recovery, _) => "sync_recovery_terminal",
(SyncStage::InventoryPage, SyncResult::Started) => "sync_inventory_page_started",
(SyncStage::InventoryPage, _) => "sync_inventory_page_terminal",
},
stage: "client",
result: result.as_str(),
sync_session_id: &self.sync_timing.session,
page_index: (stage == SyncStage::InventoryPage).then_some(self.sync_timing.page_index),
total_pages: self.sync_timing.total_pages,
started_at_ms: Some(started),
completed_at_ms: (result != SyncResult::Started).then_some(now),
elapsed: now.saturating_sub(started) as f64 / 1000.0,
..Default::default()
});
}
}
fn is_recovery_operation(context: &crate::state::CorrelationContext) -> bool {
use crate::state::CorrelationContext::*;
matches!(
context,
ScanCursors
| ScanChannelProjections
| IncrementMessageTimestampScan { .. }
| IncrementPullHttp
| IncrementPullPersist
| IncrementBatchPersist { .. }
| SyncPull { .. }
| ChannelPersist { .. }
| ChannelTerminalPersist { .. }
| MemberProjectionPersist { .. }
| TooLongReload { .. }
)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::state::CorrelationContext;
use helix_core::{Correlation, Effect, Tick};
#[test]
fn continuation_persist_remains_a_terminal_barrier() {
let mut module = ImModule::new(Default::default());
module.state.recovery_session.begin("actor");
module.observe_sync_tick(&[], 10);
module.sync_timing.inventory_committed = true;
module.state.startup_channel_projection_ready = true;
module.state.channel_sync_batch_pending = false;
module.state.recovery_session.mark_completion_published();
let corr = Correlation::from_raw(11);
module.state.corr_map.insert(
corr,
CorrelationContext::IncrementBatchPersist {
projections: vec![],
batch_id: None,
},
);
module.observe_sync_tick(&[Effect::PersistAtomic { corr, ops: vec![] }], 12);
assert!(module.sync_timing.started.is_some());
module.observe_sync_reply(&Tick::PortReply {
corr,
outcome: helix_core::tick::PortOutcome::Ok(helix_core::tick::ReplyBytes(
bytes::Bytes::new(),
)),
});
module.state.corr_map.remove(&corr);
module.observe_sync_tick(&[], 12);
assert!(module.sync_timing.started.is_none());
}
#[test]
fn old_inventory_receipt_does_not_complete_new_epoch() {
let mut module = ImModule::new(Default::default());
module.state.recovery_session.begin("actor");
module.observe_sync_tick(&[], 10);
let old = Correlation::from_raw(5);
module.state.corr_map.insert(
old,
CorrelationContext::IncrementBatchPersist {
projections: vec![],
batch_id: None,
},
);
module.observe_sync_reply(&Tick::PortReply {
corr: old,
outcome: helix_core::tick::PortOutcome::Ok(helix_core::tick::ReplyBytes(
bytes::Bytes::new(),
)),
});
assert!(!module.sync_timing.inventory_committed);
assert!(module.sync_timing.started.is_some());
}
}