use super::super::*;
use super::MissionControlApp;
use crate::output::{ActivityId, ActivityKind, ActivityMetadata, ActivityStatus};
use std::{
sync::atomic::AtomicU64,
time::{SystemTime, UNIX_EPOCH},
};
static COMPACTION_ACTIVITY_SEQUENCE: AtomicU64 = AtomicU64::new(0);
impl MissionControlApp {
pub(crate) fn start_compaction(
&mut self,
custom_instructions: Option<String>,
ui_state: &mut state::MissionControlState,
) {
let Some(session) = self.state.current_session.clone() else {
let message =
"/compact requires an active persisted session; start without --no-session or run /new"
.to_string();
let _ = send_critical(&self.events, TuiEvent::Error(message));
return;
};
let mut active_config = self
.state
.config
.clone()
.unwrap_or_else(|| self.config.clone());
active_config.model = Some(self.state.model.clone());
let settings = match crate::config::read_settings(&active_config.paths) {
Ok(settings) => settings,
Err(error) => {
let _ = send_critical(&self.events, TuiEvent::Error(error.to_string()));
return;
}
};
let refreshed_auth =
match crate::compaction::refresh_active_auth_for_compaction_if_inherited(
&mut active_config,
&settings,
&self.state.auth_state,
) {
Ok(auth_state) => auth_state,
Err(error) => {
let _ = send_critical(&self.events, TuiEvent::Error(error.to_string()));
return;
}
};
if let Some(auth_state) = refreshed_auth {
self.state.auth_state = auth_state;
}
if let Some(config) = &mut self.state.config {
config.auth = active_config.auth.clone();
config.model = active_config.model.clone();
}
self.config.auth = active_config.auth.clone();
self.config.model = active_config.model.clone();
ui_state.provider_ready = self.state.auth_state.is_ready();
ui_state.status = "compacting session…".to_string();
ui_state.record_compaction_started_transcript();
let activity_session_id = session.id().to_string();
let activity_provider = active_config.provider_id().to_string();
let activity_model = active_config
.model
.as_deref()
.unwrap_or_else(|| {
crate::providers::default_model_for_provider(active_config.provider_id())
})
.to_string();
let activity_id = compaction_activity_id(&activity_session_id);
self.active_run = true;
let cwd = self.state.cwd.clone();
let sender = self.events.clone();
let cancel = Arc::new(AtomicBool::new(false));
let outcome = Arc::new(WorkerOutcomeState::default());
let worker_cancel = Arc::clone(&cancel);
let worker_outcome = Arc::clone(&outcome);
let handle = thread::spawn(move || {
let activity_metadata = compaction_activity_metadata(
"compaction",
&activity_session_id,
&activity_provider,
&activity_model,
None,
);
let _ = send_critical(
&sender,
TuiEvent::Activity(ActivityEvent::Started {
id: activity_id.clone(),
parent_id: None,
kind: ActivityKind::Compaction,
status: ActivityStatus::Running,
metadata: activity_metadata,
}),
);
if worker_cancel.load(Ordering::SeqCst) {
let error = "compaction canceled; no checkpoint was written".to_string();
let metadata = compaction_activity_metadata(
"compaction",
&activity_session_id,
&activity_provider,
&activity_model,
Some("canceled before compaction request"),
);
send_completion(
&sender,
&worker_outcome,
Some(WorkerFinalEvent::CompactionFinished {
result: Err(error),
activity: CompactionActivityFinal {
id: activity_id,
status: ActivityStatus::Canceled,
metadata,
preview: None,
},
}),
);
return;
}
let result = crate::compaction::compact_session(crate::compaction::CompactSessionJob {
active_config,
settings,
session: session.clone(),
cwd: cwd.clone(),
cancellation: crate::agent::cancellation::AgentCancellation::new(Arc::clone(
&worker_cancel,
)),
custom_instructions,
})
.map_err(|error| error.to_string());
let (status, metadata, preview) = match &result {
Ok(Some(value)) => (
ActivityStatus::Success,
compaction_activity_metadata(
"compaction",
&value.session_id,
&value.provider,
&value.model,
None,
),
Some(crate::output::redact_sensitive_text(&value.summary)),
),
Ok(None) => (
ActivityStatus::Failed,
compaction_activity_metadata(
"compaction",
&activity_session_id,
&activity_provider,
&activity_model,
Some("empty summary; no checkpoint was written"),
),
None,
),
Err(error) if worker_cancel.load(Ordering::SeqCst) => (
ActivityStatus::Canceled,
compaction_activity_metadata(
"compaction",
&activity_session_id,
&activity_provider,
&activity_model,
Some("canceled; no checkpoint was written"),
),
Some(crate::output::redact_sensitive_text(error)),
),
Err(error) => (
ActivityStatus::Failed,
compaction_activity_metadata(
"compaction",
&activity_session_id,
&activity_provider,
&activity_model,
Some("failed; no checkpoint was written"),
),
Some(crate::output::redact_sensitive_text(error)),
),
};
if let Err(error) = &result {
let _ = record_session_event(
Some(&session),
&cwd,
SessionEventKind::Diagnostic,
serde_json::json!({"level":"error", "message": error}),
);
}
send_completion(
&sender,
&worker_outcome,
Some(WorkerFinalEvent::CompactionFinished {
result,
activity: CompactionActivityFinal {
id: activity_id,
status,
metadata,
preview,
},
}),
);
});
self.worker = Some(WorkerState {
handle,
cancel,
login_manual: None,
outcome,
outcome_reconciled: false,
steering: crate::agent::steering::AgentSteering::new(),
accepts_steering: false,
});
}
}
fn compaction_activity_id(session_id: &str) -> ActivityId {
let millis = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_millis())
.unwrap_or(0);
let sequence = COMPACTION_ACTIVITY_SEQUENCE.fetch_add(1, Ordering::SeqCst);
ActivityId::new(format!("compaction:{session_id}:{millis}:{sequence}"))
}
fn compaction_activity_metadata(
label: &str,
session_id: &str,
provider: &str,
model: &str,
note: Option<&str>,
) -> ActivityMetadata {
let mut metadata = ActivityMetadata::new(label);
metadata.fields = vec![
("session_id".to_string(), session_id.to_string()),
("provider".to_string(), provider.to_string()),
("model".to_string(), model.to_string()),
];
if let Some(note) = note {
metadata.detail = Some(note.to_string());
metadata.fields.push(("note".to_string(), note.to_string()));
}
metadata
}