use crate::store::config::duration_micros;
use crate::store::platform::clock::{mono_elapsed, Clock};
use crate::store::{
AppendReceipt, HlcPoint, Open, Store, StoreError, StoreInvariant, WatermarkKind,
};
use std::time::Duration;
struct GateWaitMeasurement {
result: Result<(), StoreError>,
waited_us: u64,
}
fn measure_gate_wait<ResolveTarget, WaitForTarget>(
clock: &dyn Clock,
resolve_target: ResolveTarget,
wait_for_target: WaitForTarget,
) -> Result<GateWaitMeasurement, StoreError>
where
ResolveTarget: FnOnce() -> Result<HlcPoint, StoreError>,
WaitForTarget: FnOnce(HlcPoint) -> Result<(), StoreError>,
{
let target = resolve_target()?;
let started_ns = clock.now_mono_ns();
let result = wait_for_target(target);
Ok(GateWaitMeasurement {
waited_us: duration_micros(mono_elapsed(clock, started_ns)),
result,
})
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct DurabilityGate {
kind: WatermarkKind,
timeout: Duration,
}
impl DurabilityGate {
#[must_use]
pub const fn new(kind: WatermarkKind, timeout: Duration) -> Self {
Self { kind, timeout }
}
pub fn kind(&self) -> WatermarkKind {
self.kind
}
#[must_use]
pub fn timeout(&self) -> Duration {
self.timeout
}
}
impl Store<Open> {
fn receipt_point(&self, receipt: &AppendReceipt) -> Result<HlcPoint, StoreError> {
use crate::id::EntityIdType;
let raw = receipt.event_id.as_u128();
self.index
.get_by_id(raw)
.map(|entry| HlcPoint {
wall_ms: entry.wall_ms,
global_sequence: entry.global_sequence,
})
.ok_or_else(|| StoreError::InvariantViolation {
kind: StoreInvariant::GateReceiptNotIndexed { event_id: raw },
})
}
pub(crate) fn wait_for_gate(
&self,
receipt: &AppendReceipt,
gate: DurabilityGate,
) -> Result<(), StoreError> {
let measurement = measure_gate_wait(
self.runtime.clock(),
|| self.receipt_point(receipt),
|target| match gate.kind {
WatermarkKind::Accepted => self.wait_for_accepted(target, gate.timeout),
WatermarkKind::Written => self.wait_for_written(target, gate.timeout),
WatermarkKind::Durable => self.wait_for_durable(target, gate.timeout),
WatermarkKind::Applied => self.wait_for_applied(target, gate.timeout),
WatermarkKind::Visible => self.wait_for_visible(target, gate.timeout),
WatermarkKind::Emitted => self.wait_for_emitted(target, gate.timeout),
},
)?;
tracing::trace!(
target: "batpak::durability_gate",
kind = ?gate.kind,
waited_us = measurement.waited_us,
ok = measurement.result.is_ok(),
"append durability gate wait completed",
);
measurement.result
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::store::platform::clock::SystemClock;
use std::sync::{
atomic::{AtomicBool, AtomicI64, Ordering},
Arc,
};
use std::time::Instant;
#[test]
fn timeout_accessor_returns_the_configured_wait_not_default() {
let gate = DurabilityGate::new(WatermarkKind::Durable, Duration::from_millis(250));
assert_eq!(
gate.timeout(),
Duration::from_millis(250),
"DurabilityGate::timeout must return the configured wait, not Duration::default()"
);
assert_ne!(
gate.timeout(),
Duration::default(),
"a non-zero configured timeout must never read back as the default zero"
);
}
#[test]
fn gate_wait_measurement_excludes_receipt_lookup_work() {
let lookup_complete = Arc::new(AtomicBool::new(false));
let wait_saw_lookup_complete = Arc::clone(&lookup_complete);
let before_lookup = Instant::now();
let measurement = measure_gate_wait(
&SystemClock::new(),
|| {
std::thread::sleep(Duration::from_millis(40));
lookup_complete.store(true, Ordering::SeqCst);
Ok(HlcPoint {
wall_ms: 1,
global_sequence: 1,
})
},
|_| {
assert!(
wait_saw_lookup_complete.load(Ordering::SeqCst),
"gate wait must start only after receipt_point lookup completes"
);
Ok(())
},
)
.expect("measurement succeeds");
let lookup_inclusive_us = duration_micros(before_lookup.elapsed());
assert!(
measurement.waited_us < lookup_inclusive_us / 2,
"waited_us must measure only the watermark wait window, not receipt lookup work; waited_us={} lookup_inclusive_us={lookup_inclusive_us}",
measurement.waited_us
);
assert!(measurement.result.is_ok());
}
struct SteppingClock {
mono_ns: AtomicI64,
step_ns: i64,
}
impl Clock for SteppingClock {
fn now_us(&self) -> i64 {
0
}
fn now_wall_ns(&self) -> i64 {
0
}
fn now_mono_ns(&self) -> i64 {
self.mono_ns.fetch_add(self.step_ns, Ordering::SeqCst)
}
fn process_boot_ns(&self) -> u64 {
0
}
}
#[test]
fn gate_wait_measurement_reads_the_injected_clock() {
let step_ns = 7_000_000_i64; let clock = SteppingClock {
mono_ns: AtomicI64::new(0),
step_ns,
};
let measurement = measure_gate_wait(
&clock,
|| {
Ok(HlcPoint {
wall_ms: 1,
global_sequence: 1,
})
},
|_| Ok(()),
)
.expect("measurement succeeds");
assert_eq!(
measurement.waited_us,
u64::try_from(step_ns).unwrap_or(u64::MAX) / 1000,
"waited_us must be derived from the injected Clock::now_mono_ns samples"
);
assert!(measurement.result.is_ok());
}
}