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 self.seat_busy(&q, now) {
205 continue;
206 }
207 let key = format!("{}:{:?}", q.thread.len(), word);
208 if matches!(self.memo.get(&q.id), Some((k, until)) if *k == key && Instant::now() < *until)
209 {
210 continue;
211 }
212 if let Err(e) = self.deliver(&q, &word, halt).await {
213 tracing::warn!("question {}: {e:#}", q.short());
214 }
215 }
216 }
217 }
218 }
219
220 fn expire(&self, q: &Question) {
223 let secs = if q.answer_timeout > 0 {
224 q.answer_timeout
225 } else {
226 self.default_timeout
227 };
228 let why = format!("no answer within {}s of asking", secs.max(1));
229 let done = self.store.update(&q.id, |r| {
230 let now = Timestamp::now();
234 let lease = self.store.read_lease(&r.id);
235 if decide(r, lease.as_ref(), false, self.default_timeout, now) != Action::Expire {
236 return Ok(false);
237 }
238 r.abandon(&why);
239 r.waiter = None;
240 Ok(true)
241 });
242 match done {
243 Ok((_, false)) => {}
244 Ok((_, true)) => {
245 self.store.drop_lease(&q.id);
246 tracing::warn!(
247 "question {} went unanswered for {secs}s with nobody waiting; \
248 it stays as the record of it",
249 q.short()
250 );
251 }
252 Err(e) => tracing::warn!("could not expire question {}: {e:#}", q.short()),
253 }
254 }
255
256 fn target(&self, q: &Question) -> std::result::Result<Target, String> {
258 let cwd = q
259 .cwd
260 .as_deref()
261 .map(PathBuf::from)
262 .filter(|p| p.is_dir())
263 .ok_or("the directory the agent was working in is gone")?;
264 let (cfg, seat, allow_write) = if q.run == crate::conduct::NODE {
265 let cfg = self
266 .conduct_cfg
267 .clone()
268 .ok_or("the conductor's configuration is not available")?;
269 let seat = crate::conduct::load_seat(&self.home)
270 .ok_or("the conductor has no recorded session")?;
271 (cfg, seat, false)
272 } else {
273 let state = RunState::load_under(&q.run, &self.home)
274 .map_err(|_| "the run that asked is gone".to_owned())?;
275 if !state.status.resumable() {
276 return Err(format!(
277 "the run that asked has already {}",
278 state.status.as_str()
279 ));
280 }
281 let seat = state
282 .seats
283 .get(&q.seat)
284 .cloned()
285 .ok_or("the run has no record of the asking seat")?;
286 let write = matches!(q.node.as_str(), "implement" | "fix");
287 (state.config, seat, write)
288 };
289 let spec = cfg
290 .agent(&seat.agent)
291 .map_err(|_| format!("agent `{}` is no longer in the roster", seat.agent))?
292 .clone();
293 if !agent::has_session(spec.kind, &seat, cfg.graph.sessions) {
294 return Err("the seat has no session to resume".to_owned());
295 }
296 Ok(Target {
297 spec,
298 seat,
299 cwd,
300 sessions: cfg.graph.sessions,
301 allow_write,
302 language: cfg.graph.language.clone(),
303 })
304 }
305
306 fn refuse(&mut self, q: &Question, key: String, why: &str) {
308 tracing::warn!("question {} cannot be delivered: {why}", q.short());
309 notices::raise_in(
310 &self.home,
311 Notice::warn(
312 &format!("question:{}", q.id),
313 format!(
314 "Question {} \"{}\": what you said cannot reach the agent \
315 ({why}). It is recorded, but nothing will read it.",
316 q.short(),
317 q.summary
318 ),
319 ),
320 );
321 let _ = self.store.update(&q.id, |r| {
322 r.waiter = None;
323 Ok(())
324 });
325 self.memo
326 .insert(q.id.clone(), (key, Instant::now() + RETRY_AFTER));
327 }
328
329 async fn deliver(
331 &mut self,
332 q: &Question,
333 word: &Word,
334 halt: &(dyn Fn() -> bool + Sync),
335 ) -> Result<()> {
336 let key = format!("{}:{:?}", q.thread.len(), word);
337 let mut target = match self.target(q) {
338 Ok(t) => t,
339 Err(why) => {
340 self.refuse(q, key, &why);
341 return Ok(());
342 }
343 };
344
345 let claim = self.store.root().join(format!("{}.claim", q.id));
346 if !take_claim(&claim) {
347 return Ok(());
348 }
349 let _release = Release(claim);
350
351 let q = self.store.get(&q.id)?;
354 let now = Timestamp::now();
355 let lease = self.store.read_lease(&q.id);
356 let Action::Deliver(word) = decide(&q, lease.as_ref(), false, self.default_timeout, now)
357 else {
358 return Ok(());
359 };
360 let snapshot = q.thread.len();
361
362 self.store.beat(&q.id, WaiterKind::Daemon);
363 self.store.update(&q.id, |r| {
364 r.waiter = Some(Note {
365 kind: WaiterKind::Daemon,
366 since: now,
367 });
368 Ok(())
369 })?;
370 tracing::info!(
371 "question {}: the asker is gone, resuming seat {} with the owner's word",
372 q.short(),
373 q.seat
374 );
375
376 let thread: Vec<(&str, &str)> = q
377 .thread
378 .iter()
379 .map(|t| {
380 (
381 if t.who == Who::Operator {
382 "operator"
383 } else {
384 "agent"
385 },
386 t.body.as_str(),
387 )
388 })
389 .collect();
390 let shown = match word {
393 Word::Said(_) => &thread[..q.delivered_turns.min(thread.len())],
394 Word::Answered(_) => &thread[..],
395 };
396 let owner_word = match &word {
397 Word::Said(s) => prompt::OwnerWord::Said(s),
398 Word::Answered(a) => prompt::OwnerWord::Answered(a),
399 };
400 let body = prompt::question_resumed(
401 &q.id,
402 &q.summary,
403 &q.detail,
404 shown,
405 &owner_word,
406 &target.language,
407 );
408
409 let artifacts = self.store.root().join(format!("{}.delivery", q.id));
410 let stem = format!("deliver-{}", now.as_second());
411 let cache_dir = None;
412 let inv = Invocation {
413 cwd: &target.cwd,
414 prompt: &body,
415 timeout: DELIVERY_TIMEOUT,
416 allow_write: target.allow_write,
417 sessions: target.sessions,
418 artifacts: &artifacts,
419 stem: &stem,
420 run: &q.run,
421 node: &q.node,
422 cache_dir,
423 attachments: &[],
424 };
425
426 let store = self.store.clone();
427 let id = q.id.clone();
428 let out = {
429 let fut = agent::invoke(&target.spec, &mut target.seat, &inv);
430 tokio::pin!(fut);
431 let mut beat = tokio::time::interval(Duration::from_secs(1));
432 let mut beats = 0u32;
433 loop {
434 tokio::select! {
435 r = &mut fut => break Some(r),
436 _ = beat.tick() => {
437 if halt() {
438 break None;
439 }
440 beats += 1;
441 if beats % 20 == 0 {
442 store.beat(&id, WaiterKind::Daemon);
443 }
444 }
445 }
446 }
447 };
448
449 let Some(out) = out else {
450 let _ = self.store.update(&q.id, |r| {
452 r.waiter = None;
453 Ok(())
454 });
455 return Ok(());
456 };
457 let out = match out {
458 Ok(o) if o.usable() => o,
459 other => {
460 let why = match other {
461 Ok(o) if o.timed_out => "the resumed turn timed out".to_owned(),
462 Ok(o) if o.quota_exhausted() => "the agent is out of quota".to_owned(),
463 Ok(o) => format!("the resumed turn failed (exit {:?})", o.exit_code),
464 Err(e) => format!("the agent could not be started: {e:#}"),
465 };
466 self.refuse(&q, key, &why);
467 return Ok(());
468 }
469 };
470
471 let text = out.text.trim().to_owned();
472 self.store.update(&q.id, |r| {
473 r.waiter = None;
474 match &word {
475 Word::Answered(_) => r.answer_delivered = true,
476 Word::Said(_) => {
477 r.delivered_turns = r.delivered_turns.max(snapshot);
478 let agent_spoke = r.thread[snapshot.min(r.thread.len())..]
483 .iter()
484 .any(|t| t.who == Who::Agent);
485 if !agent_spoke && r.thread.len() == snapshot && r.status.open() {
486 let choices = r.choices.clone();
487 r.reply(text, choices)?;
488 }
489 }
490 }
491 Ok(())
492 })?;
493 self.memo.remove(&q.id);
494 Ok(())
495 }
496}
497
498fn take_claim(path: &std::path::Path) -> bool {
502 let attempt = || {
503 std::fs::OpenOptions::new()
504 .write(true)
505 .create_new(true)
506 .open(path)
507 .is_ok()
508 };
509 if attempt() {
510 return true;
511 }
512 let stale = std::fs::metadata(path)
513 .and_then(|m| m.modified())
514 .ok()
515 .and_then(|t| t.elapsed().ok())
516 .is_some_and(|age| age > DELIVERY_TIMEOUT + Duration::from_secs(60));
517 if stale {
518 let _ = std::fs::remove_file(path);
519 return attempt();
520 }
521 false
522}
523
524struct Release(PathBuf);
526
527impl Drop for Release {
528 fn drop(&mut self) {
529 let _ = std::fs::remove_file(&self.0);
530 }
531}
532
533pub async fn run(mut waiter: Waiter, stop: crate::daemon::Stop) {
539 let halt = {
540 let stop = stop.clone();
541 move || stop.parking()
542 };
543 while !stop.stopped() {
544 waiter.tick(Timestamp::now(), &halt).await;
545 tokio::time::sleep(TICK).await;
546 }
547}
548
549#[cfg(test)]
550mod tests {
551 use super::*;
552
553 fn ts(secs: i64) -> Timestamp {
554 Timestamp::from_second(secs).unwrap()
555 }
556
557 fn asked(at: i64) -> Question {
558 let mut q = Question::new(
559 "run".into(),
560 "implement".into(),
561 "impl-A".into(),
562 "which?".into(),
563 String::new(),
564 vec![],
565 );
566 q.cwd = Some("/tmp".into());
567 q.asked_at = ts(at);
568 q.answer_timeout = 1000;
569 q
570 }
571
572 fn lease(at: i64) -> Lease {
573 Lease {
574 kind: WaiterKind::Asker,
575 pid: 1,
576 beat_at: ts(at),
577 }
578 }
579
580 #[test]
581 fn a_fresh_lease_means_nobody_else_acts() {
582 let mut q = asked(0);
583 q.say("why?").unwrap();
584 assert_eq!(
585 decide(&q, Some(&lease(95)), false, 86_400, ts(100)),
586 Action::Idle
587 );
588 }
589
590 #[test]
591 fn a_stale_lease_delivers_what_the_owner_said() {
592 let mut q = asked(0);
593 q.say("why?").unwrap();
594 assert_eq!(
595 decide(&q, Some(&lease(0)), false, 86_400, ts(500)),
596 Action::Deliver(Word::Said("why?".into()))
597 );
598 assert_eq!(
599 decide(&q, None, true, 86_400, ts(500)),
600 Action::Idle,
601 "a seat still mid-turn is not resumed"
602 );
603 }
604
605 #[test]
606 fn an_answer_nobody_read_is_delivered_once() {
607 let mut q = asked(0);
608 q.choices = vec!["A".into()];
609 q.answer(crate::ask::Answer::Choice("A".into())).unwrap();
610 assert_eq!(
611 decide(&q, None, false, 86_400, ts(10)),
612 Action::Deliver(Word::Answered("A".into()))
613 );
614 q.answer_delivered = true;
615 assert_eq!(decide(&q, None, false, 86_400, ts(10)), Action::Idle);
616 }
617
618 #[test]
619 fn only_asked_questions_and_only_before_the_deadline() {
620 let mut q = asked(0);
621 assert_eq!(decide(&q, None, false, 86_400, ts(999)), Action::Idle);
622 assert_eq!(decide(&q, None, false, 86_400, ts(1001)), Action::Expire);
623 q.cwd = None;
624 assert_eq!(decide(&q, None, false, 86_400, ts(1001)), Action::Idle);
625 }
626
627 #[test]
628 fn the_deadline_runs_from_the_last_turn_not_from_asking() {
629 let mut q = asked(0);
630 q.thread.push(crate::ask::Turn {
631 who: Who::Agent,
632 body: "context".into(),
633 at: ts(4000),
634 });
635 q.delivered_turns = 1;
636 assert_eq!(decide(&q, None, false, 86_400, ts(4500)), Action::Idle);
637 assert_eq!(decide(&q, None, false, 86_400, ts(5001)), Action::Expire);
638 }
639
640 #[test]
641 fn a_word_given_in_time_is_delivered_after_the_deadline() {
642 let mut q = asked(0);
643 q.say("wait").unwrap();
644 assert!(matches!(
645 decide(&q, None, false, 86_400, ts(5000)),
646 Action::Deliver(_)
647 ));
648 }
649
650 #[test]
651 fn an_action_answer_is_delivered_until_the_daemon_marks_it_handled() {
652 let mut q = asked(0);
653 q.choices = vec!["A".into()];
654 q.actions
655 .insert("A".into(), crate::ask::ChoiceAction::Requeue);
656 q.answer(crate::ask::Answer::Choice("A".into())).unwrap();
657 assert_eq!(
658 decide(&q, None, false, 86_400, ts(10)),
659 Action::Deliver(Word::Answered("A".into())),
660 "an action nobody applied must not swallow the answer"
661 );
662 q.answer_delivered = true;
663 assert_eq!(decide(&q, None, false, 86_400, ts(10)), Action::Idle);
664 }
665}