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)
}
}
}
}
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>,
}
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,
}
}
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
}
}
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);
let handle = thread::Builder::new()
.name("magi-fast-mode-persistence".to_string())
.spawn(move || {
let persistence_result = panic::catch_unwind(panic::AssertUnwindSafe(|| {
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}"
);
}
}
}
}