1use std::path::{Path, PathBuf};
47use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
48use std::time::Duration;
49
50use anyhow::{Context, Result, bail};
51use jiff::Timestamp;
52use serde::{Deserialize, Serialize};
53
54use crate::agent::{self, Invocation, SeatState};
55use crate::config::{AgentSpec, Config};
56use crate::queue::{Queue, Source, Task};
57
58pub const SCHEMA: u32 = 1;
60
61fn turn_timeout(cfg: &Config) -> Duration {
72 Duration::from_secs(cfg.graph.timeout_talk)
73}
74
75const SEAT: &str = "talk";
78
79const MAGI_NOTE: &str = "magi: ";
81
82#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
84#[serde(rename_all = "lowercase")]
85pub enum Who {
86 Operator,
88 Agent,
91}
92
93#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
101#[serde(deny_unknown_fields)]
102pub struct Attachment {
103 pub id: String,
105 pub name: String,
107 pub mime: String,
110 pub bytes: u64,
112}
113
114#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
116#[serde(deny_unknown_fields)]
117pub struct Turn {
118 pub who: Who,
120 pub body: String,
122 pub at: Timestamp,
124 #[serde(default)]
127 pub attachments: Vec<Attachment>,
128}
129
130#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
134#[serde(rename_all = "lowercase")]
135pub enum TalkStatus {
136 Open,
139 Closed,
141}
142
143impl TalkStatus {
144 pub fn open(self) -> bool {
146 matches!(self, Self::Open)
147 }
148
149 pub fn as_str(self) -> &'static str {
151 match self {
152 Self::Open => "open",
153 Self::Closed => "closed",
154 }
155 }
156}
157
158#[derive(Debug, Clone, Serialize, Deserialize)]
160#[serde(deny_unknown_fields)]
161pub struct Talk {
162 pub schema: u32,
164 pub id: String,
166 pub repo: PathBuf,
168 pub agent: String,
170 pub status: TalkStatus,
172 pub turns: Vec<Turn>,
174 #[serde(default)]
177 pub pending: String,
178 #[serde(default)]
180 pub pending_attachments: Vec<Attachment>,
181 pub created_at: Timestamp,
183 pub updated_at: Timestamp,
185 seat: SeatState,
190}
191
192impl Talk {
193 pub fn short(&self) -> &str {
195 short(&self.id)
196 }
197}
198
199#[derive(Debug, Clone)]
201pub struct Talks {
202 root: PathBuf,
203 lock: Arc<Mutex<()>>,
212}
213
214impl Talks {
215 pub fn open() -> Self {
217 Self::at(crate::run::home().join("talks"))
218 }
219
220 pub fn at(root: PathBuf) -> Self {
223 Self {
224 root,
225 lock: Arc::new(Mutex::new(())),
226 }
227 }
228
229 fn guard(&self) -> MutexGuard<'_, ()> {
238 self.lock.lock().unwrap_or_else(PoisonError::into_inner)
239 }
240
241 pub fn root(&self) -> &Path {
243 &self.root
244 }
245
246 pub fn path_of(&self, id: &str) -> PathBuf {
248 self.root.join(format!("{id}.json"))
249 }
250
251 pub fn artifacts_of(&self, id: &str) -> PathBuf {
254 self.root.join(format!("{id}.artifacts"))
255 }
256
257 pub fn attachments_dir(&self, id: &str) -> PathBuf {
261 self.artifacts_of(id).join("attachments")
262 }
263
264 pub fn put_attachment(
273 &self,
274 id: &str,
275 mime: &str,
276 name: &str,
277 data: &[u8],
278 ) -> Result<Attachment> {
279 let dir = self.attachments_dir(id);
280 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
281 let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
282 let att = Attachment {
283 id: new_attachment_id(),
284 name: name.to_owned(),
285 mime: mime.to_owned(),
286 bytes: data.len() as u64,
287 };
288 std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
289 .with_context(|| format!("write attachment {}", att.id))?;
290 std::fs::write(
291 dir.join(format!("{}.json", att.id)),
292 serde_json::to_string(&att).context("serialize attachment")?,
293 )
294 .with_context(|| format!("write attachment metadata {}", att.id))?;
295 Ok(att)
296 }
297
298 pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
307 if !valid_attachment_id(att_id) {
308 return Ok(None);
309 }
310 let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
311 if !meta_path.is_file() {
312 return Ok(None);
313 }
314 let att = serde_json::from_str(
315 &std::fs::read_to_string(&meta_path)
316 .with_context(|| format!("read {}", meta_path.display()))?,
317 )
318 .with_context(|| format!("parse {}", meta_path.display()))?;
319 Ok(Some(att))
320 }
321
322 pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
326 let Some(att) = self.attachment_meta(id, att_id)? else {
327 return Ok(None);
328 };
329 let ext = attachment_ext(&att.mime).with_context(|| {
330 format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
331 })?;
332 let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
333 let data =
334 std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
335 Ok(Some((att, data)))
336 }
337
338 fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
355 let ext = attachment_ext(&att.mime)?;
356 let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
357 std::path::absolute(&path).ok()
358 }
359
360 pub fn put(&self, t: &mut Talk) -> Result<()> {
369 std::fs::create_dir_all(&self.root)
370 .with_context(|| format!("create {}", self.root.display()))?;
371 t.updated_at = Timestamp::now();
372 let body = serde_json::to_string_pretty(t).context("serialize talk")?;
373 let path = self.path_of(&t.id);
374 let tmp = path.with_extension("json.tmp");
375 write_atomic(&tmp, &path, &body)
376 }
377
378 pub fn get(&self, id: &str) -> Result<Talk> {
380 let resolved = self.resolve_id(id)?;
381 read_path(&self.path_of(&resolved))
382 }
383
384 pub fn list(&self) -> Vec<Talk> {
387 let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
388 .into_iter()
389 .flatten()
390 .flatten()
391 .map(|e| e.path())
392 .filter(|p| p.extension().is_some_and(|x| x == "json"))
393 .filter_map(|p| read_path(&p).ok())
394 .collect();
395 all.sort_unstable_by(|a, b| {
396 let rank = |t: &Talk| u8::from(!t.status.open());
397 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
398 });
399 all
400 }
401
402 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
404 if self.path_of(prefix).is_file() {
405 return Ok(prefix.to_owned());
406 }
407 let hits: Vec<String> = self
408 .list()
409 .into_iter()
410 .map(|t| t.id)
411 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
412 .collect();
413 match hits.len() {
414 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
415 0 => bail!("no talk matches `{prefix}`"),
416 _ => bail!(
417 "`{prefix}` matches {} talks: {}",
418 hits.len(),
419 hits.join(", ")
420 ),
421 }
422 }
423
424 pub fn revision(&self) -> u64 {
427 std::fs::read_dir(&self.root)
428 .into_iter()
429 .flatten()
430 .flatten()
431 .filter_map(|e| e.metadata().ok())
432 .filter_map(|m| m.modified().ok())
433 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
434 .map(|d| d.as_millis() as u64)
435 .max()
436 .unwrap_or(0)
437 }
438
439 pub fn count_open(&self) -> usize {
441 self.list().iter().filter(|t| t.status.open()).count()
442 }
443
444 pub fn remove(&self, id: &str) -> Result<()> {
457 let _guard = self.guard();
458 let resolved = self.resolve_id(id)?;
459 let path = self.path_of(&resolved);
460 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
461 let artifacts = self.artifacts_of(&resolved);
462 if artifacts.is_dir() {
463 std::fs::remove_dir_all(&artifacts)
464 .with_context(|| format!("remove {}", artifacts.display()))?;
465 }
466 Ok(())
467 }
468}
469
470pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
480 let repo = repo.canonicalize().unwrap_or(repo);
483 let want = agent.or(cfg.roles.chatter.as_deref());
484 let spec = agent::pick(&cfg.agents, want, &agent::installed)?;
485
486 let now = Timestamp::now();
487 let mut talk = Talk {
488 schema: SCHEMA,
489 id: new_id(),
490 repo,
491 agent: spec.id.clone(),
492 status: TalkStatus::Open,
493 turns: Vec::new(),
494 pending: String::new(),
495 pending_attachments: Vec::new(),
496 created_at: now,
497 updated_at: now,
498 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
499 };
500 store.put(&mut talk)?;
501 Ok(talk)
502}
503
504pub fn record(
511 talk: &mut Talk,
512 store: &Talks,
513 text: &str,
514 attachments: Vec<Attachment>,
515) -> Result<String> {
516 let _guard = store.guard();
524 let Ok(fresh) = store.get(&talk.id) else {
529 bail!("talk {} was deleted", talk.short());
530 };
531 talk.status = fresh.status;
532 talk.pending = fresh.pending;
535 talk.pending_attachments = fresh.pending_attachments;
536 if !talk.status.open() {
537 bail!(
538 "talk {} is {} and takes no more turns",
539 talk.short(),
540 talk.status.as_str()
541 );
542 }
543 let text = text.trim();
544 if text.is_empty() && attachments.is_empty() {
545 bail!("nothing to say");
546 }
547 talk.turns.push(Turn {
548 who: Who::Operator,
549 body: text.to_owned(),
550 at: Timestamp::now(),
551 attachments,
552 });
553 store.put(talk)?;
554 Ok(text.to_owned())
555}
556
557pub fn queue(
559 talk: &mut Talk,
560 store: &Talks,
561 text: &str,
562 attachments: Vec<Attachment>,
563) -> Result<()> {
564 let text = text.trim();
565 if text.is_empty() && attachments.is_empty() {
566 bail!("nothing to say");
567 }
568 let _guard = store.guard();
569 let mut fresh = store
570 .get(&talk.id)
571 .with_context(|| format!("talk {} was deleted", talk.short()))?;
572 if !fresh.status.open() {
573 bail!(
574 "talk {} is {} and takes no more turns",
575 fresh.short(),
576 fresh.status.as_str()
577 );
578 }
579 if !text.is_empty() {
580 if fresh.pending.is_empty() {
581 fresh.pending = text.to_owned();
582 } else {
583 fresh.pending.push_str("\n\n");
584 fresh.pending.push_str(text);
585 }
586 }
587 fresh.pending_attachments.extend(attachments);
588 store.put(&mut fresh)?;
589 *talk = fresh;
590 Ok(())
591}
592
593pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
595 let _guard = store.guard();
596 let mut fresh = store
597 .get(&talk.id)
598 .with_context(|| format!("talk {} was deleted", talk.short()))?;
599 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
600 *talk = fresh;
601 return Ok(None);
602 }
603 let text = std::mem::take(&mut fresh.pending);
604 let attachments = std::mem::take(&mut fresh.pending_attachments);
605 fresh.turns.push(Turn {
606 who: Who::Operator,
607 body: text.clone(),
608 at: Timestamp::now(),
609 attachments,
610 });
611 store.put(&mut fresh)?;
612 *talk = fresh;
613 Ok(Some(text))
614}
615
616pub async fn say(
619 talk: &mut Talk,
620 store: &Talks,
621 cfg: &Config,
622 text: &str,
623 attachments: Vec<Attachment>,
624) -> Result<()> {
625 let text = record(talk, store, text, attachments)?;
626 turn(talk, store, cfg, &text).await
627}
628
629pub async fn respond(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
631 turn(talk, store, cfg, text).await
632}
633
634pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
652 let _guard = store.guard();
653 let mut fresh = store
654 .get(&talk.id)
655 .with_context(|| format!("talk {} was deleted", talk.short()))?;
656 fresh.status = TalkStatus::Closed;
657 fresh.pending.clear();
659 fresh.pending_attachments.clear();
660 store.put(&mut fresh)?;
661 *talk = fresh;
662 Ok(())
663}
664
665pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
676 let _guard = store.guard();
677 let mut fresh = store
678 .get(&talk.id)
679 .with_context(|| format!("talk {} was deleted", talk.short()))?;
680 fresh.status = TalkStatus::Open;
681 store.put(&mut fresh)?;
682 *talk = fresh;
683 Ok(())
684}
685
686pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
699 let _guard = store.guard();
700 let mut fresh = store
701 .get(&talk.id)
702 .with_context(|| format!("talk {} was deleted", talk.short()))?;
703 if fresh.agent == spec.id {
704 *talk = fresh;
705 return Ok(false);
706 }
707 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
708 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
709 fresh.turns.push(Turn {
710 who: Who::Agent,
711 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
712 at: Timestamp::now(),
713 attachments: Vec::new(),
714 });
715 store.put(&mut fresh)?;
716 *talk = fresh;
717 Ok(true)
718}
719
720pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
722 let _guard = store.guard();
723 let mut fresh = store
724 .get(&talk.id)
725 .with_context(|| format!("talk {} was deleted", talk.short()))?;
726 fresh.pending.clear();
727 fresh.pending_attachments.clear();
728 store.put(&mut fresh)?;
729 *talk = fresh;
730 Ok(())
731}
732
733pub fn clear_pending_if_matches(
735 talk: &mut Talk,
736 store: &Talks,
737 expected_text: &str,
738 expected_attachments: &[String],
739) -> Result<bool> {
740 let _guard = store.guard();
741 let mut fresh = store
742 .get(&talk.id)
743 .with_context(|| format!("talk {} was deleted", talk.short()))?;
744 if !pending_matches(&fresh, expected_text, expected_attachments) {
745 *talk = fresh;
746 return Ok(false);
747 }
748 fresh.pending.clear();
749 fresh.pending_attachments.clear();
750 store.put(&mut fresh)?;
751 *talk = fresh;
752 Ok(true)
753}
754
755pub fn edit_pending_text(
759 talk: &mut Talk,
760 store: &Talks,
761 text: &str,
762 expected_text: &str,
763 expected_attachments: &[String],
764) -> Result<bool> {
765 let _guard = store.guard();
766 let mut fresh = store
767 .get(&talk.id)
768 .with_context(|| format!("talk {} was deleted", talk.short()))?;
769 if !pending_matches(&fresh, expected_text, expected_attachments) {
770 *talk = fresh;
771 return Ok(false);
772 }
773 fresh.pending = text.trim().to_owned();
774 store.put(&mut fresh)?;
775 *talk = fresh;
776 Ok(true)
777}
778
779fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
780 talk.pending == expected_text
781 && talk
782 .pending_attachments
783 .iter()
784 .map(|attachment| &attachment.id)
785 .eq(expected_attachments.iter())
786}
787
788async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
795 let spec = cfg
796 .agents
797 .iter()
798 .find(|a| a.id == talk.agent)
799 .with_context(|| {
800 format!(
801 "talk {} was opened with agent `{}`, which is no longer in \
802 the roster; restore it in magi.toml or start a new \
803 conversation",
804 talk.short(),
805 talk.agent
806 )
807 })?;
808
809 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
810 let last_note = attachment_note(
814 store,
815 &talk.id,
816 talk.turns
817 .last()
818 .map_or(&[][..], |t| t.attachments.as_slice()),
819 );
820 let first_ever = talk.turns.len() <= 1;
821 let body = if talk.seat.turns == 0 && first_ever {
822 format!(
823 "{}\n\n# Operator\n\n{text}{last_note}",
824 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write)
825 )
826 } else if talk.seat.turns == 0 {
827 format!(
830 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
831 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write),
832 transcript(talk, store)
833 )
834 } else if resuming {
835 format!("{text}{last_note}")
836 } else {
837 format!("{}\n\n{text}{last_note}", transcript(talk, store))
838 };
839
840 let attachment_paths: Vec<PathBuf> = talk
846 .turns
847 .iter()
848 .flat_map(|t| t.attachments.iter())
849 .filter_map(|a| store.attachment_path(&talk.id, a))
850 .collect();
851
852 let artifacts = store.artifacts_of(&talk.id);
853 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
856 let stem = format!("turn-{}", operator_turns.max(1));
857 let cache_dir = cfg.cache_dir();
860 let inv = Invocation {
861 cwd: &talk.repo,
862 prompt: &body,
863 timeout: turn_timeout(cfg),
864 allow_write: cfg.talk.allow_write,
869 sessions: cfg.graph.sessions,
870 artifacts: &artifacts,
871 stem: &stem,
872 run: &talk.id,
875 node: "chat",
876 cache_dir: cache_dir.as_deref(),
877 attachments: &attachment_paths,
878 };
879
880 let outcome = agent::invoke(spec, &mut talk.seat, &inv).await;
881 let note = |why: String| Turn {
882 who: Who::Agent,
883 body: format!("{MAGI_NOTE}{why}"),
884 at: Timestamp::now(),
885 attachments: Vec::new(),
886 };
887 let (reply, failure) = match outcome {
888 Err(e) => (
889 note(format!("could not run agent `{}`: {e}", talk.agent)),
890 Some(format!("could not run agent `{}`: {e}", talk.agent)),
891 ),
892 Ok(out) if out.quota_exhausted() => {
893 let reset = out
894 .quota
895 .as_ref()
896 .and_then(|q| q.reset.clone())
897 .map_or_else(String::new, |r| format!(" (resets {r})"));
898 let why = format!(
899 "agent `{}` is out of quota{reset}; your message is saved, so \
900 say it again when the window reopens",
901 talk.agent
902 );
903 (note(why.clone()), Some(why))
904 }
905 Ok(out) if out.timed_out => {
906 let why = format!(
907 "agent `{}` did not answer within {}s; your message is saved",
908 talk.agent,
909 turn_timeout(cfg).as_secs()
910 );
911 (note(why.clone()), Some(why))
912 }
913 Ok(out) if !out.usable() => {
914 let why = format!(
915 "agent `{}` produced no answer (exit {}); your message is saved",
916 talk.agent,
917 out.exit_code
918 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
919 );
920 (note(why.clone()), Some(why))
921 }
922 Ok(out) => (
923 Turn {
924 who: Who::Agent,
925 body: out.text.trim().to_owned(),
926 at: Timestamp::now(),
927 attachments: Vec::new(),
928 },
929 None,
930 ),
931 };
932
933 let _guard = store.guard();
945 let Ok(fresh) = store.get(&talk.id) else {
951 return Ok(());
952 };
953 talk.status = fresh.status;
954 talk.pending = fresh.pending;
958 talk.pending_attachments = fresh.pending_attachments;
959 talk.turns.push(reply);
960 if let Err(put_err) = store.put(talk) {
961 let lost = talk.turns.pop().expect("just pushed above");
970 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
971 let why = match &stash {
972 Ok(path) => format!(
973 "agent `{}` answered, but the reply could not be saved to \
974 this conversation ({put_err:#}); the raw text was kept at \
975 {} - your message is saved, ask again",
976 talk.agent,
977 path.display()
978 ),
979 Err(stash_err) => format!(
980 "agent `{}` answered, but the reply could not be saved to \
981 this conversation ({put_err:#}), and it could not be kept \
982 anywhere else either ({stash_err:#}); your message is \
983 saved, ask again",
984 talk.agent
985 ),
986 };
987 talk.turns.push(note(why.clone()));
988 return match store.put(talk) {
995 Ok(()) => bail!("{why}"),
996 Err(note_err) => {
997 talk.turns.pop();
1017 Err(note_err).context(why)
1018 }
1019 };
1020 }
1021
1022 match failure {
1023 Some(why) => bail!("{why}"),
1024 None => Ok(()),
1025 }
1026}
1027
1028fn transcript(talk: &Talk, store: &Talks) -> String {
1031 let mut out = String::from(
1032 "This conversation cannot resume on the CLI's side, so here is \
1033 everything said so far; answer only the last message.\n",
1034 );
1035 for t in &talk.turns {
1036 let who = match t.who {
1037 Who::Operator => "operator",
1038 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1039 Who::Agent => "you",
1040 };
1041 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1042 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1043 }
1044 out
1045}
1046
1047fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1052 if attachments.is_empty() {
1053 return String::new();
1054 }
1055 let mut out = String::from(
1056 "\n\nThe operator attached the image(s) below to this message. Open \
1057 and look at each one before you answer.\n",
1058 );
1059 for att in attachments {
1060 if let Some(path) = store.attachment_path(talk_id, att) {
1061 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1062 }
1063 }
1064 out.push('\n');
1065 out
1066}
1067
1068pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
1090 let write_policy = if allow_write {
1091 "Write access is enabled for this conversation (`allow_write = \
1092 true`), so you may write files - but only a small, \
1093 already-decided edit the operator names outright in this \
1094 conversation, not an implementation. This is a permission on the \
1095 conversation as a whole, not a property of whichever repository \
1096 it happened to start in: if the operator names a different \
1097 repository for that small edit, the policy allows it there too. \
1098 Your own tool may still confine writes to the repository this \
1099 conversation started in regardless - if a write elsewhere is \
1100 refused, say so plainly rather than working around it. Once you \
1101 have made an edit, say plainly what you edited. Anything bigger, \
1102 or anything still open-ended, still goes through the queue below \
1103 rather than being done here."
1104 } else {
1105 "Do not write files. Implementing a change is not this \
1106 conversation's job; a separate, blind competition of agents does \
1107 that, and a repository this conversation has already edited would \
1108 make their diffs unjudgeable."
1109 };
1110 let mut out = format!(
1111 "You are magi's standing conversation partner for its operator, who \
1112 usually has this open on a phone. Keep replies short: no preamble, \
1113 no restating what they just said.\n\n\
1114 # Repository\n\n{repo}\n\n\
1115 You may look around: read files, run shell commands, search history, \
1116 run tests - whatever answers the question. {write_policy}\n\n\
1117 A short, command-shaped message (\"list\", \"info <id>\", \"show \
1118 3cbf\") is almost always the operator asking you to look something \
1119 up, not an instruction to file - answer it yourself with `magi \
1120 list`, `magi show <id>`, `magi task list`, or the like, the same way \
1121 you would answer any other question in this conversation.\n\n\
1122 # When the operator wants something done\n\n\
1123 Run:\n\n\
1124 magi task add --solo --repo {repo} <instruction>\n\n\
1125 and tell the operator the task id it prints, so they can follow it \
1126 from the Queue. If it refuses with a duplicate warning (the \
1127 instruction names a branch, commit or pull request that an \
1128 unfinished task, run or PR already owns), do not repeat it with \
1129 --force yourself: tell the operator what it matched and let them \
1130 decide. Write <instruction> so that an implementer who has \
1131 never seen this conversation can act on it alone - it is everything \
1132 they get. Use --solo: it runs the task through one implementer \
1133 straight into review instead of the usual multi-agent competition, \
1134 which is the right shape for a change this conversation has already \
1135 settled, rather than one still worth several independent takes.\n\n\
1136 If the operator asks for something in a different repository, \
1137 --repo does not have to be a full path: --repo owner/repo (or just \
1138 repo, when that is unambiguous) is resolved against local checkouts \
1139 the same way `magi repos` lists them. If the command fails because \
1140 nothing matches or more than one checkout shares that name, ask the \
1141 operator which repository they mean (or run `magi repos` yourself \
1142 to see the candidates) rather than guessing.\n\n\
1143 If the operator attached an image (a screenshot, say) that the task \
1144 is about, pass it with `--attach <path>`, using the absolute path \
1145 the turn's attachment note gives; repeat the flag for several. \
1146 `magi task add --solo --attach <path> <instruction>` copies the \
1147 file into the task, so the implementer receives it. Do not paste the \
1148 path into <instruction> instead: deleting this conversation deletes \
1149 its attachments, and then that path reaches no one.\n",
1150 repo = repo.display(),
1151 );
1152 out.push_str(&language_note(language));
1153 out
1154}
1155
1156fn language_note(language: &str) -> String {
1159 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
1160 String::new()
1161 } else {
1162 format!("\nHold this conversation in {language}.\n")
1163 }
1164}
1165
1166pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
1173 let mut tasks: Vec<Task> = queue
1174 .list()
1175 .into_iter()
1176 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
1177 .collect();
1178 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
1179 tasks
1180}
1181
1182fn read_path(path: &Path) -> Result<Talk> {
1183 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1184 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
1185}
1186
1187const PUT_RETRIES: u32 = 5;
1190
1191fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1203 let mut last_err = None;
1204 for attempt in 0..PUT_RETRIES {
1205 if attempt > 0 {
1206 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
1207 }
1208 match try_write_atomic(tmp, path, body) {
1209 Ok(()) => return Ok(()),
1210 Err(e) => last_err = Some(e),
1211 }
1212 }
1213 Err(last_err.expect("the loop above always runs at least once"))
1214}
1215
1216fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1217 #[cfg(test)]
1218 if failpoint::take_forced_put_failure() {
1219 bail!("simulated write failure (test)");
1220 }
1221 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
1222 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
1223 Ok(())
1224}
1225
1226fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
1231 let dir = store.artifacts_of(id);
1232 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1233 let path = dir.join(format!("{stem}-lost.txt"));
1234 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
1235 Ok(path)
1236}
1237
1238#[cfg(test)]
1245mod failpoint {
1246 use std::cell::Cell;
1247
1248 thread_local! {
1249 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
1250 }
1251
1252 pub(super) fn force_put_failures(count: u32) {
1255 FORCE_PUT_FAILURES.with(|c| c.set(count));
1256 }
1257
1258 pub(super) fn take_forced_put_failure() -> bool {
1261 FORCE_PUT_FAILURES.with(|c| {
1262 let n = c.get();
1263 if n == 0 {
1264 false
1265 } else {
1266 c.set(n - 1);
1267 true
1268 }
1269 })
1270 }
1271}
1272
1273fn short(id: &str) -> &str {
1274 id.split('-').next_back().unwrap_or(id)
1275}
1276
1277fn new_id() -> String {
1278 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1279 let seed = crate::rng::entropy();
1280 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1281}
1282
1283fn attachment_ext(mime: &str) -> Option<&'static str> {
1288 match mime {
1289 "image/png" => Some("png"),
1290 "image/jpeg" => Some("jpg"),
1291 "image/gif" => Some("gif"),
1292 "image/webp" => Some("webp"),
1293 _ => None,
1294 }
1295}
1296
1297pub fn valid_attachment_id(id: &str) -> bool {
1302 id.len() == 32
1303 && id
1304 .bytes()
1305 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
1306}
1307
1308fn new_attachment_id() -> String {
1312 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
1313 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
1314}
1315
1316#[cfg(test)]
1317mod tests {
1318 use std::collections::BTreeMap;
1319
1320 use crate::config::{AgentKind, AgentSpec, Graph};
1321 use crate::queue::{Queue, Source, Task};
1322
1323 use super::*;
1324
1325 fn store() -> (tempfile::TempDir, Talks) {
1327 let tmp = tempfile::tempdir().expect("tempdir");
1328 let talks = Talks::at(tmp.path().join("talks"));
1329 (tmp, talks)
1330 }
1331
1332 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
1336 let path = dir.join("mock-talk-agent.sh");
1337 std::fs::write(&path, script).expect("write mock");
1338 AgentSpec {
1339 id: "mock".to_owned(),
1340 kind: AgentKind::Command,
1341 model: None,
1342 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
1343 extra_args: Vec::new(),
1344 env,
1345 prompt_delivery: None,
1346 }
1347 }
1348
1349 fn config(spec: AgentSpec) -> Config {
1350 Config {
1351 agents: vec![spec],
1352 graph: Graph {
1353 language: "en".to_owned(),
1354 ..Graph::default()
1355 },
1356 ..Config::default()
1357 }
1358 }
1359
1360 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
1362
1363 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
1365
1366 const ECHO: &str = "#!/bin/sh\ncat\n";
1369
1370 fn env(reply: &str) -> BTreeMap<String, String> {
1371 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
1372 }
1373
1374 #[test]
1375 fn the_frozen_json_field_names_round_trip_through_disk() {
1376 let (tmp, talks) = store();
1377 let mut talk = Talk {
1378 schema: SCHEMA,
1379 id: "20260904-014455-ab12".to_owned(),
1380 repo: tmp.path().to_owned(),
1381 agent: "sonnet".to_owned(),
1382 status: TalkStatus::Open,
1383 turns: Vec::new(),
1384 pending: String::new(),
1385 pending_attachments: Vec::new(),
1386 created_at: Timestamp::now(),
1387 updated_at: Timestamp::now(),
1388 seat: SeatState::new(SEAT, "sonnet", 7),
1389 };
1390 talks.put(&mut talk).expect("put");
1391
1392 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
1393 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
1394 for field in [
1395 "schema",
1396 "id",
1397 "repo",
1398 "agent",
1399 "status",
1400 "turns",
1401 "created_at",
1402 "updated_at",
1403 ] {
1404 assert!(v.get(field).is_some(), "missing field `{field}`");
1405 }
1406 assert_eq!(v["schema"], 1);
1407 assert_eq!(v["status"], "open");
1408
1409 let back = talks.get(&talk.id).expect("get");
1410 assert_eq!(back.id, talk.id);
1411 assert_eq!(back.status, TalkStatus::Open);
1412 }
1413
1414 #[test]
1415 fn opening_a_talk_takes_no_agent_turn() {
1416 let (tmp, talks) = store();
1417 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1421 let cfg = config(spec);
1422
1423 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1424 assert_eq!(talk.status, TalkStatus::Open);
1425 assert!(talk.turns.is_empty(), "nothing has been said yet");
1426
1427 let on_disk = talks.get(&talk.id).expect("get");
1428 assert_eq!(on_disk.turns.len(), 0);
1429 }
1430
1431 #[test]
1439 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
1440 let (tmp, talks) = store();
1441 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1442 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1443 chatter_spec.id = "chatter-mock".to_owned();
1444
1445 let mut cfg = Config {
1446 agents: vec![first_spec.clone(), chatter_spec.clone()],
1447 graph: Graph {
1448 language: "en".to_owned(),
1449 ..Graph::default()
1450 },
1451 ..Config::default()
1452 };
1453 cfg.roles.chatter = Some(chatter_spec.id.clone());
1454
1455 let talk =
1456 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
1457 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
1458
1459 cfg.roles.chatter = None;
1460 let fallback =
1461 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
1462 assert_eq!(
1463 fallback.agent, first_spec.id,
1464 "unset chatter must fall back to agent::pick's own default order"
1465 );
1466 }
1467
1468 #[test]
1471 fn a_talk_recorded_without_attachments_still_reads() {
1472 let (tmp, talks) = store();
1473 let path = talks.path_of("20260904-014455-ab12");
1474 std::fs::create_dir_all(talks.root()).expect("talks dir");
1475 std::fs::write(
1476 &path,
1477 serde_json::json!({
1478 "schema": 1,
1479 "id": "20260904-014455-ab12",
1480 "repo": tmp.path(),
1481 "agent": "sonnet",
1482 "status": "open",
1483 "turns": [
1484 { "who": "operator", "body": "still there?",
1485 "at": Timestamp::now().to_string() },
1486 ],
1487 "created_at": Timestamp::now().to_string(),
1488 "updated_at": Timestamp::now().to_string(),
1489 "seat": SeatState::new(SEAT, "sonnet", 7),
1490 })
1491 .to_string(),
1492 )
1493 .expect("write pre-attachments talk");
1494
1495 let talk = talks.get("20260904-014455-ab12").expect("must still read");
1496 assert!(talk.turns[0].attachments.is_empty());
1497 }
1498
1499 #[test]
1500 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
1501 let (tmp, talks) = store();
1502 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1503 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1504
1505 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
1506 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
1507 let saved = talks.get(&talk.id).expect("reload queued talk");
1508 assert_eq!(saved.pending, "first\n\nsecond");
1509 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
1510
1511 let drained = drain(&mut talk, &talks).expect("drain");
1512 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
1513 let saved = talks.get(&talk.id).expect("reload drained talk");
1514 assert!(saved.pending.is_empty());
1515 assert_eq!(saved.turns.len(), 1);
1516 assert_eq!(saved.turns[0].body, "first\n\nsecond");
1517 }
1518
1519 #[test]
1520 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
1521 let (tmp, talks) = store();
1522 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1523 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1524 let attachment = Attachment {
1525 id: "a".repeat(32),
1526 name: "shot.png".to_owned(),
1527 mime: "image/png".to_owned(),
1528 bytes: 3,
1529 };
1530
1531 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
1532 assert!(
1533 edit_pending_text(
1534 &mut talk,
1535 &talks,
1536 "corrected",
1537 "first",
1538 std::slice::from_ref(&attachment.id),
1539 )
1540 .expect("edit")
1541 );
1542 let saved = talks.get(&talk.id).expect("reload edited draft");
1543 assert_eq!(saved.pending, "corrected");
1544 assert_eq!(saved.pending_attachments, vec![attachment]);
1545
1546 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
1547 assert!(
1548 !edit_pending_text(
1549 &mut talk,
1550 &talks,
1551 "stale edit",
1552 "corrected",
1553 &["a".repeat(32)],
1554 )
1555 .expect("stale edit is a conflict")
1556 );
1557 assert_eq!(
1558 talks.get(&talk.id).expect("reload after conflict").pending,
1559 "corrected\n\nlater"
1560 );
1561 assert!(
1562 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
1563 .expect("stale clear is a conflict")
1564 );
1565 assert_eq!(
1566 talks
1567 .get(&talk.id)
1568 .expect("reload after stale clear")
1569 .pending,
1570 "corrected\n\nlater"
1571 );
1572 }
1573
1574 #[tokio::test]
1575 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
1576 let (tmp, talks) = store();
1577 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
1578 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
1579 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1580 let id = running.id.clone();
1581 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
1582
1583 let response_talks = talks.clone();
1584 let response_cfg = cfg.clone();
1585 let reply = tokio::spawn(async move {
1586 respond(&mut running, &response_talks, &response_cfg, &first).await
1587 });
1588 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
1589
1590 let mut queued = talks.get(&id).expect("queued handle");
1591 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
1592 reply.await.expect("join").expect("reply");
1593
1594 let saved = talks.get(&id).expect("reload");
1595 assert_eq!(saved.pending, "next");
1596 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
1597 }
1598
1599 #[tokio::test]
1600 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
1601 let (tmp, talks) = store();
1602 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
1603 let cfg = config(spec);
1604 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1605
1606 say(
1607 &mut talk,
1608 &talks,
1609 &cfg,
1610 "what does the queue module do?",
1611 Vec::new(),
1612 )
1613 .await
1614 .expect("first turn");
1615 let first_prompt = &talk.turns[1].body;
1616 assert!(first_prompt.contains("magi task add --solo"));
1617 assert!(first_prompt.contains("what does the queue module do?"));
1618
1619 say(&mut talk, &talks, &cfg, "and how is it locked?", Vec::new())
1620 .await
1621 .expect("second turn");
1622 let second_prompt = &talk.turns[3].body;
1623 assert!(
1624 !second_prompt.contains("magi task add --solo"),
1625 "the briefing is sent once, not on every turn: {second_prompt}"
1626 );
1627 assert!(second_prompt.contains("and how is it locked?"));
1628 }
1629
1630 #[tokio::test]
1631 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
1632 let (tmp, talks) = store();
1633 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
1634 let mut b = a.clone();
1635 b.id = "other".to_owned();
1636 let mut cfg = config(a.clone());
1637 cfg.agents.push(b.clone());
1638 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
1639 say(&mut talk, &talks, &cfg, "remember the walrus", Vec::new())
1640 .await
1641 .expect("first turn");
1642 let old_session = talk.seat.claude_session.clone();
1643 assert_eq!(talk.seat.turns, 1);
1644
1645 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
1646 assert_eq!(talk.agent, "other");
1647 assert_eq!(talk.seat.turns, 0);
1648 assert_eq!(talk.seat.agent, "other");
1649 assert_ne!(talk.seat.claude_session, old_session);
1650 let note = talk.turns.last().expect("note");
1651 assert_eq!(note.who, Who::Agent);
1652 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
1653 assert!(note.body.contains("changed from"), "{}", note.body);
1654 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
1655
1656 let before = talk.turns.len();
1657 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
1658 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
1659
1660 say(&mut talk, &talks, &cfg, "what did I say?", Vec::new())
1661 .await
1662 .expect("turn after switch");
1663 let prompt = &talk.turns.last().expect("reply").body;
1664 assert!(prompt.contains("remember the walrus"), "{prompt}");
1665 assert!(prompt.contains("## magi"), "{prompt}");
1666 assert!(prompt.contains("what did I say?"), "{prompt}");
1667 }
1668
1669 #[tokio::test]
1670 async fn say_appends_the_operator_turn_then_the_agent_turn() {
1671 let (tmp, talks) = store();
1672 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
1673 let cfg = config(spec);
1674 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1675
1676 say(
1677 &mut talk,
1678 &talks,
1679 &cfg,
1680 "can I rename this function?",
1681 Vec::new(),
1682 )
1683 .await
1684 .expect("say");
1685
1686 assert_eq!(talk.turns.len(), 2);
1687 assert_eq!(talk.turns[0].who, Who::Operator);
1688 assert_eq!(talk.turns[0].body, "can I rename this function?");
1689 assert_eq!(talk.turns[1].who, Who::Agent);
1690 assert_eq!(talk.turns[1].body, "go ahead");
1691 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
1692 }
1693
1694 #[tokio::test]
1695 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
1696 let (tmp, talks) = store();
1697 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1698 let cfg = config(spec);
1699 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1700
1701 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
1702 .await
1703 .expect_err("a turn with no answer is an error");
1704 assert!(err.to_string().contains("no answer"), "{err}");
1705
1706 let on_disk = talks.get(&talk.id).expect("get");
1707 assert_eq!(on_disk.turns.len(), 2);
1708 assert_eq!(on_disk.turns[0].body, "check the tests");
1709 let note = &on_disk.turns[1];
1710 assert_eq!(note.who, Who::Agent);
1711 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
1712 assert!(note.body.contains("your message is saved"));
1713 }
1714
1715 #[tokio::test]
1721 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
1722 let (tmp, talks) = store();
1723 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
1724 let cfg = config(spec);
1725 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1726
1727 let text =
1728 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
1729 failpoint::force_put_failures(PUT_RETRIES - 1);
1732 respond(&mut talk, &talks, &cfg, &text)
1733 .await
1734 .expect("respond must survive a write failure its own retries can outlast");
1735
1736 assert_eq!(talk.turns.len(), 2);
1737 assert_eq!(talk.turns[1].who, Who::Agent);
1738 assert_eq!(talk.turns[1].body, "go ahead");
1739 let on_disk = talks.get(&talk.id).expect("get");
1740 assert_eq!(
1741 on_disk.turns, talk.turns,
1742 "the reply must reach disk despite the early write failures"
1743 );
1744 }
1745
1746 #[tokio::test]
1752 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
1753 let (tmp, talks) = store();
1754 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
1755 let cfg = config(spec);
1756 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1757
1758 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
1759 failpoint::force_put_failures(PUT_RETRIES);
1764 let err = respond(&mut talk, &talks, &cfg, &text)
1765 .await
1766 .expect_err("a reply that cannot be saved must be reported, not swallowed");
1767 assert!(err.to_string().contains("could not be saved"), "{err}");
1768
1769 let on_disk = talks.get(&talk.id).expect("get");
1770 assert_eq!(
1771 on_disk.turns.len(),
1772 2,
1773 "the operator turn plus a visible note"
1774 );
1775 assert_eq!(on_disk.turns[0].body, "check the tests");
1776 let note = &on_disk.turns[1];
1777 assert_eq!(note.who, Who::Agent);
1778 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
1779 assert!(
1780 note.body.contains("could not be saved"),
1781 "the operator must be told the reply is missing, not left staring \
1782 at a gap with no explanation: {}",
1783 note.body
1784 );
1785 assert_eq!(
1786 talk.turns, on_disk.turns,
1787 "the in-memory talk must match what actually landed on disk"
1788 );
1789
1790 let artifacts = talks.artifacts_of(&talk.id);
1793 let stash = std::fs::read_dir(&artifacts)
1794 .expect("artifacts dir")
1795 .filter_map(|e| e.ok())
1796 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
1797 .expect("a stash file for the lost reply");
1798 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
1799 assert_eq!(stashed, "go ahead");
1800
1801 assert_eq!(
1809 on_disk.seat.turns, 1,
1810 "the note's write must carry the turn the CLI actually took"
1811 );
1812 assert_eq!(
1813 on_disk.seat.claude_session, talk.seat.claude_session,
1814 "the session id handed to the CLI must survive the failed reply"
1815 );
1816 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
1817 assert!(
1818 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
1819 "the next turn must resume, not open the same session id twice"
1820 );
1821 }
1822
1823 #[tokio::test]
1828 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
1829 let (tmp, talks) = store();
1830 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
1831 let cfg = config(spec);
1832 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1833
1834 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
1835 failpoint::force_put_failures(PUT_RETRIES * 2);
1838 let err = respond(&mut talk, &talks, &cfg, &text)
1839 .await
1840 .expect_err("neither the reply nor the note could be saved");
1841 assert!(err.to_string().contains("could not be saved"), "{err}");
1842
1843 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
1844 let on_disk = talks.get(&talk.id).expect("get");
1845 assert_eq!(on_disk.turns.len(), 1);
1846
1847 assert_eq!(
1857 on_disk.seat.turns, 0,
1858 "an unwritable file cannot record the turn the CLI took"
1859 );
1860 assert_eq!(
1861 talk.seat.turns, 1,
1862 "the in-memory seat still reports the turn the CLI actually took"
1863 );
1864 assert_eq!(
1865 on_disk.seat.claude_session, talk.seat.claude_session,
1866 "the session id was minted at `begin` and never changes here"
1867 );
1868 }
1869
1870 #[tokio::test]
1874 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
1875 let (tmp, talks) = store();
1876 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
1877 let cfg = config(spec);
1878 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1879
1880 let att = talks
1881 .put_attachment(
1882 &talk.id,
1883 "image/png",
1884 "screenshot.png",
1885 b"pretend-png-bytes",
1886 )
1887 .expect("put attachment");
1888
1889 say(&mut talk, &talks, &cfg, "", vec![att.clone()])
1890 .await
1891 .expect("an empty body with an attachment is still a turn");
1892
1893 let operator_turn = &talk.turns[0];
1894 assert_eq!(operator_turn.who, Who::Operator);
1895 assert_eq!(operator_turn.body, "");
1896 assert_eq!(operator_turn.attachments, vec![att.clone()]);
1897
1898 let prompt = &talk.turns[1].body;
1899 let expected_path = talks
1900 .attachments_dir(&talk.id)
1901 .join(format!("{}.png", att.id));
1902 assert!(
1903 prompt.contains(&expected_path.display().to_string()),
1904 "the agent must be told the attachment's absolute path: {prompt}"
1905 );
1906 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
1907 }
1908
1909 #[test]
1918 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
1919 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
1920 let att = Attachment {
1921 id: "0".repeat(32),
1922 name: "shot.png".to_owned(),
1923 mime: "image/png".to_owned(),
1924 bytes: 3,
1925 };
1926 let path = talks
1927 .attachment_path("some-talk-id", &att)
1928 .expect("a supported mime always yields a path");
1929 assert!(
1930 path.is_absolute(),
1931 "must be absolute even off a relative store root: {}",
1932 path.display()
1933 );
1934 }
1935
1936 #[tokio::test]
1937 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
1938 let (tmp, talks) = store();
1943 let slow = mock_agent(
1944 tmp.path(),
1945 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
1946 BTreeMap::new(),
1947 );
1948 let mut cfg = config(slow);
1949 cfg.graph.timeout_talk = 1;
1950 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1951
1952 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
1953 .await
1954 .expect_err("a turn that never answers is an error");
1955 assert!(
1956 err.to_string().contains("did not answer within 1s"),
1957 "{err}"
1958 );
1959
1960 let on_disk = talks.get(&talk.id).expect("get");
1961 let note = on_disk.turns.last().expect("a note turn was recorded");
1962 assert!(
1963 note.body.contains("did not answer within 1s"),
1964 "the transcript must show the configured timeout: {}",
1965 note.body
1966 );
1967 }
1968
1969 #[test]
1970 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
1971 let (tmp, talks) = store();
1972 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
1973 let cfg = config(spec);
1974 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1975
1976 close(&mut talk, &talks).expect("close");
1977 assert_eq!(talk.status, TalkStatus::Closed);
1978 close(&mut talk, &talks).expect("closing twice is not an error");
1979
1980 let err =
1981 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
1982 assert!(err.to_string().contains("closed"));
1983 let _ = &cfg; }
1985
1986 #[tokio::test]
1987 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
1988 let (tmp, talks) = store();
1989 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
1990 let cfg = config(spec);
1991 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1994
1995 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
1999 close(&mut closed_elsewhere, &talks).expect("close");
2000 assert_eq!(
2001 talks.get(&in_flight.id).expect("reread").status,
2002 TalkStatus::Closed,
2003 "the close landed on disk before the turn finished"
2004 );
2005
2006 assert_eq!(in_flight.status, TalkStatus::Open);
2010 respond(&mut in_flight, &talks, &cfg, "one more question")
2011 .await
2012 .expect("the turn itself still completes");
2013
2014 let on_disk = talks.get(&in_flight.id).expect("reread");
2015 assert_eq!(
2016 on_disk.status,
2017 TalkStatus::Closed,
2018 "a close must stick even when a turn that started before it finishes after it"
2019 );
2020 assert!(
2023 on_disk.turns.iter().any(|t| t.body == "here you go"),
2024 "the in-flight turn's own reply is still recorded: {:?}",
2025 on_disk.turns
2026 );
2027 }
2028
2029 #[test]
2030 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
2031 let (tmp, talks) = store();
2032 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2033 let cfg = config(spec);
2034 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2037
2038 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
2041 close(&mut closed_elsewhere, &talks).expect("close");
2042 assert_eq!(
2043 talks.get(&stale.id).expect("reread").status,
2044 TalkStatus::Closed,
2045 "the close landed on disk before record was called"
2046 );
2047
2048 assert_eq!(stale.status, TalkStatus::Open);
2052 let err = record(&mut stale, &talks, "still there?", Vec::new())
2053 .expect_err("a close that landed first must be honored, not overwritten");
2054 assert!(err.to_string().contains("closed"));
2055
2056 let on_disk = talks.get(&stale.id).expect("reread");
2057 assert_eq!(
2058 on_disk.status,
2059 TalkStatus::Closed,
2060 "record must not resurrect a conversation closed while its snapshot was stale"
2061 );
2062 assert!(
2063 on_disk.turns.is_empty(),
2064 "the rejected turn must not have been appended: {:?}",
2065 on_disk.turns
2066 );
2067 let _ = &cfg; }
2069
2070 #[test]
2071 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
2072 let (tmp, talks) = store();
2073 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2074 let cfg = config(spec);
2075 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2076
2077 let held = talks.guard();
2081
2082 let talks2 = talks.clone();
2083 let id = talk.id.clone();
2084 let closing = std::thread::spawn(move || {
2085 let mut talk = talks2.get(&id).expect("get");
2086 close(&mut talk, &talks2).expect("close");
2087 });
2088
2089 std::thread::sleep(Duration::from_millis(50));
2090 assert!(
2091 !closing.is_finished(),
2092 "close must wait for the guard, not read and write while it is held - \
2093 a re-read alone narrows this window without closing it"
2094 );
2095
2096 drop(held);
2097 closing.join().expect("close thread panicked");
2098
2099 assert_eq!(
2100 talks.get(&talk.id).expect("reread").status,
2101 TalkStatus::Closed,
2102 "once the guard is free, close still lands"
2103 );
2104 let _ = &cfg; }
2106
2107 #[test]
2108 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
2109 let (tmp, talks) = store();
2110 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2111 let cfg = config(spec);
2112 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2113
2114 close(&mut talk, &talks).expect("close");
2115 assert_eq!(talk.status, TalkStatus::Closed);
2116
2117 reopen(&mut talk, &talks).expect("reopen");
2118 assert_eq!(talk.status, TalkStatus::Open);
2119 assert_eq!(
2120 talks.get(&talk.id).expect("reread").status,
2121 TalkStatus::Open
2122 );
2123
2124 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
2126 assert_eq!(talk.status, TalkStatus::Open);
2127
2128 record(&mut talk, &talks, "one more thing", Vec::new())
2129 .expect("a reopened talk takes turns again");
2130 let _ = &cfg; }
2132
2133 #[test]
2134 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
2135 let (tmp, talks) = store();
2136 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2137 let cfg = config(spec);
2138 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2139
2140 let artifacts = talks.artifacts_of(&talk.id);
2141 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
2142 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
2143
2144 talks.remove(&talk.id).expect("remove");
2145 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
2146 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
2147 assert!(
2148 talks.get(&talk.id).is_err(),
2149 "a removed talk cannot be read back"
2150 );
2151
2152 let err = talks
2153 .remove("nonexistent-id")
2154 .expect_err("unknown id refused");
2155 assert!(err.to_string().contains("no talk matches"), "{err}");
2156 let _ = &cfg; }
2158
2159 #[tokio::test]
2160 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2161 let (tmp, talks) = store();
2162 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2163 let cfg = config(spec);
2164 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2167
2168 talks.remove(&in_flight.id).expect("remove");
2169 assert!(
2170 talks.get(&in_flight.id).is_err(),
2171 "the delete landed on disk before the turn finished"
2172 );
2173
2174 respond(&mut in_flight, &talks, &cfg, "one more question")
2177 .await
2178 .expect("the turn itself still completes rather than erroring");
2179
2180 assert!(
2181 talks.get(&in_flight.id).is_err(),
2182 "a delete must stick even when a turn that started before it finishes after it"
2183 );
2184 }
2185
2186 #[test]
2187 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
2188 let (tmp, talks) = store();
2189 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2190 let cfg = config(spec);
2191 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2194
2195 talks.remove(&stale.id).expect("remove");
2196
2197 let err = record(&mut stale, &talks, "still there?", Vec::new())
2201 .expect_err("a delete that landed first must be honored, not overwritten");
2202 assert!(err.to_string().contains("deleted"), "{err}");
2203
2204 assert!(
2205 talks.get(&stale.id).is_err(),
2206 "record must not resurrect a conversation deleted while its snapshot was stale"
2207 );
2208 let _ = &cfg; }
2210
2211 #[test]
2212 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
2213 let (tmp, talks) = store();
2214 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2215 let cfg = config(spec);
2216 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2219
2220 talks.remove(&stale.id).expect("remove");
2221
2222 let err = close(&mut stale, &talks)
2226 .expect_err("a delete that landed first must be honored, not overwritten");
2227 assert!(err.to_string().contains("deleted"), "{err}");
2228
2229 assert!(
2230 talks.get(&stale.id).is_err(),
2231 "close must not resurrect a conversation deleted while its snapshot was stale"
2232 );
2233 let _ = &cfg; }
2235
2236 #[test]
2237 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
2238 let (tmp, talks) = store();
2239 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2240 let cfg = config(spec);
2241 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2244 close(&mut stale, &talks).expect("close");
2245
2246 talks.remove(&stale.id).expect("remove");
2247
2248 let err = reopen(&mut stale, &talks)
2252 .expect_err("a delete that landed first must be honored, not overwritten");
2253 assert!(err.to_string().contains("deleted"), "{err}");
2254
2255 assert!(
2256 talks.get(&stale.id).is_err(),
2257 "reopen must not resurrect a conversation deleted while its snapshot was stale"
2258 );
2259 let _ = &cfg; }
2261
2262 #[test]
2263 fn list_puts_open_talks_before_closed_ones() {
2264 let (tmp, talks) = store();
2265 let make = |id: &str, status: TalkStatus| {
2266 let mut t = Talk {
2267 schema: SCHEMA,
2268 id: id.to_owned(),
2269 repo: tmp.path().to_owned(),
2270 agent: "mock".to_owned(),
2271 status,
2272 turns: Vec::new(),
2273 pending: String::new(),
2274 pending_attachments: Vec::new(),
2275 created_at: Timestamp::now(),
2276 updated_at: Timestamp::now(),
2277 seat: SeatState::new(SEAT, "mock", 7),
2278 };
2279 talks.put(&mut t).expect("put");
2280 };
2281 make("20260901-000000-0001", TalkStatus::Open);
2282 make("20260902-000000-0002", TalkStatus::Open);
2283 make("20260903-000000-0003", TalkStatus::Closed);
2284
2285 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
2286 assert_eq!(
2287 ids,
2288 [
2289 "20260902-000000-0002",
2290 "20260901-000000-0001",
2291 "20260903-000000-0003"
2292 ]
2293 );
2294 assert_eq!(talks.count_open(), 2);
2295 }
2296
2297 #[test]
2298 fn tasks_of_finds_only_this_talks_own_tasks() {
2299 let dir = tempfile::tempdir().expect("tempdir");
2300 let queue = Queue::at(dir.path().join("queue"));
2301
2302 let mut mine = Task::new(
2303 "rework the loader".to_owned(),
2304 "rework the loader".to_owned(),
2305 PathBuf::from("/repo"),
2306 Source::Agent {
2307 run: "20260904-014455-ab12".to_owned(),
2308 node: "chat".to_owned(),
2309 },
2310 );
2311 queue.put(&mut mine).expect("put mine");
2312
2313 let mut theirs = Task::new(
2314 "unrelated".to_owned(),
2315 "unrelated".to_owned(),
2316 PathBuf::from("/repo"),
2317 Source::Agent {
2318 run: "20260904-090000-zz99".to_owned(),
2319 node: "implement".to_owned(),
2320 },
2321 );
2322 queue.put(&mut theirs).expect("put theirs");
2323
2324 let mut human = Task::new(
2325 "typed by hand".to_owned(),
2326 "typed by hand".to_owned(),
2327 PathBuf::from("/repo"),
2328 Source::Human,
2329 );
2330 queue.put(&mut human).expect("put human");
2331
2332 let found = tasks_of(&queue, "20260904-014455-ab12");
2333 assert_eq!(found.len(), 1);
2334 assert_eq!(found[0].id, mine.id);
2335 }
2336
2337 #[test]
2338 fn the_briefing_names_solo_task_add() {
2339 let brief = briefing(Path::new("/repo"), "en", false);
2340 assert!(brief.contains("magi task add --solo"));
2341 assert!(brief.contains("/repo"));
2342 assert!(!brief.contains("Hold this conversation in"));
2343 }
2344
2345 #[test]
2351 fn the_briefing_explains_targeting_a_different_repository_by_name() {
2352 let brief = briefing(Path::new("/repo"), "en", false);
2353 assert!(brief.contains("--repo does not have to be a full path"));
2354 assert!(brief.contains("owner/repo"));
2355 assert!(brief.contains("magi repos"));
2356 assert!(brief.contains("ask the operator"));
2357 }
2358
2359 #[test]
2360 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
2361 let brief = briefing(Path::new("/repo"), "en", false);
2362 assert!(brief.contains("--attach <path>"), "{brief}");
2363 assert!(brief.contains("deleting this conversation"), "{brief}");
2364 }
2365
2366 #[test]
2367 fn the_briefing_names_the_language_when_it_is_not_english() {
2368 let brief = briefing(Path::new("/repo"), "Japanese", false);
2369 assert!(brief.contains("Hold this conversation in Japanese"));
2370 }
2371
2372 #[test]
2373 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
2374 let read_only = briefing(Path::new("/repo"), "en", false);
2375 assert!(read_only.contains("Do not write files"));
2376 assert!(!read_only.contains("allow_write"));
2377
2378 let writable = briefing(Path::new("/repo"), "en", true);
2379 assert!(!writable.contains("Do not write files"));
2380 assert!(writable.contains("allow_write = true"));
2381 assert!(writable.contains("magi task add --solo"));
2384 assert!(writable.contains("say plainly what you"));
2385 }
2386}