1use anyhow::{Context, Result, bail};
21use jiff::Timestamp;
22
23use crate::ask::{ChatConsult, Question, Questions};
24use crate::queue::{CHAT_NODE, Source, Task};
25use crate::talk::{self, Talk, Talks};
26
27pub fn origin_talk(tasks: &[Task], talks: &[Talk], q: &Question) -> Option<Talk> {
34 if !q.status.open()
35 || q.node == crate::land::APPROVAL_NODE
36 || q.node == crate::bump::NOTICE_NODE
37 {
38 return None;
39 }
40 let task = crate::daemon::task_of_question(tasks, q)?;
41 let Source::Agent { run, .. } = &chat_origin(tasks, task)?.source else {
42 return None;
43 };
44 talks
45 .iter()
46 .find(|t| &t.id == run && t.status.open())
47 .cloned()
48}
49
50fn chat_origin<'a>(tasks: &'a [Task], start: &'a Task) -> Option<&'a Task> {
58 let mut seen = std::collections::HashSet::new();
59 let mut cur = start;
60 for _ in 0..=crate::followup::MAX_FOLLOWUP_GENERATION {
61 if !seen.insert(cur.id.as_str()) {
62 return None;
63 }
64 if matches!(&cur.source, Source::Agent { node, .. } if node == CHAT_NODE) {
65 return Some(cur);
66 }
67 let f = cur.followup.as_ref()?;
68 cur = match &f.origin_task {
69 Some(id) => tasks.iter().find(|t| &t.id == id)?,
70 None => tasks.iter().find(|t| t.runs.contains(&f.run))?,
71 };
72 }
73 None
74}
75
76pub fn pending_consults(questions: &Questions, talk_id: &str) -> bool {
83 questions
84 .list()
85 .iter()
86 .any(|q| q.status.open() && q.consult.as_ref().is_some_and(|c| c.talk == talk_id))
87}
88
89pub fn begin(questions: &Questions, talks: &Talks, q: &Question, talk: &Talk) -> Result<bool> {
95 let (q, fresh) = questions.update(&q.id, |r| {
96 if !r.status.open() {
97 bail!("question {} is already {}", r.short(), r.status.as_str());
98 }
99 if r.consult.is_some() {
100 return Ok(false);
101 }
102 r.consult = Some(ChatConsult {
103 talk: talk.id.clone(),
104 at: Timestamp::now(),
105 });
106 Ok(true)
107 })?;
108 if !fresh {
109 return Ok(false);
110 }
111 let mut talk = talk.clone();
112 if let Err(e) = talk::queue(
113 &mut talk,
114 talks,
115 &crate::prompt::chat_consult(&q),
116 Vec::new(),
117 ) {
118 let _ = questions.update(&q.id, |r| {
119 r.consult = None;
120 Ok(())
121 });
122 return Err(e).context("queue the question into the chat");
123 }
124 Ok(true)
125}
126
127#[derive(Debug, PartialEq, Eq)]
129pub enum Handled {
130 Ran(usize),
132 Idle,
134 Busy,
136}
137
138pub async fn run_turn(
146 talks: &Talks,
147 cfg: &crate::config::Config,
148 talk_id: &str,
149 wait: std::time::Duration,
150 poll: std::time::Duration,
151) -> Result<Handled> {
152 let deadline = std::time::Instant::now() + wait;
153 let lease = loop {
154 if let Some(lease) = talks.claim_turn(talk_id)? {
155 break lease;
156 }
157 if std::time::Instant::now() >= deadline {
158 return Ok(Handled::Busy);
159 }
160 tokio::time::sleep(poll).await;
161 };
162 let mut talk = talks.get(talk_id)?;
165 let mut ran = 0;
166 loop {
167 if !lease.beat()? {
168 bail!("the turn lease was taken over; the remaining drafts stay queued");
169 }
170 let Some(text) = talk::drain(&mut talk, talks)? else {
171 break;
172 };
173 lease
174 .beating(talk::respond(&mut talk, talks, cfg, &text))
175 .await??;
176 ran += 1;
177 }
178 Ok(if ran == 0 {
179 Handled::Idle
180 } else {
181 Handled::Ran(ran)
182 })
183}
184
185#[cfg(test)]
186mod tests {
187 use std::collections::BTreeMap;
188 use std::path::PathBuf;
189
190 use super::*;
191 use crate::config::{AgentKind, AgentSpec, Config};
192 use crate::queue::Source;
193 use crate::talk::TalkStatus;
194
195 const RUN: &str = "20260902-000000-beef";
196
197 fn talks() -> (tempfile::TempDir, Talks, Talk) {
198 let tmp = tempfile::tempdir().expect("tempdir");
199 let store = Talks::at(tmp.path().join("talks"));
200 let cfg = Config {
201 agents: vec![AgentSpec {
202 id: "mock".to_owned(),
203 kind: AgentKind::Command,
204 model: None,
205 command: vec!["true".to_owned()],
206 extra_args: Vec::new(),
207 env: BTreeMap::new(),
208 prompt_delivery: None,
209 }],
210 ..Config::default()
211 };
212 let talk = talk::begin(&store, &cfg, tmp.path().to_path_buf(), Some("mock")).expect("talk");
213 (tmp, store, talk)
214 }
215
216 fn task(source: Source) -> Task {
217 let mut t = Task::new(
218 "t".to_owned(),
219 "Do it".to_owned(),
220 PathBuf::from("/repo"),
221 source,
222 );
223 t.start(RUN.to_owned());
224 t
225 }
226
227 fn from_chat(talk: &Talk) -> Task {
228 task(Source::Agent {
229 run: talk.id.clone(),
230 node: CHAT_NODE.to_owned(),
231 })
232 }
233
234 fn question(node: &str) -> Question {
235 Question::new(
236 RUN.to_owned(),
237 node.to_owned(),
238 "impl-A".to_owned(),
239 "Which backend?".to_owned(),
240 "SQLite is simpler.".to_owned(),
241 vec!["SQLite".to_owned(), "Redis".to_owned()],
242 )
243 }
244
245 #[test]
246 fn origin_talk_is_decided_in_one_table() {
247 let (_tmp, _store, talk) = talks();
248 let mut closed = talk.clone();
249 closed.status = TalkStatus::Closed;
250 let chat = from_chat(&talk);
251 let q = question("implement");
252
253 let hit = |tasks: &[Task], talks: &[Talk], q: &Question| origin_talk(tasks, talks, q);
254 assert_eq!(
255 hit(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).map(|t| t.id),
256 Some(talk.id.clone())
257 );
258 assert!(hit(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
260 let other = task(Source::Agent {
261 run: talk.id.clone(),
262 node: "implement".to_owned(),
263 });
264 assert!(hit(&[other], std::slice::from_ref(&talk), &q).is_none());
265 assert!(hit(std::slice::from_ref(&chat), &[], &q).is_none());
267 assert!(hit(std::slice::from_ref(&chat), &[closed], &q).is_none());
268 assert!(hit(&[], std::slice::from_ref(&talk), &q).is_none());
270 for node in [crate::land::APPROVAL_NODE, crate::bump::NOTICE_NODE] {
272 assert!(
273 hit(
274 std::slice::from_ref(&chat),
275 std::slice::from_ref(&talk),
276 &question(node)
277 )
278 .is_none()
279 );
280 }
281 let mut answered = question("implement");
283 answered
284 .answer(crate::ask::Answer::Choice("Redis".to_owned()))
285 .unwrap();
286 assert!(hit(&[chat], &[talk], &answered).is_none());
287 }
288
289 #[test]
290 fn a_conductor_question_is_found_through_the_task_id() {
291 let (_tmp, _store, talk) = talks();
292 let chat = from_chat(&talk);
293 let mut q = question(crate::conduct::NODE);
294 q.run = chat.id.clone();
295 assert!(origin_talk(&[chat], &[talk], &q).is_some());
296 }
297
298 fn followup_of(parent: &Task, n: u32) -> Task {
299 let mut t = Task::new(
300 "f".to_owned(),
301 "Fix".to_owned(),
302 PathBuf::from("/repo"),
303 Source::Agent {
304 run: format!("merged-{n}"),
305 node: "followup".to_owned(),
306 },
307 );
308 t.start(format!("run-f{n}"));
309 t.followup = Some(crate::queue::FollowUp {
310 run: format!("merged-{n}"),
311 origin_task: Some(parent.id.clone()),
312 pr: "https://example.invalid/pr/1".to_owned(),
313 findings: Vec::new(),
314 generation: n,
315 });
316 t
317 }
318
319 fn q_for(t: &Task) -> Question {
320 let mut q = question("implement");
321 q.run = t.runs[0].clone();
322 q
323 }
324
325 #[test]
326 fn a_followup_traces_back_to_the_chat() {
327 let (_tmp, _store, talk) = talks();
328 let mut chat = from_chat(&talk);
329 chat.runs = vec!["chat-run".to_owned()];
330 let f1 = followup_of(&chat, 1);
331 let f2 = followup_of(&f1, 2);
332 let ts = [chat, f1.clone(), f2.clone()];
333 let id = Some(talk.id.clone());
334 let tk = std::slice::from_ref(&talk);
335 assert_eq!(origin_talk(&ts, tk, &q_for(&f1)).map(|t| t.id), id);
336 assert_eq!(origin_talk(&ts, tk, &q_for(&f2)).map(|t| t.id), id);
337 }
338
339 #[test]
340 fn a_followup_without_origin_task_is_found_through_runs() {
341 let (_tmp, _store, talk) = talks();
342 let mut chat = from_chat(&talk);
343 chat.runs = vec!["merged-1".to_owned()];
344 let mut f1 = followup_of(&chat, 1);
345 f1.followup.as_mut().unwrap().origin_task = None;
346 let q = q_for(&f1);
347 assert!(origin_talk(&[chat, f1], &[talk], &q).is_some());
348 }
349
350 #[test]
351 fn a_dangling_origin_task_has_no_chat() {
352 let (_tmp, _store, talk) = talks();
353 let mut chat = from_chat(&talk);
354 chat.runs = vec!["chat-run".to_owned()];
355 let mut f1 = followup_of(&chat, 1);
356 f1.followup.as_mut().unwrap().origin_task = Some("gone".to_owned());
357 let q = q_for(&f1);
358 assert!(origin_talk(&[chat, f1], &[talk], &q).is_none());
359 }
360
361 #[test]
362 fn a_followup_cycle_ends_without_a_chat() {
363 let (_tmp, _store, talk) = talks();
364 let mut a = followup_of(&from_chat(&talk), 1);
365 let mut b = followup_of(&a, 2);
366 a.followup.as_mut().unwrap().origin_task = Some(b.id.clone());
367 b.followup.as_mut().unwrap().origin_task = Some(a.id.clone());
368 let q = q_for(&a);
369 assert!(origin_talk(&[a, b], &[talk], &q).is_none());
370 }
371
372 #[test]
373 fn pending_consults_follow_the_store_not_the_turn_text() {
374 let (tmp, store, talk) = talks();
375 let questions = Questions::at(tmp.path().join("questions"));
376 assert!(!pending_consults(&questions, &talk.id), "empty store");
377
378 let mut plain = question("implement");
379 questions.put(&mut plain).unwrap();
380 assert!(!pending_consults(&questions, &talk.id), "no consult record");
381
382 let mut q = question("implement");
383 questions.put(&mut q).unwrap();
384 begin(&questions, &store, &q, &talk).unwrap();
385 assert!(pending_consults(&questions, &talk.id));
386 assert!(!pending_consults(&questions, "other-talk"));
387
388 questions
389 .update(&q.id, |q| {
390 q.abandon("test");
391 Ok(())
392 })
393 .unwrap();
394 assert!(!pending_consults(&questions, &talk.id), "closed question");
395 }
396
397 #[test]
398 fn begin_queues_once_and_leaves_the_question_open() {
399 let (tmp, store, talk) = talks();
400 let questions = Questions::at(tmp.path().join("questions"));
401 let mut q = question("implement");
402 questions.put(&mut q).unwrap();
403
404 assert!(begin(&questions, &store, &q, &talk).unwrap());
405 assert!(!begin(&questions, &store, &q, &talk).unwrap(), "idempotent");
406
407 let after = questions.get(&q.id).unwrap();
408 assert!(after.status.open());
409 assert!(after.thread.is_empty());
410 assert_eq!(after.choices, q.choices, "choices are not touched");
411 assert_eq!(
412 after.consult.as_ref().map(|c| c.talk.as_str()),
413 Some(talk.id.as_str())
414 );
415
416 let queued = store.get(&talk.id).unwrap().pending;
417 assert_eq!(
418 queued.matches(&q.id).count(),
419 2 + 1,
420 "id once per use: {queued}"
421 );
422 assert!(queued.contains("Which backend?"));
423 assert!(queued.contains("SQLite is simpler."));
424 assert!(queued.contains("- Redis"));
425 assert_eq!(
427 queued.matches(crate::prompt::CHAT_CONSULT_HEADING).count(),
428 1
429 );
430 }
431
432 #[test]
433 fn begin_withdraws_its_record_when_the_talk_cannot_take_it() {
434 let (tmp, store, talk) = talks();
435 let questions = Questions::at(tmp.path().join("questions"));
436 let mut q = question("implement");
437 questions.put(&mut q).unwrap();
438 std::fs::remove_dir_all(tmp.path().join("talks")).unwrap();
439 assert!(begin(&questions, &store, &q, &talk).is_err());
440 assert!(questions.get(&q.id).unwrap().consult.is_none());
441 }
442
443 #[test]
444 fn an_older_question_file_reads_without_a_consult() {
445 let mut q = question("implement");
446 q.schema = 5;
447 let mut v = serde_json::to_value(&q).unwrap();
448 v.as_object_mut().unwrap().remove("consult");
449 let back: Question = serde_json::from_value(v).unwrap();
450 assert!(back.consult.is_none());
451 }
452
453 #[tokio::test]
454 async fn run_turn_leaves_the_draft_when_the_lease_stays_held() {
455 let (_tmp, store, talk) = talks();
456 let mut t = store.get(&talk.id).expect("talk");
457 talk::queue(&mut t, &store, "consult", Vec::new()).expect("queue");
458 let _held = store.claim_turn(&talk.id).expect("claim").expect("free");
459 let cfg = Config::default();
460 let out = run_turn(
461 &store,
462 &cfg,
463 &talk.id,
464 std::time::Duration::from_millis(50),
465 std::time::Duration::from_millis(10),
466 )
467 .await
468 .expect("run");
469 assert_eq!(out, Handled::Busy);
470 assert!(!store.get(&talk.id).expect("talk").pending.is_empty());
471 }
472
473 #[tokio::test]
474 async fn run_turn_with_nothing_queued_is_idle() {
475 let (_tmp, store, talk) = talks();
476 let out = run_turn(
477 &store,
478 &Config::default(),
479 &talk.id,
480 std::time::Duration::ZERO,
481 std::time::Duration::from_millis(10),
482 )
483 .await
484 .expect("run");
485 assert_eq!(out, Handled::Idle);
486 assert!(store.claim_turn(&talk.id).expect("claim").is_some());
487 }
488}