use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::sync::Mutex;
use crate::app_state::{AppState, ReplayState};
use crate::record::replay::{ReplayError, Replayer};
pub struct ReplayDriver {
replayer: Replayer,
displayed_at: Instant,
last_frame_time_offset: Option<Duration>,
}
impl ReplayDriver {
pub fn open(path: PathBuf) -> Result<Self, ReplayError> {
let replayer = Replayer::open(&path)?;
Ok(Self {
replayer,
displayed_at: Instant::now(),
last_frame_time_offset: None,
})
}
pub fn total_hosts(&self) -> Vec<String> {
self.replayer
.header()
.map(|h| h.hosts.clone())
.unwrap_or_default()
}
pub async fn tick(&mut self, app_state: Arc<Mutex<AppState>>) -> Result<(), ReplayError> {
let control = {
let state = app_state.lock().await;
state.replay.clone()
};
let Some(control) = control else {
return Ok(());
};
if let Some(seek_to) = control.pending_seek {
match self.replayer.seek(seek_to) {
Ok(_) => {}
Err(e) => {
tracing::warn!(error = %e, "replay: seek failed");
}
}
self.displayed_at = Instant::now();
self.last_frame_time_offset = self.replayer.elapsed();
apply_frame_to_state(app_state.clone(), &mut self.replayer, true).await;
clear_pending_seek(app_state.clone()).await;
return Ok(());
}
if let Some(step) = control.pending_step {
if step > 0 {
let _ = self.replayer.next()?;
} else if step < 0 {
let _ = self.replayer.prev()?;
}
self.displayed_at = Instant::now();
self.last_frame_time_offset = self.replayer.elapsed();
apply_frame_to_state(app_state.clone(), &mut self.replayer, true).await;
clear_pending_step(app_state.clone()).await;
return Ok(());
}
if control.paused {
return Ok(());
}
let Some(current) = self.replayer.current() else {
return Ok(());
};
let current_offset = self.replayer.elapsed().unwrap_or(Duration::ZERO);
if self.last_frame_time_offset.is_none() {
self.last_frame_time_offset = Some(current_offset);
self.displayed_at = Instant::now();
apply_frame_to_state(app_state.clone(), &mut self.replayer, true).await;
return Ok(());
}
let prev_seq = current.seq;
let next_frame = self.replayer.next()?;
let (due, reached_end) = match next_frame {
None => {
if control.replay_loop {
self.replayer.seek(Duration::ZERO).ok();
self.displayed_at = Instant::now();
self.last_frame_time_offset = self.replayer.elapsed();
apply_frame_to_state(app_state.clone(), &mut self.replayer, true).await;
return Ok(());
}
pause_at_end(app_state.clone()).await;
return Ok(());
}
Some(next) => {
let next_offset = (next.timestamp - first_frame_ts(&self.replayer))
.to_std()
.unwrap_or(Duration::ZERO);
let delta = next_offset.saturating_sub(current_offset);
let mut speed = control.speed;
if !speed.is_finite() {
speed = 1.0;
}
let speed = speed.clamp(0.05, 16.0);
let scaled_delta = Duration::from_secs_f32(delta.as_secs_f32() / speed);
let due_time = self.displayed_at + scaled_delta;
let due_now = Instant::now() >= due_time;
(due_now, false)
}
};
if !due {
let _ = self.replayer.prev()?;
let landed_seq = self.replayer.current().map(|f| f.seq);
if landed_seq != Some(prev_seq) {
tracing::warn!(
expected = prev_seq,
got = ?landed_seq,
"replay cursor retreat landed on unexpected frame; proceeding"
);
}
return Ok(());
}
self.displayed_at = Instant::now();
self.last_frame_time_offset = self.replayer.elapsed();
apply_frame_to_state(app_state.clone(), &mut self.replayer, reached_end).await;
Ok(())
}
}
fn first_frame_ts(replayer: &Replayer) -> chrono::DateTime<chrono::Utc> {
replayer
.current()
.map(|f| f.timestamp)
.unwrap_or_else(chrono::Utc::now)
}
async fn apply_frame_to_state(
app_state: Arc<Mutex<AppState>>,
replayer: &mut Replayer,
force: bool,
) {
let Some(frame) = replayer.current() else {
return;
};
let mut state = app_state.lock().await;
let snap = &frame.snapshot;
if force || state.gpu_info.is_empty() {
state.gpu_info = snap.gpus.clone().unwrap_or_default();
} else {
let new_gpus = snap.gpus.clone().unwrap_or_default();
state.gpu_info = new_gpus;
}
state.cpu_info = snap.cpus.clone().unwrap_or_default();
state.memory_info = snap.memory.clone().unwrap_or_default();
state.chassis_info = snap.chassis.clone().unwrap_or_default();
state.process_info = snap.processes.clone().unwrap_or_default();
state.storage_info = snap.storage.clone().unwrap_or_default();
let host_label = snap
.gpus
.as_ref()
.and_then(|v| v.first())
.map(|g| g.host_id.clone())
.unwrap_or_else(|| snap.hostname.clone());
if let Some(procs) = snap.processes.as_ref() {
state.remote_process_info = procs
.iter()
.map(|p| {
crate::network::metrics_parser::ParsedProcessRow::from_local_process(p, &host_label)
})
.collect();
} else {
state.remote_process_info.clear();
}
state.loading = false;
if let Some(replay) = state.replay.as_mut() {
replay.current_seq = frame.seq;
replay.total_frames = replayer.frames_seen();
replay.at_eof = replayer.at_eof();
replay.elapsed = replayer.elapsed().unwrap_or(Duration::ZERO);
}
let mut host_ids: std::collections::BTreeSet<String> = state
.gpu_info
.iter()
.map(|g| g.host_id.clone())
.filter(|h| !h.is_empty())
.collect();
if host_ids.is_empty() {
if !snap.hostname.is_empty() {
host_ids.insert(snap.hostname.clone());
}
}
let mut tabs = vec![
"All".to_string(),
crate::ui::tabs::USERS_TAB_NAME.to_string(),
crate::ui::tabs::TOPOLOGY_TAB_NAME.to_string(),
];
tabs.extend(host_ids);
let previous_name = state.tabs.get(state.current_tab).cloned();
state.tabs = tabs;
if let Some(name) = previous_name
&& let Some(idx) = state.tabs.iter().position(|t| *t == name)
{
state.current_tab = idx;
} else if state.current_tab >= state.tabs.len() {
state.current_tab = 0;
}
if let Some(last) = state.topology_last_host_tab.as_ref()
&& !state.tabs.iter().any(|t| t == last)
{
state.topology_last_host_tab = None;
}
state.mark_collector_data_changed();
}
async fn clear_pending_seek(app_state: Arc<Mutex<AppState>>) {
let mut state = app_state.lock().await;
if let Some(r) = state.replay.as_mut() {
r.pending_seek = None;
}
}
async fn clear_pending_step(app_state: Arc<Mutex<AppState>>) {
let mut state = app_state.lock().await;
if let Some(r) = state.replay.as_mut() {
r.pending_step = None;
}
}
async fn pause_at_end(app_state: Arc<Mutex<AppState>>) {
let mut state = app_state.lock().await;
if let Some(r) = state.replay.as_mut() {
r.paused = true;
r.at_eof = true;
}
}
pub fn initial_replay_state(speed: f32, replay_loop: bool) -> ReplayState {
ReplayState {
paused: false,
speed,
current_seq: 0,
total_frames: 0,
elapsed: Duration::ZERO,
at_eof: false,
replay_loop,
pending_seek: None,
pending_step: None,
timecode_input_mode: false,
timecode_buffer: String::new(),
timecode_error: None,
}
}