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 self.list_counting_unreadable().0
529 }
530
531 pub fn list_counting_unreadable(&self) -> (Vec<Talk>, usize) {
534 let mut unreadable = 0;
535 let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
536 .into_iter()
537 .flatten()
538 .flatten()
539 .map(|e| e.path())
540 .filter(|p| p.extension().is_some_and(|x| x == "json"))
541 .filter_map(|p| {
542 let talk = read_path(&p).ok();
543 if talk.is_none() {
544 unreadable += 1;
545 }
546 talk
547 })
548 .collect();
549 all.sort_unstable_by(|a, b| {
550 let rank = |t: &Talk| u8::from(!t.status.open());
551 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
552 });
553 (all, unreadable)
554 }
555
556 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
558 if self.path_of(prefix).is_file() {
559 return Ok(prefix.to_owned());
560 }
561 let hits: Vec<String> = self
562 .list()
563 .into_iter()
564 .map(|t| t.id)
565 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
566 .collect();
567 match hits.len() {
568 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
569 0 => bail!("no talk matches `{prefix}`"),
570 _ => bail!(
571 "`{prefix}` matches {} talks: {}",
572 hits.len(),
573 hits.join(", ")
574 ),
575 }
576 }
577
578 pub fn revision(&self) -> u64 {
581 std::fs::read_dir(&self.root)
582 .into_iter()
583 .flatten()
584 .flatten()
585 .filter_map(|e| e.metadata().ok())
586 .filter_map(|m| m.modified().ok())
587 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
588 .map(|d| d.as_millis() as u64)
589 .max()
590 .unwrap_or(0)
591 }
592
593 pub fn count_open(&self) -> usize {
595 self.list().iter().filter(|t| t.status.open()).count()
596 }
597
598 pub fn remove(&self, id: &str) -> Result<()> {
611 let _guard = self.guard();
612 let resolved = self.resolve_id(id)?;
613 let path = self.path_of(&resolved);
614 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
615 let artifacts = self.artifacts_of(&resolved);
616 if artifacts.is_dir() {
617 std::fs::remove_dir_all(&artifacts)
618 .with_context(|| format!("remove {}", artifacts.display()))?;
619 }
620 Ok(())
621 }
622}
623
624pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
634 let repo = repo.canonicalize().unwrap_or(repo);
637 let spec = match agent {
640 Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
641 None => agent::pick_chain(
642 &cfg.agents,
643 cfg.roles.chatter.as_ref(),
644 &agent::installed,
645 "chatter",
646 )?
647 .remove(0),
648 };
649
650 let now = Timestamp::now();
651 let mut talk = Talk {
652 schema: SCHEMA,
653 id: new_id(),
654 repo,
655 agent: spec.id.clone(),
656 status: TalkStatus::Open,
657 turns: Vec::new(),
658 pending: String::new(),
659 pending_attachments: Vec::new(),
660 fallback: agent.is_none(),
661 created_at: now,
662 updated_at: now,
663 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
664 };
665 store.put(&mut talk)?;
666 Ok(talk)
667}
668
669pub fn record(
676 talk: &mut Talk,
677 store: &Talks,
678 text: &str,
679 attachments: Vec<Attachment>,
680) -> Result<String> {
681 let _guard = store.guard();
689 let Ok(fresh) = store.get(&talk.id) else {
694 bail!("talk {} was deleted", talk.short());
695 };
696 talk.status = fresh.status;
697 talk.pending = fresh.pending;
700 talk.pending_attachments = fresh.pending_attachments;
701 if !talk.status.open() {
702 bail!(
703 "talk {} is {} and takes no more turns",
704 talk.short(),
705 talk.status.as_str()
706 );
707 }
708 let text = text.trim();
709 if text.is_empty() && attachments.is_empty() {
710 bail!("nothing to say");
711 }
712 talk.turns.push(Turn {
713 who: Who::Operator,
714 body: text.to_owned(),
715 at: Timestamp::now(),
716 attachments,
717 usage: None,
718 });
719 store.put(talk)?;
720 Ok(text.to_owned())
721}
722
723pub fn queue(
725 talk: &mut Talk,
726 store: &Talks,
727 text: &str,
728 attachments: Vec<Attachment>,
729) -> Result<()> {
730 let text = text.trim();
731 if text.is_empty() && attachments.is_empty() {
732 bail!("nothing to say");
733 }
734 let _guard = store.guard();
735 let mut fresh = store
736 .get(&talk.id)
737 .with_context(|| format!("talk {} was deleted", talk.short()))?;
738 if !fresh.status.open() {
739 bail!(
740 "talk {} is {} and takes no more turns",
741 fresh.short(),
742 fresh.status.as_str()
743 );
744 }
745 if !text.is_empty() {
746 if fresh.pending.is_empty() {
747 fresh.pending = text.to_owned();
748 } else {
749 fresh.pending.push_str("\n\n");
750 fresh.pending.push_str(text);
751 }
752 }
753 fresh.pending_attachments.extend(attachments);
754 store.put(&mut fresh)?;
755 *talk = fresh;
756 Ok(())
757}
758
759pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
761 let _guard = store.guard();
762 let mut fresh = store
763 .get(&talk.id)
764 .with_context(|| format!("talk {} was deleted", talk.short()))?;
765 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
766 *talk = fresh;
767 return Ok(None);
768 }
769 let text = std::mem::take(&mut fresh.pending);
770 let attachments = std::mem::take(&mut fresh.pending_attachments);
771 fresh.turns.push(Turn {
772 who: Who::Operator,
773 body: text.clone(),
774 at: Timestamp::now(),
775 attachments,
776 usage: None,
777 });
778 store.put(&mut fresh)?;
779 *talk = fresh;
780 Ok(Some(text))
781}
782
783pub async fn say(
786 talk: &mut Talk,
787 store: &Talks,
788 cfg: &Config,
789 text: &str,
790 attachments: Vec<Attachment>,
791) -> Result<()> {
792 let text = record(talk, store, text, attachments)?;
793 turn(talk, store, cfg, &text).await
794}
795
796pub async fn respond(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
798 turn(talk, store, cfg, text).await
799}
800
801pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
819 let _guard = store.guard();
820 let mut fresh = store
821 .get(&talk.id)
822 .with_context(|| format!("talk {} was deleted", talk.short()))?;
823 fresh.status = TalkStatus::Closed;
824 fresh.pending.clear();
826 fresh.pending_attachments.clear();
827 store.put(&mut fresh)?;
828 *talk = fresh;
829 Ok(())
830}
831
832pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
843 let _guard = store.guard();
844 let mut fresh = store
845 .get(&talk.id)
846 .with_context(|| format!("talk {} was deleted", talk.short()))?;
847 fresh.status = TalkStatus::Open;
848 store.put(&mut fresh)?;
849 *talk = fresh;
850 Ok(())
851}
852
853pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
866 let _guard = store.guard();
867 let mut fresh = store
868 .get(&talk.id)
869 .with_context(|| format!("talk {} was deleted", talk.short()))?;
870 if fresh.agent == spec.id {
871 *talk = fresh;
872 return Ok(false);
873 }
874 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
875 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
876 fresh.fallback = false;
878 fresh.turns.push(Turn {
879 who: Who::Agent,
880 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
881 at: Timestamp::now(),
882 attachments: Vec::new(),
883 usage: None,
884 });
885 store.put(&mut fresh)?;
886 *talk = fresh;
887 Ok(true)
888}
889
890pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
892 let _guard = store.guard();
893 let mut fresh = store
894 .get(&talk.id)
895 .with_context(|| format!("talk {} was deleted", talk.short()))?;
896 fresh.pending.clear();
897 fresh.pending_attachments.clear();
898 store.put(&mut fresh)?;
899 *talk = fresh;
900 Ok(())
901}
902
903pub fn clear_pending_if_matches(
905 talk: &mut Talk,
906 store: &Talks,
907 expected_text: &str,
908 expected_attachments: &[String],
909) -> Result<bool> {
910 let _guard = store.guard();
911 let mut fresh = store
912 .get(&talk.id)
913 .with_context(|| format!("talk {} was deleted", talk.short()))?;
914 if !pending_matches(&fresh, expected_text, expected_attachments) {
915 *talk = fresh;
916 return Ok(false);
917 }
918 fresh.pending.clear();
919 fresh.pending_attachments.clear();
920 store.put(&mut fresh)?;
921 *talk = fresh;
922 Ok(true)
923}
924
925pub fn edit_pending_text(
929 talk: &mut Talk,
930 store: &Talks,
931 text: &str,
932 expected_text: &str,
933 expected_attachments: &[String],
934) -> Result<bool> {
935 let _guard = store.guard();
936 let mut fresh = store
937 .get(&talk.id)
938 .with_context(|| format!("talk {} was deleted", talk.short()))?;
939 if !pending_matches(&fresh, expected_text, expected_attachments) {
940 *talk = fresh;
941 return Ok(false);
942 }
943 fresh.pending = text.trim().to_owned();
944 store.put(&mut fresh)?;
945 *talk = fresh;
946 Ok(true)
947}
948
949fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
950 talk.pending == expected_text
951 && talk
952 .pending_attachments
953 .iter()
954 .map(|attachment| &attachment.id)
955 .eq(expected_attachments.iter())
956}
957
958async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
965 let spec = cfg
966 .agents
967 .iter()
968 .find(|a| a.id == talk.agent)
969 .with_context(|| {
970 format!(
971 "talk {} was opened with agent `{}`, which is no longer in \
972 the roster; restore it in magi.toml or start a new \
973 conversation",
974 talk.short(),
975 talk.agent
976 )
977 })?;
978
979 let last_note = attachment_note(
983 store,
984 &talk.id,
985 talk.turns
986 .last()
987 .map_or(&[][..], |t| t.attachments.as_slice()),
988 );
989
990 let attachment_paths: Vec<PathBuf> = talk
996 .turns
997 .iter()
998 .flat_map(|t| t.attachments.iter())
999 .filter_map(|a| store.attachment_path(&talk.id, a))
1000 .collect();
1001
1002 let consulted = talk
1009 .turns
1010 .last()
1011 .is_some_and(|t| t.who == Who::Operator && crate::consult::is_consult_text(&t.body));
1012 let consult_roots: Vec<PathBuf> = if consulted {
1013 vec![crate::ask::Questions::open().root().to_path_buf()]
1014 } else {
1015 Vec::new()
1016 };
1017
1018 let artifacts = store.artifacts_of(&talk.id);
1019 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
1022 let stem = format!("turn-{}", operator_turns.max(1));
1023 let cache_dir = cfg.cache_dir();
1026
1027 let mut chain = vec![spec.clone()];
1030 if let Some(choice) = cfg.roles.chatter.as_ref()
1031 && talk.fallback
1032 {
1033 for id in choice.ids() {
1034 if id == talk.agent || chain.iter().any(|s| s.id == id) {
1035 continue;
1036 }
1037 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1038 Ok(s) => chain.push(s),
1039 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1040 }
1041 }
1042 }
1043
1044 let mut outcome = None;
1045 let mut fell_back_from: Option<String> = None;
1046 let mut first_try: Option<(String, SeatState)> = None;
1049 for (n, spec) in chain.iter().enumerate() {
1050 if n > 0 {
1051 if first_try.is_none() {
1052 first_try = Some((talk.agent.clone(), talk.seat.clone()));
1053 }
1054 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1055 fell_back_from.get_or_insert_with(|| talk.agent.clone());
1058 talk.agent = spec.id.clone();
1059 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1060 }
1061 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1062 let first_ever = talk.turns.len() <= 1;
1063 let body = if talk.seat.turns == 0 && first_ever {
1064 format!(
1065 "{}\n\n# Operator\n\n{text}{last_note}",
1066 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write)
1067 )
1068 } else if talk.seat.turns == 0 {
1069 format!(
1072 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1073 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write),
1074 transcript(talk, store)
1075 )
1076 } else if resuming {
1077 format!("{text}{last_note}")
1078 } else {
1079 format!("{}\n\n{text}{last_note}", transcript(talk, store))
1080 };
1081 let attempt_stem = if n == 0 {
1082 stem.clone()
1083 } else {
1084 format!("{stem}-{}", spec.id)
1085 };
1086 let inv = Invocation {
1087 cwd: &talk.repo,
1088 prompt: &body,
1089 timeout: turn_timeout(cfg),
1090 allow_write: cfg.talk.allow_write || consulted,
1095 sessions: cfg.graph.sessions,
1096 artifacts: &artifacts,
1097 stem: &attempt_stem,
1098 run: &talk.id,
1101 node: crate::queue::CHAT_NODE,
1102 cache_dir: cache_dir.as_deref(),
1103 attachments: &attachment_paths,
1104 writable: &consult_roots,
1105 };
1106 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1107 let advance = agent::chain_advances(&result);
1108 if n == 0 || !advance {
1109 outcome = Some(result);
1110 } else {
1111 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1113 }
1114 if !advance {
1115 break;
1116 }
1117 }
1118 if outcome.as_ref().is_some_and(agent::chain_advances) {
1119 if let Some((id, seat)) = first_try {
1122 talk.agent = id;
1123 talk.seat = seat;
1124 fell_back_from = None;
1125 }
1126 }
1127 let outcome = outcome.expect("a chain holds at least one agent");
1128 let note = |why: String| Turn {
1129 who: Who::Agent,
1130 body: format!("{MAGI_NOTE}{why}"),
1131 at: Timestamp::now(),
1132 attachments: Vec::new(),
1133 usage: None,
1134 };
1135 let (reply, failure) = match outcome {
1136 Err(e) => (
1137 note(format!("could not run agent `{}`: {e}", talk.agent)),
1138 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1139 ),
1140 Ok(out) if out.quota_exhausted() => {
1141 let reset = out
1142 .quota
1143 .as_ref()
1144 .and_then(|q| q.reset.clone())
1145 .map_or_else(String::new, |r| format!(" (resets {r})"));
1146 let why = format!(
1147 "agent `{}` is out of quota{reset}; your message is saved, so \
1148 say it again when the window reopens",
1149 talk.agent
1150 );
1151 (note(why.clone()), Some(why))
1152 }
1153 Ok(out) if out.timed_out => {
1154 let why = format!(
1155 "agent `{}` did not answer within {}s; your message is saved",
1156 talk.agent,
1157 turn_timeout(cfg).as_secs()
1158 );
1159 (note(why.clone()), Some(why))
1160 }
1161 Ok(out) if !out.usable() => {
1162 let why = format!(
1163 "agent `{}` produced no answer (exit {}); your message is saved",
1164 talk.agent,
1165 out.exit_code
1166 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1167 );
1168 (note(why.clone()), Some(why))
1169 }
1170 Ok(out) => (
1171 Turn {
1172 who: Who::Agent,
1173 body: out.text.trim().to_owned(),
1174 at: Timestamp::now(),
1175 attachments: Vec::new(),
1176 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1179 context_tokens,
1180 agent: talk.agent.clone(),
1181 model: cfg
1182 .agents
1183 .iter()
1184 .find(|a| a.id == talk.agent)
1185 .and_then(|a| a.model.clone()),
1186 }),
1187 },
1188 None,
1189 ),
1190 };
1191
1192 let _guard = store.guard();
1204 let Ok(fresh) = store.get(&talk.id) else {
1210 return Ok(());
1211 };
1212 talk.status = fresh.status;
1213 talk.pending = fresh.pending;
1217 talk.pending_attachments = fresh.pending_attachments;
1218 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1219 talk.turns.push(note(format!(
1222 "agent changed from {from} to {} (fallback)",
1223 talk.agent
1224 )));
1225 }
1226 talk.turns.push(reply);
1227 if let Err(put_err) = store.put(talk) {
1228 let lost = talk.turns.pop().expect("just pushed above");
1237 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1238 let why = match &stash {
1239 Ok(path) => format!(
1240 "agent `{}` answered, but the reply could not be saved to \
1241 this conversation ({put_err:#}); the raw text was kept at \
1242 {} - your message is saved, ask again",
1243 talk.agent,
1244 path.display()
1245 ),
1246 Err(stash_err) => format!(
1247 "agent `{}` answered, but the reply could not be saved to \
1248 this conversation ({put_err:#}), and it could not be kept \
1249 anywhere else either ({stash_err:#}); your message is \
1250 saved, ask again",
1251 talk.agent
1252 ),
1253 };
1254 talk.turns.push(note(why.clone()));
1255 return match store.put(talk) {
1262 Ok(()) => bail!("{why}"),
1263 Err(note_err) => {
1264 talk.turns.pop();
1284 Err(note_err).context(why)
1285 }
1286 };
1287 }
1288
1289 match failure {
1290 Some(why) => bail!("{why}"),
1291 None => Ok(()),
1292 }
1293}
1294
1295fn transcript(talk: &Talk, store: &Talks) -> String {
1298 let mut out = String::from(
1299 "This conversation cannot resume on the CLI's side, so here is \
1300 everything said so far; answer only the last message.\n",
1301 );
1302 for t in &talk.turns {
1303 let who = match t.who {
1304 Who::Operator => "operator",
1305 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1306 Who::Agent => "you",
1307 };
1308 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1309 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1310 }
1311 out
1312}
1313
1314fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1319 if attachments.is_empty() {
1320 return String::new();
1321 }
1322 let mut out = String::from(
1323 "\n\nThe operator attached the image(s) below to this message. Open \
1324 and look at each one before you answer.\n",
1325 );
1326 for att in attachments {
1327 if let Some(path) = store.attachment_path(talk_id, att) {
1328 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1329 }
1330 }
1331 out.push('\n');
1332 out
1333}
1334
1335pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
1357 let write_policy = if allow_write {
1358 "Write access is enabled for this conversation (`allow_write = \
1359 true`), so you may write files - but only a small, \
1360 already-decided edit the operator names outright in this \
1361 conversation, not an implementation. This is a permission on the \
1362 conversation as a whole, not a property of whichever repository \
1363 it happened to start in: if the operator names a different \
1364 repository for that small edit, the policy allows it there too. \
1365 Your own tool may still confine writes to the repository this \
1366 conversation started in regardless - if a write elsewhere is \
1367 refused, say so plainly rather than working around it. Once you \
1368 have made an edit, say plainly what you edited. Anything bigger, \
1369 or anything still open-ended, still goes through the queue below \
1370 rather than being done here."
1371 } else {
1372 "Do not write files. Implementing a change is not this \
1373 conversation's job; a separate, blind competition of agents does \
1374 that, and a repository this conversation has already edited would \
1375 make their diffs unjudgeable."
1376 };
1377 let mut out = format!(
1378 "You are magi's standing conversation partner for its operator, who \
1379 usually has this open on a phone. Keep replies short: no preamble, \
1380 no restating what they just said.\n\n\
1381 # Repository\n\n{repo}\n\n\
1382 You may look around: read files, run shell commands, search history, \
1383 run tests - whatever answers the question. {write_policy}\n\n\
1384 A short, command-shaped message (\"list\", \"info <id>\", \"show \
1385 3cbf\") is almost always the operator asking you to look something \
1386 up, not an instruction to file - answer it yourself with `magi \
1387 list`, `magi show <id>`, `magi task list`, or the like, the same way \
1388 you would answer any other question in this conversation.\n\n\
1389 # When the operator wants something done\n\n\
1390 Run:\n\n\
1391 magi task add --solo --repo {repo} <instruction>\n\n\
1392 and tell the operator the task id it prints, so they can follow it \
1393 from the Queue. If it refuses with a duplicate warning (the \
1394 instruction names a branch, commit or pull request that an \
1395 unfinished task, run or PR already owns), do not repeat it with \
1396 --force yourself: tell the operator what it matched and let them \
1397 decide. Write <instruction> so that an implementer who has \
1398 never seen this conversation can act on it alone - it is everything \
1399 they get. Use --solo: it runs the task through one implementer \
1400 straight into review instead of the usual multi-agent competition, \
1401 which is the right shape for a change this conversation has already \
1402 settled, rather than one still worth several independent takes.\n\n\
1403 If the operator asks for something in a different repository, \
1404 --repo does not have to be a full path: --repo owner/repo (or just \
1405 repo, when that is unambiguous) is resolved against local checkouts \
1406 the same way `magi repos` lists them. If the command fails because \
1407 nothing matches or more than one checkout shares that name, ask the \
1408 operator which repository they mean (or run `magi repos` yourself \
1409 to see the candidates) rather than guessing.\n\n\
1410 The current state of the code is whatever origin/main holds, not \
1411 whatever a working tree shows: a primary checkout often lags \
1412 upstream, sits on a detached HEAD and carries uncommitted changes. \
1413 Before answering about code, run `git fetch origin` in that \
1414 repository if it is cheap, then read through \
1415 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
1416 If the working tree differs, say so; if the fetch fails, say that \
1417 too, so the operator knows the answer may be stale.\n\n\
1418 If the operator attached an image (a screenshot, say) that the task \
1419 is about, pass it with `--attach <path>`, using the absolute path \
1420 the turn's attachment note gives; repeat the flag for several. \
1421 `magi task add --solo --attach <path> <instruction>` copies the \
1422 file into the task, so the implementer receives it. Do not paste the \
1423 path into <instruction> instead: deleting this conversation deletes \
1424 its attachments, and then that path reaches no one.\n",
1425 repo = repo.display(),
1426 );
1427 out.push_str(&language_note(language));
1428 out
1429}
1430
1431fn language_note(language: &str) -> String {
1434 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
1435 String::new()
1436 } else {
1437 format!("\nHold this conversation in {language}.\n")
1438 }
1439}
1440
1441pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
1448 let mut tasks: Vec<Task> = queue
1449 .list()
1450 .into_iter()
1451 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
1452 .collect();
1453 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
1454 tasks
1455}
1456
1457fn read_path(path: &Path) -> Result<Talk> {
1458 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1459 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
1460}
1461
1462const PUT_RETRIES: u32 = 5;
1465
1466fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1478 let mut last_err = None;
1479 for attempt in 0..PUT_RETRIES {
1480 if attempt > 0 {
1481 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
1482 }
1483 match try_write_atomic(tmp, path, body) {
1484 Ok(()) => return Ok(()),
1485 Err(e) => last_err = Some(e),
1486 }
1487 }
1488 Err(last_err.expect("the loop above always runs at least once"))
1489}
1490
1491fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1492 #[cfg(test)]
1493 if failpoint::take_forced_put_failure() {
1494 bail!("simulated write failure (test)");
1495 }
1496 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
1497 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
1498 Ok(())
1499}
1500
1501fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
1506 let dir = store.artifacts_of(id);
1507 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1508 let path = dir.join(format!("{stem}-lost.txt"));
1509 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
1510 Ok(path)
1511}
1512
1513#[cfg(test)]
1520mod failpoint {
1521 use std::cell::Cell;
1522
1523 thread_local! {
1524 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
1525 }
1526
1527 pub(super) fn force_put_failures(count: u32) {
1530 FORCE_PUT_FAILURES.with(|c| c.set(count));
1531 }
1532
1533 pub(super) fn take_forced_put_failure() -> bool {
1536 FORCE_PUT_FAILURES.with(|c| {
1537 let n = c.get();
1538 if n == 0 {
1539 false
1540 } else {
1541 c.set(n - 1);
1542 true
1543 }
1544 })
1545 }
1546}
1547
1548fn short(id: &str) -> &str {
1549 id.split('-').next_back().unwrap_or(id)
1550}
1551
1552fn new_id() -> String {
1553 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1554 let seed = crate::rng::entropy();
1555 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1556}
1557
1558fn attachment_ext(mime: &str) -> Option<&'static str> {
1563 match mime {
1564 "image/png" => Some("png"),
1565 "image/jpeg" => Some("jpg"),
1566 "image/gif" => Some("gif"),
1567 "image/webp" => Some("webp"),
1568 _ => None,
1569 }
1570}
1571
1572pub fn valid_attachment_id(id: &str) -> bool {
1577 id.len() == 32
1578 && id
1579 .bytes()
1580 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
1581}
1582
1583fn new_attachment_id() -> String {
1587 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
1588 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
1589}
1590
1591#[cfg(test)]
1592mod tests {
1593 #[test]
1594 fn the_briefing_points_at_origin_main_not_the_working_tree() {
1595 let b = briefing(Path::new("/r"), "en", false);
1596 assert!(b.contains("origin/main"));
1597 assert!(b.contains("git show origin/main:"));
1598 }
1599 use std::collections::BTreeMap;
1600
1601 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
1602 use crate::queue::{Queue, Source, Task};
1603
1604 use super::*;
1605
1606 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
1607 AgentSpec {
1608 id: id.to_owned(),
1609 kind: AgentKind::Command,
1610 model: model.map(str::to_owned),
1611 command: Vec::new(),
1612 extra_args: Vec::new(),
1613 env: BTreeMap::new(),
1614 prompt_delivery: None,
1615 }
1616 }
1617
1618 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
1619 Talk {
1620 schema: SCHEMA,
1621 id: "20260904-014455-ab12".to_owned(),
1622 repo: PathBuf::from("."),
1623 agent: agent.to_owned(),
1624 status: TalkStatus::Open,
1625 turns,
1626 pending: String::new(),
1627 pending_attachments: Vec::new(),
1628 fallback: false,
1629 created_at: Timestamp::now(),
1630 updated_at: Timestamp::now(),
1631 seat: SeatState::new(SEAT, agent, 1),
1632 }
1633 }
1634
1635 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
1636 Turn {
1637 who: Who::Agent,
1638 body: body.to_owned(),
1639 at: Timestamp::now(),
1640 attachments: Vec::new(),
1641 usage: usage.map(|(t, a, m)| TurnUsage {
1642 context_tokens: t,
1643 agent: a.to_owned(),
1644 model: m.map(str::to_owned),
1645 }),
1646 }
1647 }
1648
1649 fn ctx_config(windows: &[(&str, u64)]) -> Config {
1650 Config {
1651 agents: vec![
1652 ctx_agent("small", Some("small-model")),
1653 ctx_agent("big", Some("big-model")),
1654 ctx_agent("plain", None),
1655 ],
1656 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
1657 ..Config::default()
1658 }
1659 }
1660
1661 #[test]
1662 fn context_usage_computes_percent_and_warns_at_eighty() {
1663 let cfg = ctx_config(&[("small-model", 1000)]);
1664 let at = |tokens| {
1665 let t = ctx_talk(
1666 "small",
1667 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
1668 );
1669 context_usage(&t, Some(&cfg))
1670 };
1671 let u = at(799);
1672 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
1673 let u = at(800);
1674 assert_eq!((u.percent, u.warn), (Some(80), true));
1675 let u = at(1500);
1676 assert_eq!((u.percent, u.warn), (Some(150), true));
1677 assert!(!u.since_switch);
1678 }
1679
1680 #[test]
1681 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
1682 let cfg = ctx_config(&[("small-model", 1000)]);
1683 let t = ctx_talk(
1684 "small",
1685 vec![
1686 reply("old", Some((900, "small", Some("small-model")))),
1687 reply("new", None),
1688 ],
1689 );
1690 let u = context_usage(&t, Some(&cfg));
1691 assert!(u.estimated);
1693 assert_ne!(u.tokens, Some(900));
1694 assert!(u.tokens.is_some());
1695 let t = ctx_talk(
1697 "small",
1698 vec![
1699 reply("old", Some((900, "small", Some("small-model")))),
1700 reply("magi: could not run agent", None),
1701 ],
1702 );
1703 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
1704 assert_eq!(
1705 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
1706 None
1707 );
1708 }
1709
1710 #[test]
1711 fn estimate_counts_chars_both_sides_and_standing_prompt() {
1712 let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
1713 assert_eq!(estimate_context_tokens(&t, 0), Some(2)); let op = Turn {
1715 who: Who::Operator,
1716 ..reply("abcdefg", None)
1717 };
1718 t.turns.push(op);
1719 assert_eq!(estimate_context_tokens(&t, 0), Some(4));
1720 assert!(
1721 estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
1722 );
1723 let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
1725 assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
1726 let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
1728 assert_eq!(estimate_context_tokens(¬e, 1000), None);
1729 assert_eq!(
1730 estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
1731 None
1732 );
1733 }
1734
1735 #[test]
1736 fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
1737 let cfg = ctx_config(&[("small-model", 1000)]);
1738 let t = ctx_talk(
1739 "small",
1740 vec![reply(
1741 &"x".repeat(5000),
1742 Some((10, "small", Some("small-model"))),
1743 )],
1744 );
1745 let u = context_usage(&t, Some(&cfg));
1746 assert_eq!((u.tokens, u.estimated), (Some(10), false));
1747 let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
1748 let u = context_usage(&t, Some(&cfg));
1749 assert!(u.estimated && !u.since_switch);
1750 assert_eq!(u.window, Some(1000));
1751 assert!(u.warn && u.percent.unwrap() >= 80);
1752 let t = ctx_talk("small", vec![reply("hi", None)]);
1753 let u = context_usage(&t, Some(&cfg));
1754 assert!(u.estimated && u.percent.is_some());
1755 }
1756
1757 #[test]
1758 fn context_usage_without_a_window_shows_tokens_only() {
1759 let cfg = ctx_config(&[]);
1760 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
1762 let u = context_usage(&t, Some(&cfg));
1763 assert_eq!(
1764 (u.tokens, u.window, u.percent, u.warn),
1765 (Some(5000), None, None, false)
1766 );
1767 let t = ctx_talk(
1768 "small",
1769 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
1770 );
1771 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
1772 assert_eq!(context_usage(&t, None).window, None);
1774 }
1775
1776 #[test]
1777 fn context_usage_switching_model_changes_the_denominator() {
1778 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
1779 let used = reply("hi", Some((900, "small", Some("small-model"))));
1780 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
1781 assert_eq!(
1782 (before.percent, before.warn, before.since_switch),
1783 (Some(90), true, false)
1784 );
1785 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
1788 assert_eq!(after.window, Some(10_000));
1789 assert_eq!(
1790 (after.percent, after.warn, after.since_switch),
1791 (Some(9), false, true)
1792 );
1793 assert_eq!(after.model.as_deref(), Some("big-model"));
1794 }
1795
1796 #[test]
1797 fn a_turn_recorded_before_usage_existed_still_reads() {
1798 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
1799 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
1800 assert!(turn.usage.is_none());
1801 let json = serde_json::to_string(&turn).expect("serialize");
1802 assert!(
1803 !json.contains("usage"),
1804 "absent usage is not written: {json}"
1805 );
1806 }
1807
1808 fn store() -> (tempfile::TempDir, Talks) {
1810 let tmp = tempfile::tempdir().expect("tempdir");
1811 let talks = Talks::at(tmp.path().join("talks"));
1812 (tmp, talks)
1813 }
1814
1815 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
1819 let path = dir.join("mock-talk-agent.sh");
1820 std::fs::write(&path, script).expect("write mock");
1821 AgentSpec {
1822 id: "mock".to_owned(),
1823 kind: AgentKind::Command,
1824 model: None,
1825 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
1826 extra_args: Vec::new(),
1827 env,
1828 prompt_delivery: None,
1829 }
1830 }
1831
1832 fn config(spec: AgentSpec) -> Config {
1833 Config {
1834 agents: vec![spec],
1835 graph: Graph {
1836 language: "en".to_owned(),
1837 ..Graph::default()
1838 },
1839 ..Config::default()
1840 }
1841 }
1842
1843 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
1845
1846 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
1848
1849 const ECHO: &str = "#!/bin/sh\ncat\n";
1852
1853 fn env(reply: &str) -> BTreeMap<String, String> {
1854 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
1855 }
1856
1857 #[test]
1858 fn the_frozen_json_field_names_round_trip_through_disk() {
1859 let (tmp, talks) = store();
1860 let mut talk = Talk {
1861 schema: SCHEMA,
1862 id: "20260904-014455-ab12".to_owned(),
1863 repo: tmp.path().to_owned(),
1864 agent: "sonnet".to_owned(),
1865 status: TalkStatus::Open,
1866 turns: Vec::new(),
1867 pending: String::new(),
1868 pending_attachments: Vec::new(),
1869 fallback: false,
1870 created_at: Timestamp::now(),
1871 updated_at: Timestamp::now(),
1872 seat: SeatState::new(SEAT, "sonnet", 7),
1873 };
1874 talks.put(&mut talk).expect("put");
1875
1876 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
1877 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
1878 for field in [
1879 "schema",
1880 "id",
1881 "repo",
1882 "agent",
1883 "status",
1884 "turns",
1885 "created_at",
1886 "updated_at",
1887 ] {
1888 assert!(v.get(field).is_some(), "missing field `{field}`");
1889 }
1890 assert_eq!(v["schema"], 1);
1891 assert_eq!(v["status"], "open");
1892
1893 let back = talks.get(&talk.id).expect("get");
1894 assert_eq!(back.id, talk.id);
1895 assert_eq!(back.status, TalkStatus::Open);
1896 }
1897
1898 #[test]
1899 fn opening_a_talk_takes_no_agent_turn() {
1900 let (tmp, talks) = store();
1901 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1905 let cfg = config(spec);
1906
1907 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1908 assert_eq!(talk.status, TalkStatus::Open);
1909 assert!(talk.turns.is_empty(), "nothing has been said yet");
1910
1911 let on_disk = talks.get(&talk.id).expect("get");
1912 assert_eq!(on_disk.turns.len(), 0);
1913 }
1914
1915 #[test]
1923 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
1924 let (tmp, talks) = store();
1925 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1926 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1927 chatter_spec.id = "chatter-mock".to_owned();
1928
1929 let mut cfg = Config {
1930 agents: vec![first_spec.clone(), chatter_spec.clone()],
1931 graph: Graph {
1932 language: "en".to_owned(),
1933 ..Graph::default()
1934 },
1935 ..Config::default()
1936 };
1937 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
1938
1939 let talk =
1940 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
1941 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
1942
1943 cfg.roles.chatter = None;
1944 let fallback =
1945 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
1946 assert_eq!(
1947 fallback.agent, first_spec.id,
1948 "unset chatter must fall back to agent::pick's own default order"
1949 );
1950 }
1951
1952 #[test]
1955 fn a_talk_recorded_without_attachments_still_reads() {
1956 let (tmp, talks) = store();
1957 let path = talks.path_of("20260904-014455-ab12");
1958 std::fs::create_dir_all(talks.root()).expect("talks dir");
1959 std::fs::write(
1960 &path,
1961 serde_json::json!({
1962 "schema": 1,
1963 "id": "20260904-014455-ab12",
1964 "repo": tmp.path(),
1965 "agent": "sonnet",
1966 "status": "open",
1967 "turns": [
1968 { "who": "operator", "body": "still there?",
1969 "at": Timestamp::now().to_string() },
1970 ],
1971 "created_at": Timestamp::now().to_string(),
1972 "updated_at": Timestamp::now().to_string(),
1973 "seat": SeatState::new(SEAT, "sonnet", 7),
1974 })
1975 .to_string(),
1976 )
1977 .expect("write pre-attachments talk");
1978
1979 let talk = talks.get("20260904-014455-ab12").expect("must still read");
1980 assert!(talk.turns[0].attachments.is_empty());
1981 }
1982
1983 #[test]
1984 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
1985 let (tmp, talks) = store();
1986 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1987 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1988
1989 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
1990 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
1991 let saved = talks.get(&talk.id).expect("reload queued talk");
1992 assert_eq!(saved.pending, "first\n\nsecond");
1993 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
1994
1995 let drained = drain(&mut talk, &talks).expect("drain");
1996 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
1997 let saved = talks.get(&talk.id).expect("reload drained talk");
1998 assert!(saved.pending.is_empty());
1999 assert_eq!(saved.turns.len(), 1);
2000 assert_eq!(saved.turns[0].body, "first\n\nsecond");
2001 }
2002
2003 #[test]
2004 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
2005 let (tmp, talks) = store();
2006 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2007 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2008 let attachment = Attachment {
2009 id: "a".repeat(32),
2010 name: "shot.png".to_owned(),
2011 mime: "image/png".to_owned(),
2012 bytes: 3,
2013 };
2014
2015 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
2016 assert!(
2017 edit_pending_text(
2018 &mut talk,
2019 &talks,
2020 "corrected",
2021 "first",
2022 std::slice::from_ref(&attachment.id),
2023 )
2024 .expect("edit")
2025 );
2026 let saved = talks.get(&talk.id).expect("reload edited draft");
2027 assert_eq!(saved.pending, "corrected");
2028 assert_eq!(saved.pending_attachments, vec![attachment]);
2029
2030 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
2031 assert!(
2032 !edit_pending_text(
2033 &mut talk,
2034 &talks,
2035 "stale edit",
2036 "corrected",
2037 &["a".repeat(32)],
2038 )
2039 .expect("stale edit is a conflict")
2040 );
2041 assert_eq!(
2042 talks.get(&talk.id).expect("reload after conflict").pending,
2043 "corrected\n\nlater"
2044 );
2045 assert!(
2046 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
2047 .expect("stale clear is a conflict")
2048 );
2049 assert_eq!(
2050 talks
2051 .get(&talk.id)
2052 .expect("reload after stale clear")
2053 .pending,
2054 "corrected\n\nlater"
2055 );
2056 }
2057
2058 #[tokio::test]
2059 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
2060 let (tmp, talks) = store();
2061 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
2062 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
2063 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2064 let id = running.id.clone();
2065 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
2066
2067 let response_talks = talks.clone();
2068 let response_cfg = cfg.clone();
2069 let reply = tokio::spawn(async move {
2070 respond(&mut running, &response_talks, &response_cfg, &first).await
2071 });
2072 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
2073
2074 let mut queued = talks.get(&id).expect("queued handle");
2075 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
2076 reply.await.expect("join").expect("reply");
2077
2078 let saved = talks.get(&id).expect("reload");
2079 assert_eq!(saved.pending, "next");
2080 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
2081 }
2082
2083 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
2086 let calls = dir.join(format!("{id}.calls"));
2087 let script = format!(
2088 "#!/bin/sh\necho x >> '{}'\n{body}\n",
2089 calls.to_string_lossy()
2090 );
2091 let path = dir.join(format!("mock-{id}.sh"));
2092 std::fs::write(&path, script).expect("write mock");
2093 AgentSpec {
2094 id: id.to_owned(),
2095 kind: AgentKind::Command,
2096 model: None,
2097 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2098 extra_args: Vec::new(),
2099 env: BTreeMap::new(),
2100 prompt_delivery: None,
2101 }
2102 }
2103
2104 fn calls(dir: &Path, id: &str) -> usize {
2105 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
2106 }
2107
2108 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
2109 let mut cfg = config(specs[0].clone());
2110 cfg.agents = specs;
2111 cfg.roles.chatter = Some(AgentChoice::Chain(
2112 ids.iter().map(|s| (*s).to_owned()).collect(),
2113 ));
2114 cfg
2115 }
2116
2117 #[tokio::test]
2118 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
2119 let (tmp, talks) = store();
2120 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2121 let b = counting_agent(tmp.path(), "b", "cat");
2122 let cfg = chain_config(vec![a, b], &["a", "b"]);
2123 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2124 assert_eq!(talk.agent, "a");
2125
2126 say(&mut talk, &talks, &cfg, "hello there", Vec::new())
2127 .await
2128 .expect("turn");
2129 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
2130 assert_eq!(calls(tmp.path(), "b"), 1);
2131 assert_eq!(talk.agent, "b", "the switch persists");
2132 assert!(talks.get(&talk.id).unwrap().agent == "b");
2133 let reply = talk.turns.last().unwrap();
2134 assert!(reply.body.contains("hello there"));
2135 assert!(
2136 reply.body.contains("magi task add --solo"),
2137 "a fresh seat gets the full briefing"
2138 );
2139 assert!(
2140 talk.turns
2141 .iter()
2142 .any(|t| t.body.contains("agent changed from a to b")),
2143 "the switch is noted"
2144 );
2145 }
2146
2147 #[tokio::test]
2148 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
2149 let (tmp, talks) = store();
2150 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2151 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
2152 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
2153 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2154
2155 let err = say(&mut talk, &talks, &cfg, "hi", Vec::new())
2156 .await
2157 .expect_err("every agent failed");
2158 assert!(err.to_string().contains("`a`"), "{err:#}");
2159 assert_eq!(calls(tmp.path(), "a"), 1);
2160 assert_eq!(calls(tmp.path(), "b"), 1);
2161 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
2162 }
2163
2164 #[test]
2165 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
2166 let (tmp, talks) = store();
2167 let b = counting_agent(tmp.path(), "b", "cat");
2168 let cfg = chain_config(vec![b], &["ghost", "b"]);
2169 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2170 assert_eq!(talk.agent, "b");
2171 }
2172
2173 #[tokio::test]
2174 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
2175 let (tmp, talks) = store();
2176 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2177 let b = counting_agent(tmp.path(), "b", "cat");
2178 let cfg = chain_config(vec![a, b], &["a", "b"]);
2179 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
2180 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2181 .await
2182 .expect_err("a alone, and it fails");
2183 assert_eq!(calls(tmp.path(), "b"), 0);
2184 assert_eq!(talk.agent, "a");
2185 }
2186
2187 #[tokio::test]
2188 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
2189 let (tmp, talks) = store();
2190 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2191 let b = counting_agent(tmp.path(), "b", "cat");
2192 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
2193 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
2194 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
2195 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2196 .await
2197 .expect_err("c alone, and it fails");
2198 assert_eq!(calls(tmp.path(), "b"), 0);
2199 }
2200
2201 #[tokio::test]
2202 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
2203 let (tmp, talks) = store();
2204 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2205 let cfg = config(spec);
2206 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2207
2208 say(
2209 &mut talk,
2210 &talks,
2211 &cfg,
2212 "what does the queue module do?",
2213 Vec::new(),
2214 )
2215 .await
2216 .expect("first turn");
2217 let first_prompt = &talk.turns[1].body;
2218 assert!(first_prompt.contains("magi task add --solo"));
2219 assert!(first_prompt.contains("what does the queue module do?"));
2220
2221 say(&mut talk, &talks, &cfg, "and how is it locked?", Vec::new())
2222 .await
2223 .expect("second turn");
2224 let second_prompt = &talk.turns[3].body;
2225 assert!(
2226 !second_prompt.contains("magi task add --solo"),
2227 "the briefing is sent once, not on every turn: {second_prompt}"
2228 );
2229 assert!(second_prompt.contains("and how is it locked?"));
2230 }
2231
2232 #[tokio::test]
2233 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
2234 let (tmp, talks) = store();
2235 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2236 let mut b = a.clone();
2237 b.id = "other".to_owned();
2238 let mut cfg = config(a.clone());
2239 cfg.agents.push(b.clone());
2240 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
2241 say(&mut talk, &talks, &cfg, "remember the walrus", Vec::new())
2242 .await
2243 .expect("first turn");
2244 let old_session = talk.seat.claude_session.clone();
2245 assert_eq!(talk.seat.turns, 1);
2246
2247 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
2248 assert_eq!(talk.agent, "other");
2249 assert_eq!(talk.seat.turns, 0);
2250 assert_eq!(talk.seat.agent, "other");
2251 assert_ne!(talk.seat.claude_session, old_session);
2252 let note = talk.turns.last().expect("note");
2253 assert_eq!(note.who, Who::Agent);
2254 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2255 assert!(note.body.contains("changed from"), "{}", note.body);
2256 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
2257
2258 let before = talk.turns.len();
2259 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
2260 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
2261
2262 say(&mut talk, &talks, &cfg, "what did I say?", Vec::new())
2263 .await
2264 .expect("turn after switch");
2265 let prompt = &talk.turns.last().expect("reply").body;
2266 assert!(prompt.contains("remember the walrus"), "{prompt}");
2267 assert!(prompt.contains("## magi"), "{prompt}");
2268 assert!(prompt.contains("what did I say?"), "{prompt}");
2269 }
2270
2271 #[tokio::test]
2272 async fn say_appends_the_operator_turn_then_the_agent_turn() {
2273 let (tmp, talks) = store();
2274 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2275 let cfg = config(spec);
2276 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2277
2278 say(
2279 &mut talk,
2280 &talks,
2281 &cfg,
2282 "can I rename this function?",
2283 Vec::new(),
2284 )
2285 .await
2286 .expect("say");
2287
2288 assert_eq!(talk.turns.len(), 2);
2289 assert_eq!(talk.turns[0].who, Who::Operator);
2290 assert_eq!(talk.turns[0].body, "can I rename this function?");
2291 assert_eq!(talk.turns[1].who, Who::Agent);
2292 assert_eq!(talk.turns[1].body, "go ahead");
2293 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
2294 }
2295
2296 #[tokio::test]
2297 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
2298 let (tmp, talks) = store();
2299 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2300 let cfg = config(spec);
2301 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2302
2303 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2304 .await
2305 .expect_err("a turn with no answer is an error");
2306 assert!(err.to_string().contains("no answer"), "{err}");
2307
2308 let on_disk = talks.get(&talk.id).expect("get");
2309 assert_eq!(on_disk.turns.len(), 2);
2310 assert_eq!(on_disk.turns[0].body, "check the tests");
2311 let note = &on_disk.turns[1];
2312 assert_eq!(note.who, Who::Agent);
2313 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2314 assert!(note.body.contains("your message is saved"));
2315 }
2316
2317 #[tokio::test]
2323 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
2324 let (tmp, talks) = store();
2325 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2326 let cfg = config(spec);
2327 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2328
2329 let text =
2330 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
2331 failpoint::force_put_failures(PUT_RETRIES - 1);
2334 respond(&mut talk, &talks, &cfg, &text)
2335 .await
2336 .expect("respond must survive a write failure its own retries can outlast");
2337
2338 assert_eq!(talk.turns.len(), 2);
2339 assert_eq!(talk.turns[1].who, Who::Agent);
2340 assert_eq!(talk.turns[1].body, "go ahead");
2341 let on_disk = talks.get(&talk.id).expect("get");
2342 assert_eq!(
2343 on_disk.turns, talk.turns,
2344 "the reply must reach disk despite the early write failures"
2345 );
2346 }
2347
2348 #[tokio::test]
2354 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
2355 let (tmp, talks) = store();
2356 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2357 let cfg = config(spec);
2358 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2359
2360 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2361 failpoint::force_put_failures(PUT_RETRIES);
2366 let err = respond(&mut talk, &talks, &cfg, &text)
2367 .await
2368 .expect_err("a reply that cannot be saved must be reported, not swallowed");
2369 assert!(err.to_string().contains("could not be saved"), "{err}");
2370
2371 let on_disk = talks.get(&talk.id).expect("get");
2372 assert_eq!(
2373 on_disk.turns.len(),
2374 2,
2375 "the operator turn plus a visible note"
2376 );
2377 assert_eq!(on_disk.turns[0].body, "check the tests");
2378 let note = &on_disk.turns[1];
2379 assert_eq!(note.who, Who::Agent);
2380 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2381 assert!(
2382 note.body.contains("could not be saved"),
2383 "the operator must be told the reply is missing, not left staring \
2384 at a gap with no explanation: {}",
2385 note.body
2386 );
2387 assert_eq!(
2388 talk.turns, on_disk.turns,
2389 "the in-memory talk must match what actually landed on disk"
2390 );
2391
2392 let artifacts = talks.artifacts_of(&talk.id);
2395 let stash = std::fs::read_dir(&artifacts)
2396 .expect("artifacts dir")
2397 .filter_map(|e| e.ok())
2398 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
2399 .expect("a stash file for the lost reply");
2400 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
2401 assert_eq!(stashed, "go ahead");
2402
2403 assert_eq!(
2411 on_disk.seat.turns, 1,
2412 "the note's write must carry the turn the CLI actually took"
2413 );
2414 assert_eq!(
2415 on_disk.seat.claude_session, talk.seat.claude_session,
2416 "the session id handed to the CLI must survive the failed reply"
2417 );
2418 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
2419 assert!(
2420 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
2421 "the next turn must resume, not open the same session id twice"
2422 );
2423 }
2424
2425 #[tokio::test]
2430 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
2431 let (tmp, talks) = store();
2432 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2433 let cfg = config(spec);
2434 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2435
2436 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2437 failpoint::force_put_failures(PUT_RETRIES * 2);
2440 let err = respond(&mut talk, &talks, &cfg, &text)
2441 .await
2442 .expect_err("neither the reply nor the note could be saved");
2443 assert!(err.to_string().contains("could not be saved"), "{err}");
2444
2445 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
2446 let on_disk = talks.get(&talk.id).expect("get");
2447 assert_eq!(on_disk.turns.len(), 1);
2448
2449 assert_eq!(
2459 on_disk.seat.turns, 0,
2460 "an unwritable file cannot record the turn the CLI took"
2461 );
2462 assert_eq!(
2463 talk.seat.turns, 1,
2464 "the in-memory seat still reports the turn the CLI actually took"
2465 );
2466 assert_eq!(
2467 on_disk.seat.claude_session, talk.seat.claude_session,
2468 "the session id was minted at `begin` and never changes here"
2469 );
2470 }
2471
2472 #[tokio::test]
2476 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
2477 let (tmp, talks) = store();
2478 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2479 let cfg = config(spec);
2480 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2481
2482 let att = talks
2483 .put_attachment(
2484 &talk.id,
2485 "image/png",
2486 "screenshot.png",
2487 b"pretend-png-bytes",
2488 )
2489 .expect("put attachment");
2490
2491 say(&mut talk, &talks, &cfg, "", vec![att.clone()])
2492 .await
2493 .expect("an empty body with an attachment is still a turn");
2494
2495 let operator_turn = &talk.turns[0];
2496 assert_eq!(operator_turn.who, Who::Operator);
2497 assert_eq!(operator_turn.body, "");
2498 assert_eq!(operator_turn.attachments, vec![att.clone()]);
2499
2500 let prompt = &talk.turns[1].body;
2501 let expected_path = talks
2502 .attachments_dir(&talk.id)
2503 .join(format!("{}.png", att.id));
2504 assert!(
2505 prompt.contains(&expected_path.display().to_string()),
2506 "the agent must be told the attachment's absolute path: {prompt}"
2507 );
2508 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
2509 }
2510
2511 #[test]
2520 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
2521 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
2522 let att = Attachment {
2523 id: "0".repeat(32),
2524 name: "shot.png".to_owned(),
2525 mime: "image/png".to_owned(),
2526 bytes: 3,
2527 };
2528 let path = talks
2529 .attachment_path("some-talk-id", &att)
2530 .expect("a supported mime always yields a path");
2531 assert!(
2532 path.is_absolute(),
2533 "must be absolute even off a relative store root: {}",
2534 path.display()
2535 );
2536 }
2537
2538 #[tokio::test]
2539 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
2540 let (tmp, talks) = store();
2545 let slow = mock_agent(
2546 tmp.path(),
2547 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
2548 BTreeMap::new(),
2549 );
2550 let mut cfg = config(slow);
2551 cfg.graph.timeout_talk = 1;
2552 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2553
2554 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2555 .await
2556 .expect_err("a turn that never answers is an error");
2557 assert!(
2558 err.to_string().contains("did not answer within 1s"),
2559 "{err}"
2560 );
2561
2562 let on_disk = talks.get(&talk.id).expect("get");
2563 let note = on_disk.turns.last().expect("a note turn was recorded");
2564 assert!(
2565 note.body.contains("did not answer within 1s"),
2566 "the transcript must show the configured timeout: {}",
2567 note.body
2568 );
2569 }
2570
2571 #[test]
2572 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
2573 let (tmp, talks) = store();
2574 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2575 let cfg = config(spec);
2576 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2577
2578 close(&mut talk, &talks).expect("close");
2579 assert_eq!(talk.status, TalkStatus::Closed);
2580 close(&mut talk, &talks).expect("closing twice is not an error");
2581
2582 let err =
2583 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
2584 assert!(err.to_string().contains("closed"));
2585 let _ = &cfg; }
2587
2588 #[tokio::test]
2589 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2590 let (tmp, talks) = store();
2591 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2592 let cfg = config(spec);
2593 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2596
2597 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
2601 close(&mut closed_elsewhere, &talks).expect("close");
2602 assert_eq!(
2603 talks.get(&in_flight.id).expect("reread").status,
2604 TalkStatus::Closed,
2605 "the close landed on disk before the turn finished"
2606 );
2607
2608 assert_eq!(in_flight.status, TalkStatus::Open);
2612 respond(&mut in_flight, &talks, &cfg, "one more question")
2613 .await
2614 .expect("the turn itself still completes");
2615
2616 let on_disk = talks.get(&in_flight.id).expect("reread");
2617 assert_eq!(
2618 on_disk.status,
2619 TalkStatus::Closed,
2620 "a close must stick even when a turn that started before it finishes after it"
2621 );
2622 assert!(
2625 on_disk.turns.iter().any(|t| t.body == "here you go"),
2626 "the in-flight turn's own reply is still recorded: {:?}",
2627 on_disk.turns
2628 );
2629 }
2630
2631 #[test]
2632 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
2633 let (tmp, talks) = store();
2634 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2635 let cfg = config(spec);
2636 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2639
2640 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
2643 close(&mut closed_elsewhere, &talks).expect("close");
2644 assert_eq!(
2645 talks.get(&stale.id).expect("reread").status,
2646 TalkStatus::Closed,
2647 "the close landed on disk before record was called"
2648 );
2649
2650 assert_eq!(stale.status, TalkStatus::Open);
2654 let err = record(&mut stale, &talks, "still there?", Vec::new())
2655 .expect_err("a close that landed first must be honored, not overwritten");
2656 assert!(err.to_string().contains("closed"));
2657
2658 let on_disk = talks.get(&stale.id).expect("reread");
2659 assert_eq!(
2660 on_disk.status,
2661 TalkStatus::Closed,
2662 "record must not resurrect a conversation closed while its snapshot was stale"
2663 );
2664 assert!(
2665 on_disk.turns.is_empty(),
2666 "the rejected turn must not have been appended: {:?}",
2667 on_disk.turns
2668 );
2669 let _ = &cfg; }
2671
2672 #[test]
2673 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
2674 let (tmp, talks) = store();
2675 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2676 let cfg = config(spec);
2677 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2678
2679 let held = talks.guard();
2683
2684 let talks2 = talks.clone();
2685 let id = talk.id.clone();
2686 let closing = std::thread::spawn(move || {
2687 let mut talk = talks2.get(&id).expect("get");
2688 close(&mut talk, &talks2).expect("close");
2689 });
2690
2691 std::thread::sleep(Duration::from_millis(50));
2692 assert!(
2693 !closing.is_finished(),
2694 "close must wait for the guard, not read and write while it is held - \
2695 a re-read alone narrows this window without closing it"
2696 );
2697
2698 drop(held);
2699 closing.join().expect("close thread panicked");
2700
2701 assert_eq!(
2702 talks.get(&talk.id).expect("reread").status,
2703 TalkStatus::Closed,
2704 "once the guard is free, close still lands"
2705 );
2706 let _ = &cfg; }
2708
2709 #[test]
2710 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
2711 let (tmp, talks) = store();
2712 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2713 let cfg = config(spec);
2714 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2715
2716 close(&mut talk, &talks).expect("close");
2717 assert_eq!(talk.status, TalkStatus::Closed);
2718
2719 reopen(&mut talk, &talks).expect("reopen");
2720 assert_eq!(talk.status, TalkStatus::Open);
2721 assert_eq!(
2722 talks.get(&talk.id).expect("reread").status,
2723 TalkStatus::Open
2724 );
2725
2726 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
2728 assert_eq!(talk.status, TalkStatus::Open);
2729
2730 record(&mut talk, &talks, "one more thing", Vec::new())
2731 .expect("a reopened talk takes turns again");
2732 let _ = &cfg; }
2734
2735 #[test]
2736 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
2737 let (tmp, talks) = store();
2738 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2739 let cfg = config(spec);
2740 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2741
2742 let artifacts = talks.artifacts_of(&talk.id);
2743 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
2744 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
2745
2746 talks.remove(&talk.id).expect("remove");
2747 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
2748 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
2749 assert!(
2750 talks.get(&talk.id).is_err(),
2751 "a removed talk cannot be read back"
2752 );
2753
2754 let err = talks
2755 .remove("nonexistent-id")
2756 .expect_err("unknown id refused");
2757 assert!(err.to_string().contains("no talk matches"), "{err}");
2758 let _ = &cfg; }
2760
2761 #[tokio::test]
2762 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2763 let (tmp, talks) = store();
2764 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2765 let cfg = config(spec);
2766 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2769
2770 talks.remove(&in_flight.id).expect("remove");
2771 assert!(
2772 talks.get(&in_flight.id).is_err(),
2773 "the delete landed on disk before the turn finished"
2774 );
2775
2776 respond(&mut in_flight, &talks, &cfg, "one more question")
2779 .await
2780 .expect("the turn itself still completes rather than erroring");
2781
2782 assert!(
2783 talks.get(&in_flight.id).is_err(),
2784 "a delete must stick even when a turn that started before it finishes after it"
2785 );
2786 }
2787
2788 #[test]
2789 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
2790 let (tmp, talks) = store();
2791 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2792 let cfg = config(spec);
2793 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2796
2797 talks.remove(&stale.id).expect("remove");
2798
2799 let err = record(&mut stale, &talks, "still there?", Vec::new())
2803 .expect_err("a delete that landed first must be honored, not overwritten");
2804 assert!(err.to_string().contains("deleted"), "{err}");
2805
2806 assert!(
2807 talks.get(&stale.id).is_err(),
2808 "record must not resurrect a conversation deleted while its snapshot was stale"
2809 );
2810 let _ = &cfg; }
2812
2813 #[test]
2814 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
2815 let (tmp, talks) = store();
2816 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2817 let cfg = config(spec);
2818 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2821
2822 talks.remove(&stale.id).expect("remove");
2823
2824 let err = close(&mut stale, &talks)
2828 .expect_err("a delete that landed first must be honored, not overwritten");
2829 assert!(err.to_string().contains("deleted"), "{err}");
2830
2831 assert!(
2832 talks.get(&stale.id).is_err(),
2833 "close must not resurrect a conversation deleted while its snapshot was stale"
2834 );
2835 let _ = &cfg; }
2837
2838 #[test]
2839 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
2840 let (tmp, talks) = store();
2841 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2842 let cfg = config(spec);
2843 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2846 close(&mut stale, &talks).expect("close");
2847
2848 talks.remove(&stale.id).expect("remove");
2849
2850 let err = reopen(&mut stale, &talks)
2854 .expect_err("a delete that landed first must be honored, not overwritten");
2855 assert!(err.to_string().contains("deleted"), "{err}");
2856
2857 assert!(
2858 talks.get(&stale.id).is_err(),
2859 "reopen must not resurrect a conversation deleted while its snapshot was stale"
2860 );
2861 let _ = &cfg; }
2863
2864 #[test]
2865 fn list_puts_open_talks_before_closed_ones() {
2866 let (tmp, talks) = store();
2867 let make = |id: &str, status: TalkStatus| {
2868 let mut t = Talk {
2869 schema: SCHEMA,
2870 id: id.to_owned(),
2871 repo: tmp.path().to_owned(),
2872 agent: "mock".to_owned(),
2873 status,
2874 turns: Vec::new(),
2875 pending: String::new(),
2876 pending_attachments: Vec::new(),
2877 fallback: false,
2878 created_at: Timestamp::now(),
2879 updated_at: Timestamp::now(),
2880 seat: SeatState::new(SEAT, "mock", 7),
2881 };
2882 talks.put(&mut t).expect("put");
2883 };
2884 make("20260901-000000-0001", TalkStatus::Open);
2885 make("20260902-000000-0002", TalkStatus::Open);
2886 make("20260903-000000-0003", TalkStatus::Closed);
2887
2888 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
2889 assert_eq!(
2890 ids,
2891 [
2892 "20260902-000000-0002",
2893 "20260901-000000-0001",
2894 "20260903-000000-0003"
2895 ]
2896 );
2897 assert_eq!(talks.count_open(), 2);
2898 }
2899
2900 #[test]
2901 fn tasks_of_finds_only_this_talks_own_tasks() {
2902 let dir = tempfile::tempdir().expect("tempdir");
2903 let queue = Queue::at(dir.path().join("queue"));
2904
2905 let mut mine = Task::new(
2906 "rework the loader".to_owned(),
2907 "rework the loader".to_owned(),
2908 PathBuf::from("/repo"),
2909 Source::Agent {
2910 run: "20260904-014455-ab12".to_owned(),
2911 node: "chat".to_owned(),
2912 },
2913 );
2914 queue.put(&mut mine).expect("put mine");
2915
2916 let mut theirs = Task::new(
2917 "unrelated".to_owned(),
2918 "unrelated".to_owned(),
2919 PathBuf::from("/repo"),
2920 Source::Agent {
2921 run: "20260904-090000-zz99".to_owned(),
2922 node: "implement".to_owned(),
2923 },
2924 );
2925 queue.put(&mut theirs).expect("put theirs");
2926
2927 let mut human = Task::new(
2928 "typed by hand".to_owned(),
2929 "typed by hand".to_owned(),
2930 PathBuf::from("/repo"),
2931 Source::Human,
2932 );
2933 queue.put(&mut human).expect("put human");
2934
2935 let found = tasks_of(&queue, "20260904-014455-ab12");
2936 assert_eq!(found.len(), 1);
2937 assert_eq!(found[0].id, mine.id);
2938 }
2939
2940 #[test]
2941 fn the_briefing_names_solo_task_add() {
2942 let brief = briefing(Path::new("/repo"), "en", false);
2943 assert!(brief.contains("magi task add --solo"));
2944 assert!(brief.contains("/repo"));
2945 assert!(!brief.contains("Hold this conversation in"));
2946 }
2947
2948 #[test]
2954 fn the_briefing_explains_targeting_a_different_repository_by_name() {
2955 let brief = briefing(Path::new("/repo"), "en", false);
2956 assert!(brief.contains("--repo does not have to be a full path"));
2957 assert!(brief.contains("owner/repo"));
2958 assert!(brief.contains("magi repos"));
2959 assert!(brief.contains("ask the operator"));
2960 }
2961
2962 #[test]
2963 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
2964 let brief = briefing(Path::new("/repo"), "en", false);
2965 assert!(brief.contains("--attach <path>"), "{brief}");
2966 assert!(brief.contains("deleting this conversation"), "{brief}");
2967 }
2968
2969 #[test]
2970 fn the_briefing_names_the_language_when_it_is_not_english() {
2971 let brief = briefing(Path::new("/repo"), "Japanese", false);
2972 assert!(brief.contains("Hold this conversation in Japanese"));
2973 }
2974
2975 #[test]
2976 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
2977 let read_only = briefing(Path::new("/repo"), "en", false);
2978 assert!(read_only.contains("Do not write files"));
2979 assert!(!read_only.contains("allow_write"));
2980
2981 let writable = briefing(Path::new("/repo"), "en", true);
2982 assert!(!writable.contains("Do not write files"));
2983 assert!(writable.contains("allow_write = true"));
2984 assert!(writable.contains("magi task add --solo"));
2987 assert!(writable.contains("say plainly what you"));
2988 }
2989}