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 #[test]
277 fn origin_talk_is_decided_in_one_table() {
278 let (_tmp, _store, talk) = talks();
279 let mut closed = talk.clone();
280 closed.status = TalkStatus::Closed;
281 let chat = from_chat(&talk);
282 let q = question("implement");
283
284 let hit = |tasks: &[Task], talks: &[Talk], q: &Question| origin_talk(tasks, talks, q);
285 assert_eq!(
286 hit(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).map(|t| t.id),
287 Some(talk.id.clone())
288 );
289 assert!(hit(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
291 let other = task(Source::Agent {
292 run: talk.id.clone(),
293 node: "implement".to_owned(),
294 });
295 assert!(hit(&[other], std::slice::from_ref(&talk), &q).is_none());
296 assert!(hit(std::slice::from_ref(&chat), &[], &q).is_none());
298 assert!(hit(std::slice::from_ref(&chat), &[closed], &q).is_none());
299 assert!(hit(&[], std::slice::from_ref(&talk), &q).is_none());
301 for node in [crate::land::APPROVAL_NODE, crate::bump::NOTICE_NODE] {
303 assert!(
304 hit(
305 std::slice::from_ref(&chat),
306 std::slice::from_ref(&talk),
307 &question(node)
308 )
309 .is_none()
310 );
311 }
312 let mut answered = question("implement");
314 answered
315 .answer(crate::ask::Answer::Choice("Redis".to_owned()))
316 .unwrap();
317 assert!(hit(&[chat], &[talk], &answered).is_none());
318 }
319
320 #[test]
321 fn a_conductor_question_is_found_through_the_task_id() {
322 let (_tmp, _store, talk) = talks();
323 let chat = from_chat(&talk);
324 let mut q = question(crate::conduct::NODE);
325 q.run = chat.id.clone();
326 assert!(origin_talk(&[chat], &[talk], &q).is_some());
327 }
328
329 fn followup_of(parent: &Task, n: u32) -> Task {
330 let mut t = Task::new(
331 "f".to_owned(),
332 "Fix".to_owned(),
333 PathBuf::from("/repo"),
334 Source::Agent {
335 run: format!("merged-{n}"),
336 node: "followup".to_owned(),
337 },
338 );
339 t.start(format!("run-f{n}"));
340 t.followup = Some(crate::queue::FollowUp {
341 run: format!("merged-{n}"),
342 origin_task: Some(parent.id.clone()),
343 pr: "https://example.invalid/pr/1".to_owned(),
344 findings: Vec::new(),
345 generation: n,
346 });
347 t
348 }
349
350 fn q_for(t: &Task) -> Question {
351 let mut q = question("implement");
352 q.run = t.runs[0].clone();
353 q
354 }
355
356 #[test]
357 fn a_followup_traces_back_to_the_chat() {
358 let (_tmp, _store, talk) = talks();
359 let mut chat = from_chat(&talk);
360 chat.runs = vec!["chat-run".to_owned()];
361 let f1 = followup_of(&chat, 1);
362 let f2 = followup_of(&f1, 2);
363 let ts = [chat, f1.clone(), f2.clone()];
364 let id = Some(talk.id.clone());
365 let tk = std::slice::from_ref(&talk);
366 assert_eq!(origin_talk(&ts, tk, &q_for(&f1)).map(|t| t.id), id);
367 assert_eq!(origin_talk(&ts, tk, &q_for(&f2)).map(|t| t.id), id);
368 }
369
370 #[test]
371 fn a_followup_without_origin_task_is_found_through_runs() {
372 let (_tmp, _store, talk) = talks();
373 let mut chat = from_chat(&talk);
374 chat.runs = vec!["merged-1".to_owned()];
375 let mut f1 = followup_of(&chat, 1);
376 f1.followup.as_mut().unwrap().origin_task = None;
377 let q = q_for(&f1);
378 assert!(origin_talk(&[chat, f1], &[talk], &q).is_some());
379 }
380
381 #[test]
382 fn a_dangling_origin_task_has_no_chat() {
383 let (_tmp, _store, talk) = talks();
384 let mut chat = from_chat(&talk);
385 chat.runs = vec!["chat-run".to_owned()];
386 let mut f1 = followup_of(&chat, 1);
387 f1.followup.as_mut().unwrap().origin_task = Some("gone".to_owned());
388 let q = q_for(&f1);
389 assert!(origin_talk(&[chat, f1], &[talk], &q).is_none());
390 }
391
392 #[test]
393 fn a_followup_cycle_ends_without_a_chat() {
394 let (_tmp, _store, talk) = talks();
395 let mut a = followup_of(&from_chat(&talk), 1);
396 let mut b = followup_of(&a, 2);
397 a.followup.as_mut().unwrap().origin_task = Some(b.id.clone());
398 b.followup.as_mut().unwrap().origin_task = Some(a.id.clone());
399 let q = q_for(&a);
400 assert!(origin_talk(&[a, b], &[talk], &q).is_none());
401 }
402
403 #[test]
404 fn pending_consults_follow_the_store_not_the_turn_text() {
405 let (tmp, store, talk) = talks();
406 let questions = Questions::at(tmp.path().join("questions"));
407 assert!(!pending_consults(&questions, &talk.id), "empty store");
408
409 let mut plain = question("implement");
410 questions.put(&mut plain).unwrap();
411 assert!(!pending_consults(&questions, &talk.id), "no consult record");
412
413 let mut q = question("implement");
414 questions.put(&mut q).unwrap();
415 begin(&questions, &store, &q, &talk).unwrap();
416 assert!(pending_consults(&questions, &talk.id));
417 assert!(!pending_consults(&questions, "other-talk"));
418
419 questions
420 .update(&q.id, |q| {
421 q.abandon("test");
422 Ok(())
423 })
424 .unwrap();
425 assert!(!pending_consults(&questions, &talk.id), "closed question");
426 }
427
428 #[test]
429 fn begin_queues_once_and_leaves_the_question_open() {
430 let (tmp, store, talk) = talks();
431 let questions = Questions::at(tmp.path().join("questions"));
432 let mut q = question("implement");
433 questions.put(&mut q).unwrap();
434
435 assert!(begin(&questions, &store, &q, &talk).unwrap());
436 assert!(!begin(&questions, &store, &q, &talk).unwrap(), "idempotent");
437
438 let after = questions.get(&q.id).unwrap();
439 assert!(after.status.open());
440 assert!(after.thread.is_empty());
441 assert_eq!(after.choices, q.choices, "choices are not touched");
442 assert_eq!(
443 after.consult.as_ref().map(|c| c.talk.as_str()),
444 Some(talk.id.as_str())
445 );
446
447 let queued = store.get(&talk.id).unwrap().pending;
448 assert_eq!(
449 queued.matches(&q.id).count(),
450 2 + 1,
451 "id once per use: {queued}"
452 );
453 assert!(queued.contains("Which backend?"));
454 assert!(queued.contains("SQLite is simpler."));
455 assert!(queued.contains("- Redis"));
456 assert_eq!(
458 queued.matches(crate::prompt::CHAT_CONSULT_HEADING).count(),
459 1
460 );
461 }
462
463 #[test]
464 fn begin_withdraws_its_record_when_the_talk_cannot_take_it() {
465 let (tmp, store, talk) = talks();
466 let questions = Questions::at(tmp.path().join("questions"));
467 let mut q = question("implement");
468 questions.put(&mut q).unwrap();
469 std::fs::remove_dir_all(tmp.path().join("talks")).unwrap();
470 assert!(begin(&questions, &store, &q, &talk).is_err());
471 assert!(questions.get(&q.id).unwrap().consult.is_none());
472 }
473
474 #[test]
475 fn an_older_question_file_reads_without_a_consult() {
476 let mut q = question("implement");
477 q.schema = 5;
478 let mut v = serde_json::to_value(&q).unwrap();
479 v.as_object_mut().unwrap().remove("consult");
480 let back: Question = serde_json::from_value(v).unwrap();
481 assert!(back.consult.is_none());
482 }
483}