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 if q.cwd.is_none() || lease.is_some_and(|l| l.fresh(now)) {
108 return Action::Idle;
109 }
110 match q.status {
111 QuestionStatus::Abandoned => Action::Idle,
112 QuestionStatus::Answered => match q.resolution() {
113 Some(a) if !q.answer_delivered && !seat_busy => Action::Deliver(Word::Answered(a)),
114 _ => Action::Idle,
115 },
116 QuestionStatus::Open => {
117 if let Some(said) = q.unread_from_owner() {
118 if seat_busy {
119 return Action::Idle;
120 }
121 return Action::Deliver(Word::Said(said));
122 }
123 let secs = if q.answer_timeout > 0 {
124 q.answer_timeout
125 } else {
126 default_timeout
127 };
128 let deadline = q.last_activity().saturating_add(secs as i64);
129 if now.as_second() > deadline {
130 Action::Expire
131 } else {
132 Action::Idle
133 }
134 }
135 }
136}
137
138struct Target {
140 spec: AgentSpec,
141 seat: SeatState,
142 cwd: PathBuf,
143 sessions: bool,
144 allow_write: bool,
145 language: String,
146}
147
148pub struct Waiter {
151 store: Questions,
152 home: PathBuf,
153 conduct_cfg: Option<Config>,
156 default_timeout: u64,
158 memo: HashMap<String, (String, Instant)>,
160}
161
162impl Waiter {
163 pub fn new(store: Questions, home: PathBuf, conduct_cfg: Option<Config>) -> Self {
165 let default_timeout = conduct_cfg
166 .as_ref()
167 .map_or(Config::default().graph.answer_timeout, |c| {
168 c.graph.answer_timeout
169 });
170 Self {
171 store,
172 home,
173 conduct_cfg,
174 default_timeout,
175 memo: HashMap::new(),
176 }
177 }
178
179 fn seat_busy(&self, q: &Question, now: Timestamp) -> bool {
181 if q.run == crate::conduct::NODE {
182 return crate::conduct::busy(&self.home);
183 }
184 RunState::load_under(&q.run, &self.home).is_ok_and(|s| {
185 s.seats_active()
186 .any(|(k, a)| *k == q.seat && a.remaining_secs(now) > 0)
187 })
188 }
189
190 pub async fn tick(&mut self, now: Timestamp, halt: &(dyn Fn() -> bool + Sync)) {
195 for q in self.store.list() {
196 if halt() {
197 return;
198 }
199 let lease = self.store.read_lease(&q.id);
200 match decide(&q, lease.as_ref(), false, self.default_timeout, now) {
201 Action::Idle => {}
202 Action::Expire => self.expire(&q),
203 Action::Deliver(word) => {
204 if q.deputy.is_some() {
208 if crate::deputy::exhausted_past_deadline(
212 &q,
213 crate::deputy::can_start(
214 self.conduct_cfg.as_ref(),
215 crate::deputy::agent_of(&q),
216 ),
217 self.default_timeout,
218 now,
219 ) {
220 self.expire(&q);
221 }
222 continue;
223 }
224 if self.seat_busy(&q, now) {
225 continue;
226 }
227 let key = format!("{}:{:?}", q.thread.len(), word);
228 if matches!(self.memo.get(&q.id), Some((k, until)) if *k == key && Instant::now() < *until)
229 {
230 continue;
231 }
232 if let Err(e) = self.deliver(&q, &word, halt).await {
233 tracing::warn!("question {}: {e:#}", q.short());
234 }
235 }
236 }
237 }
238 }
239
240 fn expire(&self, q: &Question) {
243 let secs = if q.answer_timeout > 0 {
244 q.answer_timeout
245 } else {
246 self.default_timeout
247 };
248 let startable =
249 crate::deputy::can_start(self.conduct_cfg.as_ref(), crate::deputy::agent_of(q));
250 let why = format!("no answer within {}s of asking", secs.max(1));
251 let done = self.store.update(&q.id, |r| {
252 let now = Timestamp::now();
256 let lease = self.store.read_lease(&r.id);
257 if decide(r, lease.as_ref(), false, self.default_timeout, now) != Action::Expire
258 && !(lease.as_ref().is_none_or(|l| !l.fresh(now))
259 && crate::deputy::exhausted_past_deadline(
260 r,
261 startable,
262 self.default_timeout,
263 now,
264 ))
265 {
266 return Ok(false);
267 }
268 r.abandon(&why);
269 r.waiter = None;
270 Ok(true)
271 });
272 match done {
273 Ok((_, false)) => {}
274 Ok((_, true)) => {
275 self.store.drop_lease(&q.id);
276 tracing::warn!(
277 "question {} went unanswered for {secs}s with nobody waiting; \
278 it stays as the record of it",
279 q.short()
280 );
281 }
282 Err(e) => tracing::warn!("could not expire question {}: {e:#}", q.short()),
283 }
284 }
285
286 fn target(&self, q: &Question) -> std::result::Result<Target, String> {
288 let cwd = q
289 .cwd
290 .as_deref()
291 .map(PathBuf::from)
292 .filter(|p| p.is_dir())
293 .ok_or("the directory the agent was working in is gone")?;
294 let (cfg, seat, allow_write) = if q.run == crate::conduct::NODE {
295 let cfg = self
296 .conduct_cfg
297 .clone()
298 .ok_or("the conductor's configuration is not available")?;
299 let seat = crate::conduct::load_seat(&self.home)
300 .ok_or("the conductor has no recorded session")?;
301 (cfg, seat, false)
302 } else {
303 let state = RunState::load_under(&q.run, &self.home)
304 .map_err(|_| "the run that asked is gone".to_owned())?;
305 if !state.status.resumable() {
306 return Err(format!(
307 "the run that asked has already {}",
308 state.status.as_str()
309 ));
310 }
311 let seat = state
312 .seats
313 .get(&q.seat)
314 .cloned()
315 .ok_or("the run has no record of the asking seat")?;
316 let write = matches!(q.node.as_str(), "implement" | "fix");
317 (state.config, seat, write)
318 };
319 let spec = cfg
320 .agent(&seat.agent)
321 .map_err(|_| format!("agent `{}` is no longer in the roster", seat.agent))?
322 .clone();
323 if !agent::has_session(spec.kind, &seat, cfg.graph.sessions) {
324 return Err("the seat has no session to resume".to_owned());
325 }
326 Ok(Target {
327 spec,
328 seat,
329 cwd,
330 sessions: cfg.graph.sessions,
331 allow_write,
332 language: cfg.graph.language.clone(),
333 })
334 }
335
336 fn refuse(&mut self, q: &Question, key: String, why: &str) {
338 tracing::warn!("question {} cannot be delivered: {why}", q.short());
339 notices::raise_in(
340 &self.home,
341 Notice::warn(
342 &format!("question:{}", q.id),
343 format!(
344 "Question {} \"{}\": what you said cannot reach the agent \
345 ({why}). It is recorded, but nothing will read it.",
346 q.short(),
347 q.summary
348 ),
349 ),
350 );
351 let _ = self.store.update(&q.id, |r| {
352 r.waiter = None;
353 Ok(())
354 });
355 self.memo
356 .insert(q.id.clone(), (key, Instant::now() + RETRY_AFTER));
357 }
358
359 async fn deliver(
361 &mut self,
362 q: &Question,
363 word: &Word,
364 halt: &(dyn Fn() -> bool + Sync),
365 ) -> Result<()> {
366 let key = format!("{}:{:?}", q.thread.len(), word);
367 let mut target = match self.target(q) {
368 Ok(t) => t,
369 Err(why) => {
370 self.refuse(q, key, &why);
371 return Ok(());
372 }
373 };
374
375 let claim = self.store.root().join(format!("{}.claim", q.id));
376 if !take_claim(&claim, DELIVERY_TIMEOUT + Duration::from_secs(60)) {
377 return Ok(());
378 }
379 let _release = Release(claim);
380
381 let q = self.store.get(&q.id)?;
384 let now = Timestamp::now();
385 let lease = self.store.read_lease(&q.id);
386 let Action::Deliver(word) = decide(&q, lease.as_ref(), false, self.default_timeout, now)
387 else {
388 return Ok(());
389 };
390 let snapshot = q.thread.len();
391
392 self.store.beat(&q.id, WaiterKind::Daemon);
393 self.store.update(&q.id, |r| {
394 r.waiter = Some(Note {
395 kind: WaiterKind::Daemon,
396 since: now,
397 });
398 Ok(())
399 })?;
400 tracing::info!(
401 "question {}: the asker is gone, resuming seat {} with the owner's word",
402 q.short(),
403 q.seat
404 );
405
406 let thread: Vec<(&str, &str)> = q
407 .thread
408 .iter()
409 .map(|t| {
410 (
411 if t.who == Who::Operator {
412 "operator"
413 } else {
414 "agent"
415 },
416 t.body.as_str(),
417 )
418 })
419 .collect();
420 let shown = match word {
423 Word::Said(_) => &thread[..q.delivered_turns.min(thread.len())],
424 Word::Answered(_) => &thread[..],
425 };
426 let owner_word = match &word {
427 Word::Said(s) => prompt::OwnerWord::Said(s),
428 Word::Answered(a) => prompt::OwnerWord::Answered(a),
429 };
430 let body = prompt::question_resumed(
431 &q.id,
432 &q.summary,
433 &q.detail,
434 shown,
435 &owner_word,
436 &target.language,
437 );
438
439 let artifacts = self.store.root().join(format!("{}.delivery", q.id));
440 let stem = format!("deliver-{}", now.as_second());
441 let cache_dir = None;
442 let inv = Invocation {
443 cwd: &target.cwd,
444 prompt: &body,
445 timeout: DELIVERY_TIMEOUT,
446 allow_write: target.allow_write,
447 sessions: target.sessions,
448 artifacts: &artifacts,
449 stem: &stem,
450 run: &q.run,
451 node: &q.node,
452 cache_dir,
453 attachments: &[],
454 writable: &[],
455 };
456
457 let store = self.store.clone();
458 let id = q.id.clone();
459 let out = {
460 let fut = agent::invoke(&target.spec, &mut target.seat, &inv);
461 tokio::pin!(fut);
462 let mut beat = tokio::time::interval(Duration::from_secs(1));
463 let mut beats = 0u32;
464 loop {
465 tokio::select! {
466 r = &mut fut => break Some(r),
467 _ = beat.tick() => {
468 if halt() {
469 break None;
470 }
471 beats += 1;
472 if beats % 20 == 0 {
473 store.beat(&id, WaiterKind::Daemon);
474 }
475 }
476 }
477 }
478 };
479
480 let Some(out) = out else {
481 let _ = self.store.update(&q.id, |r| {
483 r.waiter = None;
484 Ok(())
485 });
486 return Ok(());
487 };
488 let out = match out {
489 Ok(o) if o.usable() => o,
490 other => {
491 let why = match other {
492 Ok(o) if o.timed_out => "the resumed turn timed out".to_owned(),
493 Ok(o) if o.quota_exhausted() => "the agent is out of quota".to_owned(),
494 Ok(o) => format!("the resumed turn failed (exit {:?})", o.exit_code),
495 Err(e) => format!("the agent could not be started: {e:#}"),
496 };
497 self.refuse(&q, key, &why);
498 return Ok(());
499 }
500 };
501
502 let text = out.text.trim().to_owned();
503 self.store.update(&q.id, |r| {
504 r.waiter = None;
505 match &word {
506 Word::Answered(_) => r.answer_delivered = true,
507 Word::Said(_) => {
508 r.delivered_turns = r.delivered_turns.max(snapshot);
509 let agent_spoke = r.thread[snapshot.min(r.thread.len())..]
514 .iter()
515 .any(|t| t.who == Who::Agent);
516 if !agent_spoke && r.thread.len() == snapshot && r.status.open() {
517 let choices = r.choices.clone();
518 r.reply(text, choices)?;
519 }
520 }
521 }
522 Ok(())
523 })?;
524 self.memo.remove(&q.id);
525 Ok(())
526 }
527}
528
529pub(crate) fn take_claim(path: &std::path::Path, stale_after: Duration) -> bool {
533 let attempt = || {
534 std::fs::OpenOptions::new()
535 .write(true)
536 .create_new(true)
537 .open(path)
538 .is_ok()
539 };
540 if attempt() {
541 return true;
542 }
543 let stale = std::fs::metadata(path)
544 .and_then(|m| m.modified())
545 .ok()
546 .and_then(|t| t.elapsed().ok())
547 .is_some_and(|age| age > stale_after);
548 if stale {
549 let _ = std::fs::remove_file(path);
550 return attempt();
551 }
552 false
553}
554
555pub(crate) struct Release(pub(crate) PathBuf);
557
558impl Drop for Release {
559 fn drop(&mut self) {
560 let _ = std::fs::remove_file(&self.0);
561 }
562}
563
564pub async fn run(mut waiter: Waiter, stop: crate::daemon::Stop) {
570 let halt = {
571 let stop = stop.clone();
572 move || stop.parking()
573 };
574 while !stop.stopped() {
575 waiter.tick(Timestamp::now(), &halt).await;
576 tokio::time::sleep(TICK).await;
577 }
578}
579
580#[cfg(test)]
581mod tests {
582 use super::*;
583
584 fn ts(secs: i64) -> Timestamp {
585 Timestamp::from_second(secs).unwrap()
586 }
587
588 fn asked(at: i64) -> Question {
589 let mut q = Question::new(
590 "run".into(),
591 "implement".into(),
592 "impl-A".into(),
593 "which?".into(),
594 String::new(),
595 vec![],
596 );
597 q.cwd = Some("/tmp".into());
598 q.asked_at = ts(at);
599 q.answer_timeout = 1000;
600 q
601 }
602
603 fn lease(at: i64) -> Lease {
604 Lease {
605 kind: WaiterKind::Asker,
606 pid: 1,
607 beat_at: ts(at),
608 }
609 }
610
611 #[test]
612 fn a_fresh_lease_means_nobody_else_acts() {
613 let mut q = asked(0);
614 q.say("why?").unwrap();
615 assert_eq!(
616 decide(&q, Some(&lease(95)), false, 86_400, ts(100)),
617 Action::Idle
618 );
619 }
620
621 #[test]
622 fn a_stale_lease_delivers_what_the_owner_said() {
623 let mut q = asked(0);
624 q.say("why?").unwrap();
625 assert_eq!(
626 decide(&q, Some(&lease(0)), false, 86_400, ts(500)),
627 Action::Deliver(Word::Said("why?".into()))
628 );
629 assert_eq!(
630 decide(&q, None, true, 86_400, ts(500)),
631 Action::Idle,
632 "a seat still mid-turn is not resumed"
633 );
634 }
635
636 #[test]
637 fn an_answer_nobody_read_is_delivered_once() {
638 let mut q = asked(0);
639 q.choices = vec!["A".into()];
640 q.answer(crate::ask::Answer::Choice("A".into())).unwrap();
641 assert_eq!(
642 decide(&q, None, false, 86_400, ts(10)),
643 Action::Deliver(Word::Answered("A".into()))
644 );
645 q.answer_delivered = true;
646 assert_eq!(decide(&q, None, false, 86_400, ts(10)), Action::Idle);
647 }
648
649 #[test]
650 fn only_asked_questions_and_only_before_the_deadline() {
651 let mut q = asked(0);
652 assert_eq!(decide(&q, None, false, 86_400, ts(999)), Action::Idle);
653 assert_eq!(decide(&q, None, false, 86_400, ts(1001)), Action::Expire);
654 q.cwd = None;
655 assert_eq!(decide(&q, None, false, 86_400, ts(1001)), Action::Idle);
656 }
657
658 #[test]
659 fn the_deadline_runs_from_the_last_turn_not_from_asking() {
660 let mut q = asked(0);
661 q.thread.push(crate::ask::Turn {
662 who: Who::Agent,
663 body: "context".into(),
664 at: ts(4000),
665 });
666 q.delivered_turns = 1;
667 assert_eq!(decide(&q, None, false, 86_400, ts(4500)), Action::Idle);
668 assert_eq!(decide(&q, None, false, 86_400, ts(5001)), Action::Expire);
669 }
670
671 #[test]
672 fn a_word_given_in_time_is_delivered_after_the_deadline() {
673 let mut q = asked(0);
674 q.say("wait").unwrap();
675 assert!(matches!(
676 decide(&q, None, false, 86_400, ts(5000)),
677 Action::Deliver(_)
678 ));
679 }
680
681 #[test]
682 fn an_action_answer_is_delivered_until_the_daemon_marks_it_handled() {
683 let mut q = asked(0);
684 q.choices = vec!["A".into()];
685 q.actions
686 .insert("A".into(), crate::ask::ChoiceAction::Requeue);
687 q.answer(crate::ask::Answer::Choice("A".into())).unwrap();
688 assert_eq!(
689 decide(&q, None, false, 86_400, ts(10)),
690 Action::Deliver(Word::Answered("A".into())),
691 "an action nobody applied must not swallow the answer"
692 );
693 q.answer_delivered = true;
694 assert_eq!(decide(&q, None, false, 86_400, ts(10)), Action::Idle);
695 }
696}