#[cfg(test)]
mod tests;
use crate::bridge::sample_cursor::{ChunkEnd, LiveEpoch};
use indicatrix_net::client::Accumulator;
use std::{
collections::VecDeque,
sync::{Arc, Mutex},
time::{Duration, Instant},
};
pub const LIVE_CHUNK_TARGET_SECS: f64 = 1.5;
pub const LIVE_FIRST_CHUNK_SAMPLES: u32 = 8;
pub const LIVE_CHUNK_MIN_SAMPLES: u32 = 2;
pub use indicatrix_dispatch::DEFAULT_MAX_CHUNK_SAMPLES as LIVE_CHUNK_MAX_SAMPLES;
pub const MAX_CONSECUTIVE_CHUNK_FAILURES: u32 = 2;
#[must_use]
pub fn live_chunk_samples(rate_samples_per_sec: Option<f64>) -> u32 {
let Some(rate) = rate_samples_per_sec else {
return LIVE_FIRST_CHUNK_SAMPLES;
};
let target = rate * LIVE_CHUNK_TARGET_SECS;
if !target.is_finite() || target < f64::from(LIVE_CHUNK_MIN_SAMPLES) {
return LIVE_CHUNK_MIN_SAMPLES;
}
if target > f64::from(LIVE_CHUNK_MAX_SAMPLES) {
return LIVE_CHUNK_MAX_SAMPLES;
}
target.round() as u32
}
pub struct ChunkRequest {
pub request_id: u32,
pub first_sample: u32,
pub samples: u32,
pub accumulator: Arc<Mutex<Accumulator>>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ChunkVerdict {
Continue,
GaveUp,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum LaneState {
Active,
Exhausted,
GaveUp,
Abandoned,
}
struct InFlight {
request_id: u32,
dispatched_at: Instant,
first_progress: Option<(Instant, u32)>,
}
pub struct LiveLane {
epoch: Arc<LiveEpoch>,
local_lane: bool,
display_only: bool,
rate: Option<f64>,
consecutive_failures: u32,
in_flight: Option<InFlight>,
own_retry: VecDeque<(u32, u32)>,
state: LaneState,
}
impl LiveLane {
#[must_use]
pub const fn new(epoch: Arc<LiveEpoch>, local_lane: bool) -> Self {
Self {
epoch,
local_lane,
display_only: false,
rate: None,
consecutive_failures: 0,
in_flight: None,
own_retry: VecDeque::new(),
state: LaneState::Active,
}
}
#[must_use]
pub const fn display_only(epoch: Arc<LiveEpoch>) -> Self {
let mut lane = Self::new(epoch, false);
lane.display_only = true;
lane
}
#[must_use]
pub const fn is_display_only(&self) -> bool {
self.display_only
}
#[must_use]
pub const fn epoch(&self) -> &Arc<LiveEpoch> {
&self.epoch
}
#[cfg(test)]
#[must_use]
pub const fn rate(&self) -> Option<f64> {
self.rate
}
#[must_use]
pub fn is_idle(&self) -> bool {
self.state == LaneState::Active && self.in_flight.is_none()
}
#[must_use]
pub fn is_finished(&self) -> bool {
self.state != LaneState::Active && self.in_flight.is_none()
}
pub fn next_chunk(&mut self, request_id: u32, now: Instant) -> Option<ChunkRequest> {
if !self.is_idle() {
return None;
}
let want = if self.display_only {
LIVE_CHUNK_MAX_SAMPLES
} else {
live_chunk_samples(self.rate)
};
let claimed = match self.own_retry.pop_front() {
Some((start, count)) if count > want => {
self.own_retry.push_front((start + want, count - want));
Some((start, want))
}
Some(range) => Some(range),
None => self.epoch.claim_remote(want),
};
let Some((first_sample, samples)) = claimed else {
self.state = LaneState::Exhausted;
return None;
};
let (width, height) = self.epoch.dimensions();
let mut accumulator = Accumulator::new(width, height);
accumulator.begin_request_for_range(request_id, first_sample, samples);
let accumulator = Arc::new(Mutex::new(accumulator));
self.epoch
.begin_chunk(first_sample, samples, Arc::clone(&accumulator));
self.in_flight = Some(InFlight {
request_id,
dispatched_at: now,
first_progress: None,
});
Some(ChunkRequest {
request_id,
first_sample,
samples,
accumulator,
})
}
pub const fn observe_progress(&mut self, request_id: u32, samples_done: u32, now: Instant) {
if let Some(f) = self.in_flight.as_mut()
&& f.request_id == request_id
&& f.first_progress.is_none()
&& samples_done > 0
{
f.first_progress = Some((now, samples_done));
}
}
pub fn chunk_done(&mut self, request_id: u32, now: Instant) -> Option<ChunkEnd> {
let flight = self.take_in_flight(request_id)?;
let end = self.epoch.finish_chunk()?;
if self.display_only {
self.consecutive_failures = 0;
self.state = LaneState::Exhausted;
return Some(ChunkEnd {
done: end.count,
..end
});
}
if let Some(measured) = measured_rate(&flight, end.done, now) {
let blended = self
.rate
.map_or(measured, |prev| prev.mul_add(0.5, measured * 0.5));
self.rate = Some(blended);
}
if end.done == end.count {
self.consecutive_failures = 0;
} else {
self.hand_back(end);
}
Some(end)
}
pub fn chunk_paused(&mut self, request_id: u32) -> Option<ChunkEnd> {
self.take_in_flight(request_id)?;
let end = self.epoch.finish_chunk()?;
if end.done < end.count {
self.hand_back(end);
}
Some(end)
}
pub fn chunk_failed(&mut self, request_id: u32) -> Option<ChunkVerdict> {
self.take_in_flight(request_id)?;
if let Some(end) = self.epoch.finish_chunk() {
self.hand_back(end);
}
self.consecutive_failures += 1;
if self.consecutive_failures >= MAX_CONSECUTIVE_CHUNK_FAILURES {
self.state = LaneState::GaveUp;
return Some(ChunkVerdict::GaveUp);
}
Some(ChunkVerdict::Continue)
}
pub fn abandon(&mut self) -> Option<u32> {
self.state = LaneState::Abandoned;
self.epoch.abandon_chunk();
self.own_retry.clear();
self.in_flight.take().map(|f| f.request_id)
}
fn take_in_flight(&mut self, request_id: u32) -> Option<InFlight> {
if self.state == LaneState::Abandoned {
return None;
}
match &self.in_flight {
Some(f) if f.request_id == request_id => self.in_flight.take(),
_ => None,
}
}
fn hand_back(&mut self, end: ChunkEnd) {
let (start, count) = end.remainder();
if count == 0 {
return;
}
if self.local_lane {
self.epoch.return_to_local(start, count);
} else {
self.own_retry.push_back((start, count));
}
}
}
fn measured_rate(flight: &InFlight, done: u32, now: Instant) -> Option<f64> {
const MIN_SPAN: Duration = Duration::from_millis(50);
if let Some((t0, s0)) = flight.first_progress {
let span = now.saturating_duration_since(t0);
if done > s0 && span >= MIN_SPAN {
return Some(f64::from(done - s0) / span.as_secs_f64());
}
}
let span = now.saturating_duration_since(flight.dispatched_at);
if done == 0 || span.is_zero() {
return None;
}
Some(f64::from(done) / span.as_secs_f64())
}
#[cfg(test)]
mod chunk_paused_tests {
use super::*;
use crate::bridge::sample_cursor::tests::apply_frame;
fn t(ms: u64) -> Instant {
static ORIGIN: std::sync::OnceLock<Instant> = std::sync::OnceLock::new();
*ORIGIN.get_or_init(Instant::now) + Duration::from_millis(ms)
}
#[test]
fn a_fully_traced_chunk_paused_at_completion_merges_everything() {
let epoch = Arc::new(LiveEpoch::new(1, 1, 100));
let mut lane = LiveLane::new(Arc::clone(&epoch), true);
let chunk = lane.next_chunk(1, t(0)).unwrap();
apply_frame(
&chunk.accumulator,
1,
chunk.first_sample,
chunk.samples,
2.0,
);
let end = lane.chunk_paused(1).expect("in flight");
assert_eq!(end.remainder(), (chunk.first_sample + chunk.samples, 0));
assert_eq!(epoch.remote_done(), chunk.samples);
}
#[test]
fn a_chunk_paused_mid_flight_merges_the_prefix_and_hands_back_the_remainder() {
let epoch = Arc::new(LiveEpoch::new(1, 1, 100));
let mut lane = LiveLane::new(Arc::clone(&epoch), true);
let chunk = lane.next_chunk(1, t(0)).unwrap();
assert!(chunk.samples > 1, "need room for a partial trace");
apply_frame(&chunk.accumulator, 1, chunk.first_sample, 1, 3.0);
let end = lane.chunk_paused(1).expect("in flight");
assert_eq!(end.done, 1);
assert_eq!(epoch.remote_done(), 1, "the valid prefix was merged");
assert!(lane.is_idle(), "left Active and idle, not GaveUp/Exhausted");
for _ in 0..=MAX_CONSECUTIVE_CHUNK_FAILURES {
let c = lane.next_chunk(2, t(10)).expect("still claimable");
assert!(
lane.chunk_paused(c.request_id).is_some(),
"repeated pauses must never exhaust the failure budget"
);
}
assert!(
lane.is_idle(),
"chunk_paused must never drive the lane to GaveUp"
);
}
#[test]
fn a_mismatched_request_id_is_a_no_op() {
let epoch = Arc::new(LiveEpoch::new(1, 1, 100));
let mut lane = LiveLane::new(Arc::clone(&epoch), true);
let chunk = lane.next_chunk(1, t(0)).unwrap();
assert!(lane.chunk_paused(chunk.request_id + 1).is_none());
assert!(lane.chunk_paused(chunk.request_id).is_some());
}
}