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 #[serde(default, skip_serializing_if = "Option::is_none")]
133 pub usage: Option<TurnUsage>,
134}
135
136#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
144pub struct TurnUsage {
145 pub context_tokens: u64,
147 pub agent: String,
149 #[serde(default)]
151 pub model: Option<String>,
152}
153
154#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
159pub struct ContextUsage {
160 pub tokens: Option<u64>,
162 pub window: Option<u64>,
164 pub percent: Option<u64>,
166 pub warn: bool,
168 pub since_switch: bool,
172 pub model: Option<String>,
174 #[serde(default)]
177 pub estimated: bool,
178}
179
180const CONTEXT_WARN_PERCENT: u64 = 80;
182
183const STANDING_PROMPT_FALLBACK_CHARS: u64 = 7000;
186
187pub fn estimate_context_tokens(talk: &Talk, standing_chars: u64) -> Option<u64> {
200 let mut counted = false;
201 let mut chars = standing_chars;
202 for t in talk.turns.iter().filter(|t| !t.body.starts_with(MAGI_NOTE)) {
203 counted = true;
204 chars += t.body.chars().count() as u64;
205 }
206 counted.then(|| (chars * 2).div_ceil(7))
208}
209
210pub fn context_usage(talk: &Talk, cfg: Option<&Config>) -> ContextUsage {
223 let current = cfg.and_then(|c| c.agents.iter().find(|a| a.id == talk.agent));
224 let model = current.and_then(|a| a.model.clone());
225 let window = cfg
226 .zip(model.as_deref())
227 .and_then(|(c, m)| c.context_window(m))
228 .filter(|w| *w > 0);
229 let usage = talk
230 .turns
231 .iter()
232 .rev()
233 .find(|t| t.who == Who::Agent && !t.body.starts_with(MAGI_NOTE))
234 .and_then(|t| t.usage.as_ref());
235 let measured = usage.map(|u| u.context_tokens);
236 let tokens = measured.or_else(|| {
237 let standing = cfg.map_or(STANDING_PROMPT_FALLBACK_CHARS, |c| {
238 briefing(&talk.repo, &c.graph.language, c.talk.allow_write)
239 .chars()
240 .count() as u64
241 });
242 estimate_context_tokens(talk, standing)
243 });
244 let estimated = measured.is_none() && tokens.is_some();
245 let since_switch =
246 usage.is_some_and(|u| u.agent != talk.agent || (current.is_some() && u.model != model));
247 let (percent, warn) = match (tokens, window) {
248 (Some(t), Some(w)) => (
249 Some(t.saturating_mul(100) / w),
250 t.saturating_mul(100) >= w.saturating_mul(CONTEXT_WARN_PERCENT),
251 ),
252 _ => (None, false),
253 };
254 ContextUsage {
255 tokens,
256 window,
257 percent,
258 warn,
259 since_switch,
260 model,
261 estimated,
262 }
263}
264
265#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
269#[serde(rename_all = "lowercase")]
270pub enum TalkStatus {
271 Open,
274 Closed,
276}
277
278impl TalkStatus {
279 pub fn open(self) -> bool {
281 matches!(self, Self::Open)
282 }
283
284 pub fn as_str(self) -> &'static str {
286 match self {
287 Self::Open => "open",
288 Self::Closed => "closed",
289 }
290 }
291}
292
293#[derive(Debug, Clone, Serialize, Deserialize)]
295#[serde(deny_unknown_fields)]
296pub struct Talk {
297 pub schema: u32,
299 pub id: String,
301 pub repo: PathBuf,
303 pub agent: String,
305 pub status: TalkStatus,
307 pub turns: Vec<Turn>,
309 #[serde(default)]
312 pub pending: String,
313 #[serde(default)]
315 pub pending_attachments: Vec<Attachment>,
316 #[serde(default)]
321 pub fallback: bool,
322 pub created_at: Timestamp,
324 pub updated_at: Timestamp,
326 seat: SeatState,
331}
332
333impl Talk {
334 pub fn short(&self) -> &str {
336 short(&self.id)
337 }
338}
339
340#[derive(Debug, Clone)]
342pub struct Talks {
343 root: PathBuf,
344 lock: Arc<Mutex<()>>,
353}
354
355impl Talks {
356 pub fn open() -> Self {
358 Self::at(crate::run::home().join("talks"))
359 }
360
361 pub fn at(root: PathBuf) -> Self {
364 Self {
365 root,
366 lock: Arc::new(Mutex::new(())),
367 }
368 }
369
370 fn guard(&self) -> MutexGuard<'_, ()> {
379 self.lock.lock().unwrap_or_else(PoisonError::into_inner)
380 }
381
382 pub fn root(&self) -> &Path {
384 &self.root
385 }
386
387 pub fn path_of(&self, id: &str) -> PathBuf {
389 self.root.join(format!("{id}.json"))
390 }
391
392 pub fn artifacts_of(&self, id: &str) -> PathBuf {
395 self.root.join(format!("{id}.artifacts"))
396 }
397
398 pub fn attachments_dir(&self, id: &str) -> PathBuf {
402 self.artifacts_of(id).join("attachments")
403 }
404
405 pub fn put_attachment(
414 &self,
415 id: &str,
416 mime: &str,
417 name: &str,
418 data: &[u8],
419 ) -> Result<Attachment> {
420 let dir = self.attachments_dir(id);
421 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
422 let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
423 let att = Attachment {
424 id: new_attachment_id(),
425 name: name.to_owned(),
426 mime: mime.to_owned(),
427 bytes: data.len() as u64,
428 };
429 std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
430 .with_context(|| format!("write attachment {}", att.id))?;
431 std::fs::write(
432 dir.join(format!("{}.json", att.id)),
433 serde_json::to_string(&att).context("serialize attachment")?,
434 )
435 .with_context(|| format!("write attachment metadata {}", att.id))?;
436 Ok(att)
437 }
438
439 pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
448 if !valid_attachment_id(att_id) {
449 return Ok(None);
450 }
451 let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
452 if !meta_path.is_file() {
453 return Ok(None);
454 }
455 let att = serde_json::from_str(
456 &std::fs::read_to_string(&meta_path)
457 .with_context(|| format!("read {}", meta_path.display()))?,
458 )
459 .with_context(|| format!("parse {}", meta_path.display()))?;
460 Ok(Some(att))
461 }
462
463 pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
467 let Some(att) = self.attachment_meta(id, att_id)? else {
468 return Ok(None);
469 };
470 let ext = attachment_ext(&att.mime).with_context(|| {
471 format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
472 })?;
473 let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
474 let data =
475 std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
476 Ok(Some((att, data)))
477 }
478
479 fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
496 let ext = attachment_ext(&att.mime)?;
497 let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
498 std::path::absolute(&path).ok()
499 }
500
501 pub fn put(&self, t: &mut Talk) -> Result<()> {
510 std::fs::create_dir_all(&self.root)
511 .with_context(|| format!("create {}", self.root.display()))?;
512 t.updated_at = Timestamp::now();
513 let body = serde_json::to_string_pretty(t).context("serialize talk")?;
514 let path = self.path_of(&t.id);
515 let tmp = path.with_extension("json.tmp");
516 write_atomic(&tmp, &path, &body)
517 }
518
519 pub fn get(&self, id: &str) -> Result<Talk> {
521 let resolved = self.resolve_id(id)?;
522 read_path(&self.path_of(&resolved))
523 }
524
525 pub fn list(&self) -> Vec<Talk> {
528 let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
529 .into_iter()
530 .flatten()
531 .flatten()
532 .map(|e| e.path())
533 .filter(|p| p.extension().is_some_and(|x| x == "json"))
534 .filter_map(|p| read_path(&p).ok())
535 .collect();
536 all.sort_unstable_by(|a, b| {
537 let rank = |t: &Talk| u8::from(!t.status.open());
538 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
539 });
540 all
541 }
542
543 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
545 if self.path_of(prefix).is_file() {
546 return Ok(prefix.to_owned());
547 }
548 let hits: Vec<String> = self
549 .list()
550 .into_iter()
551 .map(|t| t.id)
552 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
553 .collect();
554 match hits.len() {
555 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
556 0 => bail!("no talk matches `{prefix}`"),
557 _ => bail!(
558 "`{prefix}` matches {} talks: {}",
559 hits.len(),
560 hits.join(", ")
561 ),
562 }
563 }
564
565 pub fn revision(&self) -> u64 {
568 std::fs::read_dir(&self.root)
569 .into_iter()
570 .flatten()
571 .flatten()
572 .filter_map(|e| e.metadata().ok())
573 .filter_map(|m| m.modified().ok())
574 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
575 .map(|d| d.as_millis() as u64)
576 .max()
577 .unwrap_or(0)
578 }
579
580 pub fn count_open(&self) -> usize {
582 self.list().iter().filter(|t| t.status.open()).count()
583 }
584
585 pub fn remove(&self, id: &str) -> Result<()> {
598 let _guard = self.guard();
599 let resolved = self.resolve_id(id)?;
600 let path = self.path_of(&resolved);
601 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
602 let artifacts = self.artifacts_of(&resolved);
603 if artifacts.is_dir() {
604 std::fs::remove_dir_all(&artifacts)
605 .with_context(|| format!("remove {}", artifacts.display()))?;
606 }
607 Ok(())
608 }
609}
610
611pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
621 let repo = repo.canonicalize().unwrap_or(repo);
624 let spec = match agent {
627 Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
628 None => agent::pick_chain(
629 &cfg.agents,
630 cfg.roles.chatter.as_ref(),
631 &agent::installed,
632 "chatter",
633 )?
634 .remove(0),
635 };
636
637 let now = Timestamp::now();
638 let mut talk = Talk {
639 schema: SCHEMA,
640 id: new_id(),
641 repo,
642 agent: spec.id.clone(),
643 status: TalkStatus::Open,
644 turns: Vec::new(),
645 pending: String::new(),
646 pending_attachments: Vec::new(),
647 fallback: agent.is_none(),
648 created_at: now,
649 updated_at: now,
650 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
651 };
652 store.put(&mut talk)?;
653 Ok(talk)
654}
655
656pub fn record(
663 talk: &mut Talk,
664 store: &Talks,
665 text: &str,
666 attachments: Vec<Attachment>,
667) -> Result<String> {
668 let _guard = store.guard();
676 let Ok(fresh) = store.get(&talk.id) else {
681 bail!("talk {} was deleted", talk.short());
682 };
683 talk.status = fresh.status;
684 talk.pending = fresh.pending;
687 talk.pending_attachments = fresh.pending_attachments;
688 if !talk.status.open() {
689 bail!(
690 "talk {} is {} and takes no more turns",
691 talk.short(),
692 talk.status.as_str()
693 );
694 }
695 let text = text.trim();
696 if text.is_empty() && attachments.is_empty() {
697 bail!("nothing to say");
698 }
699 talk.turns.push(Turn {
700 who: Who::Operator,
701 body: text.to_owned(),
702 at: Timestamp::now(),
703 attachments,
704 usage: None,
705 });
706 store.put(talk)?;
707 Ok(text.to_owned())
708}
709
710pub fn queue(
712 talk: &mut Talk,
713 store: &Talks,
714 text: &str,
715 attachments: Vec<Attachment>,
716) -> Result<()> {
717 let text = text.trim();
718 if text.is_empty() && attachments.is_empty() {
719 bail!("nothing to say");
720 }
721 let _guard = store.guard();
722 let mut fresh = store
723 .get(&talk.id)
724 .with_context(|| format!("talk {} was deleted", talk.short()))?;
725 if !fresh.status.open() {
726 bail!(
727 "talk {} is {} and takes no more turns",
728 fresh.short(),
729 fresh.status.as_str()
730 );
731 }
732 if !text.is_empty() {
733 if fresh.pending.is_empty() {
734 fresh.pending = text.to_owned();
735 } else {
736 fresh.pending.push_str("\n\n");
737 fresh.pending.push_str(text);
738 }
739 }
740 fresh.pending_attachments.extend(attachments);
741 store.put(&mut fresh)?;
742 *talk = fresh;
743 Ok(())
744}
745
746pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
748 let _guard = store.guard();
749 let mut fresh = store
750 .get(&talk.id)
751 .with_context(|| format!("talk {} was deleted", talk.short()))?;
752 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
753 *talk = fresh;
754 return Ok(None);
755 }
756 let text = std::mem::take(&mut fresh.pending);
757 let attachments = std::mem::take(&mut fresh.pending_attachments);
758 fresh.turns.push(Turn {
759 who: Who::Operator,
760 body: text.clone(),
761 at: Timestamp::now(),
762 attachments,
763 usage: None,
764 });
765 store.put(&mut fresh)?;
766 *talk = fresh;
767 Ok(Some(text))
768}
769
770pub async fn say(
773 talk: &mut Talk,
774 store: &Talks,
775 cfg: &Config,
776 text: &str,
777 attachments: Vec<Attachment>,
778) -> Result<()> {
779 let text = record(talk, store, text, attachments)?;
780 turn(talk, store, cfg, &text).await
781}
782
783pub async fn respond(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
785 turn(talk, store, cfg, text).await
786}
787
788pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
806 let _guard = store.guard();
807 let mut fresh = store
808 .get(&talk.id)
809 .with_context(|| format!("talk {} was deleted", talk.short()))?;
810 fresh.status = TalkStatus::Closed;
811 fresh.pending.clear();
813 fresh.pending_attachments.clear();
814 store.put(&mut fresh)?;
815 *talk = fresh;
816 Ok(())
817}
818
819pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
830 let _guard = store.guard();
831 let mut fresh = store
832 .get(&talk.id)
833 .with_context(|| format!("talk {} was deleted", talk.short()))?;
834 fresh.status = TalkStatus::Open;
835 store.put(&mut fresh)?;
836 *talk = fresh;
837 Ok(())
838}
839
840pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
853 let _guard = store.guard();
854 let mut fresh = store
855 .get(&talk.id)
856 .with_context(|| format!("talk {} was deleted", talk.short()))?;
857 if fresh.agent == spec.id {
858 *talk = fresh;
859 return Ok(false);
860 }
861 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
862 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
863 fresh.fallback = false;
865 fresh.turns.push(Turn {
866 who: Who::Agent,
867 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
868 at: Timestamp::now(),
869 attachments: Vec::new(),
870 usage: None,
871 });
872 store.put(&mut fresh)?;
873 *talk = fresh;
874 Ok(true)
875}
876
877pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
879 let _guard = store.guard();
880 let mut fresh = store
881 .get(&talk.id)
882 .with_context(|| format!("talk {} was deleted", talk.short()))?;
883 fresh.pending.clear();
884 fresh.pending_attachments.clear();
885 store.put(&mut fresh)?;
886 *talk = fresh;
887 Ok(())
888}
889
890pub fn clear_pending_if_matches(
892 talk: &mut Talk,
893 store: &Talks,
894 expected_text: &str,
895 expected_attachments: &[String],
896) -> Result<bool> {
897 let _guard = store.guard();
898 let mut fresh = store
899 .get(&talk.id)
900 .with_context(|| format!("talk {} was deleted", talk.short()))?;
901 if !pending_matches(&fresh, expected_text, expected_attachments) {
902 *talk = fresh;
903 return Ok(false);
904 }
905 fresh.pending.clear();
906 fresh.pending_attachments.clear();
907 store.put(&mut fresh)?;
908 *talk = fresh;
909 Ok(true)
910}
911
912pub fn edit_pending_text(
916 talk: &mut Talk,
917 store: &Talks,
918 text: &str,
919 expected_text: &str,
920 expected_attachments: &[String],
921) -> Result<bool> {
922 let _guard = store.guard();
923 let mut fresh = store
924 .get(&talk.id)
925 .with_context(|| format!("talk {} was deleted", talk.short()))?;
926 if !pending_matches(&fresh, expected_text, expected_attachments) {
927 *talk = fresh;
928 return Ok(false);
929 }
930 fresh.pending = text.trim().to_owned();
931 store.put(&mut fresh)?;
932 *talk = fresh;
933 Ok(true)
934}
935
936fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
937 talk.pending == expected_text
938 && talk
939 .pending_attachments
940 .iter()
941 .map(|attachment| &attachment.id)
942 .eq(expected_attachments.iter())
943}
944
945async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
952 let spec = cfg
953 .agents
954 .iter()
955 .find(|a| a.id == talk.agent)
956 .with_context(|| {
957 format!(
958 "talk {} was opened with agent `{}`, which is no longer in \
959 the roster; restore it in magi.toml or start a new \
960 conversation",
961 talk.short(),
962 talk.agent
963 )
964 })?;
965
966 let last_note = attachment_note(
970 store,
971 &talk.id,
972 talk.turns
973 .last()
974 .map_or(&[][..], |t| t.attachments.as_slice()),
975 );
976
977 let attachment_paths: Vec<PathBuf> = talk
983 .turns
984 .iter()
985 .flat_map(|t| t.attachments.iter())
986 .filter_map(|a| store.attachment_path(&talk.id, a))
987 .collect();
988
989 let artifacts = store.artifacts_of(&talk.id);
990 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
993 let stem = format!("turn-{}", operator_turns.max(1));
994 let cache_dir = cfg.cache_dir();
997
998 let mut chain = vec![spec.clone()];
1001 if let Some(choice) = cfg.roles.chatter.as_ref()
1002 && talk.fallback
1003 {
1004 for id in choice.ids() {
1005 if id == talk.agent || chain.iter().any(|s| s.id == id) {
1006 continue;
1007 }
1008 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1009 Ok(s) => chain.push(s),
1010 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1011 }
1012 }
1013 }
1014
1015 let mut outcome = None;
1016 let mut fell_back_from: Option<String> = None;
1017 let mut first_try: Option<(String, SeatState)> = None;
1020 for (n, spec) in chain.iter().enumerate() {
1021 if n > 0 {
1022 if first_try.is_none() {
1023 first_try = Some((talk.agent.clone(), talk.seat.clone()));
1024 }
1025 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1026 fell_back_from.get_or_insert_with(|| talk.agent.clone());
1029 talk.agent = spec.id.clone();
1030 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1031 }
1032 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1033 let first_ever = talk.turns.len() <= 1;
1034 let body = if talk.seat.turns == 0 && first_ever {
1035 format!(
1036 "{}\n\n# Operator\n\n{text}{last_note}",
1037 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write)
1038 )
1039 } else if talk.seat.turns == 0 {
1040 format!(
1043 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1044 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write),
1045 transcript(talk, store)
1046 )
1047 } else if resuming {
1048 format!("{text}{last_note}")
1049 } else {
1050 format!("{}\n\n{text}{last_note}", transcript(talk, store))
1051 };
1052 let attempt_stem = if n == 0 {
1053 stem.clone()
1054 } else {
1055 format!("{stem}-{}", spec.id)
1056 };
1057 let inv = Invocation {
1058 cwd: &talk.repo,
1059 prompt: &body,
1060 timeout: turn_timeout(cfg),
1061 allow_write: cfg.talk.allow_write,
1066 sessions: cfg.graph.sessions,
1067 artifacts: &artifacts,
1068 stem: &attempt_stem,
1069 run: &talk.id,
1072 node: crate::queue::CHAT_NODE,
1073 cache_dir: cache_dir.as_deref(),
1074 attachments: &attachment_paths,
1075 writable: &[],
1076 };
1077 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1078 let advance = agent::chain_advances(&result);
1079 if n == 0 || !advance {
1080 outcome = Some(result);
1081 } else {
1082 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1084 }
1085 if !advance {
1086 break;
1087 }
1088 }
1089 if outcome.as_ref().is_some_and(agent::chain_advances) {
1090 if let Some((id, seat)) = first_try {
1093 talk.agent = id;
1094 talk.seat = seat;
1095 fell_back_from = None;
1096 }
1097 }
1098 let outcome = outcome.expect("a chain holds at least one agent");
1099 let note = |why: String| Turn {
1100 who: Who::Agent,
1101 body: format!("{MAGI_NOTE}{why}"),
1102 at: Timestamp::now(),
1103 attachments: Vec::new(),
1104 usage: None,
1105 };
1106 let (reply, failure) = match outcome {
1107 Err(e) => (
1108 note(format!("could not run agent `{}`: {e}", talk.agent)),
1109 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1110 ),
1111 Ok(out) if out.quota_exhausted() => {
1112 let reset = out
1113 .quota
1114 .as_ref()
1115 .and_then(|q| q.reset.clone())
1116 .map_or_else(String::new, |r| format!(" (resets {r})"));
1117 let why = format!(
1118 "agent `{}` is out of quota{reset}; your message is saved, so \
1119 say it again when the window reopens",
1120 talk.agent
1121 );
1122 (note(why.clone()), Some(why))
1123 }
1124 Ok(out) if out.timed_out => {
1125 let why = format!(
1126 "agent `{}` did not answer within {}s; your message is saved",
1127 talk.agent,
1128 turn_timeout(cfg).as_secs()
1129 );
1130 (note(why.clone()), Some(why))
1131 }
1132 Ok(out) if !out.usable() => {
1133 let why = format!(
1134 "agent `{}` produced no answer (exit {}); your message is saved",
1135 talk.agent,
1136 out.exit_code
1137 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1138 );
1139 (note(why.clone()), Some(why))
1140 }
1141 Ok(out) => (
1142 Turn {
1143 who: Who::Agent,
1144 body: out.text.trim().to_owned(),
1145 at: Timestamp::now(),
1146 attachments: Vec::new(),
1147 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1150 context_tokens,
1151 agent: talk.agent.clone(),
1152 model: cfg
1153 .agents
1154 .iter()
1155 .find(|a| a.id == talk.agent)
1156 .and_then(|a| a.model.clone()),
1157 }),
1158 },
1159 None,
1160 ),
1161 };
1162
1163 let _guard = store.guard();
1175 let Ok(fresh) = store.get(&talk.id) else {
1181 return Ok(());
1182 };
1183 talk.status = fresh.status;
1184 talk.pending = fresh.pending;
1188 talk.pending_attachments = fresh.pending_attachments;
1189 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1190 talk.turns.push(note(format!(
1193 "agent changed from {from} to {} (fallback)",
1194 talk.agent
1195 )));
1196 }
1197 talk.turns.push(reply);
1198 if let Err(put_err) = store.put(talk) {
1199 let lost = talk.turns.pop().expect("just pushed above");
1208 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1209 let why = match &stash {
1210 Ok(path) => format!(
1211 "agent `{}` answered, but the reply could not be saved to \
1212 this conversation ({put_err:#}); the raw text was kept at \
1213 {} - your message is saved, ask again",
1214 talk.agent,
1215 path.display()
1216 ),
1217 Err(stash_err) => format!(
1218 "agent `{}` answered, but the reply could not be saved to \
1219 this conversation ({put_err:#}), and it could not be kept \
1220 anywhere else either ({stash_err:#}); your message is \
1221 saved, ask again",
1222 talk.agent
1223 ),
1224 };
1225 talk.turns.push(note(why.clone()));
1226 return match store.put(talk) {
1233 Ok(()) => bail!("{why}"),
1234 Err(note_err) => {
1235 talk.turns.pop();
1255 Err(note_err).context(why)
1256 }
1257 };
1258 }
1259
1260 match failure {
1261 Some(why) => bail!("{why}"),
1262 None => Ok(()),
1263 }
1264}
1265
1266fn transcript(talk: &Talk, store: &Talks) -> String {
1269 let mut out = String::from(
1270 "This conversation cannot resume on the CLI's side, so here is \
1271 everything said so far; answer only the last message.\n",
1272 );
1273 for t in &talk.turns {
1274 let who = match t.who {
1275 Who::Operator => "operator",
1276 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1277 Who::Agent => "you",
1278 };
1279 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1280 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1281 }
1282 out
1283}
1284
1285fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1290 if attachments.is_empty() {
1291 return String::new();
1292 }
1293 let mut out = String::from(
1294 "\n\nThe operator attached the image(s) below to this message. Open \
1295 and look at each one before you answer.\n",
1296 );
1297 for att in attachments {
1298 if let Some(path) = store.attachment_path(talk_id, att) {
1299 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1300 }
1301 }
1302 out.push('\n');
1303 out
1304}
1305
1306pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
1328 let write_policy = if allow_write {
1329 "Write access is enabled for this conversation (`allow_write = \
1330 true`), so you may write files - but only a small, \
1331 already-decided edit the operator names outright in this \
1332 conversation, not an implementation. This is a permission on the \
1333 conversation as a whole, not a property of whichever repository \
1334 it happened to start in: if the operator names a different \
1335 repository for that small edit, the policy allows it there too. \
1336 Your own tool may still confine writes to the repository this \
1337 conversation started in regardless - if a write elsewhere is \
1338 refused, say so plainly rather than working around it. Once you \
1339 have made an edit, say plainly what you edited. Anything bigger, \
1340 or anything still open-ended, still goes through the queue below \
1341 rather than being done here."
1342 } else {
1343 "Do not write files. Implementing a change is not this \
1344 conversation's job; a separate, blind competition of agents does \
1345 that, and a repository this conversation has already edited would \
1346 make their diffs unjudgeable."
1347 };
1348 let mut out = format!(
1349 "You are magi's standing conversation partner for its operator, who \
1350 usually has this open on a phone. Keep replies short: no preamble, \
1351 no restating what they just said.\n\n\
1352 # Repository\n\n{repo}\n\n\
1353 You may look around: read files, run shell commands, search history, \
1354 run tests - whatever answers the question. {write_policy}\n\n\
1355 A short, command-shaped message (\"list\", \"info <id>\", \"show \
1356 3cbf\") is almost always the operator asking you to look something \
1357 up, not an instruction to file - answer it yourself with `magi \
1358 list`, `magi show <id>`, `magi task list`, or the like, the same way \
1359 you would answer any other question in this conversation.\n\n\
1360 # When the operator wants something done\n\n\
1361 Run:\n\n\
1362 magi task add --solo --repo {repo} <instruction>\n\n\
1363 and tell the operator the task id it prints, so they can follow it \
1364 from the Queue. If it refuses with a duplicate warning (the \
1365 instruction names a branch, commit or pull request that an \
1366 unfinished task, run or PR already owns), do not repeat it with \
1367 --force yourself: tell the operator what it matched and let them \
1368 decide. Write <instruction> so that an implementer who has \
1369 never seen this conversation can act on it alone - it is everything \
1370 they get. Use --solo: it runs the task through one implementer \
1371 straight into review instead of the usual multi-agent competition, \
1372 which is the right shape for a change this conversation has already \
1373 settled, rather than one still worth several independent takes.\n\n\
1374 If the operator asks for something in a different repository, \
1375 --repo does not have to be a full path: --repo owner/repo (or just \
1376 repo, when that is unambiguous) is resolved against local checkouts \
1377 the same way `magi repos` lists them. If the command fails because \
1378 nothing matches or more than one checkout shares that name, ask the \
1379 operator which repository they mean (or run `magi repos` yourself \
1380 to see the candidates) rather than guessing.\n\n\
1381 The current state of the code is whatever origin/main holds, not \
1382 whatever a working tree shows: a primary checkout often lags \
1383 upstream, sits on a detached HEAD and carries uncommitted changes. \
1384 Before answering about code, run `git fetch origin` in that \
1385 repository if it is cheap, then read through \
1386 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
1387 If the working tree differs, say so; if the fetch fails, say that \
1388 too, so the operator knows the answer may be stale.\n\n\
1389 If the operator attached an image (a screenshot, say) that the task \
1390 is about, pass it with `--attach <path>`, using the absolute path \
1391 the turn's attachment note gives; repeat the flag for several. \
1392 `magi task add --solo --attach <path> <instruction>` copies the \
1393 file into the task, so the implementer receives it. Do not paste the \
1394 path into <instruction> instead: deleting this conversation deletes \
1395 its attachments, and then that path reaches no one.\n",
1396 repo = repo.display(),
1397 );
1398 out.push_str(&language_note(language));
1399 out
1400}
1401
1402fn language_note(language: &str) -> String {
1405 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
1406 String::new()
1407 } else {
1408 format!("\nHold this conversation in {language}.\n")
1409 }
1410}
1411
1412pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
1419 let mut tasks: Vec<Task> = queue
1420 .list()
1421 .into_iter()
1422 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
1423 .collect();
1424 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
1425 tasks
1426}
1427
1428fn read_path(path: &Path) -> Result<Talk> {
1429 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1430 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
1431}
1432
1433const PUT_RETRIES: u32 = 5;
1436
1437fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1449 let mut last_err = None;
1450 for attempt in 0..PUT_RETRIES {
1451 if attempt > 0 {
1452 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
1453 }
1454 match try_write_atomic(tmp, path, body) {
1455 Ok(()) => return Ok(()),
1456 Err(e) => last_err = Some(e),
1457 }
1458 }
1459 Err(last_err.expect("the loop above always runs at least once"))
1460}
1461
1462fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1463 #[cfg(test)]
1464 if failpoint::take_forced_put_failure() {
1465 bail!("simulated write failure (test)");
1466 }
1467 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
1468 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
1469 Ok(())
1470}
1471
1472fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
1477 let dir = store.artifacts_of(id);
1478 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1479 let path = dir.join(format!("{stem}-lost.txt"));
1480 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
1481 Ok(path)
1482}
1483
1484#[cfg(test)]
1491mod failpoint {
1492 use std::cell::Cell;
1493
1494 thread_local! {
1495 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
1496 }
1497
1498 pub(super) fn force_put_failures(count: u32) {
1501 FORCE_PUT_FAILURES.with(|c| c.set(count));
1502 }
1503
1504 pub(super) fn take_forced_put_failure() -> bool {
1507 FORCE_PUT_FAILURES.with(|c| {
1508 let n = c.get();
1509 if n == 0 {
1510 false
1511 } else {
1512 c.set(n - 1);
1513 true
1514 }
1515 })
1516 }
1517}
1518
1519fn short(id: &str) -> &str {
1520 id.split('-').next_back().unwrap_or(id)
1521}
1522
1523fn new_id() -> String {
1524 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1525 let seed = crate::rng::entropy();
1526 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1527}
1528
1529fn attachment_ext(mime: &str) -> Option<&'static str> {
1534 match mime {
1535 "image/png" => Some("png"),
1536 "image/jpeg" => Some("jpg"),
1537 "image/gif" => Some("gif"),
1538 "image/webp" => Some("webp"),
1539 _ => None,
1540 }
1541}
1542
1543pub fn valid_attachment_id(id: &str) -> bool {
1548 id.len() == 32
1549 && id
1550 .bytes()
1551 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
1552}
1553
1554fn new_attachment_id() -> String {
1558 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
1559 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
1560}
1561
1562#[cfg(test)]
1563mod tests {
1564 #[test]
1565 fn the_briefing_points_at_origin_main_not_the_working_tree() {
1566 let b = briefing(Path::new("/r"), "en", false);
1567 assert!(b.contains("origin/main"));
1568 assert!(b.contains("git show origin/main:"));
1569 }
1570 use std::collections::BTreeMap;
1571
1572 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
1573 use crate::queue::{Queue, Source, Task};
1574
1575 use super::*;
1576
1577 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
1578 AgentSpec {
1579 id: id.to_owned(),
1580 kind: AgentKind::Command,
1581 model: model.map(str::to_owned),
1582 command: Vec::new(),
1583 extra_args: Vec::new(),
1584 env: BTreeMap::new(),
1585 prompt_delivery: None,
1586 }
1587 }
1588
1589 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
1590 Talk {
1591 schema: SCHEMA,
1592 id: "20260904-014455-ab12".to_owned(),
1593 repo: PathBuf::from("."),
1594 agent: agent.to_owned(),
1595 status: TalkStatus::Open,
1596 turns,
1597 pending: String::new(),
1598 pending_attachments: Vec::new(),
1599 fallback: false,
1600 created_at: Timestamp::now(),
1601 updated_at: Timestamp::now(),
1602 seat: SeatState::new(SEAT, agent, 1),
1603 }
1604 }
1605
1606 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
1607 Turn {
1608 who: Who::Agent,
1609 body: body.to_owned(),
1610 at: Timestamp::now(),
1611 attachments: Vec::new(),
1612 usage: usage.map(|(t, a, m)| TurnUsage {
1613 context_tokens: t,
1614 agent: a.to_owned(),
1615 model: m.map(str::to_owned),
1616 }),
1617 }
1618 }
1619
1620 fn ctx_config(windows: &[(&str, u64)]) -> Config {
1621 Config {
1622 agents: vec![
1623 ctx_agent("small", Some("small-model")),
1624 ctx_agent("big", Some("big-model")),
1625 ctx_agent("plain", None),
1626 ],
1627 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
1628 ..Config::default()
1629 }
1630 }
1631
1632 #[test]
1633 fn context_usage_computes_percent_and_warns_at_eighty() {
1634 let cfg = ctx_config(&[("small-model", 1000)]);
1635 let at = |tokens| {
1636 let t = ctx_talk(
1637 "small",
1638 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
1639 );
1640 context_usage(&t, Some(&cfg))
1641 };
1642 let u = at(799);
1643 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
1644 let u = at(800);
1645 assert_eq!((u.percent, u.warn), (Some(80), true));
1646 let u = at(1500);
1647 assert_eq!((u.percent, u.warn), (Some(150), true));
1648 assert!(!u.since_switch);
1649 }
1650
1651 #[test]
1652 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
1653 let cfg = ctx_config(&[("small-model", 1000)]);
1654 let t = ctx_talk(
1655 "small",
1656 vec![
1657 reply("old", Some((900, "small", Some("small-model")))),
1658 reply("new", None),
1659 ],
1660 );
1661 let u = context_usage(&t, Some(&cfg));
1662 assert!(u.estimated);
1664 assert_ne!(u.tokens, Some(900));
1665 assert!(u.tokens.is_some());
1666 let t = ctx_talk(
1668 "small",
1669 vec![
1670 reply("old", Some((900, "small", Some("small-model")))),
1671 reply("magi: could not run agent", None),
1672 ],
1673 );
1674 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
1675 assert_eq!(
1676 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
1677 None
1678 );
1679 }
1680
1681 #[test]
1682 fn estimate_counts_chars_both_sides_and_standing_prompt() {
1683 let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
1684 assert_eq!(estimate_context_tokens(&t, 0), Some(2)); let op = Turn {
1686 who: Who::Operator,
1687 ..reply("abcdefg", None)
1688 };
1689 t.turns.push(op);
1690 assert_eq!(estimate_context_tokens(&t, 0), Some(4));
1691 assert!(
1692 estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
1693 );
1694 let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
1696 assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
1697 let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
1699 assert_eq!(estimate_context_tokens(¬e, 1000), None);
1700 assert_eq!(
1701 estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
1702 None
1703 );
1704 }
1705
1706 #[test]
1707 fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
1708 let cfg = ctx_config(&[("small-model", 1000)]);
1709 let t = ctx_talk(
1710 "small",
1711 vec![reply(
1712 &"x".repeat(5000),
1713 Some((10, "small", Some("small-model"))),
1714 )],
1715 );
1716 let u = context_usage(&t, Some(&cfg));
1717 assert_eq!((u.tokens, u.estimated), (Some(10), false));
1718 let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
1719 let u = context_usage(&t, Some(&cfg));
1720 assert!(u.estimated && !u.since_switch);
1721 assert_eq!(u.window, Some(1000));
1722 assert!(u.warn && u.percent.unwrap() >= 80);
1723 let t = ctx_talk("small", vec![reply("hi", None)]);
1724 let u = context_usage(&t, Some(&cfg));
1725 assert!(u.estimated && u.percent.is_some());
1726 }
1727
1728 #[test]
1729 fn context_usage_without_a_window_shows_tokens_only() {
1730 let cfg = ctx_config(&[]);
1731 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
1733 let u = context_usage(&t, Some(&cfg));
1734 assert_eq!(
1735 (u.tokens, u.window, u.percent, u.warn),
1736 (Some(5000), None, None, false)
1737 );
1738 let t = ctx_talk(
1739 "small",
1740 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
1741 );
1742 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
1743 assert_eq!(context_usage(&t, None).window, None);
1745 }
1746
1747 #[test]
1748 fn context_usage_switching_model_changes_the_denominator() {
1749 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
1750 let used = reply("hi", Some((900, "small", Some("small-model"))));
1751 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
1752 assert_eq!(
1753 (before.percent, before.warn, before.since_switch),
1754 (Some(90), true, false)
1755 );
1756 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
1759 assert_eq!(after.window, Some(10_000));
1760 assert_eq!(
1761 (after.percent, after.warn, after.since_switch),
1762 (Some(9), false, true)
1763 );
1764 assert_eq!(after.model.as_deref(), Some("big-model"));
1765 }
1766
1767 #[test]
1768 fn a_turn_recorded_before_usage_existed_still_reads() {
1769 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
1770 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
1771 assert!(turn.usage.is_none());
1772 let json = serde_json::to_string(&turn).expect("serialize");
1773 assert!(
1774 !json.contains("usage"),
1775 "absent usage is not written: {json}"
1776 );
1777 }
1778
1779 fn store() -> (tempfile::TempDir, Talks) {
1781 let tmp = tempfile::tempdir().expect("tempdir");
1782 let talks = Talks::at(tmp.path().join("talks"));
1783 (tmp, talks)
1784 }
1785
1786 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
1790 let path = dir.join("mock-talk-agent.sh");
1791 std::fs::write(&path, script).expect("write mock");
1792 AgentSpec {
1793 id: "mock".to_owned(),
1794 kind: AgentKind::Command,
1795 model: None,
1796 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
1797 extra_args: Vec::new(),
1798 env,
1799 prompt_delivery: None,
1800 }
1801 }
1802
1803 fn config(spec: AgentSpec) -> Config {
1804 Config {
1805 agents: vec![spec],
1806 graph: Graph {
1807 language: "en".to_owned(),
1808 ..Graph::default()
1809 },
1810 ..Config::default()
1811 }
1812 }
1813
1814 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
1816
1817 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
1819
1820 const ECHO: &str = "#!/bin/sh\ncat\n";
1823
1824 fn env(reply: &str) -> BTreeMap<String, String> {
1825 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
1826 }
1827
1828 #[test]
1829 fn the_frozen_json_field_names_round_trip_through_disk() {
1830 let (tmp, talks) = store();
1831 let mut talk = Talk {
1832 schema: SCHEMA,
1833 id: "20260904-014455-ab12".to_owned(),
1834 repo: tmp.path().to_owned(),
1835 agent: "sonnet".to_owned(),
1836 status: TalkStatus::Open,
1837 turns: Vec::new(),
1838 pending: String::new(),
1839 pending_attachments: Vec::new(),
1840 fallback: false,
1841 created_at: Timestamp::now(),
1842 updated_at: Timestamp::now(),
1843 seat: SeatState::new(SEAT, "sonnet", 7),
1844 };
1845 talks.put(&mut talk).expect("put");
1846
1847 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
1848 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
1849 for field in [
1850 "schema",
1851 "id",
1852 "repo",
1853 "agent",
1854 "status",
1855 "turns",
1856 "created_at",
1857 "updated_at",
1858 ] {
1859 assert!(v.get(field).is_some(), "missing field `{field}`");
1860 }
1861 assert_eq!(v["schema"], 1);
1862 assert_eq!(v["status"], "open");
1863
1864 let back = talks.get(&talk.id).expect("get");
1865 assert_eq!(back.id, talk.id);
1866 assert_eq!(back.status, TalkStatus::Open);
1867 }
1868
1869 #[test]
1870 fn opening_a_talk_takes_no_agent_turn() {
1871 let (tmp, talks) = store();
1872 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1876 let cfg = config(spec);
1877
1878 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1879 assert_eq!(talk.status, TalkStatus::Open);
1880 assert!(talk.turns.is_empty(), "nothing has been said yet");
1881
1882 let on_disk = talks.get(&talk.id).expect("get");
1883 assert_eq!(on_disk.turns.len(), 0);
1884 }
1885
1886 #[test]
1894 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
1895 let (tmp, talks) = store();
1896 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1897 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1898 chatter_spec.id = "chatter-mock".to_owned();
1899
1900 let mut cfg = Config {
1901 agents: vec![first_spec.clone(), chatter_spec.clone()],
1902 graph: Graph {
1903 language: "en".to_owned(),
1904 ..Graph::default()
1905 },
1906 ..Config::default()
1907 };
1908 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
1909
1910 let talk =
1911 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
1912 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
1913
1914 cfg.roles.chatter = None;
1915 let fallback =
1916 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
1917 assert_eq!(
1918 fallback.agent, first_spec.id,
1919 "unset chatter must fall back to agent::pick's own default order"
1920 );
1921 }
1922
1923 #[test]
1926 fn a_talk_recorded_without_attachments_still_reads() {
1927 let (tmp, talks) = store();
1928 let path = talks.path_of("20260904-014455-ab12");
1929 std::fs::create_dir_all(talks.root()).expect("talks dir");
1930 std::fs::write(
1931 &path,
1932 serde_json::json!({
1933 "schema": 1,
1934 "id": "20260904-014455-ab12",
1935 "repo": tmp.path(),
1936 "agent": "sonnet",
1937 "status": "open",
1938 "turns": [
1939 { "who": "operator", "body": "still there?",
1940 "at": Timestamp::now().to_string() },
1941 ],
1942 "created_at": Timestamp::now().to_string(),
1943 "updated_at": Timestamp::now().to_string(),
1944 "seat": SeatState::new(SEAT, "sonnet", 7),
1945 })
1946 .to_string(),
1947 )
1948 .expect("write pre-attachments talk");
1949
1950 let talk = talks.get("20260904-014455-ab12").expect("must still read");
1951 assert!(talk.turns[0].attachments.is_empty());
1952 }
1953
1954 #[test]
1955 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
1956 let (tmp, talks) = store();
1957 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1958 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1959
1960 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
1961 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
1962 let saved = talks.get(&talk.id).expect("reload queued talk");
1963 assert_eq!(saved.pending, "first\n\nsecond");
1964 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
1965
1966 let drained = drain(&mut talk, &talks).expect("drain");
1967 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
1968 let saved = talks.get(&talk.id).expect("reload drained talk");
1969 assert!(saved.pending.is_empty());
1970 assert_eq!(saved.turns.len(), 1);
1971 assert_eq!(saved.turns[0].body, "first\n\nsecond");
1972 }
1973
1974 #[test]
1975 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
1976 let (tmp, talks) = store();
1977 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1978 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1979 let attachment = Attachment {
1980 id: "a".repeat(32),
1981 name: "shot.png".to_owned(),
1982 mime: "image/png".to_owned(),
1983 bytes: 3,
1984 };
1985
1986 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
1987 assert!(
1988 edit_pending_text(
1989 &mut talk,
1990 &talks,
1991 "corrected",
1992 "first",
1993 std::slice::from_ref(&attachment.id),
1994 )
1995 .expect("edit")
1996 );
1997 let saved = talks.get(&talk.id).expect("reload edited draft");
1998 assert_eq!(saved.pending, "corrected");
1999 assert_eq!(saved.pending_attachments, vec![attachment]);
2000
2001 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
2002 assert!(
2003 !edit_pending_text(
2004 &mut talk,
2005 &talks,
2006 "stale edit",
2007 "corrected",
2008 &["a".repeat(32)],
2009 )
2010 .expect("stale edit is a conflict")
2011 );
2012 assert_eq!(
2013 talks.get(&talk.id).expect("reload after conflict").pending,
2014 "corrected\n\nlater"
2015 );
2016 assert!(
2017 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
2018 .expect("stale clear is a conflict")
2019 );
2020 assert_eq!(
2021 talks
2022 .get(&talk.id)
2023 .expect("reload after stale clear")
2024 .pending,
2025 "corrected\n\nlater"
2026 );
2027 }
2028
2029 #[tokio::test]
2030 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
2031 let (tmp, talks) = store();
2032 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
2033 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
2034 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2035 let id = running.id.clone();
2036 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
2037
2038 let response_talks = talks.clone();
2039 let response_cfg = cfg.clone();
2040 let reply = tokio::spawn(async move {
2041 respond(&mut running, &response_talks, &response_cfg, &first).await
2042 });
2043 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
2044
2045 let mut queued = talks.get(&id).expect("queued handle");
2046 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
2047 reply.await.expect("join").expect("reply");
2048
2049 let saved = talks.get(&id).expect("reload");
2050 assert_eq!(saved.pending, "next");
2051 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
2052 }
2053
2054 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
2057 let calls = dir.join(format!("{id}.calls"));
2058 let script = format!(
2059 "#!/bin/sh\necho x >> '{}'\n{body}\n",
2060 calls.to_string_lossy()
2061 );
2062 let path = dir.join(format!("mock-{id}.sh"));
2063 std::fs::write(&path, script).expect("write mock");
2064 AgentSpec {
2065 id: id.to_owned(),
2066 kind: AgentKind::Command,
2067 model: None,
2068 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2069 extra_args: Vec::new(),
2070 env: BTreeMap::new(),
2071 prompt_delivery: None,
2072 }
2073 }
2074
2075 fn calls(dir: &Path, id: &str) -> usize {
2076 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
2077 }
2078
2079 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
2080 let mut cfg = config(specs[0].clone());
2081 cfg.agents = specs;
2082 cfg.roles.chatter = Some(AgentChoice::Chain(
2083 ids.iter().map(|s| (*s).to_owned()).collect(),
2084 ));
2085 cfg
2086 }
2087
2088 #[tokio::test]
2089 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
2090 let (tmp, talks) = store();
2091 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2092 let b = counting_agent(tmp.path(), "b", "cat");
2093 let cfg = chain_config(vec![a, b], &["a", "b"]);
2094 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2095 assert_eq!(talk.agent, "a");
2096
2097 say(&mut talk, &talks, &cfg, "hello there", Vec::new())
2098 .await
2099 .expect("turn");
2100 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
2101 assert_eq!(calls(tmp.path(), "b"), 1);
2102 assert_eq!(talk.agent, "b", "the switch persists");
2103 assert!(talks.get(&talk.id).unwrap().agent == "b");
2104 let reply = talk.turns.last().unwrap();
2105 assert!(reply.body.contains("hello there"));
2106 assert!(
2107 reply.body.contains("magi task add --solo"),
2108 "a fresh seat gets the full briefing"
2109 );
2110 assert!(
2111 talk.turns
2112 .iter()
2113 .any(|t| t.body.contains("agent changed from a to b")),
2114 "the switch is noted"
2115 );
2116 }
2117
2118 #[tokio::test]
2119 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
2120 let (tmp, talks) = store();
2121 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2122 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
2123 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
2124 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2125
2126 let err = say(&mut talk, &talks, &cfg, "hi", Vec::new())
2127 .await
2128 .expect_err("every agent failed");
2129 assert!(err.to_string().contains("`a`"), "{err:#}");
2130 assert_eq!(calls(tmp.path(), "a"), 1);
2131 assert_eq!(calls(tmp.path(), "b"), 1);
2132 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
2133 }
2134
2135 #[test]
2136 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
2137 let (tmp, talks) = store();
2138 let b = counting_agent(tmp.path(), "b", "cat");
2139 let cfg = chain_config(vec![b], &["ghost", "b"]);
2140 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2141 assert_eq!(talk.agent, "b");
2142 }
2143
2144 #[tokio::test]
2145 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
2146 let (tmp, talks) = store();
2147 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2148 let b = counting_agent(tmp.path(), "b", "cat");
2149 let cfg = chain_config(vec![a, b], &["a", "b"]);
2150 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
2151 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2152 .await
2153 .expect_err("a alone, and it fails");
2154 assert_eq!(calls(tmp.path(), "b"), 0);
2155 assert_eq!(talk.agent, "a");
2156 }
2157
2158 #[tokio::test]
2159 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
2160 let (tmp, talks) = store();
2161 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2162 let b = counting_agent(tmp.path(), "b", "cat");
2163 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
2164 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
2165 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
2166 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2167 .await
2168 .expect_err("c alone, and it fails");
2169 assert_eq!(calls(tmp.path(), "b"), 0);
2170 }
2171
2172 #[tokio::test]
2173 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
2174 let (tmp, talks) = store();
2175 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2176 let cfg = config(spec);
2177 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2178
2179 say(
2180 &mut talk,
2181 &talks,
2182 &cfg,
2183 "what does the queue module do?",
2184 Vec::new(),
2185 )
2186 .await
2187 .expect("first turn");
2188 let first_prompt = &talk.turns[1].body;
2189 assert!(first_prompt.contains("magi task add --solo"));
2190 assert!(first_prompt.contains("what does the queue module do?"));
2191
2192 say(&mut talk, &talks, &cfg, "and how is it locked?", Vec::new())
2193 .await
2194 .expect("second turn");
2195 let second_prompt = &talk.turns[3].body;
2196 assert!(
2197 !second_prompt.contains("magi task add --solo"),
2198 "the briefing is sent once, not on every turn: {second_prompt}"
2199 );
2200 assert!(second_prompt.contains("and how is it locked?"));
2201 }
2202
2203 #[tokio::test]
2204 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
2205 let (tmp, talks) = store();
2206 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2207 let mut b = a.clone();
2208 b.id = "other".to_owned();
2209 let mut cfg = config(a.clone());
2210 cfg.agents.push(b.clone());
2211 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
2212 say(&mut talk, &talks, &cfg, "remember the walrus", Vec::new())
2213 .await
2214 .expect("first turn");
2215 let old_session = talk.seat.claude_session.clone();
2216 assert_eq!(talk.seat.turns, 1);
2217
2218 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
2219 assert_eq!(talk.agent, "other");
2220 assert_eq!(talk.seat.turns, 0);
2221 assert_eq!(talk.seat.agent, "other");
2222 assert_ne!(talk.seat.claude_session, old_session);
2223 let note = talk.turns.last().expect("note");
2224 assert_eq!(note.who, Who::Agent);
2225 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2226 assert!(note.body.contains("changed from"), "{}", note.body);
2227 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
2228
2229 let before = talk.turns.len();
2230 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
2231 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
2232
2233 say(&mut talk, &talks, &cfg, "what did I say?", Vec::new())
2234 .await
2235 .expect("turn after switch");
2236 let prompt = &talk.turns.last().expect("reply").body;
2237 assert!(prompt.contains("remember the walrus"), "{prompt}");
2238 assert!(prompt.contains("## magi"), "{prompt}");
2239 assert!(prompt.contains("what did I say?"), "{prompt}");
2240 }
2241
2242 #[tokio::test]
2243 async fn say_appends_the_operator_turn_then_the_agent_turn() {
2244 let (tmp, talks) = store();
2245 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2246 let cfg = config(spec);
2247 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2248
2249 say(
2250 &mut talk,
2251 &talks,
2252 &cfg,
2253 "can I rename this function?",
2254 Vec::new(),
2255 )
2256 .await
2257 .expect("say");
2258
2259 assert_eq!(talk.turns.len(), 2);
2260 assert_eq!(talk.turns[0].who, Who::Operator);
2261 assert_eq!(talk.turns[0].body, "can I rename this function?");
2262 assert_eq!(talk.turns[1].who, Who::Agent);
2263 assert_eq!(talk.turns[1].body, "go ahead");
2264 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
2265 }
2266
2267 #[tokio::test]
2268 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
2269 let (tmp, talks) = store();
2270 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2271 let cfg = config(spec);
2272 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2273
2274 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2275 .await
2276 .expect_err("a turn with no answer is an error");
2277 assert!(err.to_string().contains("no answer"), "{err}");
2278
2279 let on_disk = talks.get(&talk.id).expect("get");
2280 assert_eq!(on_disk.turns.len(), 2);
2281 assert_eq!(on_disk.turns[0].body, "check the tests");
2282 let note = &on_disk.turns[1];
2283 assert_eq!(note.who, Who::Agent);
2284 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2285 assert!(note.body.contains("your message is saved"));
2286 }
2287
2288 #[tokio::test]
2294 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
2295 let (tmp, talks) = store();
2296 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2297 let cfg = config(spec);
2298 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2299
2300 let text =
2301 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
2302 failpoint::force_put_failures(PUT_RETRIES - 1);
2305 respond(&mut talk, &talks, &cfg, &text)
2306 .await
2307 .expect("respond must survive a write failure its own retries can outlast");
2308
2309 assert_eq!(talk.turns.len(), 2);
2310 assert_eq!(talk.turns[1].who, Who::Agent);
2311 assert_eq!(talk.turns[1].body, "go ahead");
2312 let on_disk = talks.get(&talk.id).expect("get");
2313 assert_eq!(
2314 on_disk.turns, talk.turns,
2315 "the reply must reach disk despite the early write failures"
2316 );
2317 }
2318
2319 #[tokio::test]
2325 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
2326 let (tmp, talks) = store();
2327 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2328 let cfg = config(spec);
2329 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2330
2331 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2332 failpoint::force_put_failures(PUT_RETRIES);
2337 let err = respond(&mut talk, &talks, &cfg, &text)
2338 .await
2339 .expect_err("a reply that cannot be saved must be reported, not swallowed");
2340 assert!(err.to_string().contains("could not be saved"), "{err}");
2341
2342 let on_disk = talks.get(&talk.id).expect("get");
2343 assert_eq!(
2344 on_disk.turns.len(),
2345 2,
2346 "the operator turn plus a visible note"
2347 );
2348 assert_eq!(on_disk.turns[0].body, "check the tests");
2349 let note = &on_disk.turns[1];
2350 assert_eq!(note.who, Who::Agent);
2351 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2352 assert!(
2353 note.body.contains("could not be saved"),
2354 "the operator must be told the reply is missing, not left staring \
2355 at a gap with no explanation: {}",
2356 note.body
2357 );
2358 assert_eq!(
2359 talk.turns, on_disk.turns,
2360 "the in-memory talk must match what actually landed on disk"
2361 );
2362
2363 let artifacts = talks.artifacts_of(&talk.id);
2366 let stash = std::fs::read_dir(&artifacts)
2367 .expect("artifacts dir")
2368 .filter_map(|e| e.ok())
2369 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
2370 .expect("a stash file for the lost reply");
2371 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
2372 assert_eq!(stashed, "go ahead");
2373
2374 assert_eq!(
2382 on_disk.seat.turns, 1,
2383 "the note's write must carry the turn the CLI actually took"
2384 );
2385 assert_eq!(
2386 on_disk.seat.claude_session, talk.seat.claude_session,
2387 "the session id handed to the CLI must survive the failed reply"
2388 );
2389 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
2390 assert!(
2391 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
2392 "the next turn must resume, not open the same session id twice"
2393 );
2394 }
2395
2396 #[tokio::test]
2401 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
2402 let (tmp, talks) = store();
2403 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2404 let cfg = config(spec);
2405 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2406
2407 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2408 failpoint::force_put_failures(PUT_RETRIES * 2);
2411 let err = respond(&mut talk, &talks, &cfg, &text)
2412 .await
2413 .expect_err("neither the reply nor the note could be saved");
2414 assert!(err.to_string().contains("could not be saved"), "{err}");
2415
2416 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
2417 let on_disk = talks.get(&talk.id).expect("get");
2418 assert_eq!(on_disk.turns.len(), 1);
2419
2420 assert_eq!(
2430 on_disk.seat.turns, 0,
2431 "an unwritable file cannot record the turn the CLI took"
2432 );
2433 assert_eq!(
2434 talk.seat.turns, 1,
2435 "the in-memory seat still reports the turn the CLI actually took"
2436 );
2437 assert_eq!(
2438 on_disk.seat.claude_session, talk.seat.claude_session,
2439 "the session id was minted at `begin` and never changes here"
2440 );
2441 }
2442
2443 #[tokio::test]
2447 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
2448 let (tmp, talks) = store();
2449 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2450 let cfg = config(spec);
2451 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2452
2453 let att = talks
2454 .put_attachment(
2455 &talk.id,
2456 "image/png",
2457 "screenshot.png",
2458 b"pretend-png-bytes",
2459 )
2460 .expect("put attachment");
2461
2462 say(&mut talk, &talks, &cfg, "", vec![att.clone()])
2463 .await
2464 .expect("an empty body with an attachment is still a turn");
2465
2466 let operator_turn = &talk.turns[0];
2467 assert_eq!(operator_turn.who, Who::Operator);
2468 assert_eq!(operator_turn.body, "");
2469 assert_eq!(operator_turn.attachments, vec![att.clone()]);
2470
2471 let prompt = &talk.turns[1].body;
2472 let expected_path = talks
2473 .attachments_dir(&talk.id)
2474 .join(format!("{}.png", att.id));
2475 assert!(
2476 prompt.contains(&expected_path.display().to_string()),
2477 "the agent must be told the attachment's absolute path: {prompt}"
2478 );
2479 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
2480 }
2481
2482 #[test]
2491 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
2492 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
2493 let att = Attachment {
2494 id: "0".repeat(32),
2495 name: "shot.png".to_owned(),
2496 mime: "image/png".to_owned(),
2497 bytes: 3,
2498 };
2499 let path = talks
2500 .attachment_path("some-talk-id", &att)
2501 .expect("a supported mime always yields a path");
2502 assert!(
2503 path.is_absolute(),
2504 "must be absolute even off a relative store root: {}",
2505 path.display()
2506 );
2507 }
2508
2509 #[tokio::test]
2510 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
2511 let (tmp, talks) = store();
2516 let slow = mock_agent(
2517 tmp.path(),
2518 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
2519 BTreeMap::new(),
2520 );
2521 let mut cfg = config(slow);
2522 cfg.graph.timeout_talk = 1;
2523 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2524
2525 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2526 .await
2527 .expect_err("a turn that never answers is an error");
2528 assert!(
2529 err.to_string().contains("did not answer within 1s"),
2530 "{err}"
2531 );
2532
2533 let on_disk = talks.get(&talk.id).expect("get");
2534 let note = on_disk.turns.last().expect("a note turn was recorded");
2535 assert!(
2536 note.body.contains("did not answer within 1s"),
2537 "the transcript must show the configured timeout: {}",
2538 note.body
2539 );
2540 }
2541
2542 #[test]
2543 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
2544 let (tmp, talks) = store();
2545 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2546 let cfg = config(spec);
2547 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2548
2549 close(&mut talk, &talks).expect("close");
2550 assert_eq!(talk.status, TalkStatus::Closed);
2551 close(&mut talk, &talks).expect("closing twice is not an error");
2552
2553 let err =
2554 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
2555 assert!(err.to_string().contains("closed"));
2556 let _ = &cfg; }
2558
2559 #[tokio::test]
2560 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2561 let (tmp, talks) = store();
2562 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2563 let cfg = config(spec);
2564 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2567
2568 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
2572 close(&mut closed_elsewhere, &talks).expect("close");
2573 assert_eq!(
2574 talks.get(&in_flight.id).expect("reread").status,
2575 TalkStatus::Closed,
2576 "the close landed on disk before the turn finished"
2577 );
2578
2579 assert_eq!(in_flight.status, TalkStatus::Open);
2583 respond(&mut in_flight, &talks, &cfg, "one more question")
2584 .await
2585 .expect("the turn itself still completes");
2586
2587 let on_disk = talks.get(&in_flight.id).expect("reread");
2588 assert_eq!(
2589 on_disk.status,
2590 TalkStatus::Closed,
2591 "a close must stick even when a turn that started before it finishes after it"
2592 );
2593 assert!(
2596 on_disk.turns.iter().any(|t| t.body == "here you go"),
2597 "the in-flight turn's own reply is still recorded: {:?}",
2598 on_disk.turns
2599 );
2600 }
2601
2602 #[test]
2603 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
2604 let (tmp, talks) = store();
2605 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2606 let cfg = config(spec);
2607 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2610
2611 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
2614 close(&mut closed_elsewhere, &talks).expect("close");
2615 assert_eq!(
2616 talks.get(&stale.id).expect("reread").status,
2617 TalkStatus::Closed,
2618 "the close landed on disk before record was called"
2619 );
2620
2621 assert_eq!(stale.status, TalkStatus::Open);
2625 let err = record(&mut stale, &talks, "still there?", Vec::new())
2626 .expect_err("a close that landed first must be honored, not overwritten");
2627 assert!(err.to_string().contains("closed"));
2628
2629 let on_disk = talks.get(&stale.id).expect("reread");
2630 assert_eq!(
2631 on_disk.status,
2632 TalkStatus::Closed,
2633 "record must not resurrect a conversation closed while its snapshot was stale"
2634 );
2635 assert!(
2636 on_disk.turns.is_empty(),
2637 "the rejected turn must not have been appended: {:?}",
2638 on_disk.turns
2639 );
2640 let _ = &cfg; }
2642
2643 #[test]
2644 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
2645 let (tmp, talks) = store();
2646 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2647 let cfg = config(spec);
2648 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2649
2650 let held = talks.guard();
2654
2655 let talks2 = talks.clone();
2656 let id = talk.id.clone();
2657 let closing = std::thread::spawn(move || {
2658 let mut talk = talks2.get(&id).expect("get");
2659 close(&mut talk, &talks2).expect("close");
2660 });
2661
2662 std::thread::sleep(Duration::from_millis(50));
2663 assert!(
2664 !closing.is_finished(),
2665 "close must wait for the guard, not read and write while it is held - \
2666 a re-read alone narrows this window without closing it"
2667 );
2668
2669 drop(held);
2670 closing.join().expect("close thread panicked");
2671
2672 assert_eq!(
2673 talks.get(&talk.id).expect("reread").status,
2674 TalkStatus::Closed,
2675 "once the guard is free, close still lands"
2676 );
2677 let _ = &cfg; }
2679
2680 #[test]
2681 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
2682 let (tmp, talks) = store();
2683 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2684 let cfg = config(spec);
2685 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2686
2687 close(&mut talk, &talks).expect("close");
2688 assert_eq!(talk.status, TalkStatus::Closed);
2689
2690 reopen(&mut talk, &talks).expect("reopen");
2691 assert_eq!(talk.status, TalkStatus::Open);
2692 assert_eq!(
2693 talks.get(&talk.id).expect("reread").status,
2694 TalkStatus::Open
2695 );
2696
2697 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
2699 assert_eq!(talk.status, TalkStatus::Open);
2700
2701 record(&mut talk, &talks, "one more thing", Vec::new())
2702 .expect("a reopened talk takes turns again");
2703 let _ = &cfg; }
2705
2706 #[test]
2707 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
2708 let (tmp, talks) = store();
2709 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2710 let cfg = config(spec);
2711 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2712
2713 let artifacts = talks.artifacts_of(&talk.id);
2714 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
2715 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
2716
2717 talks.remove(&talk.id).expect("remove");
2718 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
2719 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
2720 assert!(
2721 talks.get(&talk.id).is_err(),
2722 "a removed talk cannot be read back"
2723 );
2724
2725 let err = talks
2726 .remove("nonexistent-id")
2727 .expect_err("unknown id refused");
2728 assert!(err.to_string().contains("no talk matches"), "{err}");
2729 let _ = &cfg; }
2731
2732 #[tokio::test]
2733 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2734 let (tmp, talks) = store();
2735 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2736 let cfg = config(spec);
2737 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2740
2741 talks.remove(&in_flight.id).expect("remove");
2742 assert!(
2743 talks.get(&in_flight.id).is_err(),
2744 "the delete landed on disk before the turn finished"
2745 );
2746
2747 respond(&mut in_flight, &talks, &cfg, "one more question")
2750 .await
2751 .expect("the turn itself still completes rather than erroring");
2752
2753 assert!(
2754 talks.get(&in_flight.id).is_err(),
2755 "a delete must stick even when a turn that started before it finishes after it"
2756 );
2757 }
2758
2759 #[test]
2760 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
2761 let (tmp, talks) = store();
2762 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2763 let cfg = config(spec);
2764 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2767
2768 talks.remove(&stale.id).expect("remove");
2769
2770 let err = record(&mut stale, &talks, "still there?", Vec::new())
2774 .expect_err("a delete that landed first must be honored, not overwritten");
2775 assert!(err.to_string().contains("deleted"), "{err}");
2776
2777 assert!(
2778 talks.get(&stale.id).is_err(),
2779 "record must not resurrect a conversation deleted while its snapshot was stale"
2780 );
2781 let _ = &cfg; }
2783
2784 #[test]
2785 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
2786 let (tmp, talks) = store();
2787 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2788 let cfg = config(spec);
2789 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2792
2793 talks.remove(&stale.id).expect("remove");
2794
2795 let err = close(&mut stale, &talks)
2799 .expect_err("a delete that landed first must be honored, not overwritten");
2800 assert!(err.to_string().contains("deleted"), "{err}");
2801
2802 assert!(
2803 talks.get(&stale.id).is_err(),
2804 "close must not resurrect a conversation deleted while its snapshot was stale"
2805 );
2806 let _ = &cfg; }
2808
2809 #[test]
2810 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
2811 let (tmp, talks) = store();
2812 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2813 let cfg = config(spec);
2814 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2817 close(&mut stale, &talks).expect("close");
2818
2819 talks.remove(&stale.id).expect("remove");
2820
2821 let err = reopen(&mut stale, &talks)
2825 .expect_err("a delete that landed first must be honored, not overwritten");
2826 assert!(err.to_string().contains("deleted"), "{err}");
2827
2828 assert!(
2829 talks.get(&stale.id).is_err(),
2830 "reopen must not resurrect a conversation deleted while its snapshot was stale"
2831 );
2832 let _ = &cfg; }
2834
2835 #[test]
2836 fn list_puts_open_talks_before_closed_ones() {
2837 let (tmp, talks) = store();
2838 let make = |id: &str, status: TalkStatus| {
2839 let mut t = Talk {
2840 schema: SCHEMA,
2841 id: id.to_owned(),
2842 repo: tmp.path().to_owned(),
2843 agent: "mock".to_owned(),
2844 status,
2845 turns: Vec::new(),
2846 pending: String::new(),
2847 pending_attachments: Vec::new(),
2848 fallback: false,
2849 created_at: Timestamp::now(),
2850 updated_at: Timestamp::now(),
2851 seat: SeatState::new(SEAT, "mock", 7),
2852 };
2853 talks.put(&mut t).expect("put");
2854 };
2855 make("20260901-000000-0001", TalkStatus::Open);
2856 make("20260902-000000-0002", TalkStatus::Open);
2857 make("20260903-000000-0003", TalkStatus::Closed);
2858
2859 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
2860 assert_eq!(
2861 ids,
2862 [
2863 "20260902-000000-0002",
2864 "20260901-000000-0001",
2865 "20260903-000000-0003"
2866 ]
2867 );
2868 assert_eq!(talks.count_open(), 2);
2869 }
2870
2871 #[test]
2872 fn tasks_of_finds_only_this_talks_own_tasks() {
2873 let dir = tempfile::tempdir().expect("tempdir");
2874 let queue = Queue::at(dir.path().join("queue"));
2875
2876 let mut mine = Task::new(
2877 "rework the loader".to_owned(),
2878 "rework the loader".to_owned(),
2879 PathBuf::from("/repo"),
2880 Source::Agent {
2881 run: "20260904-014455-ab12".to_owned(),
2882 node: "chat".to_owned(),
2883 },
2884 );
2885 queue.put(&mut mine).expect("put mine");
2886
2887 let mut theirs = Task::new(
2888 "unrelated".to_owned(),
2889 "unrelated".to_owned(),
2890 PathBuf::from("/repo"),
2891 Source::Agent {
2892 run: "20260904-090000-zz99".to_owned(),
2893 node: "implement".to_owned(),
2894 },
2895 );
2896 queue.put(&mut theirs).expect("put theirs");
2897
2898 let mut human = Task::new(
2899 "typed by hand".to_owned(),
2900 "typed by hand".to_owned(),
2901 PathBuf::from("/repo"),
2902 Source::Human,
2903 );
2904 queue.put(&mut human).expect("put human");
2905
2906 let found = tasks_of(&queue, "20260904-014455-ab12");
2907 assert_eq!(found.len(), 1);
2908 assert_eq!(found[0].id, mine.id);
2909 }
2910
2911 #[test]
2912 fn the_briefing_names_solo_task_add() {
2913 let brief = briefing(Path::new("/repo"), "en", false);
2914 assert!(brief.contains("magi task add --solo"));
2915 assert!(brief.contains("/repo"));
2916 assert!(!brief.contains("Hold this conversation in"));
2917 }
2918
2919 #[test]
2925 fn the_briefing_explains_targeting_a_different_repository_by_name() {
2926 let brief = briefing(Path::new("/repo"), "en", false);
2927 assert!(brief.contains("--repo does not have to be a full path"));
2928 assert!(brief.contains("owner/repo"));
2929 assert!(brief.contains("magi repos"));
2930 assert!(brief.contains("ask the operator"));
2931 }
2932
2933 #[test]
2934 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
2935 let brief = briefing(Path::new("/repo"), "en", false);
2936 assert!(brief.contains("--attach <path>"), "{brief}");
2937 assert!(brief.contains("deleting this conversation"), "{brief}");
2938 }
2939
2940 #[test]
2941 fn the_briefing_names_the_language_when_it_is_not_english() {
2942 let brief = briefing(Path::new("/repo"), "Japanese", false);
2943 assert!(brief.contains("Hold this conversation in Japanese"));
2944 }
2945
2946 #[test]
2947 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
2948 let read_only = briefing(Path::new("/repo"), "en", false);
2949 assert!(read_only.contains("Do not write files"));
2950 assert!(!read_only.contains("allow_write"));
2951
2952 let writable = briefing(Path::new("/repo"), "en", true);
2953 assert!(!writable.contains("Do not write files"));
2954 assert!(writable.contains("allow_write = true"));
2955 assert!(writable.contains("magi task add --solo"));
2958 assert!(writable.contains("say plainly what you"));
2959 }
2960}