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#[cfg(test)]
105mod tests {
106 use std::collections::BTreeMap;
107 use std::path::PathBuf;
108
109 use super::*;
110 use crate::config::{AgentKind, AgentSpec, Config};
111 use crate::queue::Source;
112 use crate::talk::TalkStatus;
113
114 const RUN: &str = "20260902-000000-beef";
115
116 fn talks() -> (tempfile::TempDir, Talks, Talk) {
117 let tmp = tempfile::tempdir().expect("tempdir");
118 let store = Talks::at(tmp.path().join("talks"));
119 let cfg = Config {
120 agents: vec![AgentSpec {
121 id: "mock".to_owned(),
122 kind: AgentKind::Command,
123 model: None,
124 command: vec!["true".to_owned()],
125 extra_args: Vec::new(),
126 env: BTreeMap::new(),
127 prompt_delivery: None,
128 }],
129 ..Config::default()
130 };
131 let talk = talk::begin(&store, &cfg, tmp.path().to_path_buf(), Some("mock")).expect("talk");
132 (tmp, store, talk)
133 }
134
135 fn task(source: Source) -> Task {
136 let mut t = Task::new(
137 "t".to_owned(),
138 "Do it".to_owned(),
139 PathBuf::from("/repo"),
140 source,
141 );
142 t.start(RUN.to_owned());
143 t
144 }
145
146 fn from_chat(talk: &Talk) -> Task {
147 task(Source::Agent {
148 run: talk.id.clone(),
149 node: CHAT_NODE.to_owned(),
150 })
151 }
152
153 fn question(node: &str) -> Question {
154 Question::new(
155 RUN.to_owned(),
156 node.to_owned(),
157 "impl-A".to_owned(),
158 "Which backend?".to_owned(),
159 "SQLite is simpler.".to_owned(),
160 vec!["SQLite".to_owned(), "Redis".to_owned()],
161 )
162 }
163
164 #[test]
165 fn origin_talk_is_decided_in_one_table() {
166 let (_tmp, _store, talk) = talks();
167 let mut closed = talk.clone();
168 closed.status = TalkStatus::Closed;
169 let chat = from_chat(&talk);
170 let q = question("implement");
171
172 let hit = |tasks: &[Task], talks: &[Talk], q: &Question| origin_talk(tasks, talks, q);
173 assert_eq!(
174 hit(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).map(|t| t.id),
175 Some(talk.id.clone())
176 );
177 assert!(hit(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
179 let other = task(Source::Agent {
180 run: talk.id.clone(),
181 node: "implement".to_owned(),
182 });
183 assert!(hit(&[other], std::slice::from_ref(&talk), &q).is_none());
184 assert!(hit(std::slice::from_ref(&chat), &[], &q).is_none());
186 assert!(hit(std::slice::from_ref(&chat), &[closed], &q).is_none());
187 assert!(hit(&[], std::slice::from_ref(&talk), &q).is_none());
189 for node in [crate::land::APPROVAL_NODE, crate::bump::NOTICE_NODE] {
191 assert!(
192 hit(
193 std::slice::from_ref(&chat),
194 std::slice::from_ref(&talk),
195 &question(node)
196 )
197 .is_none()
198 );
199 }
200 let mut answered = question("implement");
202 answered
203 .answer(crate::ask::Answer::Choice("Redis".to_owned()))
204 .unwrap();
205 assert!(hit(&[chat], &[talk], &answered).is_none());
206 }
207
208 #[test]
209 fn a_conductor_question_is_found_through_the_task_id() {
210 let (_tmp, _store, talk) = talks();
211 let chat = from_chat(&talk);
212 let mut q = question(crate::conduct::NODE);
213 q.run = chat.id.clone();
214 assert!(origin_talk(&[chat], &[talk], &q).is_some());
215 }
216
217 #[test]
218 fn pending_consults_follow_the_store_not_the_turn_text() {
219 let (tmp, store, talk) = talks();
220 let questions = Questions::at(tmp.path().join("questions"));
221 assert!(!pending_consults(&questions, &talk.id), "empty store");
222
223 let mut plain = question("implement");
224 questions.put(&mut plain).unwrap();
225 assert!(!pending_consults(&questions, &talk.id), "no consult record");
226
227 let mut q = question("implement");
228 questions.put(&mut q).unwrap();
229 begin(&questions, &store, &q, &talk).unwrap();
230 assert!(pending_consults(&questions, &talk.id));
231 assert!(!pending_consults(&questions, "other-talk"));
232
233 questions
234 .update(&q.id, |q| {
235 q.abandon("test");
236 Ok(())
237 })
238 .unwrap();
239 assert!(!pending_consults(&questions, &talk.id), "closed question");
240 }
241
242 #[test]
243 fn begin_queues_once_and_leaves_the_question_open() {
244 let (tmp, store, talk) = talks();
245 let questions = Questions::at(tmp.path().join("questions"));
246 let mut q = question("implement");
247 questions.put(&mut q).unwrap();
248
249 assert!(begin(&questions, &store, &q, &talk).unwrap());
250 assert!(!begin(&questions, &store, &q, &talk).unwrap(), "idempotent");
251
252 let after = questions.get(&q.id).unwrap();
253 assert!(after.status.open());
254 assert!(after.thread.is_empty());
255 assert_eq!(after.choices, q.choices, "choices are not touched");
256 assert_eq!(
257 after.consult.as_ref().map(|c| c.talk.as_str()),
258 Some(talk.id.as_str())
259 );
260
261 let queued = store.get(&talk.id).unwrap().pending;
262 assert_eq!(
263 queued.matches(&q.id).count(),
264 2 + 1,
265 "id once per use: {queued}"
266 );
267 assert!(queued.contains("Which backend?"));
268 assert!(queued.contains("SQLite is simpler."));
269 assert!(queued.contains("- Redis"));
270 assert_eq!(
272 queued.matches(crate::prompt::CHAT_CONSULT_HEADING).count(),
273 1
274 );
275 }
276
277 #[test]
278 fn begin_withdraws_its_record_when_the_talk_cannot_take_it() {
279 let (tmp, store, talk) = talks();
280 let questions = Questions::at(tmp.path().join("questions"));
281 let mut q = question("implement");
282 questions.put(&mut q).unwrap();
283 std::fs::remove_dir_all(tmp.path().join("talks")).unwrap();
284 assert!(begin(&questions, &store, &q, &talk).is_err());
285 assert!(questions.get(&q.id).unwrap().consult.is_none());
286 }
287
288 #[test]
289 fn an_older_question_file_reads_without_a_consult() {
290 let mut q = question("implement");
291 q.schema = 5;
292 let mut v = serde_json::to_value(&q).unwrap();
293 v.as_object_mut().unwrap().remove("consult");
294 let back: Question = serde_json::from_value(v).unwrap();
295 assert!(back.consult.is_none());
296 }
297}