use super::super::*;
use super::MissionControlApp;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct PendingFastModePersistence {
request_id: u64,
enabled: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum FastModePersistenceDelivery {
Delivered,
ReceiverDisconnected,
Shutdown,
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum RetainedFastModePersistenceFailureKind {
Persistence(String),
EventDeliveryDisconnected,
WorkerJoinPanic,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct RetainedFastModePersistenceFailure {
request_id: u64,
enabled: bool,
kind: RetainedFastModePersistenceFailureKind,
}
impl RetainedFastModePersistenceFailure {
fn from_worker_result(
request_id: u64,
enabled: bool,
result: &FastModePersistenceWorkerResult,
) -> Option<Self> {
match (&result.persistence_result, result.delivery) {
(Err(error), _) => Some(Self {
request_id,
enabled,
kind: RetainedFastModePersistenceFailureKind::Persistence(error.clone()),
}),
(
Ok(()),
FastModePersistenceDelivery::Delivered | FastModePersistenceDelivery::Shutdown,
) => None,
(Ok(()), FastModePersistenceDelivery::ReceiverDisconnected) => Some(Self {
request_id,
enabled,
kind: RetainedFastModePersistenceFailureKind::EventDeliveryDisconnected,
}),
}
}
fn worker_join_panic(request_id: u64, enabled: bool) -> Self {
Self {
request_id,
enabled,
kind: RetainedFastModePersistenceFailureKind::WorkerJoinPanic,
}
}
fn cleanup_message(&self) -> String {
match &self.kind {
RetainedFastModePersistenceFailureKind::Persistence(error) => format!(
"fast mode persistence request {} failed: {error}",
self.request_id
),
RetainedFastModePersistenceFailureKind::EventDeliveryDisconnected => format!(
"fast mode persistence request {} failed: TUI event channel closed",
self.request_id
),
RetainedFastModePersistenceFailureKind::WorkerJoinPanic => {
format!("fast mode persistence worker {} panicked", self.request_id)
}
}
}
}
#[cfg(test)]
#[derive(Debug)]
enum FastModePersistenceTestInjection {
PersistenceError(String),
Panic,
}
pub(super) const FAST_MODE_PERSISTENCE_RETRY_BLOCKED_STATUS: &str =
"fast mode persistence failed; restart Mission Control before retrying.";
#[derive(Debug)]
struct FastModePersistenceWorker {
request_id: u64,
handle: JoinHandle<FastModePersistenceWorkerResult>,
}
#[derive(Debug)]
struct FastModePersistenceWorkerResult {
persistence_result: Result<(), String>,
delivery: FastModePersistenceDelivery,
}
#[derive(Debug)]
pub(super) struct FastModePersistenceRuntime {
next_request_id: u64,
pending: Option<PendingFastModePersistence>,
shutdown: Arc<AtomicBool>,
workers: Vec<FastModePersistenceWorker>,
retained_failure: Option<RetainedFastModePersistenceFailure>,
#[cfg(test)]
test_injection: Option<FastModePersistenceTestInjection>,
}
impl FastModePersistenceRuntime {
pub(super) fn new() -> Self {
Self {
next_request_id: 0,
pending: None,
shutdown: Arc::new(AtomicBool::new(false)),
workers: Vec::new(),
retained_failure: None,
#[cfg(test)]
test_injection: None,
}
}
pub(super) fn is_pending(&self) -> bool {
self.pending.is_some()
}
pub(super) fn has_retained_failure(&self) -> bool {
self.retained_failure.is_some()
}
pub(super) fn request_shutdown(&self) {
self.shutdown.store(true, Ordering::SeqCst);
}
}
fn send_fast_mode_persisted(
sender: &Sender<TuiEvent>,
shutdown: &AtomicBool,
request_id: u64,
enabled: bool,
persistence_result: &Result<(), String>,
) -> FastModePersistenceDelivery {
let mut event = Some(TuiEvent::FastModePersisted {
request_id,
enabled,
result: persistence_result.clone(),
});
loop {
if shutdown.load(Ordering::SeqCst) {
return FastModePersistenceDelivery::Shutdown;
}
let Some(next) = event.take() else {
return FastModePersistenceDelivery::Delivered;
};
match sender.send_timeout(next, CRITICAL_EVENT_TIMEOUT) {
Ok(()) => return FastModePersistenceDelivery::Delivered,
Err(crossbeam_channel::SendTimeoutError::Disconnected(_)) => {
return FastModePersistenceDelivery::ReceiverDisconnected;
}
Err(crossbeam_channel::SendTimeoutError::Timeout(next)) => event = Some(next),
}
}
}
impl FastModePersistenceRuntime {
fn reap(&mut self, ui_state: &mut state::MissionControlState) -> bool {
let mut changed = false;
let mut active = Vec::with_capacity(self.workers.len());
for worker in std::mem::take(&mut self.workers) {
if !worker.handle.is_finished() {
active.push(worker);
continue;
}
let pending = self.pending.filter(|p| p.request_id == worker.request_id);
match worker.handle.join() {
Ok(result) => {
if let Some(pending) = pending {
if let Some(failure) =
RetainedFastModePersistenceFailure::from_worker_result(
worker.request_id,
pending.enabled,
&result,
)
{
self.retained_failure = Some(failure);
}
if matches!(
result.delivery,
FastModePersistenceDelivery::ReceiverDisconnected
| FastModePersistenceDelivery::Shutdown
) {
self.pending = None;
ui_state.status =
"failed to save fast mode: TUI event channel closed".to_string();
changed = true;
}
}
}
Err(_) => {
if let Some(pending) = pending {
self.retained_failure =
Some(RetainedFastModePersistenceFailure::worker_join_panic(
worker.request_id,
pending.enabled,
));
self.pending = None;
ui_state.status =
"failed to save fast mode: fast mode persistence worker panicked"
.to_string();
changed = true;
}
}
}
}
self.workers = active;
changed
}
pub(super) fn cleanup(&mut self) -> Vec<String> {
self.shutdown.store(true, Ordering::SeqCst);
let retained_request = self.retained_failure.as_ref().map(|f| f.request_id);
let mut errors = self
.retained_failure
.as_ref()
.map(|f| f.cleanup_message())
.into_iter()
.collect::<Vec<_>>();
for worker in std::mem::take(&mut self.workers) {
let pending = self.pending.filter(|p| p.request_id == worker.request_id);
let failure = match worker.handle.join() {
Ok(result) => pending.and_then(|p| {
RetainedFastModePersistenceFailure::from_worker_result(
worker.request_id,
p.enabled,
&result,
)
}),
Err(_) => pending.map(|p| {
RetainedFastModePersistenceFailure::worker_join_panic(
worker.request_id,
p.enabled,
)
}),
};
if let Some(failure) = failure {
if retained_request != Some(failure.request_id) {
errors.push(failure.cleanup_message());
}
self.retained_failure = Some(failure);
}
}
self.pending = None;
errors
}
#[cfg(test)]
fn set_injection(&mut self, injection: FastModePersistenceTestInjection) {
self.test_injection = Some(injection);
}
}
impl MissionControlApp {
pub(super) fn reap_fast_mode_persistence_workers(
&mut self,
ui_state: &mut state::MissionControlState,
) -> bool {
self.fast_mode_persistence.reap(ui_state)
}
pub(in crate::tui) fn handle_fast_mode_persistence_drain(
&mut self,
ui_state: &mut state::MissionControlState,
drain_result: &DrainResult,
) -> bool {
let mut changed = false;
for (request_id, enabled, result) in &drain_result.fast_mode_persisted {
let pending_matches = self
.fast_mode_persistence
.pending
.is_some_and(|p| p.request_id == *request_id && p.enabled == *enabled);
let other = self
.fast_mode_persistence
.pending
.is_some_and(|p| p.request_id != *request_id || p.enabled != *enabled);
if self
.fast_mode_persistence
.retained_failure
.as_ref()
.is_some_and(|f| f.request_id == *request_id && f.enabled == *enabled)
&& !other
{
self.fast_mode_persistence.retained_failure = None;
}
if !pending_matches {
continue;
}
self.fast_mode_persistence.pending = None;
match result {
Ok(()) => {
self.settings.fast.enabled = *enabled;
ui_state.set_fast_mode_enabled(*enabled);
if !*enabled {
ui_state.clear_fast_observations();
}
self.refresh_fast_mode_state(ui_state);
ui_state.status = self.fast_mode_status(ui_state);
}
Err(error) => ui_state.status = format!("failed to save fast mode: {error}"),
}
changed = true;
}
changed
}
pub(super) fn handle_fast_mode_command(
&mut self,
command: crate::commands::FastModeCommand,
ui_state: &mut state::MissionControlState,
) {
self.reap_fast_mode_persistence_workers(ui_state);
let target = match command {
crate::commands::FastModeCommand::Status => {
ui_state.status = if self.fast_mode_persistence.has_retained_failure() {
FAST_MODE_PERSISTENCE_RETRY_BLOCKED_STATUS.to_string()
} else {
self.fast_mode_status(ui_state)
};
return;
}
crate::commands::FastModeCommand::Toggle => !self.settings.fast.enabled,
crate::commands::FastModeCommand::On => true,
crate::commands::FastModeCommand::Off => false,
};
if self.fast_mode_persistence.is_pending() {
ui_state.status = "fast mode persistence already in progress".to_string();
return;
}
if self.fast_mode_persistence.has_retained_failure() {
ui_state.status = FAST_MODE_PERSISTENCE_RETRY_BLOCKED_STATUS.to_string();
return;
}
let runtime = &mut self.fast_mode_persistence;
runtime.next_request_id = runtime.next_request_id.saturating_add(1);
let request_id = runtime.next_request_id;
runtime.pending = Some(PendingFastModePersistence {
request_id,
enabled: target,
});
ui_state.status = "saving fast mode setting…".to_string();
let paths = self.config.paths.clone();
let sender = self.events.clone();
let shutdown = Arc::clone(&runtime.shutdown);
#[cfg(test)]
let test_injection = runtime.test_injection.take();
let handle = thread::Builder::new()
.name("magi-fast-mode-persistence".to_string())
.spawn(move || {
let persistence_result = panic::catch_unwind(panic::AssertUnwindSafe(|| {
#[cfg(test)]
if let Some(injection) = test_injection {
return match injection {
FastModePersistenceTestInjection::PersistenceError(error) => Err(error),
FastModePersistenceTestInjection::Panic => {
panic!("injected fast mode persistence panic")
}
};
}
crate::config::set_fast_mode(&paths, target)
.map(|_| ())
.map_err(|e| e.to_string())
}))
.unwrap_or_else(|_| Err("fast mode persistence worker panicked".to_string()));
let delivery = send_fast_mode_persisted(
&sender,
&shutdown,
request_id,
target,
&persistence_result,
);
FastModePersistenceWorkerResult {
persistence_result,
delivery,
}
});
match handle {
Ok(handle) => runtime
.workers
.push(FastModePersistenceWorker { request_id, handle }),
Err(error) => {
runtime.pending = None;
ui_state.status = format!(
"failed to save fast mode: could not start persistence worker: {error}"
);
}
}
}
}
#[cfg(test)]
mod tests {
use super::super::tests::{test_app, test_area};
use super::super::*;
use super::*;
#[test]
fn cleanup_after_run_reports_fast_mode_persistence_error() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, receiver) = bounded::<TuiEvent>(4);
let mut app = test_app(&temp, sender);
app.fast_mode_persistence.set_injection(
FastModePersistenceTestInjection::PersistenceError(
"injected persistence failure".to_string(),
),
);
let mut ui_state = state::MissionControlState::default();
assert!(!app.submit(
"/fast on".to_string(),
&mut ui_state,
&receiver,
test_area()
));
drop(receiver);
let error = app.cleanup_after_run().unwrap_err().to_string();
assert!(error.contains("injected persistence failure"), "{error}");
assert!(!app.settings.fast.enabled);
assert!(app.fast_mode_persistence.pending.is_none());
}
#[test]
fn cleanup_after_run_reports_caught_fast_mode_persistence_panic() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, receiver) = bounded::<TuiEvent>(4);
let mut app = test_app(&temp, sender);
app.fast_mode_persistence
.set_injection(FastModePersistenceTestInjection::Panic);
let mut ui_state = state::MissionControlState::default();
assert!(!app.submit(
"/fast on".to_string(),
&mut ui_state,
&receiver,
test_area()
));
drop(receiver);
let error = app.cleanup_after_run().unwrap_err().to_string();
assert!(
error.contains("fast mode persistence worker panicked"),
"{error}"
);
assert!(!app.settings.fast.enabled);
assert!(app.fast_mode_persistence.pending.is_none());
}
#[test]
fn true_fast_mode_worker_join_panic_blocks_retry_and_cleanup_reports_first_failure() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, receiver) = bounded::<TuiEvent>(4);
let mut app = test_app(&temp, sender);
let mut ui_state = state::MissionControlState::default();
app.fast_mode_persistence.pending = Some(PendingFastModePersistence {
request_id: 1,
enabled: true,
});
let handle = thread::spawn(|| -> FastModePersistenceWorkerResult {
panic!("true fast mode persistence worker panic");
});
while !handle.is_finished() {
thread::yield_now();
}
app.fast_mode_persistence
.workers
.push(FastModePersistenceWorker {
request_id: 1,
handle,
});
assert!(app.reap_fast_mode_persistence_workers(&mut ui_state));
assert!(app.fast_mode_persistence.workers.is_empty());
assert!(app.fast_mode_persistence.pending.is_none());
assert!(matches!(
app.fast_mode_persistence
.retained_failure
.as_ref()
.map(|failure| &failure.kind),
Some(RetainedFastModePersistenceFailureKind::WorkerJoinPanic)
));
app.fast_mode_persistence.set_injection(
FastModePersistenceTestInjection::PersistenceError(
"second persistence failure".to_string(),
),
);
assert!(!app.submit(
"/fast on".to_string(),
&mut ui_state,
&receiver,
test_area()
));
assert_eq!(ui_state.status, FAST_MODE_PERSISTENCE_RETRY_BLOCKED_STATUS);
assert!(app.fast_mode_persistence.workers.is_empty());
assert!(
matches!(app.fast_mode_persistence.test_injection.as_ref(), Some(FastModePersistenceTestInjection::PersistenceError(error)) if error == "second persistence failure")
);
app.submit(
"/fast invalid".to_string(),
&mut ui_state,
&receiver,
test_area(),
);
assert_eq!(ui_state.status, crate::commands::FAST_MODE_USAGE);
app.submit("/help".to_string(), &mut ui_state, &receiver, test_area());
assert!(ui_state.show_help);
assert!(!app.submit("/quit".to_string(), &mut ui_state, &receiver, test_area()));
app.submit(
"/fast status".to_string(),
&mut ui_state,
&receiver,
test_area(),
);
assert_eq!(ui_state.status, FAST_MODE_PERSISTENCE_RETRY_BLOCKED_STATUS);
let error = app.cleanup_after_run().unwrap_err().to_string();
assert!(
error.contains("fast mode persistence worker 1 panicked"),
"{error}"
);
assert!(!error.contains("second persistence failure"), "{error}");
}
fn wait_for_fast_mode_persistence_delivery(
app: &MissionControlApp,
receiver: &Receiver<TuiEvent>,
) {
let deadline = Instant::now() + Duration::from_secs(1);
while Instant::now() < deadline {
let event_delivered = receiver.len() == 1;
let worker_finished = app
.fast_mode_persistence
.workers
.iter()
.all(|worker| worker.handle.is_finished());
if event_delivered && worker_finished {
return;
}
thread::yield_now();
}
panic!(
"fast mode persistence did not deliver and finish: events={}, workers={}",
receiver.len(),
app.fast_mode_persistence.workers.len()
);
}
#[test]
fn normal_reaping_retains_fast_mode_failure_until_cleanup() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, receiver) = bounded::<TuiEvent>(4);
let mut app = test_app(&temp, sender);
app.fast_mode_persistence.set_injection(
FastModePersistenceTestInjection::PersistenceError(
"injected persistence failure".to_string(),
),
);
let mut ui_state = state::MissionControlState::default();
assert!(!app.submit("/fast".to_string(), &mut ui_state, &receiver, test_area()));
wait_for_fast_mode_persistence_delivery(&app, &receiver);
assert_eq!(receiver.len(), 1);
assert!(!app.reap_fast_mode_persistence_workers(&mut ui_state));
assert!(app.fast_mode_persistence.workers.is_empty());
assert!(app.fast_mode_persistence.retained_failure.is_some());
let error = app.cleanup_after_run().unwrap_err().to_string();
assert!(error.contains("injected persistence failure"), "{error}");
}
#[test]
fn handled_fast_mode_failure_is_not_re_reported_by_cleanup() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, receiver) = bounded::<TuiEvent>(4);
let mut app = test_app(&temp, sender);
app.fast_mode_persistence.set_injection(
FastModePersistenceTestInjection::PersistenceError(
"injected persistence failure".to_string(),
),
);
let mut ui_state = state::MissionControlState::default();
assert!(!app.submit("/fast".to_string(), &mut ui_state, &receiver, test_area()));
wait_for_fast_mode_persistence_delivery(&app, &receiver);
assert!(!app.reap_fast_mode_persistence_workers(&mut ui_state));
let event = receiver.try_recv().unwrap();
let mut drain_result = DrainResult::default();
apply_control_event_to_state(&mut ui_state, event, &mut drain_result);
assert!(app.handle_fast_mode_persistence_drain(&mut ui_state, &drain_result));
assert!(app.fast_mode_persistence.retained_failure.is_none());
assert!(app.fast_mode_persistence.pending.is_none());
assert!(app.cleanup_after_run().is_ok());
}
#[test]
fn stale_or_mismatched_completion_is_ignored() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, receiver) = bounded::<TuiEvent>(4);
let mut app = test_app(&temp, sender);
let mut ui_state = state::MissionControlState::default();
app.submit(
"/fast on".to_string(),
&mut ui_state,
&receiver,
test_area(),
);
let mut stale = DrainResult::default();
stale.fast_mode_persisted.push((999, true, Ok(())));
assert!(!app.handle_fast_mode_persistence_drain(&mut ui_state, &stale));
assert!(!app.settings.fast.enabled);
assert!(!ui_state.fast_mode_enabled);
let mut mismatched = DrainResult::default();
mismatched.fast_mode_persisted.push((1, false, Ok(())));
assert!(!app.handle_fast_mode_persistence_drain(&mut ui_state, &mismatched));
assert!(!app.settings.fast.enabled);
assert!(!ui_state.fast_mode_enabled);
let event = receiver
.recv_timeout(Duration::from_secs(1))
.expect("fast mode persistence completion");
let mut result = DrainResult::default();
apply_control_event_to_state(&mut ui_state, event, &mut result);
assert_eq!(result.fast_mode_persisted.len(), 1);
assert!(app.handle_fast_mode_persistence_drain(&mut ui_state, &result));
assert!(app.settings.fast.enabled);
assert!(ui_state.fast_mode_enabled);
}
}