use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, Instant};
use super::evaluator::{
check_cadence, project_evaluation, CadenceRefusal, EvaluationRequest, ReadinessEvaluation,
DEFAULT_ATTESTATION_CADENCE_FLOOR,
};
use super::identity::{
AudienceScopeCommitment, CanonicalConstraints, CapabilityId, CapabilityInterestKey, Digest256,
InterestSpec, WorkLatencyEnvelope,
};
use super::incarnation::Incarnation;
use super::wire::UnsignedAttestation;
pub const MAX_LIVE_SENSING_STREAMS: usize = 1024;
pub const MAX_STREAM_SLOTS: usize = 8192;
pub const STREAM_SLOTS_LOW_WATER: usize = 6144;
const OVERFLOW_PARK: Duration = Duration::from_secs(60 * 60 * 24 * 365);
fn schedule_after(base: Instant, delta: Duration) -> Instant {
base.checked_add(delta)
.or_else(|| base.checked_add(OVERFLOW_PARK))
.unwrap_or(base)
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum StreamRefusal {
Cadence(CadenceRefusal),
AtCapacity,
}
#[derive(Debug)]
struct CompiledPredicate {
capability_id: CapabilityId,
constraints: CanonicalConstraints,
work_latency: WorkLatencyEnvelope,
}
#[derive(Debug)]
struct LiveStream {
key: CapabilityInterestKey,
predicate: Arc<CompiledPredicate>,
audience: AudienceScopeCommitment,
promised_cadence: Duration,
due_at: Instant,
last_emitted_at: Option<Instant>,
registered_stamp: u64,
}
#[derive(Debug)]
struct StreamSlot {
next_seq: u64,
touched: u64,
live: Option<LiveStream>,
}
#[derive(Debug)]
pub struct DueBeat {
key: CapabilityInterestKey,
predicate: Arc<CompiledPredicate>,
audience: AudienceScopeCommitment,
origin: u64,
incarnation: Incarnation,
generation: u64,
seq: u64,
promised_cadence: Duration,
stamp: u64,
}
impl DueBeat {
pub fn key(&self) -> &CapabilityInterestKey {
&self.key
}
pub fn stamp(&self) -> u64 {
self.stamp
}
pub fn request(&self) -> EvaluationRequest<'_> {
EvaluationRequest {
capability_id: &self.predicate.capability_id,
constraints: &self.predicate.constraints,
work_latency: &self.predicate.work_latency,
}
}
pub fn into_unsigned(self, evaluation: Option<ReadinessEvaluation>) -> UnsignedAttestation {
let evaluation = evaluation.unwrap_or(ReadinessEvaluation::TemporarilyUnevaluable);
let (status, status_reason) = project_evaluation(&evaluation);
let estimated_start = match evaluation {
ReadinessEvaluation::Ready { estimated_start } => estimated_start,
_ => None,
};
UnsignedAttestation {
interest_digest: self.key.interest_digest,
origin: self.origin,
origin_incarnation: self.incarnation,
capability_id: self.key.capability_id,
capability_generation: self.generation,
status,
status_reason,
estimated_start,
seq: self.seq,
promised_cadence: self.promised_cadence,
audience_scope: self.audience,
}
}
}
#[derive(Debug)]
pub struct OriginEmitter {
origin: u64,
incarnation: Incarnation,
cadence_floor: Duration,
slots: HashMap<Digest256, StreamSlot>,
live_count: usize,
touch_counter: u64,
}
impl OriginEmitter {
pub fn new(origin: u64, incarnation: Incarnation, cadence_floor: Duration) -> Self {
let cadence_floor = if cadence_floor.is_zero() {
DEFAULT_ATTESTATION_CADENCE_FLOOR
} else {
cadence_floor
};
Self {
origin,
incarnation,
cadence_floor,
slots: HashMap::new(),
live_count: 0,
touch_counter: 0,
}
}
fn touch(&mut self) -> u64 {
self.touch_counter += 1;
self.touch_counter
}
pub fn stamp(&self) -> u64 {
self.touch_counter
}
pub fn register(
&mut self,
spec: &InterestSpec,
strictest: Duration,
now: Instant,
) -> Result<(), StreamRefusal> {
check_cadence(strictest, self.cadence_floor).map_err(StreamRefusal::Cadence)?;
let promised_cadence = (strictest / 2).max(self.cadence_floor);
let key = CapabilityInterestKey::for_spec(spec);
let digest = key.interest_digest;
let already_live = self
.slots
.get(&digest)
.is_some_and(|slot| slot.live.is_some());
if !already_live && self.live_count >= MAX_LIVE_SENSING_STREAMS {
return Err(StreamRefusal::AtCapacity);
}
let stamp = self.touch();
let slot = self.slots.entry(digest).or_insert(StreamSlot {
next_seq: 0,
touched: 0,
live: None,
});
slot.touched = stamp;
match &mut slot.live {
Some(stream) => {
stream.registered_stamp = stamp;
if stream.promised_cadence != promised_cadence {
stream.promised_cadence = promised_cadence;
stream.due_at = match stream.last_emitted_at {
Some(last) => schedule_after(last, promised_cadence).max(now),
None => now,
};
}
}
None => {
slot.live = Some(LiveStream {
key,
predicate: Arc::new(CompiledPredicate {
capability_id: spec.capability_id.clone(),
constraints: spec.constraints.clone(),
work_latency: spec.work_latency,
}),
audience: spec.audience,
promised_cadence,
due_at: now,
last_emitted_at: None,
registered_stamp: stamp,
});
self.live_count += 1;
}
}
self.evict_retired();
Ok(())
}
pub fn retire(&mut self, digest: &Digest256) {
if let Some(slot) = self.slots.get_mut(digest) {
if slot.live.take().is_some() {
self.live_count -= 1;
}
}
}
pub fn retire_if_stale(&mut self, digest: &Digest256, seen_stamp: u64) -> bool {
let Some(slot) = self.slots.get_mut(digest) else {
return false;
};
let Some(stream) = &slot.live else {
return false;
};
if stream.registered_stamp > seen_stamp {
return false;
}
slot.live = None;
self.live_count -= 1;
true
}
pub fn poke(&mut self, capability_id: &CapabilityId, now: Instant) -> bool {
let floor = self.cadence_floor;
let mut moved = false;
for slot in self.slots.values_mut() {
let Some(stream) = &mut slot.live else {
continue;
};
if stream.key.capability_id != *capability_id {
continue;
}
let earliest = match stream.last_emitted_at {
Some(last) => schedule_after(last, floor).max(now),
None => now,
};
if earliest < stream.due_at {
stream.due_at = earliest;
moved = true;
}
}
moved
}
pub fn next_due(&self) -> Option<Instant> {
self.slots
.values()
.filter_map(|slot| slot.live.as_ref().map(|stream| stream.due_at))
.min()
}
pub fn collect_due(&mut self, now: Instant, generation: u64) -> Vec<DueBeat> {
let mut due = Vec::new();
let stamp = self.touch();
for slot in self.slots.values_mut() {
let Some(stream) = &mut slot.live else {
continue;
};
if stream.due_at > now {
continue;
}
let seq = slot.next_seq;
slot.next_seq += 1;
slot.touched = stamp;
stream.last_emitted_at = Some(now);
stream.due_at = schedule_after(now, stream.promised_cadence);
due.push(DueBeat {
key: stream.key.clone(),
predicate: stream.predicate.clone(),
audience: stream.audience,
origin: self.origin,
incarnation: self.incarnation,
generation,
seq,
promised_cadence: stream.promised_cadence,
stamp,
});
}
due
}
pub fn refusal_beat(
&mut self,
spec: &InterestSpec,
refusal: CadenceRefusal,
generation: u64,
) -> UnsignedAttestation {
let key = CapabilityInterestKey::for_spec(spec);
let digest = key.interest_digest;
let stamp = self.touch();
let slot = self.slots.entry(digest).or_insert(StreamSlot {
next_seq: 0,
touched: 0,
live: None,
});
slot.touched = stamp;
let seq = slot.next_seq;
slot.next_seq += 1;
let (status, status_reason) = refusal.as_status();
let beat = UnsignedAttestation {
interest_digest: digest,
origin: self.origin,
origin_incarnation: self.incarnation,
capability_id: key.capability_id,
capability_generation: generation,
status,
status_reason,
estimated_start: None,
seq,
promised_cadence: refusal.minimum_supported,
audience_scope: spec.audience,
};
self.evict_retired();
beat
}
fn evict_retired(&mut self) {
if self.slots.len() <= MAX_STREAM_SLOTS {
return;
}
let mut retired: Vec<(Digest256, u64)> = self
.slots
.iter()
.filter(|(_, slot)| slot.live.is_none())
.map(|(digest, slot)| (*digest, slot.touched))
.collect();
retired.sort_by_key(|(_, touched)| *touched);
let excess = self.slots.len().saturating_sub(STREAM_SLOTS_LOW_WATER);
for (digest, _) in retired.into_iter().take(excess) {
self.slots.remove(&digest);
}
}
pub fn live_streams(&self) -> usize {
self.live_count
}
pub fn slot_count(&self) -> usize {
self.slots.len()
}
pub fn stream_cadence(&self, digest: &Digest256) -> Option<Duration> {
self.slots
.get(digest)
.and_then(|slot| slot.live.as_ref())
.map(|stream| stream.promised_cadence)
}
}
#[cfg(test)]
mod tests {
use std::time::{Duration, Instant};
use super::super::continuity::AttestedStatus;
use super::super::evaluator::StatusReason;
use super::super::identity::{DisclosureClass, ProviderSelector, ResultMode};
use super::*;
const FLOOR: Duration = DEFAULT_ATTESTATION_CADENCE_FLOOR;
fn spec(capability: &str, marker: &str) -> InterestSpec {
InterestSpec {
capability_id: CapabilityId::new(capability),
constraints: CanonicalConstraints::from_entries([("marker", marker)]).unwrap(),
work_latency: WorkLatencyEnvelope::start_within(Duration::from_millis(100)),
providers: ProviderSelector::AnyAuthorized,
result_mode: ResultMode::Any,
disclosure_class: DisclosureClass::Owner,
audience: AudienceScopeCommitment::from_bytes([7u8; 32]),
}
}
fn ready() -> Option<ReadinessEvaluation> {
Some(ReadinessEvaluation::Ready {
estimated_start: Some(Duration::from_millis(3)),
})
}
fn beats(
emitter: &mut OriginEmitter,
now: Instant,
generation: u64,
) -> Vec<(CapabilityInterestKey, UnsignedAttestation)> {
emitter
.collect_due(now, generation)
.into_iter()
.map(|beat| (beat.key().clone(), beat.into_unsigned(ready())))
.collect()
}
#[test]
fn first_beat_immediate_then_cadence_spacing_and_monotonic_seq() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(1), FLOOR);
let spec = spec("job.run", "a");
emitter
.register(&spec, Duration::from_millis(200), t0)
.unwrap();
assert_eq!(emitter.next_due(), Some(t0));
let out = beats(&mut emitter, t0, 5);
assert_eq!(out.len(), 1);
let (key, beat) = &out[0];
assert_eq!(key.interest_digest, spec.interest_digest());
assert_eq!(beat.seq, 0);
assert_eq!(beat.capability_generation, 5);
assert_eq!(beat.status, AttestedStatus::Ready);
assert_eq!(beat.promised_cadence, Duration::from_millis(100));
assert_eq!(beat.origin, 11);
assert!(beats(&mut emitter, t0 + Duration::from_millis(99), 6).is_empty());
let out = beats(&mut emitter, t0 + Duration::from_millis(100), 6);
assert_eq!(out.len(), 1);
assert_eq!(out[0].1.seq, 1);
assert_eq!(out[0].1.capability_generation, 6);
}
#[test]
fn two_phase_split_reserves_under_lock_and_seals_pure() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(2), FLOOR);
let spec = spec("job.run", "a");
emitter
.register(&spec, Duration::from_millis(200), t0)
.unwrap();
let due = emitter.collect_due(t0, 9);
assert_eq!(due.len(), 1);
assert_eq!(
emitter.next_due(),
Some(t0 + Duration::from_millis(100)),
"schedule re-armed before evaluation",
);
assert!(
emitter.collect_due(t0, 9).is_empty(),
"seq/schedule reserved exactly once",
);
let beat = due.into_iter().next().unwrap();
assert_eq!(beat.request().capability_id.as_str(), "job.run");
let unsigned = beat.into_unsigned(None);
assert_eq!(unsigned.status, AttestedStatus::ProviderUnknown);
assert_eq!(unsigned.status_reason, StatusReason::TemporarilyUnevaluable);
assert_eq!(unsigned.seq, 0);
assert_eq!(unsigned.origin_incarnation, Incarnation::new(2));
}
#[test]
fn refresh_does_not_starve_cadence_but_tightening_reschedules() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(1), FLOOR);
let spec = spec("job.run", "a");
emitter
.register(&spec, Duration::from_millis(400), t0)
.unwrap();
assert_eq!(
emitter.stream_cadence(&spec.interest_digest()),
Some(Duration::from_millis(200)),
);
let _ = beats(&mut emitter, t0, 1);
assert_eq!(emitter.next_due(), Some(t0 + Duration::from_millis(200)));
emitter
.register(
&spec,
Duration::from_millis(400),
t0 + Duration::from_millis(50),
)
.unwrap();
assert_eq!(emitter.next_due(), Some(t0 + Duration::from_millis(200)));
emitter
.register(
&spec,
Duration::from_millis(120),
t0 + Duration::from_millis(50),
)
.unwrap();
assert_eq!(
emitter.stream_cadence(&spec.interest_digest()),
Some(Duration::from_millis(60)),
);
assert_eq!(emitter.next_due(), Some(t0 + Duration::from_millis(60)));
emitter
.register(&spec, FLOOR, t0 + Duration::from_millis(50))
.unwrap();
assert_eq!(emitter.stream_cadence(&spec.interest_digest()), Some(FLOOR));
}
#[test]
fn poke_pulls_forward_with_floor_min_gap() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(1), FLOOR);
let spec = spec("job.run", "a");
emitter
.register(&spec, Duration::from_millis(400), t0)
.unwrap();
let _ = beats(&mut emitter, t0, 1);
assert!(emitter.poke(
&CapabilityId::new("job.run"),
t0 + Duration::from_millis(10)
));
assert_eq!(emitter.next_due(), Some(t0 + FLOOR));
let late = t0 + Duration::from_millis(150);
let _ = beats(&mut emitter, t0 + FLOOR, 1);
assert!(emitter.poke(&CapabilityId::new("job.run"), late));
assert_eq!(emitter.next_due(), Some(late));
assert!(!emitter.poke(&CapabilityId::new("other.cap"), late));
}
#[test]
fn refusal_beat_carries_floor_in_promised_cadence_and_shares_seq_space() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(1), FLOOR);
let spec = spec("job.run", "a");
let refused = emitter
.register(&spec, Duration::from_millis(10), t0)
.unwrap_err();
let StreamRefusal::Cadence(refusal) = refused else {
panic!("expected a cadence refusal, got {refused:?}");
};
assert_eq!(refusal.minimum_supported, FLOOR);
assert_eq!(emitter.live_streams(), 0);
let beat = emitter.refusal_beat(&spec, refusal, 9);
assert_eq!(beat.status, AttestedStatus::ProviderUnknown);
assert_eq!(
beat.status_reason,
StatusReason::SamplingIntervalUnsupported
);
assert_eq!(beat.promised_cadence, FLOOR);
assert_eq!(beat.estimated_start, None);
assert_eq!(beat.seq, 0);
emitter
.register(&spec, Duration::from_millis(200), t0)
.unwrap();
let out = beats(&mut emitter, t0, 9);
assert_eq!(out.len(), 1);
assert_eq!(out[0].1.seq, 1);
}
#[test]
fn retire_stops_emission_and_resurrection_keeps_seq() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(1), FLOOR);
let spec = spec("job.run", "a");
emitter
.register(&spec, Duration::from_millis(200), t0)
.unwrap();
let _ = beats(&mut emitter, t0, 1);
emitter.retire(&spec.interest_digest());
assert_eq!(emitter.live_streams(), 0);
assert_eq!(emitter.next_due(), None);
assert!(beats(&mut emitter, t0 + Duration::from_secs(5), 1).is_empty());
emitter
.register(
&spec,
Duration::from_millis(200),
t0 + Duration::from_secs(6),
)
.unwrap();
let out = beats(&mut emitter, t0 + Duration::from_secs(6), 1);
assert_eq!(out[0].1.seq, 1);
}
#[test]
fn compile_once_per_distinct_digest_and_per_interest_streams() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(1), FLOOR);
let a = spec("job.run", "a");
let b = spec("job.run", "b");
emitter
.register(&a, Duration::from_millis(200), t0)
.unwrap();
emitter
.register(&b, Duration::from_millis(400), t0)
.unwrap();
assert_eq!(emitter.live_streams(), 2);
let out = beats(&mut emitter, t0, 1);
assert_eq!(out.len(), 2);
assert!(out.iter().all(|(_, beat)| beat.seq == 0));
let out = beats(&mut emitter, t0 + Duration::from_millis(100), 1);
assert_eq!(out.len(), 1);
assert_eq!(out[0].0.interest_digest, a.interest_digest());
}
#[test]
fn live_capacity_refuses_new_digests_never_evicts_and_mints_no_slot() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(1), FLOOR);
let mut specs = Vec::new();
for i in 0..MAX_LIVE_SENSING_STREAMS {
let s = spec("job.run", &format!("live-{i}"));
emitter
.register(&s, Duration::from_millis(200), t0)
.unwrap();
specs.push(s);
}
assert_eq!(emitter.live_streams(), MAX_LIVE_SENSING_STREAMS);
let slots_at_cap = emitter.slot_count();
let overflow = spec("job.run", "overflow");
assert_eq!(
emitter.register(&overflow, Duration::from_millis(200), t0),
Err(StreamRefusal::AtCapacity),
);
assert_eq!(emitter.live_streams(), MAX_LIVE_SENSING_STREAMS);
assert_eq!(emitter.slot_count(), slots_at_cap, "no slot minted");
assert!(emitter
.register(&specs[0], Duration::from_millis(120), t0)
.is_ok());
assert_eq!(
emitter.stream_cadence(&specs[0].interest_digest()),
Some(Duration::from_millis(60)),
);
emitter.retire(&specs[1].interest_digest());
emitter
.register(&overflow, Duration::from_millis(200), t0)
.expect("one live slot freed");
assert_eq!(emitter.live_streams(), MAX_LIVE_SENSING_STREAMS);
assert_eq!(
emitter.register(&specs[1], Duration::from_millis(200), t0),
Err(StreamRefusal::AtCapacity),
"resurrection counts against the live cap",
);
}
#[test]
fn stamped_retire_skips_streams_registered_after_the_snapshot() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(1), FLOOR);
let spec = spec("job.run", "a");
emitter
.register(&spec, Duration::from_millis(200), t0)
.unwrap();
let digest = spec.interest_digest();
let seen = emitter.stamp();
emitter
.register(
&spec,
Duration::from_millis(200),
t0 + Duration::from_millis(5),
)
.unwrap();
assert!(!emitter.retire_if_stale(&digest, seen));
assert_eq!(emitter.live_streams(), 1);
let seen = emitter.stamp();
assert!(emitter.retire_if_stale(&digest, seen));
assert_eq!(emitter.live_streams(), 0);
assert!(!emitter.retire_if_stale(&digest, seen), "idempotent");
}
#[test]
fn absurd_durations_never_panic_and_park_instead_of_hot_looping() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(1), Duration::ZERO);
let spec = spec("job.run", "a");
assert!(matches!(
emitter.register(&spec, Duration::from_millis(1), t0),
Err(StreamRefusal::Cadence(_)),
));
emitter.register(&spec, Duration::MAX, t0).unwrap();
let out = beats(&mut emitter, t0, 1);
assert_eq!(out.len(), 1, "first beat still fires");
let next = emitter.next_due().expect("stream live");
assert!(
next > t0 + Duration::from_secs(60 * 60 * 24 * 30),
"overflowed schedule parks far in the future",
);
assert!(
beats(&mut emitter, t0 + Duration::from_secs(1), 1).is_empty(),
"no hot loop",
);
}
#[test]
fn retired_slots_evict_oldest_first_live_never() {
let t0 = Instant::now();
let mut emitter = OriginEmitter::new(11, Incarnation::new(1), FLOOR);
let live = spec("job.run", "live");
emitter
.register(&live, Duration::from_millis(200), t0)
.unwrap();
for i in 0..MAX_STREAM_SLOTS {
let s = spec("job.run", &format!("retired-{i}"));
emitter
.register(&s, Duration::from_millis(200), t0)
.unwrap();
emitter.retire(&s.interest_digest());
}
assert!(emitter.slot_count() <= STREAM_SLOTS_LOW_WATER + 1);
assert_eq!(emitter.live_streams(), 1);
assert_eq!(
emitter.stream_cadence(&live.interest_digest()),
Some(Duration::from_millis(100)),
);
}
}