1use anyhow::{Context, Result, bail};
21use jiff::Timestamp;
22
23use crate::ask::{ChatConsult, Question, Questions};
24use crate::config::Config;
25use crate::queue::{CHAT_NODE, Source, Task};
26use crate::talk::{self, Talk, Talks};
27
28pub fn origin_talk(tasks: &[Task], talks: &[Talk], q: &Question) -> Option<Talk> {
35 if !q.status.open()
36 || q.node == crate::land::APPROVAL_NODE
37 || q.node == crate::bump::NOTICE_NODE
38 {
39 return None;
40 }
41 let task = crate::daemon::task_of_question(tasks, q)?;
42 let Source::Agent { run, .. } = &chat_origin(tasks, task)?.source else {
43 return None;
44 };
45 talks
46 .iter()
47 .find(|t| &t.id == run && t.status.open())
48 .cloned()
49}
50
51fn chat_origin<'a>(tasks: &'a [Task], start: &'a Task) -> Option<&'a Task> {
59 let mut seen = std::collections::HashSet::new();
60 let mut cur = start;
61 for _ in 0..=crate::followup::MAX_FOLLOWUP_GENERATION {
62 if !seen.insert(cur.id.as_str()) {
63 return None;
64 }
65 if matches!(&cur.source, Source::Agent { node, .. } if node == CHAT_NODE) {
66 return Some(cur);
67 }
68 let f = cur.followup.as_ref()?;
69 cur = match &f.origin_task {
70 Some(id) => tasks.iter().find(|t| &t.id == id)?,
71 None => tasks.iter().find(|t| t.runs.contains(&f.run))?,
72 };
73 }
74 None
75}
76
77pub fn pending_consults(questions: &Questions, talk_id: &str) -> bool {
84 questions
85 .list()
86 .iter()
87 .any(|q| q.status.open() && q.consult.as_ref().is_some_and(|c| c.talk == talk_id))
88}
89
90pub fn begin(questions: &Questions, talks: &Talks, q: &Question, talk: &Talk) -> Result<bool> {
96 let (q, fresh) = questions.update(&q.id, |r| {
97 if !r.status.open() {
98 bail!("question {} is already {}", r.short(), r.status.as_str());
99 }
100 if r.consult.is_some() {
101 return Ok(false);
102 }
103 r.consult = Some(ChatConsult {
104 talk: talk.id.clone(),
105 at: Timestamp::now(),
106 });
107 Ok(true)
108 })?;
109 if !fresh {
110 return Ok(false);
111 }
112 let mut talk = talk.clone();
113 if let Err(e) = talk::queue(
114 &mut talk,
115 talks,
116 &crate::prompt::chat_consult(&q),
117 Vec::new(),
118 ) {
119 let _ = questions.update(&q.id, |r| {
120 r.consult = None;
121 Ok(())
122 });
123 return Err(e).context("queue the question into the chat");
124 }
125 Ok(true)
126}
127
128#[derive(Debug, PartialEq, Eq)]
130pub enum Started {
131 Answered,
133 Busy,
135 Nothing,
137}
138
139pub async fn start_turn(
146 questions: &Questions,
147 talks: &Talks,
148 q: &Question,
149 talk: &Talk,
150 cfg: &Config,
151) -> Result<Started> {
152 if q.consult.is_some() {
153 return Ok(Started::Nothing);
154 }
155 let Some(mut lease) = talks.claim_turn(&talk.id)? else {
156 return Ok(Started::Busy);
157 };
158 if !begin(questions, talks, q, talk)? {
159 return Ok(Started::Nothing);
160 }
161 let mut talk = talks.get(&talk.id)?;
162 let mut failed = None;
167 loop {
168 while let Some(text) = talk::drain(&mut talk, talks)? {
169 if let Err(e) = talk::respond(&lease, &mut talk, talks, cfg, &text).await {
170 failed.get_or_insert(e);
171 }
172 if !lease.beat()? {
173 bail!("the turn lease for chat {} was lost", talk.short());
174 }
175 }
176 drop(lease);
181 talk = talks.get(&talk.id)?;
182 let owed = talk.status.open()
183 && (!talk.pending.is_empty() || !talk.pending_attachments.is_empty());
184 if !owed {
185 break;
186 }
187 let Some(again) = talks.claim_turn(&talk.id)? else {
189 break;
190 };
191 lease = again;
192 }
193 if let Some(e) = failed {
194 return Err(e);
195 }
196 Ok(Started::Answered)
197}
198
199#[cfg(test)]
200mod tests {
201 use std::collections::BTreeMap;
202 use std::path::PathBuf;
203
204 use super::*;
205 use crate::config::{AgentKind, AgentSpec, Config};
206 use crate::queue::Source;
207 use crate::talk::TalkStatus;
208
209 const RUN: &str = "20260902-000000-beef";
210
211 fn talks() -> (tempfile::TempDir, Talks, Talk) {
212 let tmp = tempfile::tempdir().expect("tempdir");
213 let store = Talks::at(tmp.path().join("talks"));
214 let cfg = Config {
215 agents: vec![AgentSpec {
216 id: "mock".to_owned(),
217 kind: AgentKind::Command,
218 model: None,
219 command: vec!["true".to_owned()],
220 extra_args: Vec::new(),
221 env: BTreeMap::new(),
222 prompt_delivery: None,
223 }],
224 ..Config::default()
225 };
226 let talk = talk::begin(&store, &cfg, tmp.path().to_path_buf(), Some("mock")).expect("talk");
227 (tmp, store, talk)
228 }
229
230 fn task(source: Source) -> Task {
231 let mut t = Task::new(
232 "t".to_owned(),
233 "Do it".to_owned(),
234 PathBuf::from("/repo"),
235 source,
236 );
237 t.start(RUN.to_owned());
238 t
239 }
240
241 fn from_chat(talk: &Talk) -> Task {
242 task(Source::Agent {
243 run: talk.id.clone(),
244 node: CHAT_NODE.to_owned(),
245 })
246 }
247
248 fn question(node: &str) -> Question {
249 Question::new(
250 RUN.to_owned(),
251 node.to_owned(),
252 "impl-A".to_owned(),
253 "Which backend?".to_owned(),
254 "SQLite is simpler.".to_owned(),
255 vec!["SQLite".to_owned(), "Redis".to_owned()],
256 )
257 }
258
259 #[tokio::test]
260 async fn start_turn_is_busy_and_writes_nothing_while_the_lease_is_held() {
261 let (tmp, store, talk) = talks();
262 let questions = Questions::at(tmp.path().join("questions"));
263 let mut q = question("implement");
264 questions.put(&mut q).expect("put");
265 let held = store.claim_turn(&talk.id).expect("claim").expect("first");
266
267 let got = start_turn(&questions, &store, &q, &talk, &Config::default())
268 .await
269 .expect("start");
270 assert_eq!(got, Started::Busy);
271 assert!(questions.get(&q.id).unwrap().consult.is_none());
272 assert!(store.get(&talk.id).unwrap().pending.is_empty());
273 drop(held);
274 }
275
276 fn scripted(tmp: &std::path::Path, script: &str, lease: &std::path::Path) -> Config {
279 let path = tmp.join("mock-consult-agent.sh");
280 std::fs::write(&path, script).expect("write mock");
281 Config {
282 agents: vec![AgentSpec {
283 id: "mock".to_owned(),
284 kind: AgentKind::Command,
285 model: None,
286 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
287 extra_args: Vec::new(),
288 env: BTreeMap::from([("LEASE".to_owned(), lease.to_string_lossy().into_owned())]),
289 prompt_delivery: None,
290 }],
291 ..Config::default()
292 }
293 }
294
295 async fn run_turn(
298 script: &str,
299 ) -> (
300 tempfile::TempDir,
301 Result<Started>,
302 Questions,
303 Question,
304 Talk,
305 ) {
306 crate::run::set_home(crate::run::test_home());
307 let (tmp, store, talk) = talks();
308 let questions = Questions::at(tmp.path().join("questions"));
309 let mut q = question("implement");
310 questions.put(&mut q).expect("put");
311 let cfg = scripted(tmp.path(), script, &store.turn_path(&talk.id));
312 let got = start_turn(&questions, &store, &q, &talk, &cfg).await;
313 (tmp, got, questions, q, talk)
314 }
315
316 #[tokio::test]
317 async fn the_lease_is_held_during_the_turn_and_free_once_it_ends() {
318 let (tmp, got, _questions, _q, talk) =
320 run_turn("#!/bin/sh\ncat >/dev/null\n[ -f \"$LEASE\" ] || exit 7\nprintf 'ok\\n'\n")
321 .await;
322 assert_eq!(got.expect("turn"), Started::Answered);
323 let other = Talks::at(tmp.path().join("talks"));
324 assert!(
325 other.claim_turn(&talk.id).expect("claim").is_some(),
326 "free once the turn ended"
327 );
328 }
329
330 #[tokio::test]
331 async fn the_lease_is_free_again_after_a_failed_turn() {
332 let (tmp, got, questions, q, talk) = run_turn("#!/bin/sh\ncat >/dev/null\nexit 3\n").await;
333 assert!(got.is_err(), "the agent failed");
334 assert!(
335 questions.get(&q.id).unwrap().consult.is_some(),
336 "the turn got as far as running"
337 );
338 let other = Talks::at(tmp.path().join("talks"));
339 assert!(
340 other.claim_turn(&talk.id).expect("claim").is_some(),
341 "free after the failed turn"
342 );
343 }
344
345 #[test]
346 fn origin_talk_is_decided_in_one_table() {
347 let (_tmp, _store, talk) = talks();
348 let mut closed = talk.clone();
349 closed.status = TalkStatus::Closed;
350 let chat = from_chat(&talk);
351 let q = question("implement");
352
353 let hit = |tasks: &[Task], talks: &[Talk], q: &Question| origin_talk(tasks, talks, q);
354 assert_eq!(
355 hit(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).map(|t| t.id),
356 Some(talk.id.clone())
357 );
358 assert!(hit(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
360 let other = task(Source::Agent {
361 run: talk.id.clone(),
362 node: "implement".to_owned(),
363 });
364 assert!(hit(&[other], std::slice::from_ref(&talk), &q).is_none());
365 assert!(hit(std::slice::from_ref(&chat), &[], &q).is_none());
367 assert!(hit(std::slice::from_ref(&chat), &[closed], &q).is_none());
368 assert!(hit(&[], std::slice::from_ref(&talk), &q).is_none());
370 for node in [crate::land::APPROVAL_NODE, crate::bump::NOTICE_NODE] {
372 assert!(
373 hit(
374 std::slice::from_ref(&chat),
375 std::slice::from_ref(&talk),
376 &question(node)
377 )
378 .is_none()
379 );
380 }
381 let mut answered = question("implement");
383 answered
384 .answer(crate::ask::Answer::Choice("Redis".to_owned()))
385 .unwrap();
386 assert!(hit(&[chat], &[talk], &answered).is_none());
387 }
388
389 #[test]
390 fn a_conductor_question_is_found_through_the_task_id() {
391 let (_tmp, _store, talk) = talks();
392 let chat = from_chat(&talk);
393 let mut q = question(crate::conduct::NODE);
394 q.run = chat.id.clone();
395 assert!(origin_talk(&[chat], &[talk], &q).is_some());
396 }
397
398 fn followup_of(parent: &Task, n: u32) -> Task {
399 let mut t = Task::new(
400 "f".to_owned(),
401 "Fix".to_owned(),
402 PathBuf::from("/repo"),
403 Source::Agent {
404 run: format!("merged-{n}"),
405 node: "followup".to_owned(),
406 },
407 );
408 t.start(format!("run-f{n}"));
409 t.followup = Some(crate::queue::FollowUp {
410 run: format!("merged-{n}"),
411 origin_task: Some(parent.id.clone()),
412 pr: "https://example.invalid/pr/1".to_owned(),
413 findings: Vec::new(),
414 generation: n,
415 });
416 t
417 }
418
419 fn q_for(t: &Task) -> Question {
420 let mut q = question("implement");
421 q.run = t.runs[0].clone();
422 q
423 }
424
425 #[test]
426 fn a_followup_traces_back_to_the_chat() {
427 let (_tmp, _store, talk) = talks();
428 let mut chat = from_chat(&talk);
429 chat.runs = vec!["chat-run".to_owned()];
430 let f1 = followup_of(&chat, 1);
431 let f2 = followup_of(&f1, 2);
432 let ts = [chat, f1.clone(), f2.clone()];
433 let id = Some(talk.id.clone());
434 let tk = std::slice::from_ref(&talk);
435 assert_eq!(origin_talk(&ts, tk, &q_for(&f1)).map(|t| t.id), id);
436 assert_eq!(origin_talk(&ts, tk, &q_for(&f2)).map(|t| t.id), id);
437 }
438
439 #[test]
440 fn a_followup_without_origin_task_is_found_through_runs() {
441 let (_tmp, _store, talk) = talks();
442 let mut chat = from_chat(&talk);
443 chat.runs = vec!["merged-1".to_owned()];
444 let mut f1 = followup_of(&chat, 1);
445 f1.followup.as_mut().unwrap().origin_task = None;
446 let q = q_for(&f1);
447 assert!(origin_talk(&[chat, f1], &[talk], &q).is_some());
448 }
449
450 #[test]
451 fn a_dangling_origin_task_has_no_chat() {
452 let (_tmp, _store, talk) = talks();
453 let mut chat = from_chat(&talk);
454 chat.runs = vec!["chat-run".to_owned()];
455 let mut f1 = followup_of(&chat, 1);
456 f1.followup.as_mut().unwrap().origin_task = Some("gone".to_owned());
457 let q = q_for(&f1);
458 assert!(origin_talk(&[chat, f1], &[talk], &q).is_none());
459 }
460
461 #[test]
462 fn a_followup_cycle_ends_without_a_chat() {
463 let (_tmp, _store, talk) = talks();
464 let mut a = followup_of(&from_chat(&talk), 1);
465 let mut b = followup_of(&a, 2);
466 a.followup.as_mut().unwrap().origin_task = Some(b.id.clone());
467 b.followup.as_mut().unwrap().origin_task = Some(a.id.clone());
468 let q = q_for(&a);
469 assert!(origin_talk(&[a, b], &[talk], &q).is_none());
470 }
471
472 #[test]
473 fn pending_consults_follow_the_store_not_the_turn_text() {
474 let (tmp, store, talk) = talks();
475 let questions = Questions::at(tmp.path().join("questions"));
476 assert!(!pending_consults(&questions, &talk.id), "empty store");
477
478 let mut plain = question("implement");
479 questions.put(&mut plain).unwrap();
480 assert!(!pending_consults(&questions, &talk.id), "no consult record");
481
482 let mut q = question("implement");
483 questions.put(&mut q).unwrap();
484 begin(&questions, &store, &q, &talk).unwrap();
485 assert!(pending_consults(&questions, &talk.id));
486 assert!(!pending_consults(&questions, "other-talk"));
487
488 questions
489 .update(&q.id, |q| {
490 q.abandon("test");
491 Ok(())
492 })
493 .unwrap();
494 assert!(!pending_consults(&questions, &talk.id), "closed question");
495 }
496
497 #[test]
498 fn begin_queues_once_and_leaves_the_question_open() {
499 let (tmp, store, talk) = talks();
500 let questions = Questions::at(tmp.path().join("questions"));
501 let mut q = question("implement");
502 questions.put(&mut q).unwrap();
503
504 assert!(begin(&questions, &store, &q, &talk).unwrap());
505 assert!(!begin(&questions, &store, &q, &talk).unwrap(), "idempotent");
506
507 let after = questions.get(&q.id).unwrap();
508 assert!(after.status.open());
509 assert!(after.thread.is_empty());
510 assert_eq!(after.choices, q.choices, "choices are not touched");
511 assert_eq!(
512 after.consult.as_ref().map(|c| c.talk.as_str()),
513 Some(talk.id.as_str())
514 );
515
516 let queued = store.get(&talk.id).unwrap().pending;
517 assert_eq!(
518 queued.matches(&q.id).count(),
519 2 + 1,
520 "id once per use: {queued}"
521 );
522 assert!(queued.contains("Which backend?"));
523 assert!(queued.contains("SQLite is simpler."));
524 assert!(queued.contains("- Redis"));
525 assert_eq!(
527 queued.matches(crate::prompt::CHAT_CONSULT_HEADING).count(),
528 1
529 );
530 }
531
532 #[test]
533 fn begin_withdraws_its_record_when_the_talk_cannot_take_it() {
534 let (tmp, store, talk) = talks();
535 let questions = Questions::at(tmp.path().join("questions"));
536 let mut q = question("implement");
537 questions.put(&mut q).unwrap();
538 std::fs::remove_dir_all(tmp.path().join("talks")).unwrap();
539 assert!(begin(&questions, &store, &q, &talk).is_err());
540 assert!(questions.get(&q.id).unwrap().consult.is_none());
541 }
542
543 #[test]
544 fn an_older_question_file_reads_without_a_consult() {
545 let mut q = question("implement");
546 q.schema = 5;
547 let mut v = serde_json::to_value(&q).unwrap();
548 v.as_object_mut().unwrap().remove("consult");
549 let back: Question = serde_json::from_value(v).unwrap();
550 assert!(back.consult.is_none());
551 }
552}