use anyhow::{Context, Result, bail};
use jiff::Timestamp;
use crate::ask::{ChatConsult, Question, Questions};
use crate::queue::{CHAT_NODE, Source, Task};
use crate::talk::{self, Talk, Talks};
pub fn origin_talk(tasks: &[Task], talks: &[Talk], q: &Question) -> Option<Talk> {
if !q.status.open()
|| q.node == crate::land::APPROVAL_NODE
|| q.node == crate::bump::NOTICE_NODE
{
return None;
}
let task = crate::daemon::task_of_question(tasks, q)?;
let Source::Agent { run, node } = &task.source else {
return None;
};
if node != CHAT_NODE {
return None;
}
talks
.iter()
.find(|t| &t.id == run && t.status.open())
.cloned()
}
pub fn pending_consults(questions: &Questions, talk_id: &str) -> bool {
questions
.list()
.iter()
.any(|q| q.status.open() && q.consult.as_ref().is_some_and(|c| c.talk == talk_id))
}
pub fn begin(questions: &Questions, talks: &Talks, q: &Question, talk: &Talk) -> Result<bool> {
let (q, fresh) = questions.update(&q.id, |r| {
if !r.status.open() {
bail!("question {} is already {}", r.short(), r.status.as_str());
}
if r.consult.is_some() {
return Ok(false);
}
r.consult = Some(ChatConsult {
talk: talk.id.clone(),
at: Timestamp::now(),
});
Ok(true)
})?;
if !fresh {
return Ok(false);
}
let mut talk = talk.clone();
if let Err(e) = talk::queue(
&mut talk,
talks,
&crate::prompt::chat_consult(&q),
Vec::new(),
) {
let _ = questions.update(&q.id, |r| {
r.consult = None;
Ok(())
});
return Err(e).context("queue the question into the chat");
}
Ok(true)
}
#[cfg(test)]
mod tests {
use std::collections::BTreeMap;
use std::path::PathBuf;
use super::*;
use crate::config::{AgentKind, AgentSpec, Config};
use crate::queue::Source;
use crate::talk::TalkStatus;
const RUN: &str = "20260902-000000-beef";
fn talks() -> (tempfile::TempDir, Talks, Talk) {
let tmp = tempfile::tempdir().expect("tempdir");
let store = Talks::at(tmp.path().join("talks"));
let cfg = Config {
agents: vec![AgentSpec {
id: "mock".to_owned(),
kind: AgentKind::Command,
model: None,
command: vec!["true".to_owned()],
extra_args: Vec::new(),
env: BTreeMap::new(),
prompt_delivery: None,
}],
..Config::default()
};
let talk = talk::begin(&store, &cfg, tmp.path().to_path_buf(), Some("mock")).expect("talk");
(tmp, store, talk)
}
fn task(source: Source) -> Task {
let mut t = Task::new(
"t".to_owned(),
"Do it".to_owned(),
PathBuf::from("/repo"),
source,
);
t.start(RUN.to_owned());
t
}
fn from_chat(talk: &Talk) -> Task {
task(Source::Agent {
run: talk.id.clone(),
node: CHAT_NODE.to_owned(),
})
}
fn question(node: &str) -> Question {
Question::new(
RUN.to_owned(),
node.to_owned(),
"impl-A".to_owned(),
"Which backend?".to_owned(),
"SQLite is simpler.".to_owned(),
vec!["SQLite".to_owned(), "Redis".to_owned()],
)
}
#[test]
fn origin_talk_is_decided_in_one_table() {
let (_tmp, _store, talk) = talks();
let mut closed = talk.clone();
closed.status = TalkStatus::Closed;
let chat = from_chat(&talk);
let q = question("implement");
let hit = |tasks: &[Task], talks: &[Talk], q: &Question| origin_talk(tasks, talks, q);
assert_eq!(
hit(std::slice::from_ref(&chat), std::slice::from_ref(&talk), &q).map(|t| t.id),
Some(talk.id.clone())
);
assert!(hit(&[task(Source::Human)], std::slice::from_ref(&talk), &q).is_none());
let other = task(Source::Agent {
run: talk.id.clone(),
node: "implement".to_owned(),
});
assert!(hit(&[other], std::slice::from_ref(&talk), &q).is_none());
assert!(hit(std::slice::from_ref(&chat), &[], &q).is_none());
assert!(hit(std::slice::from_ref(&chat), &[closed], &q).is_none());
assert!(hit(&[], std::slice::from_ref(&talk), &q).is_none());
for node in [crate::land::APPROVAL_NODE, crate::bump::NOTICE_NODE] {
assert!(
hit(
std::slice::from_ref(&chat),
std::slice::from_ref(&talk),
&question(node)
)
.is_none()
);
}
let mut answered = question("implement");
answered
.answer(crate::ask::Answer::Choice("Redis".to_owned()))
.unwrap();
assert!(hit(&[chat], &[talk], &answered).is_none());
}
#[test]
fn a_conductor_question_is_found_through_the_task_id() {
let (_tmp, _store, talk) = talks();
let chat = from_chat(&talk);
let mut q = question(crate::conduct::NODE);
q.run = chat.id.clone();
assert!(origin_talk(&[chat], &[talk], &q).is_some());
}
#[test]
fn pending_consults_follow_the_store_not_the_turn_text() {
let (tmp, store, talk) = talks();
let questions = Questions::at(tmp.path().join("questions"));
assert!(!pending_consults(&questions, &talk.id), "empty store");
let mut plain = question("implement");
questions.put(&mut plain).unwrap();
assert!(!pending_consults(&questions, &talk.id), "no consult record");
let mut q = question("implement");
questions.put(&mut q).unwrap();
begin(&questions, &store, &q, &talk).unwrap();
assert!(pending_consults(&questions, &talk.id));
assert!(!pending_consults(&questions, "other-talk"));
questions
.update(&q.id, |q| {
q.abandon("test");
Ok(())
})
.unwrap();
assert!(!pending_consults(&questions, &talk.id), "closed question");
}
#[test]
fn begin_queues_once_and_leaves_the_question_open() {
let (tmp, store, talk) = talks();
let questions = Questions::at(tmp.path().join("questions"));
let mut q = question("implement");
questions.put(&mut q).unwrap();
assert!(begin(&questions, &store, &q, &talk).unwrap());
assert!(!begin(&questions, &store, &q, &talk).unwrap(), "idempotent");
let after = questions.get(&q.id).unwrap();
assert!(after.status.open());
assert!(after.thread.is_empty());
assert_eq!(after.choices, q.choices, "choices are not touched");
assert_eq!(
after.consult.as_ref().map(|c| c.talk.as_str()),
Some(talk.id.as_str())
);
let queued = store.get(&talk.id).unwrap().pending;
assert_eq!(
queued.matches(&q.id).count(),
2 + 1,
"id once per use: {queued}"
);
assert!(queued.contains("Which backend?"));
assert!(queued.contains("SQLite is simpler."));
assert!(queued.contains("- Redis"));
assert_eq!(
queued.matches(crate::prompt::CHAT_CONSULT_HEADING).count(),
1
);
}
#[test]
fn begin_withdraws_its_record_when_the_talk_cannot_take_it() {
let (tmp, store, talk) = talks();
let questions = Questions::at(tmp.path().join("questions"));
let mut q = question("implement");
questions.put(&mut q).unwrap();
std::fs::remove_dir_all(tmp.path().join("talks")).unwrap();
assert!(begin(&questions, &store, &q, &talk).is_err());
assert!(questions.get(&q.id).unwrap().consult.is_none());
}
#[test]
fn an_older_question_file_reads_without_a_consult() {
let mut q = question("implement");
q.schema = 5;
let mut v = serde_json::to_value(&q).unwrap();
v.as_object_mut().unwrap().remove("consult");
let back: Question = serde_json::from_value(v).unwrap();
assert!(back.consult.is_none());
}
}