1use std::collections::HashMap;
45use std::path::PathBuf;
46use std::time::{Duration, Instant};
47
48use anyhow::Result;
49use jiff::Timestamp;
50
51use crate::agent::{self, Invocation, SeatState};
52use crate::ask::{Lease, Question, QuestionStatus, Questions, Waiter as Note, WaiterKind, Who};
53use crate::config::{AgentSpec, Config};
54use crate::notices::{self, Notice};
55use crate::prompt;
56use crate::run::RunState;
57
58pub const TICK: Duration = Duration::from_secs(5);
60
61const DELIVERY_TIMEOUT: Duration = Duration::from_secs(15 * 60);
65
66const RETRY_AFTER: Duration = Duration::from_secs(10 * 60);
70
71#[derive(Debug, Clone, PartialEq, Eq)]
73pub enum Word {
74 Said(String),
76 Answered(String),
78}
79
80#[derive(Debug, Clone, PartialEq, Eq)]
82pub enum Action {
83 Idle,
85 Expire,
87 Deliver(Word),
89}
90
91pub fn decide(
101 q: &Question,
102 lease: Option<&Lease>,
103 seat_busy: bool,
104 default_timeout: u64,
105 now: Timestamp,
106) -> Action {
107 decide_owned(q, lease, seat_busy, default_timeout, now, false)
108}
109
110pub fn decide_owned(
120 q: &Question,
121 lease: Option<&Lease>,
122 seat_busy: bool,
123 default_timeout: u64,
124 now: Timestamp,
125 daemon_owns_action: bool,
126) -> Action {
127 if q.cwd.is_none() || lease.is_some_and(|l| l.fresh(now)) {
128 return Action::Idle;
129 }
130 match q.status {
131 QuestionStatus::Abandoned => Action::Idle,
132 QuestionStatus::Answered => match q.resolution() {
133 _ if daemon_owns_action && q.chosen_action().is_some() => Action::Idle,
134 Some(a) if !q.answer_delivered && !seat_busy => Action::Deliver(Word::Answered(a)),
135 _ => Action::Idle,
136 },
137 QuestionStatus::Open => {
138 if let Some(said) = q.unread_from_owner() {
139 if seat_busy {
140 return Action::Idle;
141 }
142 return Action::Deliver(Word::Said(said));
143 }
144 let secs = if q.answer_timeout > 0 {
145 q.answer_timeout
146 } else {
147 default_timeout
148 };
149 let deadline = q.last_activity().saturating_add(secs as i64);
150 if now.as_second() > deadline {
151 Action::Expire
152 } else {
153 Action::Idle
154 }
155 }
156 }
157}
158
159struct Target {
161 spec: AgentSpec,
162 seat: SeatState,
163 cwd: PathBuf,
164 sessions: bool,
165 allow_write: bool,
166 language: String,
167}
168
169pub struct Waiter {
172 store: Questions,
173 home: PathBuf,
174 conduct_cfg: Option<Config>,
177 default_timeout: u64,
179 memo: HashMap<String, (String, Instant)>,
181}
182
183impl Waiter {
184 pub fn new(store: Questions, home: PathBuf, conduct_cfg: Option<Config>) -> Self {
186 let default_timeout = conduct_cfg
187 .as_ref()
188 .map_or(Config::default().graph.answer_timeout, |c| {
189 c.graph.answer_timeout
190 });
191 Self {
192 store,
193 home,
194 conduct_cfg,
195 default_timeout,
196 memo: HashMap::new(),
197 }
198 }
199
200 fn seat_busy(&self, q: &Question, now: Timestamp) -> bool {
202 if q.run == crate::conduct::NODE {
203 return crate::conduct::busy(&self.home);
204 }
205 RunState::load_under(&q.run, &self.home).is_ok_and(|s| {
206 s.seats_active()
207 .any(|(k, a)| *k == q.seat && a.remaining_secs(now) > 0)
208 })
209 }
210
211 fn daemon_owns_action(&self, q: &Question) -> bool {
214 crate::daemon::task_of_question(&crate::queue::Queue::at(self.home.join("queue")).list(), q)
215 .is_some_and(|t| crate::daemon::action_standing(t, q).daemon_owns())
216 }
217
218 pub async fn tick(&mut self, now: Timestamp, halt: &(dyn Fn() -> bool + Sync)) {
223 for q in self.store.list() {
224 if halt() {
225 return;
226 }
227 let lease = self.store.read_lease(&q.id);
228 let owned = self.daemon_owns_action(&q);
229 match decide_owned(&q, lease.as_ref(), false, self.default_timeout, now, owned) {
230 Action::Idle => {}
231 Action::Expire => self.expire(&q),
232 Action::Deliver(word) => {
233 if q.deputy.is_some() {
237 if crate::deputy::exhausted_past_deadline(
241 &q,
242 crate::deputy::can_start(
243 self.conduct_cfg.as_ref(),
244 crate::deputy::agent_of(&q),
245 ),
246 self.default_timeout,
247 now,
248 ) {
249 self.expire(&q);
250 }
251 continue;
252 }
253 if self.seat_busy(&q, now) {
254 continue;
255 }
256 let key = format!("{}:{:?}", q.thread.len(), word);
257 if matches!(self.memo.get(&q.id), Some((k, until)) if *k == key && Instant::now() < *until)
258 {
259 continue;
260 }
261 if let Err(e) = self.deliver(&q, &word, halt).await {
262 tracing::warn!("question {}: {e:#}", q.short());
263 }
264 }
265 }
266 }
267 }
268
269 fn expire(&self, q: &Question) {
272 let secs = if q.answer_timeout > 0 {
273 q.answer_timeout
274 } else {
275 self.default_timeout
276 };
277 let startable =
278 crate::deputy::can_start(self.conduct_cfg.as_ref(), crate::deputy::agent_of(q));
279 let why = format!("no answer within {}s of asking", secs.max(1));
280 let done = self.store.update(&q.id, |r| {
281 let now = Timestamp::now();
285 let lease = self.store.read_lease(&r.id);
286 if decide(r, lease.as_ref(), false, self.default_timeout, now) != Action::Expire
287 && !(lease.as_ref().is_none_or(|l| !l.fresh(now))
288 && crate::deputy::exhausted_past_deadline(
289 r,
290 startable,
291 self.default_timeout,
292 now,
293 ))
294 {
295 return Ok(false);
296 }
297 r.abandon(&why);
298 r.waiter = None;
299 Ok(true)
300 });
301 match done {
302 Ok((_, false)) => {}
303 Ok((_, true)) => {
304 self.store.drop_lease(&q.id);
305 tracing::warn!(
306 "question {} went unanswered for {secs}s with nobody waiting; \
307 it stays as the record of it",
308 q.short()
309 );
310 }
311 Err(e) => tracing::warn!("could not expire question {}: {e:#}", q.short()),
312 }
313 }
314
315 fn target(&self, q: &Question) -> std::result::Result<Target, String> {
317 let cwd = q
318 .cwd
319 .as_deref()
320 .map(PathBuf::from)
321 .filter(|p| p.is_dir())
322 .ok_or("the directory the agent was working in is gone")?;
323 let (cfg, seat, allow_write) = if q.run == crate::conduct::NODE {
324 let cfg = self
325 .conduct_cfg
326 .clone()
327 .ok_or("the conductor's configuration is not available")?;
328 let seat = crate::conduct::load_seat(&self.home)
329 .ok_or("the conductor has no recorded session")?;
330 (cfg, seat, false)
331 } else {
332 let state = RunState::load_under(&q.run, &self.home)
333 .map_err(|_| "the run that asked is gone".to_owned())?;
334 if !state.status.resumable() {
335 return Err(format!(
336 "the run that asked has already {}",
337 state.status.as_str()
338 ));
339 }
340 let seat = state
341 .seats
342 .get(&q.seat)
343 .cloned()
344 .ok_or("the run has no record of the asking seat")?;
345 let write = matches!(q.node.as_str(), "implement" | "fix");
346 (state.config, seat, write)
347 };
348 let spec = cfg
349 .agent(&seat.agent)
350 .map_err(|_| format!("agent `{}` is no longer in the roster", seat.agent))?
351 .clone();
352 if !agent::has_session(spec.kind, &seat, cfg.graph.sessions) {
353 return Err("the seat has no session to resume".to_owned());
354 }
355 Ok(Target {
356 spec,
357 seat,
358 cwd,
359 sessions: cfg.graph.sessions,
360 allow_write,
361 language: cfg.graph.language.clone(),
362 })
363 }
364
365 fn refuse(&mut self, q: &Question, key: String, why: &str) {
367 tracing::warn!("question {} cannot be delivered: {why}", q.short());
368 notices::raise_in(
369 &self.home,
370 Notice::warn(
371 &format!("question:{}", q.id),
372 format!(
373 "Question {} \"{}\": what you said cannot reach the agent \
374 ({why}). It is recorded, but nothing will read it.",
375 q.short(),
376 q.summary
377 ),
378 ),
379 );
380 let _ = self.store.update(&q.id, |r| {
381 r.waiter = None;
382 Ok(())
383 });
384 self.memo
385 .insert(q.id.clone(), (key, Instant::now() + RETRY_AFTER));
386 }
387
388 async fn deliver(
390 &mut self,
391 q: &Question,
392 word: &Word,
393 halt: &(dyn Fn() -> bool + Sync),
394 ) -> Result<()> {
395 let key = format!("{}:{:?}", q.thread.len(), word);
396 let mut target = match self.target(q) {
397 Ok(t) => t,
398 Err(why) => {
399 self.refuse(q, key, &why);
400 return Ok(());
401 }
402 };
403
404 let claim = self.store.root().join(format!("{}.claim", q.id));
405 if !take_claim(&claim, DELIVERY_TIMEOUT + Duration::from_secs(60)) {
406 return Ok(());
407 }
408 let _release = Release(claim);
409
410 let q = self.store.get(&q.id)?;
413 let now = Timestamp::now();
414 let lease = self.store.read_lease(&q.id);
415 let owned = self.daemon_owns_action(&q);
416 let Action::Deliver(word) =
417 decide_owned(&q, lease.as_ref(), false, self.default_timeout, now, owned)
418 else {
419 return Ok(());
420 };
421 let snapshot = q.thread.len();
422
423 let queue = crate::queue::Queue::at(self.home.join("queue"));
429 let task_id = crate::daemon::task_of_question(&queue.list(), &q).map(|t| t.id.clone());
430 let _task_claim = match (&task_id, q.chosen_action()) {
434 (Some(id), Some(_)) => match queue.claim(id) {
435 Ok(c) => Some(c),
436 Err(e) => {
437 tracing::debug!(
438 "question {}: task claim not available ({e:#}), deferring delivery",
439 q.short()
440 );
441 return Ok(());
442 }
443 },
444 _ => None,
445 };
446 let task_still_ours = || match &task_id {
449 Some(id) => queue
450 .get(id)
451 .map(|t| !crate::daemon::action_standing(&t, &q).daemon_owns())
452 .unwrap_or(true),
453 None => true,
454 };
455 if !task_still_ours() {
456 return Ok(());
457 }
458
459 self.store.beat(&q.id, WaiterKind::Daemon);
460 if !task_still_ours() {
463 return Ok(());
464 }
465 self.store.update(&q.id, |r| {
466 r.waiter = Some(Note {
467 kind: WaiterKind::Daemon,
468 since: now,
469 });
470 Ok(())
471 })?;
472 tracing::info!(
473 "question {}: the asker is gone, resuming seat {} with the owner's word",
474 q.short(),
475 q.seat
476 );
477
478 let thread: Vec<(&str, &str)> = q
479 .thread
480 .iter()
481 .map(|t| {
482 (
483 if t.who == Who::Operator {
484 "operator"
485 } else {
486 "agent"
487 },
488 t.body.as_str(),
489 )
490 })
491 .collect();
492 let shown = match word {
495 Word::Said(_) => &thread[..q.delivered_turns.min(thread.len())],
496 Word::Answered(_) => &thread[..],
497 };
498 let owner_word = match &word {
499 Word::Said(s) => prompt::OwnerWord::Said(s),
500 Word::Answered(a) => prompt::OwnerWord::Answered(a),
501 };
502 let body = prompt::question_resumed(
503 &q.id,
504 &q.summary,
505 &q.detail,
506 shown,
507 &owner_word,
508 &target.language,
509 );
510
511 let artifacts = self.store.root().join(format!("{}.delivery", q.id));
512 let stem = format!("deliver-{}", now.as_second());
513 let cache_dir = None;
514 let inv = Invocation {
515 cwd: &target.cwd,
516 prompt: &body,
517 timeout: DELIVERY_TIMEOUT,
518 allow_write: target.allow_write,
519 sessions: target.sessions,
520 artifacts: &artifacts,
521 stem: &stem,
522 run: &q.run,
523 node: &q.node,
524 cache_dir,
525 attachments: &[],
526 writable: &[],
527 };
528
529 let store = self.store.clone();
530 let id = q.id.clone();
531 let out = {
532 let fut = agent::invoke(&target.spec, &mut target.seat, &inv);
533 tokio::pin!(fut);
534 let mut beat = tokio::time::interval(Duration::from_secs(1));
535 let mut beats = 0u32;
536 loop {
537 tokio::select! {
538 r = &mut fut => break Some(r),
539 _ = beat.tick() => {
540 if halt() {
541 break None;
542 }
543 beats += 1;
544 if beats % 20 == 0 {
545 store.beat(&id, WaiterKind::Daemon);
546 }
547 }
548 }
549 }
550 };
551
552 let Some(out) = out else {
553 let _ = self.store.update(&q.id, |r| {
555 r.waiter = None;
556 Ok(())
557 });
558 return Ok(());
559 };
560 let out = match out {
561 Ok(o) if o.usable() => o,
562 other => {
563 let why = match other {
564 Ok(o) if o.timed_out => "the resumed turn timed out".to_owned(),
565 Ok(o) if o.quota_exhausted() => "the agent is out of quota".to_owned(),
566 Ok(o) => format!("the resumed turn failed (exit {:?})", o.exit_code),
567 Err(e) => format!("the agent could not be started: {e:#}"),
568 };
569 self.refuse(&q, key, &why);
570 return Ok(());
571 }
572 };
573
574 let text = out.text.trim().to_owned();
575 self.store.update(&q.id, |r| {
576 r.waiter = None;
577 match &word {
578 Word::Answered(_) => r.answer_delivered = true,
579 Word::Said(_) => {
580 r.delivered_turns = r.delivered_turns.max(snapshot);
581 let agent_spoke = r.thread[snapshot.min(r.thread.len())..]
586 .iter()
587 .any(|t| t.who == Who::Agent);
588 if !agent_spoke && r.thread.len() == snapshot && r.status.open() {
589 let choices = r.choices.clone();
590 r.reply(text, choices)?;
591 }
592 }
593 }
594 Ok(())
595 })?;
596 self.memo.remove(&q.id);
597 Ok(())
598 }
599}
600
601pub(crate) fn take_claim(path: &std::path::Path, stale_after: Duration) -> bool {
605 let attempt = || {
606 std::fs::OpenOptions::new()
607 .write(true)
608 .create_new(true)
609 .open(path)
610 .is_ok()
611 };
612 if attempt() {
613 return true;
614 }
615 let stale = std::fs::metadata(path)
616 .and_then(|m| m.modified())
617 .ok()
618 .and_then(|t| t.elapsed().ok())
619 .is_some_and(|age| age > stale_after);
620 if stale {
621 let _ = std::fs::remove_file(path);
622 return attempt();
623 }
624 false
625}
626
627pub(crate) struct Release(pub(crate) PathBuf);
629
630impl Drop for Release {
631 fn drop(&mut self) {
632 let _ = std::fs::remove_file(&self.0);
633 }
634}
635
636pub async fn run(mut waiter: Waiter, stop: crate::daemon::Stop) {
642 let halt = {
643 let stop = stop.clone();
644 move || stop.parking()
645 };
646 while !stop.stopped() {
647 waiter.tick(Timestamp::now(), &halt).await;
648 tokio::time::sleep(TICK).await;
649 }
650}
651
652#[cfg(test)]
653mod tests {
654 use super::*;
655
656 fn ts(secs: i64) -> Timestamp {
657 Timestamp::from_second(secs).unwrap()
658 }
659
660 fn asked(at: i64) -> Question {
661 let mut q = Question::new(
662 "run".into(),
663 "implement".into(),
664 "impl-A".into(),
665 "which?".into(),
666 String::new(),
667 vec![],
668 );
669 q.cwd = Some("/tmp".into());
670 q.asked_at = ts(at);
671 q.answer_timeout = 1000;
672 q
673 }
674
675 fn lease(at: i64) -> Lease {
676 Lease {
677 kind: WaiterKind::Asker,
678 pid: 1,
679 beat_at: ts(at),
680 }
681 }
682
683 #[test]
684 fn a_fresh_lease_means_nobody_else_acts() {
685 let mut q = asked(0);
686 q.say("why?").unwrap();
687 assert_eq!(
688 decide(&q, Some(&lease(95)), false, 86_400, ts(100)),
689 Action::Idle
690 );
691 }
692
693 #[test]
694 fn a_stale_lease_delivers_what_the_owner_said() {
695 let mut q = asked(0);
696 q.say("why?").unwrap();
697 assert_eq!(
698 decide(&q, Some(&lease(0)), false, 86_400, ts(500)),
699 Action::Deliver(Word::Said("why?".into()))
700 );
701 assert_eq!(
702 decide(&q, None, true, 86_400, ts(500)),
703 Action::Idle,
704 "a seat still mid-turn is not resumed"
705 );
706 }
707
708 #[test]
709 fn an_answer_nobody_read_is_delivered_once() {
710 let mut q = asked(0);
711 q.choices = vec!["A".into()];
712 q.answer(crate::ask::Answer::Choice("A".into())).unwrap();
713 assert_eq!(
714 decide(&q, None, false, 86_400, ts(10)),
715 Action::Deliver(Word::Answered("A".into()))
716 );
717 q.answer_delivered = true;
718 assert_eq!(decide(&q, None, false, 86_400, ts(10)), Action::Idle);
719 }
720
721 #[test]
722 fn only_asked_questions_and_only_before_the_deadline() {
723 let mut q = asked(0);
724 assert_eq!(decide(&q, None, false, 86_400, ts(999)), Action::Idle);
725 assert_eq!(decide(&q, None, false, 86_400, ts(1001)), Action::Expire);
726 q.cwd = None;
727 assert_eq!(decide(&q, None, false, 86_400, ts(1001)), Action::Idle);
728 }
729
730 #[test]
731 fn the_deadline_runs_from_the_last_turn_not_from_asking() {
732 let mut q = asked(0);
733 q.thread.push(crate::ask::Turn {
734 who: Who::Agent,
735 body: "context".into(),
736 at: ts(4000),
737 });
738 q.delivered_turns = 1;
739 assert_eq!(decide(&q, None, false, 86_400, ts(4500)), Action::Idle);
740 assert_eq!(decide(&q, None, false, 86_400, ts(5001)), Action::Expire);
741 }
742
743 #[test]
744 fn a_word_given_in_time_is_delivered_after_the_deadline() {
745 let mut q = asked(0);
746 q.say("wait").unwrap();
747 assert!(matches!(
748 decide(&q, None, false, 86_400, ts(5000)),
749 Action::Deliver(_)
750 ));
751 }
752
753 #[test]
754 fn an_action_answer_is_left_to_the_daemon() {
755 let mut q = asked(0);
756 q.choices = vec!["A".into()];
757 q.actions
758 .insert("A".into(), crate::ask::ChoiceAction::Requeue);
759 q.answer(crate::ask::Answer::Choice("A".into())).unwrap();
760 assert_eq!(
761 decide_owned(&q, None, false, 86_400, ts(10), true),
762 Action::Idle,
763 "the daemon applies the action; the waiter must not also resume the seat"
764 );
765 q.answer_delivered = true;
767 assert_eq!(
768 decide_owned(&q, None, false, 86_400, ts(10), true),
769 Action::Idle
770 );
771 q.answer_delivered = false;
773 assert_eq!(
774 decide_owned(&q, None, false, 86_400, ts(10), false),
775 Action::Deliver(Word::Answered("A".into()))
776 );
777 let mut plain = asked(0);
779 plain.choices = vec!["A".into()];
780 plain
781 .answer(crate::ask::Answer::Choice("A".into()))
782 .unwrap();
783 assert_eq!(
784 decide_owned(&plain, None, false, 86_400, ts(10), true),
785 Action::Deliver(Word::Answered("A".into()))
786 );
787 }
788}