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, node } = &task.source else {
42 return None;
43 };
44 if node != CHAT_NODE {
45 return None;
46 }
47 talks
48 .iter()
49 .find(|t| &t.id == run && t.status.open())
50 .cloned()
51}
52
53pub fn pending_consults(questions: &Questions, talk_id: &str) -> bool {
60 questions
61 .list()
62 .iter()
63 .any(|q| q.status.open() && q.consult.as_ref().is_some_and(|c| c.talk == talk_id))
64}
65
66pub fn begin(questions: &Questions, talks: &Talks, q: &Question, talk: &Talk) -> Result<bool> {
72 let (q, fresh) = questions.update(&q.id, |r| {
73 if !r.status.open() {
74 bail!("question {} is already {}", r.short(), r.status.as_str());
75 }
76 if r.consult.is_some() {
77 return Ok(false);
78 }
79 r.consult = Some(ChatConsult {
80 talk: talk.id.clone(),
81 at: Timestamp::now(),
82 });
83 Ok(true)
84 })?;
85 if !fresh {
86 return Ok(false);
87 }
88 let mut talk = talk.clone();
89 if let Err(e) = talk::queue(
90 &mut talk,
91 talks,
92 &crate::prompt::chat_consult(&q),
93 Vec::new(),
94 ) {
95 let _ = questions.update(&q.id, |r| {
96 r.consult = None;
97 Ok(())
98 });
99 return Err(e).context("queue the question into the chat");
100 }
101 Ok(true)
102}
103
104#[derive(Debug, PartialEq, Eq)]
106pub enum Handled {
107 Ran(usize),
109 Idle,
111 Busy,
113}
114
115pub async fn run_turn(
123 talks: &Talks,
124 cfg: &crate::config::Config,
125 talk_id: &str,
126 wait: std::time::Duration,
127 poll: std::time::Duration,
128) -> Result<Handled> {
129 let deadline = std::time::Instant::now() + wait;
130 let lease = loop {
131 if let Some(lease) = talks.claim_turn(talk_id)? {
132 break lease;
133 }
134 if std::time::Instant::now() >= deadline {
135 return Ok(Handled::Busy);
136 }
137 tokio::time::sleep(poll).await;
138 };
139 let mut talk = talks.get(talk_id)?;
142 let mut ran = 0;
143 loop {
144 if !lease.beat()? {
145 bail!("the turn lease was taken over; the remaining drafts stay queued");
146 }
147 let Some(text) = talk::drain(&mut talk, talks)? else {
148 break;
149 };
150 lease
151 .beating(talk::respond(&mut talk, talks, cfg, &text))
152 .await??;
153 ran += 1;
154 }
155 Ok(if ran == 0 {
156 Handled::Idle
157 } else {
158 Handled::Ran(ran)
159 })
160}
161
162#[cfg(test)]
163mod tests {
164 use std::collections::BTreeMap;
165 use std::path::PathBuf;
166
167 use super::*;
168 use crate::config::{AgentKind, AgentSpec, Config};
169 use crate::queue::Source;
170 use crate::talk::TalkStatus;
171
172 const RUN: &str = "20260902-000000-beef";
173
174 fn talks() -> (tempfile::TempDir, Talks, Talk) {
175 let tmp = tempfile::tempdir().expect("tempdir");
176 let store = Talks::at(tmp.path().join("talks"));
177 let cfg = Config {
178 agents: vec![AgentSpec {
179 id: "mock".to_owned(),
180 kind: AgentKind::Command,
181 model: None,
182 command: vec!["true".to_owned()],
183 extra_args: Vec::new(),
184 env: BTreeMap::new(),
185 prompt_delivery: None,
186 }],
187 ..Config::default()
188 };
189 let talk = talk::begin(&store, &cfg, tmp.path().to_path_buf(), Some("mock")).expect("talk");
190 (tmp, store, talk)
191 }
192
193 fn task(source: Source) -> Task {
194 let mut t = Task::new(
195 "t".to_owned(),
196 "Do it".to_owned(),
197 PathBuf::from("/repo"),
198 source,
199 );
200 t.start(RUN.to_owned());
201 t
202 }
203
204 fn from_chat(talk: &Talk) -> Task {
205 task(Source::Agent {
206 run: talk.id.clone(),
207 node: CHAT_NODE.to_owned(),
208 })
209 }
210
211 fn question(node: &str) -> Question {
212 Question::new(
213 RUN.to_owned(),
214 node.to_owned(),
215 "impl-A".to_owned(),
216 "Which backend?".to_owned(),
217 "SQLite is simpler.".to_owned(),
218 vec!["SQLite".to_owned(), "Redis".to_owned()],
219 )
220 }
221
222 #[test]
223 fn origin_talk_is_decided_in_one_table() {
224 let (_tmp, _store, talk) = talks();
225 let mut closed = talk.clone();
226 closed.status = TalkStatus::Closed;
227 let chat = from_chat(&talk);
228 let q = question("implement");
229
230 let hit = |tasks: &[Task], talks: &[Talk], q: &Question| origin_talk(tasks, talks, q);
231 assert_eq!(
232 hit(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).map(|t| t.id),
233 Some(talk.id.clone())
234 );
235 assert!(hit(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
237 let other = task(Source::Agent {
238 run: talk.id.clone(),
239 node: "implement".to_owned(),
240 });
241 assert!(hit(&[other], std::slice::from_ref(&talk), &q).is_none());
242 assert!(hit(std::slice::from_ref(&chat), &[], &q).is_none());
244 assert!(hit(std::slice::from_ref(&chat), &[closed], &q).is_none());
245 assert!(hit(&[], std::slice::from_ref(&talk), &q).is_none());
247 for node in [crate::land::APPROVAL_NODE, crate::bump::NOTICE_NODE] {
249 assert!(
250 hit(
251 std::slice::from_ref(&chat),
252 std::slice::from_ref(&talk),
253 &question(node)
254 )
255 .is_none()
256 );
257 }
258 let mut answered = question("implement");
260 answered
261 .answer(crate::ask::Answer::Choice("Redis".to_owned()))
262 .unwrap();
263 assert!(hit(&[chat], &[talk], &answered).is_none());
264 }
265
266 #[test]
267 fn a_conductor_question_is_found_through_the_task_id() {
268 let (_tmp, _store, talk) = talks();
269 let chat = from_chat(&talk);
270 let mut q = question(crate::conduct::NODE);
271 q.run = chat.id.clone();
272 assert!(origin_talk(&[chat], &[talk], &q).is_some());
273 }
274
275 #[test]
276 fn pending_consults_follow_the_store_not_the_turn_text() {
277 let (tmp, store, talk) = talks();
278 let questions = Questions::at(tmp.path().join("questions"));
279 assert!(!pending_consults(&questions, &talk.id), "empty store");
280
281 let mut plain = question("implement");
282 questions.put(&mut plain).unwrap();
283 assert!(!pending_consults(&questions, &talk.id), "no consult record");
284
285 let mut q = question("implement");
286 questions.put(&mut q).unwrap();
287 begin(&questions, &store, &q, &talk).unwrap();
288 assert!(pending_consults(&questions, &talk.id));
289 assert!(!pending_consults(&questions, "other-talk"));
290
291 questions
292 .update(&q.id, |q| {
293 q.abandon("test");
294 Ok(())
295 })
296 .unwrap();
297 assert!(!pending_consults(&questions, &talk.id), "closed question");
298 }
299
300 #[test]
301 fn begin_queues_once_and_leaves_the_question_open() {
302 let (tmp, store, talk) = talks();
303 let questions = Questions::at(tmp.path().join("questions"));
304 let mut q = question("implement");
305 questions.put(&mut q).unwrap();
306
307 assert!(begin(&questions, &store, &q, &talk).unwrap());
308 assert!(!begin(&questions, &store, &q, &talk).unwrap(), "idempotent");
309
310 let after = questions.get(&q.id).unwrap();
311 assert!(after.status.open());
312 assert!(after.thread.is_empty());
313 assert_eq!(after.choices, q.choices, "choices are not touched");
314 assert_eq!(
315 after.consult.as_ref().map(|c| c.talk.as_str()),
316 Some(talk.id.as_str())
317 );
318
319 let queued = store.get(&talk.id).unwrap().pending;
320 assert_eq!(
321 queued.matches(&q.id).count(),
322 2 + 1,
323 "id once per use: {queued}"
324 );
325 assert!(queued.contains("Which backend?"));
326 assert!(queued.contains("SQLite is simpler."));
327 assert!(queued.contains("- Redis"));
328 assert_eq!(
330 queued.matches(crate::prompt::CHAT_CONSULT_HEADING).count(),
331 1
332 );
333 }
334
335 #[test]
336 fn begin_withdraws_its_record_when_the_talk_cannot_take_it() {
337 let (tmp, store, talk) = talks();
338 let questions = Questions::at(tmp.path().join("questions"));
339 let mut q = question("implement");
340 questions.put(&mut q).unwrap();
341 std::fs::remove_dir_all(tmp.path().join("talks")).unwrap();
342 assert!(begin(&questions, &store, &q, &talk).is_err());
343 assert!(questions.get(&q.id).unwrap().consult.is_none());
344 }
345
346 #[test]
347 fn an_older_question_file_reads_without_a_consult() {
348 let mut q = question("implement");
349 q.schema = 5;
350 let mut v = serde_json::to_value(&q).unwrap();
351 v.as_object_mut().unwrap().remove("consult");
352 let back: Question = serde_json::from_value(v).unwrap();
353 assert!(back.consult.is_none());
354 }
355
356 #[tokio::test]
357 async fn run_turn_leaves_the_draft_when_the_lease_stays_held() {
358 let (_tmp, store, talk) = talks();
359 let mut t = store.get(&talk.id).expect("talk");
360 talk::queue(&mut t, &store, "consult", Vec::new()).expect("queue");
361 let _held = store.claim_turn(&talk.id).expect("claim").expect("free");
362 let cfg = Config::default();
363 let out = run_turn(
364 &store,
365 &cfg,
366 &talk.id,
367 std::time::Duration::from_millis(50),
368 std::time::Duration::from_millis(10),
369 )
370 .await
371 .expect("run");
372 assert_eq!(out, Handled::Busy);
373 assert!(!store.get(&talk.id).expect("talk").pending.is_empty());
374 }
375
376 #[tokio::test]
377 async fn run_turn_with_nothing_queued_is_idle() {
378 let (_tmp, store, talk) = talks();
379 let out = run_turn(
380 &store,
381 &Config::default(),
382 &talk.id,
383 std::time::Duration::ZERO,
384 std::time::Duration::from_millis(10),
385 )
386 .await
387 .expect("run");
388 assert_eq!(out, Handled::Idle);
389 assert!(store.claim_turn(&talk.id).expect("claim").is_some());
390 }
391}