use std::sync::Arc;
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use crate::app::events::AppEvent;
use crate::entities::message::{Message, MessageRole};
use crate::entities::subagent::{RunOutcome, SubagentRun};
use crate::features::tools::ChatEffect;
use super::generation::{self, BackgroundMessage, BackgroundSpawn, BackgroundStart, TurnProgress};
use super::{InflightChild, Orchestrator};
pub(super) struct PendingLanding {
generation: Uuid,
run: SubagentRun,
result: String,
effects: Vec<ChatEffect>,
}
pub(super) struct BackgroundRun {
pub(super) generation: Uuid,
pub(super) run_id: Uuid,
pub(super) chat: Uuid,
pub(super) cancel: CancellationToken,
pub(super) child: InflightChild,
}
impl BackgroundRun {
pub(super) fn is_out(&self) -> bool {
self.child.run.outcome.is_none()
}
}
impl Orchestrator {
pub(super) fn spawn_background_run(&mut self, start: BackgroundStart, chat: Uuid) {
let generation = Uuid::new_v4();
let run_id = start.run_id();
let placeholder = start.placeholder();
let sessions = self.session_budget();
let cancel = generation::spawn_background_run(
start,
BackgroundSpawn {
generation,
sessions,
evt_tx: self.evt_tx.clone(),
bg_tx: self.bg_run_tx.clone(),
},
);
self.background_runs.push(BackgroundRun {
generation,
run_id,
chat,
cancel,
child: InflightChild {
run: placeholder,
stream: Uuid::new_v4(),
partial: Default::default(),
line_role: MessageRole::Assistant,
position: None,
},
});
self.emit_chat_list();
self.emit_background_runs();
}
pub(super) fn handle_background_message(&mut self, message: BackgroundMessage) {
match message {
BackgroundMessage::Progress {
generation,
progress,
} => {
let Some(seat) = self
.background_runs
.iter_mut()
.find(|b| b.generation == generation)
else {
return;
};
match progress {
TurnProgress::ChildStarted(run) => {
let mut run = *run;
run.background = true;
run.title = seat.child.run.title.clone();
run.renamed_manually = seat.child.run.renamed_manually;
seat.child.run = run;
self.emit_chat_list();
}
other => self.handle_child_progress(other),
}
}
BackgroundMessage::Done {
generation,
run,
result,
effects,
} => self.land_background_run(generation, *run, result, effects),
}
}
fn land_background_run(
&mut self,
generation: Uuid,
run: SubagentRun,
result: String,
effects: Vec<ChatEffect>,
) {
let Some(pos) = self
.background_runs
.iter()
.position(|b| b.generation == generation)
else {
return;
};
let chat_id = self.background_runs[pos].chat;
if self.inflight.as_ref().is_some_and(|t| t.chat == chat_id) {
let mirror = &mut self.background_runs[pos].child.run;
mirror.outcome = run.outcome;
mirror.finished_at = run.finished_at;
mirror.tokens = run.tokens;
mirror.messages = run.messages.clone();
self.pending_landings.push(PendingLanding {
generation,
run,
result,
effects,
});
self.emit_background_runs();
self.emit_chat_list();
return;
}
let seat = self.background_runs.remove(pos);
self.finish_landing(seat, run, result, effects);
}
pub(super) fn land_pending_runs(&mut self, chat_id: Uuid) {
let (pending, kept): (Vec<_>, Vec<_>) = std::mem::take(&mut self.pending_landings)
.into_iter()
.partition(|landing| {
self.background_runs
.iter()
.any(|b| b.generation == landing.generation && b.chat == chat_id)
});
self.pending_landings = kept;
for landing in pending {
let Some(pos) = self
.background_runs
.iter()
.position(|b| b.generation == landing.generation)
else {
continue;
};
let seat = self.background_runs.remove(pos);
self.finish_landing(seat, landing.run, landing.result, landing.effects);
}
}
fn finish_landing(
&mut self,
seat: BackgroundRun,
mut run: SubagentRun,
result: String,
effects: Vec<ChatEffect>,
) {
run.background = true;
if seat.child.run.renamed_manually {
run.title = seat.child.run.title.clone();
run.renamed_manually = true;
}
let chat_id = seat.chat;
let run_id = run.id;
let kind = run.kind;
let label = run.name.clone().unwrap_or_else(|| run.title.clone());
let mut attached = Vec::new();
let mut stored = Vec::new();
let mut live = false;
let mut landed = false;
if let Some(chat) = self.chat_mut(chat_id) {
live = chat.child(run_id).is_some();
if let Some(record) = chat.child_mut_including_deleted(run_id) {
*record = run;
landed = true;
}
for effect in effects {
match effect {
ChatEffect::AddAttachment(a) => attached.push(*a),
ChatEffect::AddChatFile(f) => stored.push(*f),
ChatEffect::SetSystemMessage(_) | ChatEffect::SetSamplingOverride(_) => {}
}
}
}
self.emit_background_runs();
if !landed {
self.emit_chat_list();
return;
}
self.mark_dirty(chat_id);
for a in attached {
self.insert_attachment(chat_id, a, None);
}
self.list_stored_files(chat_id, stored);
self.emit_chat_list();
self.maybe_auto_title_run(run_id);
if !live {
return;
}
let loc = self.profile_locale_of(chat_id);
let key = match kind {
crate::entities::subagent::RunKind::Dialogue => "tool.start_dialogue.notification",
_ => "tool.start_subagent.notification",
};
let text = loc.tf(
key,
&[
("name", &label),
("address", &crate::features::chat_links::uri(run_id)),
("body", &result),
],
);
let message = Message::notification(run_id, text);
self.deliver_notification(chat_id, message);
}
fn deliver_notification(&mut self, chat_id: Uuid, message: Message) {
let open = self.active_id == Some(chat_id);
let Some(chat) = self.chat_mut(chat_id) else {
return;
};
chat.push_message(message);
if !open {
chat.unread = true;
}
self.mark_dirty(chat_id);
if open {
self.activate(chat_id);
}
self.emit_chat_list();
self.maybe_wake(chat_id);
}
fn maybe_wake(&mut self, chat_id: Uuid) {
if !self.config.tools.subagent_background_wake
|| self.active_id != Some(chat_id)
|| !self.gen_state.is_idle()
{
return;
}
let Ok(backend) = self.engines.backend_if_ready(self.ui_locale()) else {
return;
};
self.start_woken_generation(chat_id, backend);
}
pub(super) fn handle_stop_subagent_run(&mut self, id: Uuid) {
if let Some(seat) = self.background_runs.iter().find(|b| b.run_id == id) {
seat.cancel.cancel();
}
}
pub(super) fn cancel_orphaned_background_runs(&mut self, chat_id: Uuid) {
let live: Vec<Uuid> = self
.chats
.iter()
.find(|c| c.id == chat_id)
.map(|c| c.children().map(|r| r.id).collect())
.unwrap_or_default();
for seat in self
.background_runs
.iter()
.filter(|b| b.chat == chat_id && !live.contains(&b.run_id))
{
seat.cancel.cancel();
}
}
pub(super) fn cancel_background_runs_of(&mut self, chat_id: Uuid) {
for seat in self.background_runs.iter().filter(|b| b.chat == chat_id) {
seat.cancel.cancel();
}
}
pub(super) fn stop_all_background_runs(&mut self) {
let seats = std::mem::take(&mut self.background_runs);
for seat in seats {
seat.cancel.cancel();
let mut run = seat.child.run.clone();
run.background = true;
run.outcome = Some(RunOutcome::Cancelled);
run.finished_at = Some(chrono::Utc::now());
if let Some(record) = self
.chat_mut(seat.chat)
.and_then(|c| c.child_mut_including_deleted(run.id))
{
*record = run;
self.mark_dirty(seat.chat);
}
}
}
pub(super) fn background_child(&self, run: Uuid) -> Option<&InflightChild> {
self.background_runs
.iter()
.find(|b| b.run_id == run)
.map(|b| &b.child)
}
pub(super) fn background_run(&self, run: Uuid) -> Option<&BackgroundRun> {
self.background_runs.iter().find(|b| b.run_id == run)
}
pub(super) fn child_any(&self, run: Uuid) -> Option<&InflightChild> {
self.inflight
.as_ref()
.and_then(|t| t.child(run))
.or_else(|| self.background_child(run))
}
pub(super) fn child_mut_any(&mut self, run: Uuid) -> Option<&mut InflightChild> {
if self
.inflight
.as_ref()
.is_some_and(|t| t.child(run).is_some())
{
return self.inflight.as_mut().and_then(|t| t.child_mut(run));
}
self.background_runs
.iter_mut()
.find(|b| b.run_id == run)
.map(|b| &mut b.child)
}
pub(super) fn emit_background_runs(&self) {
let out = self.background_runs.iter().filter(|b| b.is_out()).count() as u32;
let _ = self.evt_tx.send(AppEvent::BackgroundRuns { out });
}
fn profile_locale_of(&self, chat_id: Uuid) -> &'static crate::shared::i18n::Locale {
let lang = self
.chats
.iter()
.find(|c| c.id == chat_id)
.and_then(|c| self.profiles.iter().find(|p| p.id == c.profile_id))
.map(|p| p.language)
.unwrap_or_default();
crate::shared::i18n::locale(lang)
}
}
pub(super) type BudgetKey = (crate::shared::config::ServerMode, u32, Option<u64>);
impl Orchestrator {
pub(super) fn session_budget(&mut self) -> Arc<crate::shared::session_budget::SessionBudget> {
let key: BudgetKey = (
self.config.engine.mode,
self.config.engine.active_sessions(),
self.session_pool(),
);
if let Some((k, budget)) = &self.session_budget_memo
&& *k == key
{
return budget.clone();
}
let budget = Arc::new(crate::shared::session_budget::SessionBudget::new(
key.1, key.2,
));
self.session_budget_memo = Some((key, budget.clone()));
budget
}
}