1use std::collections::{BTreeMap, HashMap, HashSet};
49use std::path::PathBuf;
50use std::sync::Arc;
51use std::time::{Duration, Instant};
52
53use anyhow::{Context as _, Result};
54use jiff::Timestamp;
55use tokio::task::JoinSet;
56
57use crate::agent::{self, Invocation, SeatState};
58use crate::ask::{ChoiceAction, Deputy, Question, Questions, Waiter as Note, WaiterKind, Who};
59use crate::config::Config;
60use crate::notices::{self, Notice};
61use crate::prompt;
62
63pub const NODE: &str = "deputy";
66
67pub const TICK: Duration = Duration::from_secs(5);
69
70pub const MAX_STARTS: u32 = 3;
74
75const RESTART_AFTER: Duration = Duration::from_secs(30);
77
78const SLACK_SECS: u64 = 120;
81
82const CLAIM_STALE: Duration = Duration::from_secs(60);
84
85const HANDOVER_TIMEOUT: Duration = Duration::from_secs(10 * 60);
87
88pub type Halt = Arc<dyn Fn() -> bool + Send + Sync>;
90
91pub fn brief(
98 task_id: &str,
99 reason: &str,
100 choices: &[String],
101 actions: &BTreeMap<String, ChoiceAction>,
102) -> String {
103 let reason = reason.trim();
104 let mut s = format!(
105 "Task {task_id} (`magi task show {task_id}`) is blocked on this question. \
106 The conductor asked because: {}\n\n\
107 When the owner answers, magi records the question and the answer on the \
108 task and unblocks it, and the answer becomes part of the instructions of \
109 the task's next attempt.",
110 if reason.is_empty() {
111 "(it recorded no reasoning)"
112 } else {
113 reason
114 }
115 );
116 if !choices.is_empty() {
117 s.push_str("\n\nWhat each option does:");
118 for c in choices {
119 match actions.get(c) {
120 Some(a) => s.push_str(&format!("\n- `{c}`: also {}", a.describe())),
121 None => s.push_str(&format!(
122 "\n- `{c}`: recorded on the task as the answer and the task is unblocked"
123 )),
124 }
125 }
126 }
127 s
128}
129
130#[derive(Debug, Clone, Copy, PartialEq, Eq)]
132pub enum Kind {
133 Conduct,
135 Land,
137}
138
139pub fn kind_of(q: &Question) -> Option<Kind> {
141 match q.node.as_str() {
142 crate::conduct::NODE => Some(Kind::Conduct),
143 crate::land::APPROVAL_NODE => Some(Kind::Land),
144 _ => None,
145 }
146}
147
148pub fn deadline(q: &Question, default_timeout: u64) -> i64 {
156 let secs = if q.answer_timeout > 0 {
157 q.answer_timeout
158 } else {
159 default_timeout
160 };
161 let from = match kind_of(q) {
162 Some(Kind::Land) => q.asked_at.as_second(),
163 _ => q.last_activity(),
164 };
165 from.saturating_add(secs as i64)
166}
167
168pub fn can_start(cfg: Option<&Config>, agent: &str) -> bool {
176 cfg.is_some_and(|c| {
177 c.daemon.max_deputies > 0
178 && ((!agent.is_empty() && c.agent(agent).is_ok()) || c.resolve_roles().is_ok())
179 })
180}
181
182pub fn agent_of(q: &Question) -> &str {
184 q.deputy.as_ref().map_or("", |d| d.agent.as_str())
185}
186
187pub fn exhausted_past_deadline(
194 q: &Question,
195 startable: bool,
196 default_timeout: u64,
197 now: Timestamp,
198) -> bool {
199 q.status.open()
200 && q.deputy
201 .as_ref()
202 .is_some_and(|d| d.starts >= MAX_STARTS || !startable)
203 && now.as_second() > deadline(q, default_timeout)
204}
205
206pub struct Deputies {
208 store: Questions,
209 home: PathBuf,
210 cfg: Option<Config>,
211 fallback_repo: PathBuf,
213 max: usize,
214 halt: Halt,
215 tasks: JoinSet<String>,
216 inflight: HashSet<String>,
217 memo: HashMap<String, Instant>,
219}
220
221impl Deputies {
222 pub fn new(
224 store: Questions,
225 home: PathBuf,
226 cfg: Option<Config>,
227 fallback_repo: PathBuf,
228 max: usize,
229 halt: Halt,
230 ) -> Self {
231 Self {
232 store,
233 home,
234 cfg,
235 fallback_repo,
236 max,
237 halt,
238 tasks: JoinSet::new(),
239 inflight: HashSet::new(),
240 memo: HashMap::new(),
241 }
242 }
243
244 fn default_timeout(&self) -> u64 {
245 self.cfg
246 .as_ref()
247 .map_or(Config::default().graph.answer_timeout, |c| {
248 c.graph.answer_timeout
249 })
250 }
251
252 fn reap(&mut self) {
253 while let Some(done) = self.tasks.try_join_next() {
254 if let Ok(id) = done {
255 self.inflight.remove(&id);
256 }
257 }
258 if self.tasks.is_empty() {
259 self.inflight.clear();
260 }
261 }
262
263 fn attach(&self, q: &Question, kind: Kind) -> Option<Question> {
271 let default_timeout = self.default_timeout();
272 let repo = self.fallback_repo.to_string_lossy().into_owned();
273 let state = match kind {
274 Kind::Land => crate::run::RunState::load(&q.run).ok(),
275 Kind::Conduct => None,
276 };
277 let timeout = state
279 .as_ref()
280 .map_or(default_timeout, |s| s.config.graph.answer_timeout);
281 self.store
282 .update(&q.id, |r| {
283 if r.deputy.is_none() {
284 r.deputy = Some(Deputy::new(match kind {
285 Kind::Conduct => brief(&r.run, &r.detail, &r.choices, &r.actions),
286 Kind::Land => crate::land::deputy_brief(r, state.as_ref()),
287 }));
288 }
289 if kind == Kind::Conduct && r.cwd.is_none() {
290 r.cwd = Some(repo.clone());
291 }
292 if r.answer_timeout == 0 {
293 r.answer_timeout = match kind {
294 Kind::Conduct => default_timeout,
295 Kind::Land => timeout,
296 };
297 }
298 Ok(())
299 })
300 .map(|(r, ())| r)
301 .map_err(|e| tracing::warn!("question {}: cannot attach a deputy: {e:#}", q.short()))
302 .ok()
303 }
304
305 pub fn tick(&mut self, now: Timestamp) {
308 self.reap();
309 for q in self.store.list() {
310 if (self.halt)() {
311 return;
312 }
313 let Some(kind) = kind_of(&q) else {
314 continue;
315 };
316 if !q.status.open() {
317 continue;
318 }
319 let needs_cwd = kind == Kind::Conduct && q.cwd.is_none();
320 let q = if q.deputy.is_none() || needs_cwd || q.answer_timeout == 0 {
321 match self.attach(&q, kind) {
322 Some(q) => q,
323 None => continue,
324 }
325 } else {
326 q
327 };
328 let Some(dep) = q.deputy.as_ref() else {
329 continue;
330 };
331 if self.inflight.contains(&q.id)
332 || self.store.read_lease(&q.id).is_some_and(|l| l.fresh(now))
333 {
334 continue;
335 }
336 if dep.starts >= MAX_STARTS {
337 self.give_up(&q);
338 continue;
339 }
340 if now.as_second() > deadline(&q, self.default_timeout())
343 && q.unread_from_owner().is_none()
344 {
345 continue;
346 }
347 if self.inflight.len() >= self.max || !can_start(self.cfg.as_ref(), dep.agent.as_str())
348 {
349 continue;
350 }
351 if matches!(self.memo.get(&q.id), Some(until) if Instant::now() < *until) {
352 continue;
353 }
354 self.memo
355 .insert(q.id.clone(), Instant::now() + RESTART_AFTER);
356 self.inflight.insert(q.id.clone());
357 let job = Job {
358 store: self.store.clone(),
359 home: self.home.clone(),
360 cfg: self.cfg.clone(),
361 fallback_repo: self.fallback_repo.clone(),
362 halt: Arc::clone(&self.halt),
363 };
364 let id = q.id.clone();
365 self.tasks.spawn(async move {
366 if let Err(e) = job.turn(&id).await {
367 tracing::warn!("deputy for question {}: {e:#}", crate::ask::short_id(&id));
368 }
369 id
370 });
371 }
372 }
373
374 pub async fn drain(&mut self) {
377 while let Some(done) = self.tasks.join_next().await {
378 if let Ok(id) = done {
379 self.inflight.remove(&id);
380 }
381 }
382 self.inflight.clear();
383 }
384
385 fn give_up(&self, q: &Question) {
387 notices::raise_in(
388 &self.home,
389 Notice::warn(
390 &format!("deputy:{}", q.id),
391 format!(
392 "Question {} \"{}\": its follow-up agent ended {MAX_STARTS} times \
393 without an answer and is not restarted. What you say is recorded \
394 but nothing will read it; answer with one of the choices instead.",
395 q.short(),
396 q.summary
397 ),
398 ),
399 );
400 }
401}
402
403struct Job {
405 store: Questions,
406 home: PathBuf,
407 cfg: Option<Config>,
408 fallback_repo: PathBuf,
409 halt: Halt,
410}
411
412impl Job {
413 async fn drive(
416 &self,
417 spec: &crate::config::AgentSpec,
418 seat: &mut SeatState,
419 inv: &Invocation<'_>,
420 id: &str,
421 ) -> Option<Result<agent::AgentOutput>> {
422 let fut = agent::invoke(spec, seat, inv);
423 tokio::pin!(fut);
424 let mut beat = tokio::time::interval(Duration::from_secs(1));
425 let mut beats = 0u32;
426 loop {
427 tokio::select! {
428 r = &mut fut => break Some(r),
429 _ = beat.tick() => {
430 if (self.halt)() {
431 break None;
432 }
433 beats += 1;
434 if beats % 20 == 0 {
435 self.store.beat(id, WaiterKind::Deputy);
436 }
437 }
438 }
439 }
440 }
441
442 fn park_quietly(&self, id: &str) {
445 let _ = self.store.update(id, |r| {
446 r.waiter = None;
447 Ok(())
448 });
449 self.store.drop_lease(id);
450 }
451
452 fn park(&self, id: &str) {
455 let _ = self.store.update(id, |r| {
456 if let Some(d) = r.deputy.as_mut() {
457 d.starts = d.starts.saturating_sub(1);
458 }
459 r.waiter = None;
460 Ok(())
461 });
462 self.store.drop_lease(id);
463 }
464
465 async fn turn(&self, id: &str) -> Result<()> {
466 let claim = self.store.root().join(format!("{id}.deputy-claim"));
467 if !crate::waiter::take_claim(&claim, CLAIM_STALE) {
471 return Ok(());
472 }
473 let release = crate::waiter::Release(claim);
474
475 let q = self.store.get(id)?;
477 let now = Timestamp::now();
478 let Some(dep) = q.deputy.clone() else {
479 return Ok(());
480 };
481 if !q.status.open()
482 || dep.starts >= MAX_STARTS
483 || self.store.read_lease(&q.id).is_some_and(|l| l.fresh(now))
484 {
485 return Ok(());
486 }
487 let cfg = self
488 .cfg
489 .as_ref()
490 .context("the deputy's configuration is not available")?;
491 let spec = match cfg.agent(&dep.agent) {
492 Ok(s) if !dep.agent.is_empty() => s.clone(),
493 _ => {
494 cfg.resolve_roles()
495 .context("resolving the deputy's agent")?
496 .conductor
497 }
498 };
499 let key = crate::ask::deputy_seat_key(&q.id);
502 let (mut seat, resumed) = match dep.seat.clone() {
503 Some(s)
504 if s.agent == spec.id && agent::has_session(spec.kind, &s, cfg.graph.sessions) =>
505 {
506 (s, true)
507 }
508 _ => (SeatState::new(&key, &spec.id, crate::rng::entropy()), false),
509 };
510 let cwd = q
511 .cwd
512 .as_deref()
513 .map(PathBuf::from)
514 .filter(|p| p.is_dir())
515 .unwrap_or_else(|| self.fallback_repo.clone());
516
517 let thread: Vec<(&str, &str)> = q
518 .thread
519 .iter()
520 .map(|t| {
521 (
522 if t.who == Who::Operator {
523 "operator"
524 } else {
525 "agent"
526 },
527 t.body.as_str(),
528 )
529 })
530 .collect();
531 let read = &thread[..q.delivered_turns.min(thread.len())];
532 let unread = q.unread_from_owner();
533 let snapshot = q.thread.len();
534 let body = prompt::deputy(&prompt::DeputyPrompt {
535 id: &q.id,
536 summary: &q.summary,
537 detail: &q.detail,
538 brief: &dep.brief,
539 choices: &q.choices,
540 thread: read,
541 unread: unread.as_deref(),
542 resumed,
543 handover: false,
544 land: kind_of(&q) == Some(Kind::Land),
545 language: &cfg.graph.language,
546 });
547
548 self.store.beat(&q.id, WaiterKind::Deputy);
551 let starts = dep.starts + 1;
552 let first = seat.clone();
553 self.store.update(&q.id, |r| {
554 if let Some(d) = r.deputy.as_mut() {
555 d.agent = spec.id.clone();
556 d.seat = Some(first);
557 d.starts = starts;
558 }
559 r.waiter = Some(Note {
560 kind: WaiterKind::Deputy,
561 since: now,
562 });
563 Ok(())
564 })?;
565 drop(release);
566 tracing::info!(
567 "question {}: deputy seat {} {} (start {starts}/{MAX_STARTS})",
568 q.short(),
569 seat.key,
570 if resumed { "resuming" } else { "starting" }
571 );
572
573 let left = (deadline(&q, cfg.graph.answer_timeout) - now.as_second()).max(0) as u64;
574 let artifacts = self.store.root().join(format!("{}.deputy", q.id));
575 let cache_dir = cfg.cache_dir();
576 let allow_write = true;
580 let writable = [self.store.root().to_path_buf()];
582 macro_rules! invocation {
583 ($prompt:expr, $stem:expr, $timeout:expr) => {
584 Invocation {
585 cwd: &cwd,
586 prompt: $prompt,
587 timeout: $timeout,
588 allow_write,
589 sessions: cfg.graph.sessions,
590 artifacts: &artifacts,
591 stem: $stem,
592 run: &q.run,
596 node: NODE,
597 cache_dir: cache_dir.as_deref(),
598 attachments: &[],
599 writable: &writable,
600 }
601 };
602 }
603
604 let mut early = None;
609 if !resumed && cfg.graph.sessions {
610 let hbody = prompt::deputy(&prompt::DeputyPrompt {
611 id: &q.id,
612 summary: &q.summary,
613 detail: &q.detail,
614 brief: &dep.brief,
615 choices: &q.choices,
616 thread: read,
617 unread: None,
618 resumed: false,
619 handover: true,
620 land: kind_of(&q) == Some(Kind::Land),
621 language: &cfg.graph.language,
622 });
623 let hstem = format!("handover-{starts}");
624 let hlimit = HANDOVER_TIMEOUT.min(Duration::from_secs(left.max(1)));
627 let hinv = invocation!(&hbody, &hstem, hlimit);
628 let Some(done) = self.drive(&spec, &mut seat, &hinv, &q.id).await else {
629 self.park(&q.id);
630 return Ok(());
631 };
632 let kept = seat.clone();
633 self.store.update(&q.id, |r| {
634 if let Some(d) = r.deputy.as_mut() {
635 d.seat = Some(kept);
636 }
637 Ok(())
638 })?;
639 if !matches!(&done, Ok(o) if o.usable()) {
640 early = Some(done);
641 }
642 }
643
644 let left = if early.is_none() && !resumed && cfg.graph.sessions {
647 let now = Timestamp::now();
648 let again = self.store.get(&q.id)?;
649 let left = (deadline(&again, cfg.graph.answer_timeout) - now.as_second()).max(0) as u64;
650 if !again.status.open() || (left == 0 && again.unread_from_owner().is_none()) {
651 self.park_quietly(&q.id);
652 return Ok(());
653 }
654 left
655 } else {
656 left
657 };
658 let out = match early {
659 Some(done) => done,
660 None => {
661 let stem = format!("turn-{starts}");
662 let timeout = Duration::from_secs(left.max(60) + SLACK_SECS);
663 let inv = invocation!(&body, &stem, timeout);
664 match self.drive(&spec, &mut seat, &inv, &q.id).await {
665 Some(out) => out,
666 None => {
667 self.park(&q.id);
668 return Ok(());
669 }
670 }
671 }
672 };
673
674 let (text, why) = match out {
675 Ok(o) if o.usable() => (Some(o.text.trim().to_owned()), None),
676 Ok(o) if o.timed_out => (None, Some("its turn timed out".to_owned())),
677 Ok(o) if o.quota_exhausted() => (None, Some("the agent is out of quota".to_owned())),
678 Ok(o) => (
679 None,
680 Some(format!("its turn failed (exit {:?})", o.exit_code)),
681 ),
682 Err(e) => (None, Some(format!("the agent could not be started: {e:#}"))),
683 };
684 let kept = seat.clone();
685 self.store.update(&q.id, |r| {
686 if let Some(d) = r.deputy.as_mut() {
687 d.seat = Some(kept);
688 }
689 r.waiter = None;
690 if let Some(text) = text
694 && unread.is_some()
695 && !text.is_empty()
696 && r.status.open()
697 && r.thread.len() == snapshot
698 {
699 r.delivered_turns = r.delivered_turns.max(snapshot);
700 let choices = r.choices.clone();
701 r.reply(text, choices)?;
702 }
703 Ok(())
704 })?;
705 self.store.drop_lease(&q.id);
706 if let Some(why) = why {
707 tracing::warn!("question {}: the deputy ended: {why}", q.short());
708 notices::raise_in(
709 &self.home,
710 Notice::warn(
711 &format!("deputy-turn:{}", q.id),
712 format!(
713 "Question {} \"{}\": its follow-up agent stopped ({why}).",
714 q.short(),
715 q.summary
716 ),
717 ),
718 );
719 }
720 Ok(())
721 }
722}
723
724pub async fn run(mut deputies: Deputies, stop: crate::daemon::Stop) {
726 while !stop.stopped() {
727 deputies.tick(Timestamp::now());
728 tokio::time::sleep(TICK).await;
729 }
730}