use std::sync::Mutex;
use std::time::Duration;
use uuid::Uuid;
use super::types::{MessageEnqueueCallback, QueuedUserMessage};
pub const ROUTE_GRACE: Duration = Duration::from_secs(30);
static PARKED: Mutex<Vec<(Uuid, QueuedUserMessage)>> = Mutex::new(Vec::new());
pub fn deliver_or_park(session_id: Uuid, msg: QueuedUserMessage) -> bool {
if let Some(route) = super::background_tasks::session_route(session_id) {
route(session_id, msg);
return true;
}
match PARKED.lock() {
Ok(mut parked) => parked.push((session_id, msg)),
Err(e) => {
tracing::error!(
target: "background_task",
"Could not park restart report for session {session_id}, it is lost: {e}"
);
}
}
false
}
pub fn claim_session(session_id: Uuid, route: &MessageEnqueueCallback) -> usize {
let mine = match PARKED.lock() {
Ok(mut parked) => {
let mut mine = Vec::new();
parked.retain(|(id, msg)| {
if *id == session_id {
mine.push(msg.clone());
false
} else {
true
}
});
mine
}
Err(e) => {
tracing::error!(
target: "background_task",
"Could not read parked restart reports for session {session_id}: {e}"
);
return 0;
}
};
let count = mine.len();
for msg in mine {
route(session_id, msg);
}
if count > 0 {
tracing::info!(
target: "background_task",
"Delivered {count} parked restart report(s) to session {session_id}"
);
}
count
}
pub fn flush_parked(local: &MessageEnqueueCallback) -> usize {
let remaining = match PARKED.lock() {
Ok(mut parked) => std::mem::take(&mut *parked),
Err(e) => {
tracing::error!(
target: "background_task",
"Could not flush parked restart reports: {e}"
);
return 0;
}
};
let count = remaining.len();
for (session_id, msg) in remaining {
local(session_id, msg);
}
if count > 0 {
tracing::info!(
target: "background_task",
"No route claimed {count} restart report(s) within the grace period, delivered locally"
);
}
count
}
pub fn schedule_flush(local: MessageEnqueueCallback) {
tokio::spawn(async move {
tokio::time::sleep(ROUTE_GRACE).await;
flush_parked(&local);
});
}
pub fn parked_count() -> usize {
PARKED.lock().map(|p| p.len()).unwrap_or(0)
}
#[cfg(test)]
pub(crate) fn clear_parked_for_test() {
if let Ok(mut parked) = PARKED.lock() {
parked.clear();
}
}
pub async fn recover(local: MessageEnqueueCallback) -> usize {
let orphans = crate::brain::tools::subagent::reconcile::reconcile_orphaned_agents();
let mut reported = 0usize;
for orphan in orphans {
match Uuid::parse_str(&orphan.parent_session_id) {
Ok(session_id) => {
deliver_or_park(session_id, subagent_interrupted_message(&orphan));
reported += 1;
}
Err(e) => {
tracing::error!(
target: "background_task",
"Sub-agent '{}' has an unparseable parent session '{}', its interruption \
cannot be reported: {e}",
orphan.label,
orphan.parent_session_id
);
}
}
}
reported += super::background_tasks::report_interrupted().await;
if reported > 0 {
tracing::info!(
target: "background_task",
"Recovered {reported} interrupted item(s) from a previous run, {} waiting for a \
session route",
parked_count()
);
}
schedule_flush(local);
reported
}
fn subagent_interrupted_message(
status: &crate::brain::tools::subagent::status::AgentStatus,
) -> QueuedUserMessage {
let context_text = format!(
"[SUB-AGENT INTERRUPTED] The sub-agent `{}` (id {}) was still running when OpenCrabs \
restarted, so it was killed and produced no result. Its task was:\n\n```\n{}\n```\n\nIt \
did NOT complete. Decide whether to spawn it again based on what you were doing; do not \
assume it succeeded or failed.",
status.label, status.id, status.prompt
);
QueuedUserMessage {
context_text,
display_text: format!("⚠️ Sub-agent interrupted by restart: {}", status.label),
}
}