use core::time::Duration;
use std::sync::{Arc, Mutex};
use std::time::Instant;
use keel_core_api::policy::NondeterminismResponse;
use keel_core_api::{
ENVELOPE_VERSION, ErrorClass, ErrorCode, KeelError, Outcome, OutcomeError, Request,
};
use keel_journal::{
Clock, FlowId, FlowStatus, Journal, NewFlow, ProcessId, StepKey, StepKind, StepOutcome,
StepStatus,
};
use serde_json::{Value, json};
use tracing::{debug, warn};
const BRANCH_SEQ_BASE: u64 = 1_000_000;
const ATTEMPT_SEQ: u64 = 0;
const ATTEMPT_KEY: &str = "flow:attempt";
const IDEMPOTENCY_KEY_FIELD: &str = "idempotency_key";
use crate::engine::Engine;
#[derive(Debug, Clone)]
pub struct FlowDescriptor {
pub entrypoint: String,
pub args_hash: String,
pub explicit_key: Option<String>,
pub code_hash: Option<String>,
}
impl FlowDescriptor {
#[must_use]
pub fn flow_id(&self) -> FlowId {
FlowId::new(format!(
"{}#{}#{}",
self.entrypoint,
self.args_hash,
self.explicit_key.as_deref().unwrap_or("")
))
}
}
#[derive(Debug, Clone, Copy)]
pub struct FlowConfig {
pub lease_ttl: Duration,
pub max_attempts: u32,
}
impl Default for FlowConfig {
fn default() -> Self {
Self {
lease_ttl: Duration::from_secs(30),
max_attempts: 3,
}
}
}
pub struct FlowManager {
engine: Arc<Engine>,
journal: Arc<dyn Journal>,
clock: Arc<dyn Clock>,
holder: ProcessId,
config: FlowConfig,
}
impl core::fmt::Debug for FlowManager {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("FlowManager")
.field("holder", &self.holder)
.field("config", &self.config)
.finish_non_exhaustive()
}
}
impl FlowManager {
#[must_use]
pub fn new(
engine: Arc<Engine>,
journal: Arc<dyn Journal>,
clock: Arc<dyn Clock>,
holder: ProcessId,
) -> Self {
Self::with_config(engine, journal, clock, holder, FlowConfig::default())
}
#[must_use]
pub fn with_config(
engine: Arc<Engine>,
journal: Arc<dyn Journal>,
clock: Arc<dyn Clock>,
holder: ProcessId,
config: FlowConfig,
) -> Self {
Self {
engine,
journal,
clock,
holder,
config,
}
}
pub fn enter_flow(&self, desc: &FlowDescriptor) -> Result<FlowHandle, KeelError> {
self.enter(
&desc.flow_id(),
&desc.entrypoint,
&desc.args_hash,
desc.code_hash.as_deref(),
)
}
pub fn resume_flow(
&self,
flow: &keel_journal::FlowDescriptor,
current_code_hash: Option<&str>,
) -> Result<FlowHandle, KeelError> {
self.enter(
&flow.flow_id,
&flow.entrypoint,
&flow.args_hash,
current_code_hash,
)
}
fn bump_attempt_or_kill(&self, flow_id: &FlowId, entrypoint: &str) -> Result<u32, KeelError> {
let prior = self
.journal
.step_at(flow_id, ATTEMPT_SEQ)
.map_err(|e| internal(format!("attempt lookup failed: {e}")))?
.map_or(0, |(_, o)| o.attempt);
let attempt = prior.saturating_add(1);
if attempt >= 2 {
crate::metrics::record_flow_resume(entrypoint);
}
if attempt > self.config.max_attempts {
self.journal
.complete_flow(flow_id, FlowStatus::Dead)
.map_err(|e| internal(format!("mark-dead failed: {e}")))?;
return Err(KeelError {
code: ErrorCode::FlowDead,
message: format!(
"flow {flow_id} exceeded its {} attempt cap; marked dead (KEEL-E032)",
self.config.max_attempts
),
});
}
Ok(attempt)
}
fn enter(
&self,
flow_id: &FlowId,
entrypoint: &str,
args_hash: &str,
current_code_hash: Option<&str>,
) -> Result<FlowHandle, KeelError> {
self.journal
.begin_flow(&NewFlow {
flow_id: flow_id.clone(),
entrypoint: entrypoint.to_owned(),
args_hash: args_hash.to_owned(),
code_hash: current_code_hash.map(str::to_owned),
})
.map_err(|e| internal(format!("begin_flow failed: {e}")))?;
let existing = self
.journal
.get_flow(flow_id)
.map_err(|e| internal(format!("get_flow failed: {e}")))?;
let status = existing.as_ref().map_or(FlowStatus::Running, |f| f.status);
if status == FlowStatus::Dead {
return Err(KeelError {
code: ErrorCode::FlowDead,
message: format!("flow {flow_id} is dead; refusing to resume (KEEL-E032)"),
});
}
if status == FlowStatus::Completed {
return Ok(self.new_handle(
flow_id.clone(),
status,
false,
None,
LeaseHeartbeatMonitor::new(),
true,
));
}
let attempt = self.bump_attempt_or_kill(flow_id, entrypoint)?;
if status == FlowStatus::Failed {
self.journal
.complete_flow(flow_id, FlowStatus::Running)
.map_err(|e| internal(format!("reset-to-running failed: {e}")))?;
}
let acquired = self
.journal
.acquire_lease(flow_id, &self.holder, self.config.lease_ttl)
.map_err(|e| internal(format!("acquire_lease failed: {e}")))?;
if !acquired {
return Err(KeelError {
code: ErrorCode::FlowLeaseHeld,
message: format!(
"flow {flow_id} is leased by another holder; not resuming (KEEL-E030)"
),
});
}
let now = self.clock.now_ms();
let marker = StepOutcome {
kind: StepKind::Marker,
attempt,
status: StepStatus::Ok,
payload: None,
error_class: None,
started_at: now,
ended_at: Some(now),
};
if let Err(e) =
self.journal
.record_step(flow_id, ATTEMPT_SEQ, &StepKey::new(ATTEMPT_KEY), &marker)
{
warn!(flow = %flow_id, error = %e, "attempt-counter record failed");
}
let code_hash_fenced = match (existing.and_then(|f| f.code_hash), current_code_hash) {
(Some(recorded), Some(current)) => recorded != current,
_ => false,
};
let lease_monitor = LeaseHeartbeatMonitor::new();
let heartbeat = spawn_heartbeat(
Arc::clone(&self.journal),
flow_id.clone(),
self.holder.clone(),
self.config.lease_ttl,
lease_monitor.clone(),
);
Ok(self.new_handle(
flow_id.clone(),
status,
code_hash_fenced,
heartbeat,
lease_monitor,
false,
))
}
fn new_handle(
&self,
flow_id: FlowId,
entry_status: FlowStatus,
code_hash_fenced: bool,
heartbeat: Option<HeartbeatHandle>,
lease_monitor: LeaseHeartbeatMonitor,
replay_only: bool,
) -> FlowHandle {
FlowHandle {
engine: Arc::clone(&self.engine),
journal: Arc::clone(&self.journal),
clock: Arc::clone(&self.clock),
holder: self.holder.clone(),
flow_id,
seq: 0,
code_hash_fenced,
replay_abandoned: false,
replay_only,
completed: replay_only,
entry_status,
heartbeat,
lease_ttl: self.config.lease_ttl,
lease_monitor,
}
}
}
#[expect(
clippy::struct_excessive_bools,
reason = "each flag is an independent per-handle replay/lease predicate, not a \
packed state enum: code_hash_fenced, replay_abandoned, replay_only, completed"
)]
pub struct FlowHandle {
engine: Arc<Engine>,
journal: Arc<dyn Journal>,
clock: Arc<dyn Clock>,
holder: ProcessId,
flow_id: FlowId,
seq: u64,
code_hash_fenced: bool,
replay_abandoned: bool,
replay_only: bool,
entry_status: FlowStatus,
completed: bool,
heartbeat: Option<HeartbeatHandle>,
lease_ttl: Duration,
lease_monitor: LeaseHeartbeatMonitor,
}
impl core::fmt::Debug for FlowHandle {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("FlowHandle")
.field("flow_id", &self.flow_id)
.field("seq", &self.seq)
.field("entry_status", &self.entry_status)
.field("replay_only", &self.replay_only)
.field("completed", &self.completed)
.finish_non_exhaustive()
}
}
enum StepPlan {
Live,
Replay(StepOutcome),
Diverged { recorded: StepKey },
}
struct Divergence<'a> {
seq: u64,
recorded: &'a StepKey,
observed: &'a StepKey,
mode: &'static str,
preserve: bool,
}
impl FlowHandle {
#[must_use]
pub fn flow_id(&self) -> &FlowId {
&self.flow_id
}
#[must_use]
pub fn entry_status(&self) -> FlowStatus {
self.entry_status
}
#[must_use]
pub fn is_replay_only(&self) -> bool {
self.replay_only
}
fn now(&self) -> i64 {
self.clock.now_ms()
}
fn lease_lost(&self) -> bool {
let still_recorded_holder = match self.journal.get_flow(&self.flow_id) {
Ok(Some(flow)) => flow.lease_holder.as_ref() == Some(&self.holder),
Ok(None) | Err(_) => true,
};
!still_recorded_holder || self.lease_monitor.elapsed() >= self.lease_ttl
}
fn plan_step(&self, seq: u64, key: &StepKey) -> StepPlan {
if self.replay_abandoned {
return StepPlan::Live;
}
match self.journal.step_at(&self.flow_id, seq) {
Ok(None) => StepPlan::Live,
Ok(Some((recorded_key, outcome))) => {
if recorded_key == *key {
match outcome.status {
StepStatus::Running => StepPlan::Live,
StepStatus::Ok | StepStatus::Error => StepPlan::Replay(outcome),
}
} else {
StepPlan::Diverged {
recorded: recorded_key,
}
}
}
Err(e) => {
warn!(flow = %self.flow_id, seq, error = %e, "step_at failed; executing live");
StepPlan::Live
}
}
}
fn step_key(request: &Request) -> StepKey {
StepKey::new(format!(
"{}#{}",
request.target,
request.args_hash.as_deref().unwrap_or("-")
))
}
pub async fn execute_step<F>(&mut self, request: &Request, effect: F) -> Outcome
where
F: AsyncFnMut(u32) -> keel_core_api::AttemptResult,
{
self.execute_step_with_idempotency_key(request, None, effect)
.await
}
pub async fn execute_step_with_idempotency_key<F>(
&mut self,
request: &Request,
idempotency_key: Option<&str>,
effect: F,
) -> Outcome
where
F: AsyncFnMut(u32) -> keel_core_api::AttemptResult,
{
self.seq += 1;
let seq = self.seq;
let key = Self::step_key(request);
let plan = self.plan_step(seq, &key);
if self.replay_only {
return match plan {
StepPlan::Replay(outcome) => replay_outcome(&self.flow_id, seq, &outcome),
_ => replay_miss_outcome(&self.flow_id, seq, &key),
};
}
match plan {
StepPlan::Replay(outcome) => replay_outcome(&self.flow_id, seq, &outcome),
StepPlan::Diverged { recorded } => {
self.on_divergence(seq, &recorded, &key, request, idempotency_key, effect)
.await
}
StepPlan::Live => {
self.run_live(
seq,
&key,
StepKind::Effect,
request,
idempotency_key,
effect,
)
.await
}
}
}
#[must_use]
pub fn recorded_idempotency_key(&self, step_key: &str) -> Option<String> {
if self.replay_abandoned || self.replay_only {
return None;
}
let key = StepKey::new(step_key);
match self.journal.step_at(&self.flow_id, self.seq + 1) {
Ok(Some((recorded, outcome)))
if recorded == key && outcome.status == StepStatus::Running =>
{
outcome
.payload
.as_deref()
.and_then(decode_payload)
.and_then(|v| {
v.get(IDEMPOTENCY_KEY_FIELD)
.and_then(Value::as_str)
.map(str::to_owned)
})
}
_ => None,
}
}
fn effective_response(&self) -> NondeterminismResponse {
let configured = self.engine.nondeterminism_response();
if self.code_hash_fenced && configured == NondeterminismResponse::Fail {
NondeterminismResponse::Warn
} else {
configured
}
}
async fn on_divergence<F>(
&mut self,
seq: u64,
recorded: &StepKey,
observed: &StepKey,
request: &Request,
idempotency_key: Option<&str>,
effect: F,
) -> Outcome
where
F: AsyncFnMut(u32) -> keel_core_api::AttemptResult,
{
let (mode, preserve) = match self.effective_response() {
NondeterminismResponse::Fail => {
return diverged_outcome(&self.flow_id, seq, recorded, observed);
}
NondeterminismResponse::Warn => ("warn", false),
NondeterminismResponse::Branch => ("branch", true),
};
let div = Divergence {
seq,
recorded,
observed,
mode,
preserve,
};
self.branch_and_continue(div, request, idempotency_key, effect)
.await
}
async fn branch_and_continue<F>(
&mut self,
div: Divergence<'_>,
request: &Request,
idempotency_key: Option<&str>,
effect: F,
) -> Outcome
where
F: AsyncFnMut(u32) -> keel_core_api::AttemptResult,
{
self.journal_branch_marker(&div);
let live_seq = self.seq;
self.run_live(
live_seq,
div.observed,
StepKind::Effect,
request,
idempotency_key,
effect,
)
.await
}
fn journal_branch_marker(&mut self, div: &Divergence<'_>) {
warn!(
flow = %self.flow_id, seq = div.seq, mode = div.mode,
expected = %div.recorded, observed = %div.observed,
"flow nondeterminism; abandoning replay"
);
self.replay_abandoned = true;
let marker_seq = if div.preserve {
BRANCH_SEQ_BASE + div.seq
} else {
div.seq
};
let now = self.now();
self.record(
marker_seq,
&StepKey::new(format!("flow:branch:{}", div.mode)),
&StepOutcome {
kind: StepKind::Marker,
attempt: 0,
status: StepStatus::Ok,
payload: encode_payload(&json!({
"mode": div.mode,
"expected": div.recorded.as_str(),
"observed": div.observed.as_str(),
})),
error_class: None,
started_at: now,
ended_at: Some(now),
},
);
self.seq = marker_seq + 1;
}
pub fn journal_time(&mut self, key: &str, now_ms: i64) -> Result<i64, KeelError> {
let bytes = encode_payload(&json!(now_ms)).unwrap_or_default();
let recorded = self.resolve_value_step(key, StepKind::Time, bytes)?;
Ok(decode_payload(&recorded)
.and_then(|v| v.as_i64())
.unwrap_or(now_ms))
}
pub fn journal_random(&mut self, key: &str, bytes: Vec<u8>) -> Result<Vec<u8>, KeelError> {
self.resolve_value_step(key, StepKind::Random, bytes)
}
fn resolve_value_step(
&mut self,
key_str: &str,
kind: StepKind,
live_bytes: Vec<u8>,
) -> Result<Vec<u8>, KeelError> {
self.seq += 1;
let seq = self.seq;
let key = StepKey::new(key_str);
let plan = self.plan_step(seq, &key);
if self.replay_only {
return match plan {
StepPlan::Replay(outcome) => Ok(outcome.payload.unwrap_or_default()),
_ => Err(replay_miss_error(&self.flow_id, seq, &key)),
};
}
match plan {
StepPlan::Replay(outcome) => Ok(outcome.payload.unwrap_or_default()),
StepPlan::Live => {
self.record_value(seq, &key, kind, &live_bytes);
Ok(live_bytes)
}
StepPlan::Diverged { recorded } => {
let (mode, preserve) = match self.effective_response() {
NondeterminismResponse::Fail => {
return Err(diverged_error(&self.flow_id, seq, &recorded, &key));
}
NondeterminismResponse::Warn => ("warn", false),
NondeterminismResponse::Branch => ("branch", true),
};
let div = Divergence {
seq,
recorded: &recorded,
observed: &key,
mode,
preserve,
};
self.journal_branch_marker(&div);
let live_seq = self.seq;
self.record_value(live_seq, &key, kind, &live_bytes);
Ok(live_bytes)
}
}
}
fn record_value(&self, seq: u64, key: &StepKey, kind: StepKind, bytes: &[u8]) {
let now = self.now();
self.record(
seq,
key,
&StepOutcome {
kind,
attempt: 0,
status: StepStatus::Ok,
payload: Some(bytes.to_vec()),
error_class: None,
started_at: now,
ended_at: Some(now),
},
);
}
async fn run_live<F>(
&mut self,
seq: u64,
key: &StepKey,
kind: StepKind,
request: &Request,
idempotency_key: Option<&str>,
effect: F,
) -> Outcome
where
F: AsyncFnMut(u32) -> keel_core_api::AttemptResult,
{
if self.lease_lost() {
warn!(flow = %self.flow_id, seq, "lease lost before live step; refusing to double-execute (KEEL-E030)");
return lease_lost_outcome(&self.flow_id, seq);
}
let started_at = self.now();
self.record(
seq,
key,
&StepOutcome {
kind,
attempt: 0,
status: StepStatus::Running,
payload: idempotency_key
.and_then(|k| encode_payload(&json!({ (IDEMPOTENCY_KEY_FIELD): k }))),
error_class: None,
started_at,
ended_at: None,
},
);
let outcome = self.engine.execute(request, effect).await;
let ended_at = self.now();
let (status, payload, error_class) = if outcome.result == "ok" {
(
StepStatus::Ok,
outcome.payload.as_ref().and_then(encode_payload),
None,
)
} else {
(
StepStatus::Error,
None,
outcome.error.as_ref().map(|e| e.class),
)
};
self.record(
seq,
key,
&StepOutcome {
kind,
attempt: outcome.attempts,
status,
payload,
error_class,
started_at,
ended_at: Some(ended_at),
},
);
outcome
}
fn record(&self, seq: u64, key: &StepKey, outcome: &StepOutcome) {
if let Err(e) = self.journal.record_step(&self.flow_id, seq, key, outcome) {
warn!(flow = %self.flow_id, seq, error = %e, "record_step failed; step not journaled");
}
}
pub fn complete(&mut self, status: FlowStatus) {
if let Err(e) = self.journal.complete_flow(&self.flow_id, status) {
warn!(flow = %self.flow_id, error = %e, "complete_flow failed");
}
self.completed = true;
}
pub fn complete_success(&mut self) {
self.complete(FlowStatus::Completed);
}
pub fn complete_failed(&mut self) {
self.complete(FlowStatus::Failed);
}
}
impl Drop for FlowHandle {
fn drop(&mut self) {
self.heartbeat = None;
if !self.completed {
debug!(flow = %self.flow_id, "flow handle dropped uncompleted; left running for recovery");
}
}
}
#[derive(Debug, Clone)]
struct LeaseHeartbeatMonitor(Arc<Mutex<Instant>>);
impl LeaseHeartbeatMonitor {
fn new() -> Self {
Self(Arc::new(Mutex::new(Instant::now())))
}
fn mark_renewed(&self) {
let mut at = self
.0
.lock()
.expect("lease heartbeat monitor lock poisoned");
*at = Instant::now();
}
fn elapsed(&self) -> Duration {
self.0
.lock()
.expect("lease heartbeat monitor lock poisoned")
.elapsed()
}
}
struct HeartbeatHandle {
stop: Option<std::sync::mpsc::Sender<()>>,
thread: Option<std::thread::JoinHandle<()>>,
}
impl Drop for HeartbeatHandle {
fn drop(&mut self) {
drop(self.stop.take());
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
}
}
fn spawn_heartbeat(
journal: Arc<dyn Journal>,
flow: FlowId,
holder: ProcessId,
ttl: Duration,
monitor: LeaseHeartbeatMonitor,
) -> Option<HeartbeatHandle> {
use std::sync::mpsc::{RecvTimeoutError, channel};
let period = (ttl / 2).max(Duration::from_millis(1));
let (stop, rx) = channel::<()>();
let thread = std::thread::Builder::new()
.name(String::from("keel-lease-heartbeat"))
.spawn(move || {
while let Err(RecvTimeoutError::Timeout) = rx.recv_timeout(period) {
match journal.acquire_lease(&flow, &holder, ttl) {
Ok(true) => monitor.mark_renewed(),
Ok(false) => {
warn!(flow = %flow, "lease heartbeat: lease lost to another holder");
}
Err(e) => {
warn!(flow = %flow, error = %e, "lease heartbeat renewal failed");
}
}
}
})
.ok()?;
Some(HeartbeatHandle {
stop: Some(stop),
thread: Some(thread),
})
}
fn internal(message: String) -> KeelError {
KeelError {
code: ErrorCode::Internal,
message,
}
}
const STEP_PAYLOAD_SCHEMA: &str = "keel.step/v1";
#[derive(serde::Serialize)]
struct StepPayloadRef<'a> {
schema: &'a str,
payload: &'a Value,
}
#[derive(serde::Deserialize)]
struct StepPayloadOwned {
schema: String,
payload: Value,
}
fn encode_payload(value: &Value) -> Option<Vec<u8>> {
rmp_serde::to_vec_named(&StepPayloadRef {
schema: STEP_PAYLOAD_SCHEMA,
payload: value,
})
.ok()
}
fn decode_payload(bytes: &[u8]) -> Option<Value> {
if let Ok(envelope) = rmp_serde::from_slice::<StepPayloadOwned>(bytes)
&& envelope.schema == STEP_PAYLOAD_SCHEMA
{
return Some(envelope.payload);
}
rmp_serde::from_slice(bytes).ok()
}
fn step_trace(flow: &FlowId, seq: u64) -> String {
format!("flow-{flow}-s{seq}")
}
fn replay_outcome(flow: &FlowId, seq: u64, step: &StepOutcome) -> Outcome {
let mut outcome = base_outcome(step_trace(flow, seq));
outcome.attempts = step.attempt;
match step.status {
StepStatus::Ok | StepStatus::Running => {
outcome.result = String::from("ok");
outcome.payload = step.payload.as_deref().and_then(decode_payload);
}
StepStatus::Error => {
outcome.error = Some(OutcomeError {
code: ErrorCode::NonRetryableError,
class: step.error_class.unwrap_or(ErrorClass::Other),
http_status: None,
message: String::from("replayed failed step"),
original: None,
});
}
}
outcome
}
fn divergence_message(flow: &FlowId, seq: u64, recorded: &StepKey, observed: &StepKey) -> String {
format!("flow {flow} diverged at step {seq}: expected {recorded}, got {observed} (KEEL-E031)")
}
fn diverged_outcome(flow: &FlowId, seq: u64, recorded: &StepKey, observed: &StepKey) -> Outcome {
let mut outcome = base_outcome(step_trace(flow, seq));
outcome.error = Some(OutcomeError {
code: ErrorCode::FlowNondeterminism,
class: ErrorClass::Other,
http_status: None,
message: divergence_message(flow, seq, recorded, observed),
original: None,
});
outcome
}
fn lease_lost_outcome(flow: &FlowId, seq: u64) -> Outcome {
let mut outcome = base_outcome(step_trace(flow, seq));
outcome.error = Some(OutcomeError {
code: ErrorCode::FlowLeaseHeld,
class: ErrorClass::Other,
http_status: None,
message: format!(
"flow {flow} lost its lease before step {seq}; another holder may be resuming it \
(KEEL-E030). Refusing to run the effect to avoid double execution."
),
original: None,
});
outcome
}
fn diverged_error(flow: &FlowId, seq: u64, recorded: &StepKey, observed: &StepKey) -> KeelError {
KeelError {
code: ErrorCode::FlowNondeterminism,
message: divergence_message(flow, seq, recorded, observed),
}
}
fn replay_miss_message(flow: &FlowId, seq: u64, observed: &StepKey) -> String {
format!(
"flow {flow} replay reached unrecorded step {seq} ({observed}); \
the completed flow's code changed (KEEL-E031)"
)
}
fn replay_miss_outcome(flow: &FlowId, seq: u64, observed: &StepKey) -> Outcome {
let mut outcome = base_outcome(step_trace(flow, seq));
outcome.error = Some(OutcomeError {
code: ErrorCode::FlowNondeterminism,
class: ErrorClass::Other,
http_status: None,
message: replay_miss_message(flow, seq, observed),
original: None,
});
outcome
}
fn replay_miss_error(flow: &FlowId, seq: u64, observed: &StepKey) -> KeelError {
KeelError {
code: ErrorCode::FlowNondeterminism,
message: replay_miss_message(flow, seq, observed),
}
}
fn base_outcome(trace_id: String) -> Outcome {
Outcome {
v: ENVELOPE_VERSION,
result: String::from("error"),
payload: None,
error: None,
attempts: 0,
from_cache: false,
waits_ms: Vec::new(),
throttled: false,
throttle_wait_ms: 0,
breaker: keel_core_api::BreakerState::Closed,
trace_id,
}
}
#[cfg(test)]
mod tests {
use super::{LeaseHeartbeatMonitor, decode_payload, encode_payload};
use core::time::Duration;
use serde_json::json;
#[test]
fn lease_heartbeat_monitor_tracks_real_elapsed_time_since_the_last_renewal() {
let monitor = LeaseHeartbeatMonitor::new();
assert!(
monitor.elapsed() < Duration::from_millis(50),
"freshly constructed: elapsed should be ~0"
);
std::thread::sleep(Duration::from_millis(60));
assert!(
monitor.elapsed() >= Duration::from_millis(60),
"elapsed grows with real time when nothing renews"
);
monitor.mark_renewed();
assert!(
monitor.elapsed() < Duration::from_millis(50),
"a renewal resets elapsed back to ~0"
);
}
#[test]
fn payload_round_trips_through_the_schema_tag() {
let value = json!({ "rows": 120, "nested": [1, 2, 3], "ok": true });
let bytes = encode_payload(&value).expect("encodes");
assert_eq!(decode_payload(&bytes), Some(value));
}
#[test]
fn legacy_bare_messagepack_still_decodes() {
let map = json!({ "rows": 120 });
let bare_map = rmp_serde::to_vec_named(&map).expect("bare encodes");
assert_eq!(decode_payload(&bare_map), Some(map));
let num = json!(1_783_727_616u64);
let bare_num = rmp_serde::to_vec_named(&num).expect("bare encodes");
assert_eq!(bare_num, vec![0xCE, 0x6A, 0x51, 0x86, 0x00]);
assert_eq!(decode_payload(&bare_num), Some(num));
}
}