use super::{TuiEvent, send_critical, state};
use crossbeam_channel::Sender;
use std::{
sync::{
Arc, Mutex,
atomic::{AtomicBool, Ordering},
mpsc,
},
thread::JoinHandle,
};
#[derive(Debug, Clone)]
pub(crate) struct CompactionActivityFinal {
pub(crate) id: crate::output::ActivityId,
pub(crate) status: crate::output::ActivityStatus,
pub(crate) metadata: crate::output::ActivityMetadata,
pub(crate) preview: Option<String>,
}
#[derive(Debug, Clone)]
pub(crate) enum WorkerFinalEvent {
Done,
RunCanceled {
prompt: String,
},
Error(String),
CompactionFinished {
result: Result<Option<crate::compaction::CompactionResult>, String>,
activity: CompactionActivityFinal,
},
LoginFinished {
provider_id: String,
result: Result<String, String>,
},
}
impl WorkerFinalEvent {
fn into_tui_event(self) -> TuiEvent {
match self {
Self::Done => TuiEvent::Done,
Self::RunCanceled { prompt } => TuiEvent::RunCanceled { prompt },
Self::Error(error) => TuiEvent::Error(error),
Self::CompactionFinished { result, activity } => {
TuiEvent::CompactionFinishedAuthoritative { result, activity }
}
Self::LoginFinished {
provider_id,
result,
} => TuiEvent::LoginFinished {
provider_id,
result,
},
}
}
}
#[derive(Debug, Default)]
pub(super) struct WorkerOutcomeState {
final_event: Mutex<Option<WorkerFinalEvent>>,
final_event_sent: AtomicBool,
}
impl WorkerOutcomeState {
fn store_final_event(&self, event: Option<WorkerFinalEvent>) {
let mut final_event = self
.final_event
.lock()
.unwrap_or_else(|error| error.into_inner());
*final_event = event;
}
pub(super) fn final_event(&self) -> Option<WorkerFinalEvent> {
self.final_event
.lock()
.unwrap_or_else(|error| error.into_inner())
.clone()
}
fn mark_final_event_sent(&self) {
self.final_event_sent.store(true, Ordering::SeqCst);
}
pub(super) fn final_event_sent(&self) -> bool {
self.final_event_sent.load(Ordering::SeqCst)
}
}
pub(super) struct WorkerState {
pub(super) handle: JoinHandle<()>,
pub(super) cancel: Arc<AtomicBool>,
pub(super) login_manual: Option<mpsc::Sender<String>>,
pub(super) outcome: Arc<WorkerOutcomeState>,
pub(super) outcome_reconciled: bool,
pub(super) steering: crate::agent::steering::AgentSteering,
pub(super) accepts_steering: bool,
}
pub(super) fn send_completion(
sender: &Sender<TuiEvent>,
outcome: &WorkerOutcomeState,
final_event: Option<WorkerFinalEvent>,
) {
outcome.store_final_event(final_event.clone());
let final_event_sent = final_event.is_some();
let final_event = final_event.map(|event| Box::new(event.into_tui_event()));
match send_critical(sender, TuiEvent::RunFinished { final_event }) {
Ok(()) => {
if final_event_sent {
outcome.mark_final_event_sent();
}
}
Err(error) => {
let _ = sender.try_send(TuiEvent::Error(error.to_string()));
}
}
}
pub(super) fn apply_worker_final_event(
state: &mut state::MissionControlState,
event: WorkerFinalEvent,
) {
match event {
WorkerFinalEvent::Done => {
state.finish_assistant_streaming();
state.status = "run complete".to_string();
}
WorkerFinalEvent::RunCanceled { prompt } => {
state.finish_assistant_streaming();
let preview = crate::tui::transcript::sanitize_preview(&prompt);
state.mark_running_prompt_canceled();
state.record_canceled_transcript(&preview);
state.status = format!("canceled: {preview}");
}
WorkerFinalEvent::Error(error) => {
state.finish_assistant_streaming();
state.record_error_transcript(&error);
state.status = error;
}
WorkerFinalEvent::CompactionFinished { result, activity } => {
state.apply_activity_event(crate::output::ActivityEvent::Started {
id: activity.id.clone(),
parent_id: None,
kind: crate::output::ActivityKind::Compaction,
status: crate::output::ActivityStatus::Running,
metadata: activity.metadata.clone(),
});
if let Some(preview) = activity.preview {
state.apply_activity_event(crate::output::ActivityEvent::FinalPreview {
id: activity.id.clone(),
preview,
metadata: Some(activity.metadata.clone()),
status: Some(activity.status),
});
}
state.apply_activity_event(crate::output::ActivityEvent::Finished {
id: activity.id,
status: activity.status,
metadata: Some(activity.metadata),
});
match result {
Ok(Some(result)) => {
state.record_compaction_complete_transcript(Some(&result.summary));
state.status = crate::compaction::compaction_success_message(&result);
}
Ok(None) => {
state.record_compaction_complete_transcript(None);
state.status =
"compaction produced empty summary; no checkpoint was written".to_string();
}
Err(error) => {
state.record_error_transcript(&error);
state.status = error;
}
}
}
WorkerFinalEvent::LoginFinished {
provider_id,
result,
} => match result {
Ok(message) => {
state.clear_active_login();
if state.provider == provider_id {
state.provider_ready = true;
}
state.status = message;
}
Err(error) => {
state.clear_active_login();
state.record_error_transcript(&error);
state.status = error;
}
},
}
}
#[cfg(test)]
mod tests {
use super::*;
use crossbeam_channel::bounded;
#[test]
fn send_completion_emits_single_composite_event() {
let (sender, receiver) = bounded::<TuiEvent>(1);
let outcome = WorkerOutcomeState::default();
send_completion(&sender, &outcome, Some(WorkerFinalEvent::Done));
assert!(outcome.final_event_sent());
assert!(matches!(
receiver.try_recv().unwrap(),
TuiEvent::RunFinished {
final_event: Some(event)
} if matches!(*event, TuiEvent::Done)
));
assert!(receiver.try_recv().is_err());
}
#[test]
fn final_event_lock_recovers_after_poison() {
let outcome = WorkerOutcomeState::default();
let _ = std::panic::catch_unwind(|| {
let _guard = outcome.final_event.lock().unwrap();
panic!("poison final event lock");
});
outcome.store_final_event(Some(WorkerFinalEvent::Error("kept".to_string())));
assert!(matches!(
outcome.final_event(),
Some(WorkerFinalEvent::Error(message)) if message == "kept"
));
}
#[test]
fn apply_worker_final_event_records_compaction_summary_transcript() {
let mut state = state::MissionControlState::default();
apply_worker_final_event(
&mut state,
WorkerFinalEvent::CompactionFinished {
result: Ok(Some(crate::compaction::CompactionResult {
session_id: "session-1".to_string(),
summary: "worker summary".to_string(),
provider: "local".to_string(),
model: "compact-model".to_string(),
})),
activity: CompactionActivityFinal {
id: crate::output::ActivityId::new("compaction:test"),
status: crate::output::ActivityStatus::Success,
metadata: crate::output::ActivityMetadata::new("compaction"),
preview: Some("worker summary".to_string()),
},
},
);
assert_eq!(state.transcript, vec!["compact: complete • worker summary"]);
assert!(
state
.status
.contains("compaction complete: session session-1")
);
}
#[test]
fn compaction_finished_success_reconciles_one_authoritative_activity() {
let (sender, receiver) = bounded::<TuiEvent>(1);
let outcome = WorkerOutcomeState::default();
send_completion(
&sender,
&outcome,
Some(WorkerFinalEvent::CompactionFinished {
result: Ok(Some(crate::compaction::CompactionResult {
session_id: "session-1".to_string(),
summary: "authoritative summary".to_string(),
provider: "local".to_string(),
model: "compact-model".to_string(),
})),
activity: CompactionActivityFinal {
id: crate::output::ActivityId::new("compaction:success"),
status: crate::output::ActivityStatus::Success,
metadata: crate::output::ActivityMetadata::new("compaction"),
preview: Some("authoritative summary".to_string()),
},
}),
);
let TuiEvent::RunFinished {
final_event: Some(event),
} = receiver.try_recv().unwrap()
else {
panic!("missing authoritative completion event");
};
let mut state = state::MissionControlState::default();
apply_worker_final_event(
&mut state,
match *event {
TuiEvent::CompactionFinishedAuthoritative { result, activity } => {
WorkerFinalEvent::CompactionFinished { result, activity }
}
_ => panic!("unexpected completion event"),
},
);
let id = crate::output::ActivityId::new("compaction:success");
assert_eq!(state.roots, vec![id.clone()]);
assert_eq!(state.nodes.len(), 1);
assert_eq!(
state.nodes[&id].status,
crate::output::ActivityStatus::Success
);
assert_eq!(state.nodes[&id].preview, "authoritative summary");
assert_eq!(
state.transcript,
vec!["compact: complete • authoritative summary"]
);
}
#[test]
fn compaction_finished_fallback_reconciles_failed_and_canceled_activity() {
let cases: [(crate::output::ActivityStatus, &str, Option<&str>); 2] = [
(
crate::output::ActivityStatus::Failed,
"compaction failed",
Some("compaction failed"),
),
(
crate::output::ActivityStatus::Canceled,
"compaction canceled",
None,
),
];
for (status, error, preview) in cases {
let (sender, _receiver) = bounded::<TuiEvent>(1);
sender.send(TuiEvent::Done).unwrap();
let outcome = WorkerOutcomeState::default();
let id = crate::output::ActivityId::new(format!("compaction:{status:?}"));
send_completion(
&sender,
&outcome,
Some(WorkerFinalEvent::CompactionFinished {
result: Err(error.to_string()),
activity: CompactionActivityFinal {
id: id.clone(),
status,
metadata: crate::output::ActivityMetadata::new("compaction"),
preview: preview.map(str::to_string),
},
}),
);
let mut state = state::MissionControlState::default();
apply_worker_final_event(&mut state, outcome.final_event().unwrap());
assert_eq!(state.roots, vec![id.clone()]);
assert_eq!(state.nodes.len(), 1);
assert_eq!(state.nodes[&id].status, status);
assert_eq!(state.nodes[&id].metadata.label, "compaction");
assert_eq!(state.nodes[&id].preview, preview.unwrap_or_default());
assert!(state.transcript.iter().any(|line| line.contains(error)));
}
}
#[test]
fn compaction_finished_fallback_reconciles_empty_summary_and_redacts_failure() {
let cases = [
(
Ok(None),
crate::output::ActivityStatus::Failed,
None,
"empty summary; no checkpoint was written",
None,
"compact: complete",
),
(
Err("provider failed api_key=secret-value".to_string()),
crate::output::ActivityStatus::Failed,
Some("provider failed api_key=secret-value"),
"failed; no checkpoint was written",
Some("provider failed api_key=<redacted>"),
"error: provider failed api_key=<redacted>",
),
];
for (result, status, preview_input, note, expected_preview, transcript) in cases {
let (sender, _receiver) = bounded::<TuiEvent>(1);
sender.send(TuiEvent::Done).unwrap();
let outcome = WorkerOutcomeState::default();
let id = crate::output::ActivityId::new(format!("compaction:fallback:{status:?}"));
let metadata = crate::output::ActivityMetadata {
label: "compaction".to_string(),
detail: Some(note.to_string()),
fields: vec![("note".to_string(), note.to_string())],
};
let preview = preview_input.map(crate::output::redact_sensitive_text);
send_completion(
&sender,
&outcome,
Some(WorkerFinalEvent::CompactionFinished {
result,
activity: CompactionActivityFinal {
id: id.clone(),
status,
metadata,
preview,
},
}),
);
let mut state = state::MissionControlState::default();
apply_worker_final_event(&mut state, outcome.final_event().unwrap());
assert_eq!(state.nodes.len(), 1);
assert_eq!(state.roots, vec![id.clone()]);
assert_eq!(
state.nodes[&id].kind,
crate::output::ActivityKind::Compaction
);
assert_eq!(state.nodes[&id].status, status);
assert_eq!(state.nodes[&id].metadata.detail.as_deref(), Some(note));
assert_eq!(
state.nodes[&id].preview,
expected_preview.unwrap_or_default()
);
assert!(state.transcript.iter().any(|line| line == transcript));
assert!(
state
.transcript
.iter()
.all(|line| !line.contains("secret-value"))
);
}
}
}