use std::sync::Mutex;
use std::time::Duration;
use uuid::Uuid;
use super::types::{MessageEnqueueCallback, PushOrigin, QueuedUserMessage};
pub const ROUTE_GRACE: Duration = Duration::from_secs(30);
static PARKED: Mutex<Vec<(Uuid, QueuedUserMessage)>> = Mutex::new(Vec::new());
static AWAITING_CHANNEL: Mutex<Option<std::collections::HashSet<Uuid>>> = Mutex::new(None);
pub fn expect_channel_route(session_id: Uuid) {
match AWAITING_CHANNEL.lock() {
Ok(mut guard) => {
guard
.get_or_insert_with(std::collections::HashSet::new)
.insert(session_id);
}
Err(e) => {
tracing::error!(
target: "background_task",
"Could not mark session {session_id} as channel-owned: {e}"
);
}
}
}
pub fn awaits_channel_route(session_id: Uuid) -> bool {
match AWAITING_CHANNEL.lock() {
Ok(guard) => guard.as_ref().is_some_and(|s| s.contains(&session_id)),
Err(e) => {
tracing::error!(
target: "background_task",
"Could not read the channel-owned mark for session {session_id}: {e}"
);
false
}
}
}
fn clear_channel_expectation(session_id: Uuid) {
match AWAITING_CHANNEL.lock() {
Ok(mut guard) => {
if let Some(set) = guard.as_mut() {
set.remove(&session_id);
}
}
Err(e) => {
tracing::warn!(
target: "background_task",
"Could not clear the channel-owned mark for session {session_id}: {e}"
);
}
}
}
pub fn deliver_or_park(session_id: Uuid, msg: QueuedUserMessage) -> bool {
if let Some(route) = super::session_routes::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 {
clear_channel_expectation(session_id);
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.clone());
}
if count > 0 {
tracing::info!(
target: "background_task",
"Delivered {count} parked restart report(s) to session {session_id}"
);
for msg in &mine {
super::notify_queue::clear_on_delivery(session_id, msg);
}
}
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.clone());
super::notify_queue::clear_on_delivery(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)
}
pub fn parking_route() -> MessageEnqueueCallback {
std::sync::Arc::new(|session_id, msg: QueuedUserMessage| {
tracing::warn!(
target: "background_task",
"No surface claims session {session_id} and this process has no local one; \
parking until its channel claims it: {}",
msg.display_text
);
super::notify_queue::persist(session_id, &msg);
deliver_or_park(session_id, msg);
})
}
#[cfg(test)]
pub(crate) fn clear_parked_for_test() {
if let Ok(mut parked) = PARKED.lock() {
parked.clear();
}
if let Ok(mut awaiting) = AWAITING_CHANNEL.lock()
&& let Some(set) = awaiting.as_mut()
{
set.clear();
}
}
#[cfg(test)]
pub(crate) fn test_guard() -> std::sync::MutexGuard<'static, ()> {
static TEST_LOCK: Mutex<()> = Mutex::new(());
let guard = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
clear_parked_for_test();
guard
}
pub async fn report_interrupted() -> usize {
let Some(repo) = super::background_tasks::task_repo() else {
return 0;
};
let rows = match repo.all().await {
Ok(rows) => rows,
Err(e) => {
tracing::error!(target: "background_task", "Failed to read background tasks: {e:#}");
return 0;
}
};
if rows.is_empty() {
return 0;
}
let mut count = 0usize;
for row in rows {
tracing::warn!(
target: "background_task",
"Background task '{}' for session {} was interrupted by a restart",
row.label,
row.session_id
);
deliver_or_park(row.session_id, interrupted_message(&row));
count += 1;
if let Err(e) = repo.clear(row.id).await {
tracing::error!(
target: "background_task",
"Failed to clear background task '{}' after reporting it: {e:#}",
row.label
);
}
}
count
}
fn interrupted_message(row: &crate::db::BackgroundTaskRow) -> QueuedUserMessage {
let context_text = format!(
"[BACKGROUND TASK INTERRUPTED] `{}` was still running when OpenCrabs restarted, so it \
was killed and produced no result. The command was:\n\n```\n{}\n```\n\nIt did NOT \
complete. Decide whether to run it again based on what you were doing; do not assume \
it passed or failed.",
row.label, row.command
);
QueuedUserMessage {
context_text,
display_text: format!("⚠️ Background task interrupted by restart: {}", row.label),
origin: PushOrigin::Recovery,
bg_meta: None,
}
}
pub async fn recover(local: Option<MessageEnqueueCallback>) -> usize {
let orphans = crate::brain::tools::subagent::reconcile::reconcile_orphaned_agents();
let mut reported = 0usize;
for mut orphan in orphans {
match interrupted_report_target(&orphan) {
Some(session_id) => {
deliver_or_park(session_id, subagent_interrupted_message(&orphan));
reported += 1;
}
None => {
tracing::error!(
target: "background_task",
"Sub-agent '{}' has an unparseable parent session '{}' or '{}', its \
interruption cannot be reported",
orphan.label,
orphan.parent_session_id.as_deref().unwrap_or("<none>"),
orphan.session_id
);
orphan.mark_interrupted().ok();
}
}
}
reported += report_interrupted().await;
let redelivered = super::notify_queue::redeliver_persisted().await;
if redelivered > 0 {
tracing::info!(
target: "background_task",
"Boot notify queue: redelivered={redelivered}"
);
}
if reported > 0 {
tracing::info!(
target: "background_task",
"Recovered {reported} interrupted item(s) from a previous run, {} waiting for a \
session route",
parked_count()
);
}
match local {
Some(local) => schedule_flush(local),
None => tracing::info!(
target: "background_task",
"No local surface to flush to; parked reports wait for a channel to claim their \
session"
),
}
reported
}
fn subagent_interrupted_message(
status: &crate::brain::agent::service::work_status::WorkStatus,
) -> 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.task
);
QueuedUserMessage {
context_text,
display_text: format!("⚠️ Sub-agent interrupted by restart: {}", status.label),
origin: PushOrigin::Recovery,
bg_meta: None,
}
}
fn interrupted_report_target(
status: &crate::brain::agent::service::work_status::WorkStatus,
) -> Option<Uuid> {
status
.parent_session_id
.as_deref()
.or(Some(status.session_id.as_str()))
.and_then(|p| Uuid::parse_str(p).ok())
}
pub(crate) fn deliver_revived_agent_outcome(
status: &mut crate::brain::agent::service::work_status::WorkStatus,
outcome: std::result::Result<&str, &str>,
) -> bool {
let Some(parent) = status
.parent_session_id
.as_deref()
.and_then(|p| Uuid::parse_str(p).ok())
else {
tracing::warn!(
target: "background_task",
"Revived sub-agent '{}' has no parent session recorded; finalizing status only",
status.id
);
finalize_revived_status(status, outcome);
return false;
};
let agent_id = status.id.clone();
let label = status.label.clone();
let msg = crate::brain::tools::subagent::spawn::completion_message(&label, &agent_id, outcome);
match crate::brain::agent::service::session_routes::deliver_to_session(
parent,
msg.clone(),
true,
) {
crate::brain::agent::service::session_routes::Delivery::Delivered => {
tracing::info!(
target: "background_task",
"Revived sub-agent {agent_id}'s result delivered to parent session {parent}"
);
}
crate::brain::agent::service::session_routes::Delivery::NoRoute => {
tracing::warn!(
target: "background_task",
"Revived sub-agent {agent_id}'s result had no route for parent {parent}; parked"
);
deliver_or_park(parent, msg);
}
_ => {
tracing::info!(
target: "background_task",
"Revived sub-agent {agent_id}'s result for parent {parent} handled by the \
delivery machinery (parked or redirected)"
);
}
}
finalize_revived_status(status, outcome);
true
}
fn finalize_revived_status(
status: &mut crate::brain::agent::service::work_status::WorkStatus,
outcome: std::result::Result<&str, &str>,
) {
match outcome {
Ok(output) => {
if let Err(e) = status.mark_completed(output.to_string()) {
tracing::warn!(
target: "background_task",
"Could not finalize revived sub-agent status '{}': {e}",
status.id
);
}
}
Err(error) => {
if let Err(e) = status.mark_failed(error.to_string()) {
tracing::warn!(
target: "background_task",
"Could not finalize revived sub-agent status '{}': {e}",
status.id
);
}
}
}
}