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