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}
175
176const CONTEXT_WARN_PERCENT: u64 = 80;
178
179pub fn context_usage(talk: &Talk, cfg: Option<&Config>) -> ContextUsage {
192 let current = cfg.and_then(|c| c.agents.iter().find(|a| a.id == talk.agent));
193 let model = current.and_then(|a| a.model.clone());
194 let window = cfg
195 .zip(model.as_deref())
196 .and_then(|(c, m)| c.context_window(m))
197 .filter(|w| *w > 0);
198 let usage = talk
199 .turns
200 .iter()
201 .rev()
202 .find(|t| t.who == Who::Agent && !t.body.starts_with(MAGI_NOTE))
203 .and_then(|t| t.usage.as_ref());
204 let tokens = usage.map(|u| u.context_tokens);
205 let since_switch =
206 usage.is_some_and(|u| u.agent != talk.agent || (current.is_some() && u.model != model));
207 let (percent, warn) = match (tokens, window) {
208 (Some(t), Some(w)) => (
209 Some(t.saturating_mul(100) / w),
210 t.saturating_mul(100) >= w.saturating_mul(CONTEXT_WARN_PERCENT),
211 ),
212 _ => (None, false),
213 };
214 ContextUsage {
215 tokens,
216 window,
217 percent,
218 warn,
219 since_switch,
220 model,
221 }
222}
223
224#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
228#[serde(rename_all = "lowercase")]
229pub enum TalkStatus {
230 Open,
233 Closed,
235}
236
237impl TalkStatus {
238 pub fn open(self) -> bool {
240 matches!(self, Self::Open)
241 }
242
243 pub fn as_str(self) -> &'static str {
245 match self {
246 Self::Open => "open",
247 Self::Closed => "closed",
248 }
249 }
250}
251
252#[derive(Debug, Clone, Serialize, Deserialize)]
254#[serde(deny_unknown_fields)]
255pub struct Talk {
256 pub schema: u32,
258 pub id: String,
260 pub repo: PathBuf,
262 pub agent: String,
264 pub status: TalkStatus,
266 pub turns: Vec<Turn>,
268 #[serde(default)]
271 pub pending: String,
272 #[serde(default)]
274 pub pending_attachments: Vec<Attachment>,
275 #[serde(default)]
280 pub fallback: bool,
281 pub created_at: Timestamp,
283 pub updated_at: Timestamp,
285 seat: SeatState,
290}
291
292impl Talk {
293 pub fn short(&self) -> &str {
295 short(&self.id)
296 }
297}
298
299#[derive(Debug, Clone)]
301pub struct Talks {
302 root: PathBuf,
303 lock: Arc<Mutex<()>>,
312}
313
314impl Talks {
315 pub fn open() -> Self {
317 Self::at(crate::run::home().join("talks"))
318 }
319
320 pub fn at(root: PathBuf) -> Self {
323 Self {
324 root,
325 lock: Arc::new(Mutex::new(())),
326 }
327 }
328
329 fn guard(&self) -> MutexGuard<'_, ()> {
338 self.lock.lock().unwrap_or_else(PoisonError::into_inner)
339 }
340
341 pub fn root(&self) -> &Path {
343 &self.root
344 }
345
346 pub fn path_of(&self, id: &str) -> PathBuf {
348 self.root.join(format!("{id}.json"))
349 }
350
351 pub fn artifacts_of(&self, id: &str) -> PathBuf {
354 self.root.join(format!("{id}.artifacts"))
355 }
356
357 pub fn attachments_dir(&self, id: &str) -> PathBuf {
361 self.artifacts_of(id).join("attachments")
362 }
363
364 pub fn put_attachment(
373 &self,
374 id: &str,
375 mime: &str,
376 name: &str,
377 data: &[u8],
378 ) -> Result<Attachment> {
379 let dir = self.attachments_dir(id);
380 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
381 let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
382 let att = Attachment {
383 id: new_attachment_id(),
384 name: name.to_owned(),
385 mime: mime.to_owned(),
386 bytes: data.len() as u64,
387 };
388 std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
389 .with_context(|| format!("write attachment {}", att.id))?;
390 std::fs::write(
391 dir.join(format!("{}.json", att.id)),
392 serde_json::to_string(&att).context("serialize attachment")?,
393 )
394 .with_context(|| format!("write attachment metadata {}", att.id))?;
395 Ok(att)
396 }
397
398 pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
407 if !valid_attachment_id(att_id) {
408 return Ok(None);
409 }
410 let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
411 if !meta_path.is_file() {
412 return Ok(None);
413 }
414 let att = serde_json::from_str(
415 &std::fs::read_to_string(&meta_path)
416 .with_context(|| format!("read {}", meta_path.display()))?,
417 )
418 .with_context(|| format!("parse {}", meta_path.display()))?;
419 Ok(Some(att))
420 }
421
422 pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
426 let Some(att) = self.attachment_meta(id, att_id)? else {
427 return Ok(None);
428 };
429 let ext = attachment_ext(&att.mime).with_context(|| {
430 format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
431 })?;
432 let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
433 let data =
434 std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
435 Ok(Some((att, data)))
436 }
437
438 fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
455 let ext = attachment_ext(&att.mime)?;
456 let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
457 std::path::absolute(&path).ok()
458 }
459
460 pub fn put(&self, t: &mut Talk) -> Result<()> {
469 std::fs::create_dir_all(&self.root)
470 .with_context(|| format!("create {}", self.root.display()))?;
471 t.updated_at = Timestamp::now();
472 let body = serde_json::to_string_pretty(t).context("serialize talk")?;
473 let path = self.path_of(&t.id);
474 let tmp = path.with_extension("json.tmp");
475 write_atomic(&tmp, &path, &body)
476 }
477
478 pub fn get(&self, id: &str) -> Result<Talk> {
480 let resolved = self.resolve_id(id)?;
481 read_path(&self.path_of(&resolved))
482 }
483
484 pub fn list(&self) -> Vec<Talk> {
487 let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
488 .into_iter()
489 .flatten()
490 .flatten()
491 .map(|e| e.path())
492 .filter(|p| p.extension().is_some_and(|x| x == "json"))
493 .filter_map(|p| read_path(&p).ok())
494 .collect();
495 all.sort_unstable_by(|a, b| {
496 let rank = |t: &Talk| u8::from(!t.status.open());
497 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
498 });
499 all
500 }
501
502 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
504 if self.path_of(prefix).is_file() {
505 return Ok(prefix.to_owned());
506 }
507 let hits: Vec<String> = self
508 .list()
509 .into_iter()
510 .map(|t| t.id)
511 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
512 .collect();
513 match hits.len() {
514 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
515 0 => bail!("no talk matches `{prefix}`"),
516 _ => bail!(
517 "`{prefix}` matches {} talks: {}",
518 hits.len(),
519 hits.join(", ")
520 ),
521 }
522 }
523
524 pub fn revision(&self) -> u64 {
527 std::fs::read_dir(&self.root)
528 .into_iter()
529 .flatten()
530 .flatten()
531 .filter_map(|e| e.metadata().ok())
532 .filter_map(|m| m.modified().ok())
533 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
534 .map(|d| d.as_millis() as u64)
535 .max()
536 .unwrap_or(0)
537 }
538
539 pub fn count_open(&self) -> usize {
541 self.list().iter().filter(|t| t.status.open()).count()
542 }
543
544 pub fn remove(&self, id: &str) -> Result<()> {
557 let _guard = self.guard();
558 let resolved = self.resolve_id(id)?;
559 let path = self.path_of(&resolved);
560 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
561 let artifacts = self.artifacts_of(&resolved);
562 if artifacts.is_dir() {
563 std::fs::remove_dir_all(&artifacts)
564 .with_context(|| format!("remove {}", artifacts.display()))?;
565 }
566 Ok(())
567 }
568}
569
570pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
580 let repo = repo.canonicalize().unwrap_or(repo);
583 let spec = match agent {
586 Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
587 None => agent::pick_chain(
588 &cfg.agents,
589 cfg.roles.chatter.as_ref(),
590 &agent::installed,
591 "chatter",
592 )?
593 .remove(0),
594 };
595
596 let now = Timestamp::now();
597 let mut talk = Talk {
598 schema: SCHEMA,
599 id: new_id(),
600 repo,
601 agent: spec.id.clone(),
602 status: TalkStatus::Open,
603 turns: Vec::new(),
604 pending: String::new(),
605 pending_attachments: Vec::new(),
606 fallback: agent.is_none(),
607 created_at: now,
608 updated_at: now,
609 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
610 };
611 store.put(&mut talk)?;
612 Ok(talk)
613}
614
615pub fn record(
622 talk: &mut Talk,
623 store: &Talks,
624 text: &str,
625 attachments: Vec<Attachment>,
626) -> Result<String> {
627 let _guard = store.guard();
635 let Ok(fresh) = store.get(&talk.id) else {
640 bail!("talk {} was deleted", talk.short());
641 };
642 talk.status = fresh.status;
643 talk.pending = fresh.pending;
646 talk.pending_attachments = fresh.pending_attachments;
647 if !talk.status.open() {
648 bail!(
649 "talk {} is {} and takes no more turns",
650 talk.short(),
651 talk.status.as_str()
652 );
653 }
654 let text = text.trim();
655 if text.is_empty() && attachments.is_empty() {
656 bail!("nothing to say");
657 }
658 talk.turns.push(Turn {
659 who: Who::Operator,
660 body: text.to_owned(),
661 at: Timestamp::now(),
662 attachments,
663 usage: None,
664 });
665 store.put(talk)?;
666 Ok(text.to_owned())
667}
668
669pub fn queue(
671 talk: &mut Talk,
672 store: &Talks,
673 text: &str,
674 attachments: Vec<Attachment>,
675) -> Result<()> {
676 let text = text.trim();
677 if text.is_empty() && attachments.is_empty() {
678 bail!("nothing to say");
679 }
680 let _guard = store.guard();
681 let mut fresh = store
682 .get(&talk.id)
683 .with_context(|| format!("talk {} was deleted", talk.short()))?;
684 if !fresh.status.open() {
685 bail!(
686 "talk {} is {} and takes no more turns",
687 fresh.short(),
688 fresh.status.as_str()
689 );
690 }
691 if !text.is_empty() {
692 if fresh.pending.is_empty() {
693 fresh.pending = text.to_owned();
694 } else {
695 fresh.pending.push_str("\n\n");
696 fresh.pending.push_str(text);
697 }
698 }
699 fresh.pending_attachments.extend(attachments);
700 store.put(&mut fresh)?;
701 *talk = fresh;
702 Ok(())
703}
704
705pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
707 let _guard = store.guard();
708 let mut fresh = store
709 .get(&talk.id)
710 .with_context(|| format!("talk {} was deleted", talk.short()))?;
711 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
712 *talk = fresh;
713 return Ok(None);
714 }
715 let text = std::mem::take(&mut fresh.pending);
716 let attachments = std::mem::take(&mut fresh.pending_attachments);
717 fresh.turns.push(Turn {
718 who: Who::Operator,
719 body: text.clone(),
720 at: Timestamp::now(),
721 attachments,
722 usage: None,
723 });
724 store.put(&mut fresh)?;
725 *talk = fresh;
726 Ok(Some(text))
727}
728
729pub async fn say(
732 talk: &mut Talk,
733 store: &Talks,
734 cfg: &Config,
735 text: &str,
736 attachments: Vec<Attachment>,
737) -> Result<()> {
738 let text = record(talk, store, text, attachments)?;
739 turn(talk, store, cfg, &text).await
740}
741
742pub async fn respond(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
744 turn(talk, store, cfg, text).await
745}
746
747pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
765 let _guard = store.guard();
766 let mut fresh = store
767 .get(&talk.id)
768 .with_context(|| format!("talk {} was deleted", talk.short()))?;
769 fresh.status = TalkStatus::Closed;
770 fresh.pending.clear();
772 fresh.pending_attachments.clear();
773 store.put(&mut fresh)?;
774 *talk = fresh;
775 Ok(())
776}
777
778pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
789 let _guard = store.guard();
790 let mut fresh = store
791 .get(&talk.id)
792 .with_context(|| format!("talk {} was deleted", talk.short()))?;
793 fresh.status = TalkStatus::Open;
794 store.put(&mut fresh)?;
795 *talk = fresh;
796 Ok(())
797}
798
799pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
812 let _guard = store.guard();
813 let mut fresh = store
814 .get(&talk.id)
815 .with_context(|| format!("talk {} was deleted", talk.short()))?;
816 if fresh.agent == spec.id {
817 *talk = fresh;
818 return Ok(false);
819 }
820 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
821 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
822 fresh.fallback = false;
824 fresh.turns.push(Turn {
825 who: Who::Agent,
826 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
827 at: Timestamp::now(),
828 attachments: Vec::new(),
829 usage: None,
830 });
831 store.put(&mut fresh)?;
832 *talk = fresh;
833 Ok(true)
834}
835
836pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
838 let _guard = store.guard();
839 let mut fresh = store
840 .get(&talk.id)
841 .with_context(|| format!("talk {} was deleted", talk.short()))?;
842 fresh.pending.clear();
843 fresh.pending_attachments.clear();
844 store.put(&mut fresh)?;
845 *talk = fresh;
846 Ok(())
847}
848
849pub fn clear_pending_if_matches(
851 talk: &mut Talk,
852 store: &Talks,
853 expected_text: &str,
854 expected_attachments: &[String],
855) -> Result<bool> {
856 let _guard = store.guard();
857 let mut fresh = store
858 .get(&talk.id)
859 .with_context(|| format!("talk {} was deleted", talk.short()))?;
860 if !pending_matches(&fresh, expected_text, expected_attachments) {
861 *talk = fresh;
862 return Ok(false);
863 }
864 fresh.pending.clear();
865 fresh.pending_attachments.clear();
866 store.put(&mut fresh)?;
867 *talk = fresh;
868 Ok(true)
869}
870
871pub fn edit_pending_text(
875 talk: &mut Talk,
876 store: &Talks,
877 text: &str,
878 expected_text: &str,
879 expected_attachments: &[String],
880) -> Result<bool> {
881 let _guard = store.guard();
882 let mut fresh = store
883 .get(&talk.id)
884 .with_context(|| format!("talk {} was deleted", talk.short()))?;
885 if !pending_matches(&fresh, expected_text, expected_attachments) {
886 *talk = fresh;
887 return Ok(false);
888 }
889 fresh.pending = text.trim().to_owned();
890 store.put(&mut fresh)?;
891 *talk = fresh;
892 Ok(true)
893}
894
895fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
896 talk.pending == expected_text
897 && talk
898 .pending_attachments
899 .iter()
900 .map(|attachment| &attachment.id)
901 .eq(expected_attachments.iter())
902}
903
904async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
911 let spec = cfg
912 .agents
913 .iter()
914 .find(|a| a.id == talk.agent)
915 .with_context(|| {
916 format!(
917 "talk {} was opened with agent `{}`, which is no longer in \
918 the roster; restore it in magi.toml or start a new \
919 conversation",
920 talk.short(),
921 talk.agent
922 )
923 })?;
924
925 let last_note = attachment_note(
929 store,
930 &talk.id,
931 talk.turns
932 .last()
933 .map_or(&[][..], |t| t.attachments.as_slice()),
934 );
935
936 let attachment_paths: Vec<PathBuf> = talk
942 .turns
943 .iter()
944 .flat_map(|t| t.attachments.iter())
945 .filter_map(|a| store.attachment_path(&talk.id, a))
946 .collect();
947
948 let artifacts = store.artifacts_of(&talk.id);
949 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
952 let stem = format!("turn-{}", operator_turns.max(1));
953 let cache_dir = cfg.cache_dir();
956
957 let mut chain = vec![spec.clone()];
960 if let Some(choice) = cfg.roles.chatter.as_ref()
961 && talk.fallback
962 {
963 for id in choice.ids() {
964 if id == talk.agent || chain.iter().any(|s| s.id == id) {
965 continue;
966 }
967 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
968 Ok(s) => chain.push(s),
969 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
970 }
971 }
972 }
973
974 let mut outcome = None;
975 let mut fell_back_from: Option<String> = None;
976 let mut first_try: Option<(String, SeatState)> = None;
979 for (n, spec) in chain.iter().enumerate() {
980 if n > 0 {
981 if first_try.is_none() {
982 first_try = Some((talk.agent.clone(), talk.seat.clone()));
983 }
984 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
985 fell_back_from.get_or_insert_with(|| talk.agent.clone());
988 talk.agent = spec.id.clone();
989 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
990 }
991 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
992 let first_ever = talk.turns.len() <= 1;
993 let body = if talk.seat.turns == 0 && first_ever {
994 format!(
995 "{}\n\n# Operator\n\n{text}{last_note}",
996 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write)
997 )
998 } else if talk.seat.turns == 0 {
999 format!(
1002 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1003 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write),
1004 transcript(talk, store)
1005 )
1006 } else if resuming {
1007 format!("{text}{last_note}")
1008 } else {
1009 format!("{}\n\n{text}{last_note}", transcript(talk, store))
1010 };
1011 let attempt_stem = if n == 0 {
1012 stem.clone()
1013 } else {
1014 format!("{stem}-{}", spec.id)
1015 };
1016 let inv = Invocation {
1017 cwd: &talk.repo,
1018 prompt: &body,
1019 timeout: turn_timeout(cfg),
1020 allow_write: cfg.talk.allow_write,
1025 sessions: cfg.graph.sessions,
1026 artifacts: &artifacts,
1027 stem: &attempt_stem,
1028 run: &talk.id,
1031 node: "chat",
1032 cache_dir: cache_dir.as_deref(),
1033 attachments: &attachment_paths,
1034 writable: &[],
1035 };
1036 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1037 let advance = agent::chain_advances(&result);
1038 if n == 0 || !advance {
1039 outcome = Some(result);
1040 } else {
1041 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1043 }
1044 if !advance {
1045 break;
1046 }
1047 }
1048 if outcome.as_ref().is_some_and(agent::chain_advances) {
1049 if let Some((id, seat)) = first_try {
1052 talk.agent = id;
1053 talk.seat = seat;
1054 fell_back_from = None;
1055 }
1056 }
1057 let outcome = outcome.expect("a chain holds at least one agent");
1058 let note = |why: String| Turn {
1059 who: Who::Agent,
1060 body: format!("{MAGI_NOTE}{why}"),
1061 at: Timestamp::now(),
1062 attachments: Vec::new(),
1063 usage: None,
1064 };
1065 let (reply, failure) = match outcome {
1066 Err(e) => (
1067 note(format!("could not run agent `{}`: {e}", talk.agent)),
1068 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1069 ),
1070 Ok(out) if out.quota_exhausted() => {
1071 let reset = out
1072 .quota
1073 .as_ref()
1074 .and_then(|q| q.reset.clone())
1075 .map_or_else(String::new, |r| format!(" (resets {r})"));
1076 let why = format!(
1077 "agent `{}` is out of quota{reset}; your message is saved, so \
1078 say it again when the window reopens",
1079 talk.agent
1080 );
1081 (note(why.clone()), Some(why))
1082 }
1083 Ok(out) if out.timed_out => {
1084 let why = format!(
1085 "agent `{}` did not answer within {}s; your message is saved",
1086 talk.agent,
1087 turn_timeout(cfg).as_secs()
1088 );
1089 (note(why.clone()), Some(why))
1090 }
1091 Ok(out) if !out.usable() => {
1092 let why = format!(
1093 "agent `{}` produced no answer (exit {}); your message is saved",
1094 talk.agent,
1095 out.exit_code
1096 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1097 );
1098 (note(why.clone()), Some(why))
1099 }
1100 Ok(out) => (
1101 Turn {
1102 who: Who::Agent,
1103 body: out.text.trim().to_owned(),
1104 at: Timestamp::now(),
1105 attachments: Vec::new(),
1106 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1109 context_tokens,
1110 agent: talk.agent.clone(),
1111 model: cfg
1112 .agents
1113 .iter()
1114 .find(|a| a.id == talk.agent)
1115 .and_then(|a| a.model.clone()),
1116 }),
1117 },
1118 None,
1119 ),
1120 };
1121
1122 let _guard = store.guard();
1134 let Ok(fresh) = store.get(&talk.id) else {
1140 return Ok(());
1141 };
1142 talk.status = fresh.status;
1143 talk.pending = fresh.pending;
1147 talk.pending_attachments = fresh.pending_attachments;
1148 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1149 talk.turns.push(note(format!(
1152 "agent changed from {from} to {} (fallback)",
1153 talk.agent
1154 )));
1155 }
1156 talk.turns.push(reply);
1157 if let Err(put_err) = store.put(talk) {
1158 let lost = talk.turns.pop().expect("just pushed above");
1167 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1168 let why = match &stash {
1169 Ok(path) => format!(
1170 "agent `{}` answered, but the reply could not be saved to \
1171 this conversation ({put_err:#}); the raw text was kept at \
1172 {} - your message is saved, ask again",
1173 talk.agent,
1174 path.display()
1175 ),
1176 Err(stash_err) => format!(
1177 "agent `{}` answered, but the reply could not be saved to \
1178 this conversation ({put_err:#}), and it could not be kept \
1179 anywhere else either ({stash_err:#}); your message is \
1180 saved, ask again",
1181 talk.agent
1182 ),
1183 };
1184 talk.turns.push(note(why.clone()));
1185 return match store.put(talk) {
1192 Ok(()) => bail!("{why}"),
1193 Err(note_err) => {
1194 talk.turns.pop();
1214 Err(note_err).context(why)
1215 }
1216 };
1217 }
1218
1219 match failure {
1220 Some(why) => bail!("{why}"),
1221 None => Ok(()),
1222 }
1223}
1224
1225fn transcript(talk: &Talk, store: &Talks) -> String {
1228 let mut out = String::from(
1229 "This conversation cannot resume on the CLI's side, so here is \
1230 everything said so far; answer only the last message.\n",
1231 );
1232 for t in &talk.turns {
1233 let who = match t.who {
1234 Who::Operator => "operator",
1235 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1236 Who::Agent => "you",
1237 };
1238 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1239 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1240 }
1241 out
1242}
1243
1244fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1249 if attachments.is_empty() {
1250 return String::new();
1251 }
1252 let mut out = String::from(
1253 "\n\nThe operator attached the image(s) below to this message. Open \
1254 and look at each one before you answer.\n",
1255 );
1256 for att in attachments {
1257 if let Some(path) = store.attachment_path(talk_id, att) {
1258 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1259 }
1260 }
1261 out.push('\n');
1262 out
1263}
1264
1265pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
1287 let write_policy = if allow_write {
1288 "Write access is enabled for this conversation (`allow_write = \
1289 true`), so you may write files - but only a small, \
1290 already-decided edit the operator names outright in this \
1291 conversation, not an implementation. This is a permission on the \
1292 conversation as a whole, not a property of whichever repository \
1293 it happened to start in: if the operator names a different \
1294 repository for that small edit, the policy allows it there too. \
1295 Your own tool may still confine writes to the repository this \
1296 conversation started in regardless - if a write elsewhere is \
1297 refused, say so plainly rather than working around it. Once you \
1298 have made an edit, say plainly what you edited. Anything bigger, \
1299 or anything still open-ended, still goes through the queue below \
1300 rather than being done here."
1301 } else {
1302 "Do not write files. Implementing a change is not this \
1303 conversation's job; a separate, blind competition of agents does \
1304 that, and a repository this conversation has already edited would \
1305 make their diffs unjudgeable."
1306 };
1307 let mut out = format!(
1308 "You are magi's standing conversation partner for its operator, who \
1309 usually has this open on a phone. Keep replies short: no preamble, \
1310 no restating what they just said.\n\n\
1311 # Repository\n\n{repo}\n\n\
1312 You may look around: read files, run shell commands, search history, \
1313 run tests - whatever answers the question. {write_policy}\n\n\
1314 A short, command-shaped message (\"list\", \"info <id>\", \"show \
1315 3cbf\") is almost always the operator asking you to look something \
1316 up, not an instruction to file - answer it yourself with `magi \
1317 list`, `magi show <id>`, `magi task list`, or the like, the same way \
1318 you would answer any other question in this conversation.\n\n\
1319 # When the operator wants something done\n\n\
1320 Run:\n\n\
1321 magi task add --solo --repo {repo} <instruction>\n\n\
1322 and tell the operator the task id it prints, so they can follow it \
1323 from the Queue. If it refuses with a duplicate warning (the \
1324 instruction names a branch, commit or pull request that an \
1325 unfinished task, run or PR already owns), do not repeat it with \
1326 --force yourself: tell the operator what it matched and let them \
1327 decide. Write <instruction> so that an implementer who has \
1328 never seen this conversation can act on it alone - it is everything \
1329 they get. Use --solo: it runs the task through one implementer \
1330 straight into review instead of the usual multi-agent competition, \
1331 which is the right shape for a change this conversation has already \
1332 settled, rather than one still worth several independent takes.\n\n\
1333 If the operator asks for something in a different repository, \
1334 --repo does not have to be a full path: --repo owner/repo (or just \
1335 repo, when that is unambiguous) is resolved against local checkouts \
1336 the same way `magi repos` lists them. If the command fails because \
1337 nothing matches or more than one checkout shares that name, ask the \
1338 operator which repository they mean (or run `magi repos` yourself \
1339 to see the candidates) rather than guessing.\n\n\
1340 If the operator attached an image (a screenshot, say) that the task \
1341 is about, pass it with `--attach <path>`, using the absolute path \
1342 the turn's attachment note gives; repeat the flag for several. \
1343 `magi task add --solo --attach <path> <instruction>` copies the \
1344 file into the task, so the implementer receives it. Do not paste the \
1345 path into <instruction> instead: deleting this conversation deletes \
1346 its attachments, and then that path reaches no one.\n",
1347 repo = repo.display(),
1348 );
1349 out.push_str(&language_note(language));
1350 out
1351}
1352
1353fn language_note(language: &str) -> String {
1356 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
1357 String::new()
1358 } else {
1359 format!("\nHold this conversation in {language}.\n")
1360 }
1361}
1362
1363pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
1370 let mut tasks: Vec<Task> = queue
1371 .list()
1372 .into_iter()
1373 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
1374 .collect();
1375 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
1376 tasks
1377}
1378
1379fn read_path(path: &Path) -> Result<Talk> {
1380 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1381 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
1382}
1383
1384const PUT_RETRIES: u32 = 5;
1387
1388fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1400 let mut last_err = None;
1401 for attempt in 0..PUT_RETRIES {
1402 if attempt > 0 {
1403 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
1404 }
1405 match try_write_atomic(tmp, path, body) {
1406 Ok(()) => return Ok(()),
1407 Err(e) => last_err = Some(e),
1408 }
1409 }
1410 Err(last_err.expect("the loop above always runs at least once"))
1411}
1412
1413fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1414 #[cfg(test)]
1415 if failpoint::take_forced_put_failure() {
1416 bail!("simulated write failure (test)");
1417 }
1418 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
1419 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
1420 Ok(())
1421}
1422
1423fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
1428 let dir = store.artifacts_of(id);
1429 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1430 let path = dir.join(format!("{stem}-lost.txt"));
1431 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
1432 Ok(path)
1433}
1434
1435#[cfg(test)]
1442mod failpoint {
1443 use std::cell::Cell;
1444
1445 thread_local! {
1446 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
1447 }
1448
1449 pub(super) fn force_put_failures(count: u32) {
1452 FORCE_PUT_FAILURES.with(|c| c.set(count));
1453 }
1454
1455 pub(super) fn take_forced_put_failure() -> bool {
1458 FORCE_PUT_FAILURES.with(|c| {
1459 let n = c.get();
1460 if n == 0 {
1461 false
1462 } else {
1463 c.set(n - 1);
1464 true
1465 }
1466 })
1467 }
1468}
1469
1470fn short(id: &str) -> &str {
1471 id.split('-').next_back().unwrap_or(id)
1472}
1473
1474fn new_id() -> String {
1475 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1476 let seed = crate::rng::entropy();
1477 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1478}
1479
1480fn attachment_ext(mime: &str) -> Option<&'static str> {
1485 match mime {
1486 "image/png" => Some("png"),
1487 "image/jpeg" => Some("jpg"),
1488 "image/gif" => Some("gif"),
1489 "image/webp" => Some("webp"),
1490 _ => None,
1491 }
1492}
1493
1494pub fn valid_attachment_id(id: &str) -> bool {
1499 id.len() == 32
1500 && id
1501 .bytes()
1502 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
1503}
1504
1505fn new_attachment_id() -> String {
1509 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
1510 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
1511}
1512
1513#[cfg(test)]
1514mod tests {
1515 use std::collections::BTreeMap;
1516
1517 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
1518 use crate::queue::{Queue, Source, Task};
1519
1520 use super::*;
1521
1522 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
1523 AgentSpec {
1524 id: id.to_owned(),
1525 kind: AgentKind::Command,
1526 model: model.map(str::to_owned),
1527 command: Vec::new(),
1528 extra_args: Vec::new(),
1529 env: BTreeMap::new(),
1530 prompt_delivery: None,
1531 }
1532 }
1533
1534 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
1535 Talk {
1536 schema: SCHEMA,
1537 id: "20260904-014455-ab12".to_owned(),
1538 repo: PathBuf::from("."),
1539 agent: agent.to_owned(),
1540 status: TalkStatus::Open,
1541 turns,
1542 pending: String::new(),
1543 pending_attachments: Vec::new(),
1544 fallback: false,
1545 created_at: Timestamp::now(),
1546 updated_at: Timestamp::now(),
1547 seat: SeatState::new(SEAT, agent, 1),
1548 }
1549 }
1550
1551 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
1552 Turn {
1553 who: Who::Agent,
1554 body: body.to_owned(),
1555 at: Timestamp::now(),
1556 attachments: Vec::new(),
1557 usage: usage.map(|(t, a, m)| TurnUsage {
1558 context_tokens: t,
1559 agent: a.to_owned(),
1560 model: m.map(str::to_owned),
1561 }),
1562 }
1563 }
1564
1565 fn ctx_config(windows: &[(&str, u64)]) -> Config {
1566 Config {
1567 agents: vec![
1568 ctx_agent("small", Some("small-model")),
1569 ctx_agent("big", Some("big-model")),
1570 ctx_agent("plain", None),
1571 ],
1572 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
1573 ..Config::default()
1574 }
1575 }
1576
1577 #[test]
1578 fn context_usage_computes_percent_and_warns_at_eighty() {
1579 let cfg = ctx_config(&[("small-model", 1000)]);
1580 let at = |tokens| {
1581 let t = ctx_talk(
1582 "small",
1583 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
1584 );
1585 context_usage(&t, Some(&cfg))
1586 };
1587 let u = at(799);
1588 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
1589 let u = at(800);
1590 assert_eq!((u.percent, u.warn), (Some(80), true));
1591 let u = at(1500);
1592 assert_eq!((u.percent, u.warn), (Some(150), true));
1593 assert!(!u.since_switch);
1594 }
1595
1596 #[test]
1597 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
1598 let cfg = ctx_config(&[("small-model", 1000)]);
1599 let t = ctx_talk(
1600 "small",
1601 vec![
1602 reply("old", Some((900, "small", Some("small-model")))),
1603 reply("new", None),
1604 ],
1605 );
1606 let u = context_usage(&t, Some(&cfg));
1607 assert_eq!((u.tokens, u.percent, u.warn), (None, None, false));
1608 let t = ctx_talk(
1610 "small",
1611 vec![
1612 reply("old", Some((900, "small", Some("small-model")))),
1613 reply("magi: could not run agent", None),
1614 ],
1615 );
1616 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
1617 assert_eq!(
1618 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
1619 None
1620 );
1621 }
1622
1623 #[test]
1624 fn context_usage_without_a_window_shows_tokens_only() {
1625 let cfg = ctx_config(&[]);
1626 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
1628 let u = context_usage(&t, Some(&cfg));
1629 assert_eq!(
1630 (u.tokens, u.window, u.percent, u.warn),
1631 (Some(5000), None, None, false)
1632 );
1633 let t = ctx_talk(
1634 "small",
1635 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
1636 );
1637 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
1638 assert_eq!(context_usage(&t, None).window, None);
1640 }
1641
1642 #[test]
1643 fn context_usage_switching_model_changes_the_denominator() {
1644 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
1645 let used = reply("hi", Some((900, "small", Some("small-model"))));
1646 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
1647 assert_eq!(
1648 (before.percent, before.warn, before.since_switch),
1649 (Some(90), true, false)
1650 );
1651 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
1654 assert_eq!(after.window, Some(10_000));
1655 assert_eq!(
1656 (after.percent, after.warn, after.since_switch),
1657 (Some(9), false, true)
1658 );
1659 assert_eq!(after.model.as_deref(), Some("big-model"));
1660 }
1661
1662 #[test]
1663 fn a_turn_recorded_before_usage_existed_still_reads() {
1664 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
1665 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
1666 assert!(turn.usage.is_none());
1667 let json = serde_json::to_string(&turn).expect("serialize");
1668 assert!(
1669 !json.contains("usage"),
1670 "absent usage is not written: {json}"
1671 );
1672 }
1673
1674 fn store() -> (tempfile::TempDir, Talks) {
1676 let tmp = tempfile::tempdir().expect("tempdir");
1677 let talks = Talks::at(tmp.path().join("talks"));
1678 (tmp, talks)
1679 }
1680
1681 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
1685 let path = dir.join("mock-talk-agent.sh");
1686 std::fs::write(&path, script).expect("write mock");
1687 AgentSpec {
1688 id: "mock".to_owned(),
1689 kind: AgentKind::Command,
1690 model: None,
1691 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
1692 extra_args: Vec::new(),
1693 env,
1694 prompt_delivery: None,
1695 }
1696 }
1697
1698 fn config(spec: AgentSpec) -> Config {
1699 Config {
1700 agents: vec![spec],
1701 graph: Graph {
1702 language: "en".to_owned(),
1703 ..Graph::default()
1704 },
1705 ..Config::default()
1706 }
1707 }
1708
1709 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
1711
1712 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
1714
1715 const ECHO: &str = "#!/bin/sh\ncat\n";
1718
1719 fn env(reply: &str) -> BTreeMap<String, String> {
1720 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
1721 }
1722
1723 #[test]
1724 fn the_frozen_json_field_names_round_trip_through_disk() {
1725 let (tmp, talks) = store();
1726 let mut talk = Talk {
1727 schema: SCHEMA,
1728 id: "20260904-014455-ab12".to_owned(),
1729 repo: tmp.path().to_owned(),
1730 agent: "sonnet".to_owned(),
1731 status: TalkStatus::Open,
1732 turns: Vec::new(),
1733 pending: String::new(),
1734 pending_attachments: Vec::new(),
1735 fallback: false,
1736 created_at: Timestamp::now(),
1737 updated_at: Timestamp::now(),
1738 seat: SeatState::new(SEAT, "sonnet", 7),
1739 };
1740 talks.put(&mut talk).expect("put");
1741
1742 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
1743 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
1744 for field in [
1745 "schema",
1746 "id",
1747 "repo",
1748 "agent",
1749 "status",
1750 "turns",
1751 "created_at",
1752 "updated_at",
1753 ] {
1754 assert!(v.get(field).is_some(), "missing field `{field}`");
1755 }
1756 assert_eq!(v["schema"], 1);
1757 assert_eq!(v["status"], "open");
1758
1759 let back = talks.get(&talk.id).expect("get");
1760 assert_eq!(back.id, talk.id);
1761 assert_eq!(back.status, TalkStatus::Open);
1762 }
1763
1764 #[test]
1765 fn opening_a_talk_takes_no_agent_turn() {
1766 let (tmp, talks) = store();
1767 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1771 let cfg = config(spec);
1772
1773 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1774 assert_eq!(talk.status, TalkStatus::Open);
1775 assert!(talk.turns.is_empty(), "nothing has been said yet");
1776
1777 let on_disk = talks.get(&talk.id).expect("get");
1778 assert_eq!(on_disk.turns.len(), 0);
1779 }
1780
1781 #[test]
1789 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
1790 let (tmp, talks) = store();
1791 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1792 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1793 chatter_spec.id = "chatter-mock".to_owned();
1794
1795 let mut cfg = Config {
1796 agents: vec![first_spec.clone(), chatter_spec.clone()],
1797 graph: Graph {
1798 language: "en".to_owned(),
1799 ..Graph::default()
1800 },
1801 ..Config::default()
1802 };
1803 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
1804
1805 let talk =
1806 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
1807 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
1808
1809 cfg.roles.chatter = None;
1810 let fallback =
1811 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
1812 assert_eq!(
1813 fallback.agent, first_spec.id,
1814 "unset chatter must fall back to agent::pick's own default order"
1815 );
1816 }
1817
1818 #[test]
1821 fn a_talk_recorded_without_attachments_still_reads() {
1822 let (tmp, talks) = store();
1823 let path = talks.path_of("20260904-014455-ab12");
1824 std::fs::create_dir_all(talks.root()).expect("talks dir");
1825 std::fs::write(
1826 &path,
1827 serde_json::json!({
1828 "schema": 1,
1829 "id": "20260904-014455-ab12",
1830 "repo": tmp.path(),
1831 "agent": "sonnet",
1832 "status": "open",
1833 "turns": [
1834 { "who": "operator", "body": "still there?",
1835 "at": Timestamp::now().to_string() },
1836 ],
1837 "created_at": Timestamp::now().to_string(),
1838 "updated_at": Timestamp::now().to_string(),
1839 "seat": SeatState::new(SEAT, "sonnet", 7),
1840 })
1841 .to_string(),
1842 )
1843 .expect("write pre-attachments talk");
1844
1845 let talk = talks.get("20260904-014455-ab12").expect("must still read");
1846 assert!(talk.turns[0].attachments.is_empty());
1847 }
1848
1849 #[test]
1850 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
1851 let (tmp, talks) = store();
1852 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1853 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1854
1855 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
1856 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
1857 let saved = talks.get(&talk.id).expect("reload queued talk");
1858 assert_eq!(saved.pending, "first\n\nsecond");
1859 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
1860
1861 let drained = drain(&mut talk, &talks).expect("drain");
1862 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
1863 let saved = talks.get(&talk.id).expect("reload drained talk");
1864 assert!(saved.pending.is_empty());
1865 assert_eq!(saved.turns.len(), 1);
1866 assert_eq!(saved.turns[0].body, "first\n\nsecond");
1867 }
1868
1869 #[test]
1870 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
1871 let (tmp, talks) = store();
1872 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1873 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1874 let attachment = Attachment {
1875 id: "a".repeat(32),
1876 name: "shot.png".to_owned(),
1877 mime: "image/png".to_owned(),
1878 bytes: 3,
1879 };
1880
1881 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
1882 assert!(
1883 edit_pending_text(
1884 &mut talk,
1885 &talks,
1886 "corrected",
1887 "first",
1888 std::slice::from_ref(&attachment.id),
1889 )
1890 .expect("edit")
1891 );
1892 let saved = talks.get(&talk.id).expect("reload edited draft");
1893 assert_eq!(saved.pending, "corrected");
1894 assert_eq!(saved.pending_attachments, vec![attachment]);
1895
1896 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
1897 assert!(
1898 !edit_pending_text(
1899 &mut talk,
1900 &talks,
1901 "stale edit",
1902 "corrected",
1903 &["a".repeat(32)],
1904 )
1905 .expect("stale edit is a conflict")
1906 );
1907 assert_eq!(
1908 talks.get(&talk.id).expect("reload after conflict").pending,
1909 "corrected\n\nlater"
1910 );
1911 assert!(
1912 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
1913 .expect("stale clear is a conflict")
1914 );
1915 assert_eq!(
1916 talks
1917 .get(&talk.id)
1918 .expect("reload after stale clear")
1919 .pending,
1920 "corrected\n\nlater"
1921 );
1922 }
1923
1924 #[tokio::test]
1925 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
1926 let (tmp, talks) = store();
1927 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
1928 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
1929 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1930 let id = running.id.clone();
1931 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
1932
1933 let response_talks = talks.clone();
1934 let response_cfg = cfg.clone();
1935 let reply = tokio::spawn(async move {
1936 respond(&mut running, &response_talks, &response_cfg, &first).await
1937 });
1938 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
1939
1940 let mut queued = talks.get(&id).expect("queued handle");
1941 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
1942 reply.await.expect("join").expect("reply");
1943
1944 let saved = talks.get(&id).expect("reload");
1945 assert_eq!(saved.pending, "next");
1946 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
1947 }
1948
1949 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
1952 let calls = dir.join(format!("{id}.calls"));
1953 let script = format!(
1954 "#!/bin/sh\necho x >> '{}'\n{body}\n",
1955 calls.to_string_lossy()
1956 );
1957 let path = dir.join(format!("mock-{id}.sh"));
1958 std::fs::write(&path, script).expect("write mock");
1959 AgentSpec {
1960 id: id.to_owned(),
1961 kind: AgentKind::Command,
1962 model: None,
1963 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
1964 extra_args: Vec::new(),
1965 env: BTreeMap::new(),
1966 prompt_delivery: None,
1967 }
1968 }
1969
1970 fn calls(dir: &Path, id: &str) -> usize {
1971 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
1972 }
1973
1974 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
1975 let mut cfg = config(specs[0].clone());
1976 cfg.agents = specs;
1977 cfg.roles.chatter = Some(AgentChoice::Chain(
1978 ids.iter().map(|s| (*s).to_owned()).collect(),
1979 ));
1980 cfg
1981 }
1982
1983 #[tokio::test]
1984 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
1985 let (tmp, talks) = store();
1986 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
1987 let b = counting_agent(tmp.path(), "b", "cat");
1988 let cfg = chain_config(vec![a, b], &["a", "b"]);
1989 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1990 assert_eq!(talk.agent, "a");
1991
1992 say(&mut talk, &talks, &cfg, "hello there", Vec::new())
1993 .await
1994 .expect("turn");
1995 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
1996 assert_eq!(calls(tmp.path(), "b"), 1);
1997 assert_eq!(talk.agent, "b", "the switch persists");
1998 assert!(talks.get(&talk.id).unwrap().agent == "b");
1999 let reply = talk.turns.last().unwrap();
2000 assert!(reply.body.contains("hello there"));
2001 assert!(
2002 reply.body.contains("magi task add --solo"),
2003 "a fresh seat gets the full briefing"
2004 );
2005 assert!(
2006 talk.turns
2007 .iter()
2008 .any(|t| t.body.contains("agent changed from a to b")),
2009 "the switch is noted"
2010 );
2011 }
2012
2013 #[tokio::test]
2014 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
2015 let (tmp, talks) = store();
2016 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2017 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
2018 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
2019 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2020
2021 let err = say(&mut talk, &talks, &cfg, "hi", Vec::new())
2022 .await
2023 .expect_err("every agent failed");
2024 assert!(err.to_string().contains("`a`"), "{err:#}");
2025 assert_eq!(calls(tmp.path(), "a"), 1);
2026 assert_eq!(calls(tmp.path(), "b"), 1);
2027 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
2028 }
2029
2030 #[test]
2031 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
2032 let (tmp, talks) = store();
2033 let b = counting_agent(tmp.path(), "b", "cat");
2034 let cfg = chain_config(vec![b], &["ghost", "b"]);
2035 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2036 assert_eq!(talk.agent, "b");
2037 }
2038
2039 #[tokio::test]
2040 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
2041 let (tmp, talks) = store();
2042 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2043 let b = counting_agent(tmp.path(), "b", "cat");
2044 let cfg = chain_config(vec![a, b], &["a", "b"]);
2045 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
2046 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2047 .await
2048 .expect_err("a alone, and it fails");
2049 assert_eq!(calls(tmp.path(), "b"), 0);
2050 assert_eq!(talk.agent, "a");
2051 }
2052
2053 #[tokio::test]
2054 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
2055 let (tmp, talks) = store();
2056 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2057 let b = counting_agent(tmp.path(), "b", "cat");
2058 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
2059 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
2060 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
2061 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2062 .await
2063 .expect_err("c alone, and it fails");
2064 assert_eq!(calls(tmp.path(), "b"), 0);
2065 }
2066
2067 #[tokio::test]
2068 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
2069 let (tmp, talks) = store();
2070 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2071 let cfg = config(spec);
2072 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2073
2074 say(
2075 &mut talk,
2076 &talks,
2077 &cfg,
2078 "what does the queue module do?",
2079 Vec::new(),
2080 )
2081 .await
2082 .expect("first turn");
2083 let first_prompt = &talk.turns[1].body;
2084 assert!(first_prompt.contains("magi task add --solo"));
2085 assert!(first_prompt.contains("what does the queue module do?"));
2086
2087 say(&mut talk, &talks, &cfg, "and how is it locked?", Vec::new())
2088 .await
2089 .expect("second turn");
2090 let second_prompt = &talk.turns[3].body;
2091 assert!(
2092 !second_prompt.contains("magi task add --solo"),
2093 "the briefing is sent once, not on every turn: {second_prompt}"
2094 );
2095 assert!(second_prompt.contains("and how is it locked?"));
2096 }
2097
2098 #[tokio::test]
2099 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
2100 let (tmp, talks) = store();
2101 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2102 let mut b = a.clone();
2103 b.id = "other".to_owned();
2104 let mut cfg = config(a.clone());
2105 cfg.agents.push(b.clone());
2106 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
2107 say(&mut talk, &talks, &cfg, "remember the walrus", Vec::new())
2108 .await
2109 .expect("first turn");
2110 let old_session = talk.seat.claude_session.clone();
2111 assert_eq!(talk.seat.turns, 1);
2112
2113 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
2114 assert_eq!(talk.agent, "other");
2115 assert_eq!(talk.seat.turns, 0);
2116 assert_eq!(talk.seat.agent, "other");
2117 assert_ne!(talk.seat.claude_session, old_session);
2118 let note = talk.turns.last().expect("note");
2119 assert_eq!(note.who, Who::Agent);
2120 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2121 assert!(note.body.contains("changed from"), "{}", note.body);
2122 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
2123
2124 let before = talk.turns.len();
2125 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
2126 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
2127
2128 say(&mut talk, &talks, &cfg, "what did I say?", Vec::new())
2129 .await
2130 .expect("turn after switch");
2131 let prompt = &talk.turns.last().expect("reply").body;
2132 assert!(prompt.contains("remember the walrus"), "{prompt}");
2133 assert!(prompt.contains("## magi"), "{prompt}");
2134 assert!(prompt.contains("what did I say?"), "{prompt}");
2135 }
2136
2137 #[tokio::test]
2138 async fn say_appends_the_operator_turn_then_the_agent_turn() {
2139 let (tmp, talks) = store();
2140 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2141 let cfg = config(spec);
2142 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2143
2144 say(
2145 &mut talk,
2146 &talks,
2147 &cfg,
2148 "can I rename this function?",
2149 Vec::new(),
2150 )
2151 .await
2152 .expect("say");
2153
2154 assert_eq!(talk.turns.len(), 2);
2155 assert_eq!(talk.turns[0].who, Who::Operator);
2156 assert_eq!(talk.turns[0].body, "can I rename this function?");
2157 assert_eq!(talk.turns[1].who, Who::Agent);
2158 assert_eq!(talk.turns[1].body, "go ahead");
2159 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
2160 }
2161
2162 #[tokio::test]
2163 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
2164 let (tmp, talks) = store();
2165 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2166 let cfg = config(spec);
2167 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2168
2169 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2170 .await
2171 .expect_err("a turn with no answer is an error");
2172 assert!(err.to_string().contains("no answer"), "{err}");
2173
2174 let on_disk = talks.get(&talk.id).expect("get");
2175 assert_eq!(on_disk.turns.len(), 2);
2176 assert_eq!(on_disk.turns[0].body, "check the tests");
2177 let note = &on_disk.turns[1];
2178 assert_eq!(note.who, Who::Agent);
2179 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2180 assert!(note.body.contains("your message is saved"));
2181 }
2182
2183 #[tokio::test]
2189 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
2190 let (tmp, talks) = store();
2191 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2192 let cfg = config(spec);
2193 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2194
2195 let text =
2196 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
2197 failpoint::force_put_failures(PUT_RETRIES - 1);
2200 respond(&mut talk, &talks, &cfg, &text)
2201 .await
2202 .expect("respond must survive a write failure its own retries can outlast");
2203
2204 assert_eq!(talk.turns.len(), 2);
2205 assert_eq!(talk.turns[1].who, Who::Agent);
2206 assert_eq!(talk.turns[1].body, "go ahead");
2207 let on_disk = talks.get(&talk.id).expect("get");
2208 assert_eq!(
2209 on_disk.turns, talk.turns,
2210 "the reply must reach disk despite the early write failures"
2211 );
2212 }
2213
2214 #[tokio::test]
2220 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
2221 let (tmp, talks) = store();
2222 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2223 let cfg = config(spec);
2224 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2225
2226 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2227 failpoint::force_put_failures(PUT_RETRIES);
2232 let err = respond(&mut talk, &talks, &cfg, &text)
2233 .await
2234 .expect_err("a reply that cannot be saved must be reported, not swallowed");
2235 assert!(err.to_string().contains("could not be saved"), "{err}");
2236
2237 let on_disk = talks.get(&talk.id).expect("get");
2238 assert_eq!(
2239 on_disk.turns.len(),
2240 2,
2241 "the operator turn plus a visible note"
2242 );
2243 assert_eq!(on_disk.turns[0].body, "check the tests");
2244 let note = &on_disk.turns[1];
2245 assert_eq!(note.who, Who::Agent);
2246 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2247 assert!(
2248 note.body.contains("could not be saved"),
2249 "the operator must be told the reply is missing, not left staring \
2250 at a gap with no explanation: {}",
2251 note.body
2252 );
2253 assert_eq!(
2254 talk.turns, on_disk.turns,
2255 "the in-memory talk must match what actually landed on disk"
2256 );
2257
2258 let artifacts = talks.artifacts_of(&talk.id);
2261 let stash = std::fs::read_dir(&artifacts)
2262 .expect("artifacts dir")
2263 .filter_map(|e| e.ok())
2264 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
2265 .expect("a stash file for the lost reply");
2266 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
2267 assert_eq!(stashed, "go ahead");
2268
2269 assert_eq!(
2277 on_disk.seat.turns, 1,
2278 "the note's write must carry the turn the CLI actually took"
2279 );
2280 assert_eq!(
2281 on_disk.seat.claude_session, talk.seat.claude_session,
2282 "the session id handed to the CLI must survive the failed reply"
2283 );
2284 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
2285 assert!(
2286 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
2287 "the next turn must resume, not open the same session id twice"
2288 );
2289 }
2290
2291 #[tokio::test]
2296 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
2297 let (tmp, talks) = store();
2298 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2299 let cfg = config(spec);
2300 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2301
2302 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2303 failpoint::force_put_failures(PUT_RETRIES * 2);
2306 let err = respond(&mut talk, &talks, &cfg, &text)
2307 .await
2308 .expect_err("neither the reply nor the note could be saved");
2309 assert!(err.to_string().contains("could not be saved"), "{err}");
2310
2311 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
2312 let on_disk = talks.get(&talk.id).expect("get");
2313 assert_eq!(on_disk.turns.len(), 1);
2314
2315 assert_eq!(
2325 on_disk.seat.turns, 0,
2326 "an unwritable file cannot record the turn the CLI took"
2327 );
2328 assert_eq!(
2329 talk.seat.turns, 1,
2330 "the in-memory seat still reports the turn the CLI actually took"
2331 );
2332 assert_eq!(
2333 on_disk.seat.claude_session, talk.seat.claude_session,
2334 "the session id was minted at `begin` and never changes here"
2335 );
2336 }
2337
2338 #[tokio::test]
2342 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
2343 let (tmp, talks) = store();
2344 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2345 let cfg = config(spec);
2346 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2347
2348 let att = talks
2349 .put_attachment(
2350 &talk.id,
2351 "image/png",
2352 "screenshot.png",
2353 b"pretend-png-bytes",
2354 )
2355 .expect("put attachment");
2356
2357 say(&mut talk, &talks, &cfg, "", vec![att.clone()])
2358 .await
2359 .expect("an empty body with an attachment is still a turn");
2360
2361 let operator_turn = &talk.turns[0];
2362 assert_eq!(operator_turn.who, Who::Operator);
2363 assert_eq!(operator_turn.body, "");
2364 assert_eq!(operator_turn.attachments, vec![att.clone()]);
2365
2366 let prompt = &talk.turns[1].body;
2367 let expected_path = talks
2368 .attachments_dir(&talk.id)
2369 .join(format!("{}.png", att.id));
2370 assert!(
2371 prompt.contains(&expected_path.display().to_string()),
2372 "the agent must be told the attachment's absolute path: {prompt}"
2373 );
2374 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
2375 }
2376
2377 #[test]
2386 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
2387 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
2388 let att = Attachment {
2389 id: "0".repeat(32),
2390 name: "shot.png".to_owned(),
2391 mime: "image/png".to_owned(),
2392 bytes: 3,
2393 };
2394 let path = talks
2395 .attachment_path("some-talk-id", &att)
2396 .expect("a supported mime always yields a path");
2397 assert!(
2398 path.is_absolute(),
2399 "must be absolute even off a relative store root: {}",
2400 path.display()
2401 );
2402 }
2403
2404 #[tokio::test]
2405 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
2406 let (tmp, talks) = store();
2411 let slow = mock_agent(
2412 tmp.path(),
2413 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
2414 BTreeMap::new(),
2415 );
2416 let mut cfg = config(slow);
2417 cfg.graph.timeout_talk = 1;
2418 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2419
2420 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2421 .await
2422 .expect_err("a turn that never answers is an error");
2423 assert!(
2424 err.to_string().contains("did not answer within 1s"),
2425 "{err}"
2426 );
2427
2428 let on_disk = talks.get(&talk.id).expect("get");
2429 let note = on_disk.turns.last().expect("a note turn was recorded");
2430 assert!(
2431 note.body.contains("did not answer within 1s"),
2432 "the transcript must show the configured timeout: {}",
2433 note.body
2434 );
2435 }
2436
2437 #[test]
2438 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
2439 let (tmp, talks) = store();
2440 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2441 let cfg = config(spec);
2442 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2443
2444 close(&mut talk, &talks).expect("close");
2445 assert_eq!(talk.status, TalkStatus::Closed);
2446 close(&mut talk, &talks).expect("closing twice is not an error");
2447
2448 let err =
2449 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
2450 assert!(err.to_string().contains("closed"));
2451 let _ = &cfg; }
2453
2454 #[tokio::test]
2455 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2456 let (tmp, talks) = store();
2457 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2458 let cfg = config(spec);
2459 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2462
2463 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
2467 close(&mut closed_elsewhere, &talks).expect("close");
2468 assert_eq!(
2469 talks.get(&in_flight.id).expect("reread").status,
2470 TalkStatus::Closed,
2471 "the close landed on disk before the turn finished"
2472 );
2473
2474 assert_eq!(in_flight.status, TalkStatus::Open);
2478 respond(&mut in_flight, &talks, &cfg, "one more question")
2479 .await
2480 .expect("the turn itself still completes");
2481
2482 let on_disk = talks.get(&in_flight.id).expect("reread");
2483 assert_eq!(
2484 on_disk.status,
2485 TalkStatus::Closed,
2486 "a close must stick even when a turn that started before it finishes after it"
2487 );
2488 assert!(
2491 on_disk.turns.iter().any(|t| t.body == "here you go"),
2492 "the in-flight turn's own reply is still recorded: {:?}",
2493 on_disk.turns
2494 );
2495 }
2496
2497 #[test]
2498 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
2499 let (tmp, talks) = store();
2500 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2501 let cfg = config(spec);
2502 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2505
2506 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
2509 close(&mut closed_elsewhere, &talks).expect("close");
2510 assert_eq!(
2511 talks.get(&stale.id).expect("reread").status,
2512 TalkStatus::Closed,
2513 "the close landed on disk before record was called"
2514 );
2515
2516 assert_eq!(stale.status, TalkStatus::Open);
2520 let err = record(&mut stale, &talks, "still there?", Vec::new())
2521 .expect_err("a close that landed first must be honored, not overwritten");
2522 assert!(err.to_string().contains("closed"));
2523
2524 let on_disk = talks.get(&stale.id).expect("reread");
2525 assert_eq!(
2526 on_disk.status,
2527 TalkStatus::Closed,
2528 "record must not resurrect a conversation closed while its snapshot was stale"
2529 );
2530 assert!(
2531 on_disk.turns.is_empty(),
2532 "the rejected turn must not have been appended: {:?}",
2533 on_disk.turns
2534 );
2535 let _ = &cfg; }
2537
2538 #[test]
2539 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
2540 let (tmp, talks) = store();
2541 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2542 let cfg = config(spec);
2543 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2544
2545 let held = talks.guard();
2549
2550 let talks2 = talks.clone();
2551 let id = talk.id.clone();
2552 let closing = std::thread::spawn(move || {
2553 let mut talk = talks2.get(&id).expect("get");
2554 close(&mut talk, &talks2).expect("close");
2555 });
2556
2557 std::thread::sleep(Duration::from_millis(50));
2558 assert!(
2559 !closing.is_finished(),
2560 "close must wait for the guard, not read and write while it is held - \
2561 a re-read alone narrows this window without closing it"
2562 );
2563
2564 drop(held);
2565 closing.join().expect("close thread panicked");
2566
2567 assert_eq!(
2568 talks.get(&talk.id).expect("reread").status,
2569 TalkStatus::Closed,
2570 "once the guard is free, close still lands"
2571 );
2572 let _ = &cfg; }
2574
2575 #[test]
2576 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
2577 let (tmp, talks) = store();
2578 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2579 let cfg = config(spec);
2580 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2581
2582 close(&mut talk, &talks).expect("close");
2583 assert_eq!(talk.status, TalkStatus::Closed);
2584
2585 reopen(&mut talk, &talks).expect("reopen");
2586 assert_eq!(talk.status, TalkStatus::Open);
2587 assert_eq!(
2588 talks.get(&talk.id).expect("reread").status,
2589 TalkStatus::Open
2590 );
2591
2592 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
2594 assert_eq!(talk.status, TalkStatus::Open);
2595
2596 record(&mut talk, &talks, "one more thing", Vec::new())
2597 .expect("a reopened talk takes turns again");
2598 let _ = &cfg; }
2600
2601 #[test]
2602 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
2603 let (tmp, talks) = store();
2604 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2605 let cfg = config(spec);
2606 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2607
2608 let artifacts = talks.artifacts_of(&talk.id);
2609 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
2610 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
2611
2612 talks.remove(&talk.id).expect("remove");
2613 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
2614 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
2615 assert!(
2616 talks.get(&talk.id).is_err(),
2617 "a removed talk cannot be read back"
2618 );
2619
2620 let err = talks
2621 .remove("nonexistent-id")
2622 .expect_err("unknown id refused");
2623 assert!(err.to_string().contains("no talk matches"), "{err}");
2624 let _ = &cfg; }
2626
2627 #[tokio::test]
2628 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2629 let (tmp, talks) = store();
2630 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2631 let cfg = config(spec);
2632 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2635
2636 talks.remove(&in_flight.id).expect("remove");
2637 assert!(
2638 talks.get(&in_flight.id).is_err(),
2639 "the delete landed on disk before the turn finished"
2640 );
2641
2642 respond(&mut in_flight, &talks, &cfg, "one more question")
2645 .await
2646 .expect("the turn itself still completes rather than erroring");
2647
2648 assert!(
2649 talks.get(&in_flight.id).is_err(),
2650 "a delete must stick even when a turn that started before it finishes after it"
2651 );
2652 }
2653
2654 #[test]
2655 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
2656 let (tmp, talks) = store();
2657 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2658 let cfg = config(spec);
2659 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2662
2663 talks.remove(&stale.id).expect("remove");
2664
2665 let err = record(&mut stale, &talks, "still there?", Vec::new())
2669 .expect_err("a delete that landed first must be honored, not overwritten");
2670 assert!(err.to_string().contains("deleted"), "{err}");
2671
2672 assert!(
2673 talks.get(&stale.id).is_err(),
2674 "record must not resurrect a conversation deleted while its snapshot was stale"
2675 );
2676 let _ = &cfg; }
2678
2679 #[test]
2680 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
2681 let (tmp, talks) = store();
2682 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2683 let cfg = config(spec);
2684 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2687
2688 talks.remove(&stale.id).expect("remove");
2689
2690 let err = close(&mut stale, &talks)
2694 .expect_err("a delete that landed first must be honored, not overwritten");
2695 assert!(err.to_string().contains("deleted"), "{err}");
2696
2697 assert!(
2698 talks.get(&stale.id).is_err(),
2699 "close must not resurrect a conversation deleted while its snapshot was stale"
2700 );
2701 let _ = &cfg; }
2703
2704 #[test]
2705 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
2706 let (tmp, talks) = store();
2707 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2708 let cfg = config(spec);
2709 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2712 close(&mut stale, &talks).expect("close");
2713
2714 talks.remove(&stale.id).expect("remove");
2715
2716 let err = reopen(&mut stale, &talks)
2720 .expect_err("a delete that landed first must be honored, not overwritten");
2721 assert!(err.to_string().contains("deleted"), "{err}");
2722
2723 assert!(
2724 talks.get(&stale.id).is_err(),
2725 "reopen must not resurrect a conversation deleted while its snapshot was stale"
2726 );
2727 let _ = &cfg; }
2729
2730 #[test]
2731 fn list_puts_open_talks_before_closed_ones() {
2732 let (tmp, talks) = store();
2733 let make = |id: &str, status: TalkStatus| {
2734 let mut t = Talk {
2735 schema: SCHEMA,
2736 id: id.to_owned(),
2737 repo: tmp.path().to_owned(),
2738 agent: "mock".to_owned(),
2739 status,
2740 turns: Vec::new(),
2741 pending: String::new(),
2742 pending_attachments: Vec::new(),
2743 fallback: false,
2744 created_at: Timestamp::now(),
2745 updated_at: Timestamp::now(),
2746 seat: SeatState::new(SEAT, "mock", 7),
2747 };
2748 talks.put(&mut t).expect("put");
2749 };
2750 make("20260901-000000-0001", TalkStatus::Open);
2751 make("20260902-000000-0002", TalkStatus::Open);
2752 make("20260903-000000-0003", TalkStatus::Closed);
2753
2754 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
2755 assert_eq!(
2756 ids,
2757 [
2758 "20260902-000000-0002",
2759 "20260901-000000-0001",
2760 "20260903-000000-0003"
2761 ]
2762 );
2763 assert_eq!(talks.count_open(), 2);
2764 }
2765
2766 #[test]
2767 fn tasks_of_finds_only_this_talks_own_tasks() {
2768 let dir = tempfile::tempdir().expect("tempdir");
2769 let queue = Queue::at(dir.path().join("queue"));
2770
2771 let mut mine = Task::new(
2772 "rework the loader".to_owned(),
2773 "rework the loader".to_owned(),
2774 PathBuf::from("/repo"),
2775 Source::Agent {
2776 run: "20260904-014455-ab12".to_owned(),
2777 node: "chat".to_owned(),
2778 },
2779 );
2780 queue.put(&mut mine).expect("put mine");
2781
2782 let mut theirs = Task::new(
2783 "unrelated".to_owned(),
2784 "unrelated".to_owned(),
2785 PathBuf::from("/repo"),
2786 Source::Agent {
2787 run: "20260904-090000-zz99".to_owned(),
2788 node: "implement".to_owned(),
2789 },
2790 );
2791 queue.put(&mut theirs).expect("put theirs");
2792
2793 let mut human = Task::new(
2794 "typed by hand".to_owned(),
2795 "typed by hand".to_owned(),
2796 PathBuf::from("/repo"),
2797 Source::Human,
2798 );
2799 queue.put(&mut human).expect("put human");
2800
2801 let found = tasks_of(&queue, "20260904-014455-ab12");
2802 assert_eq!(found.len(), 1);
2803 assert_eq!(found[0].id, mine.id);
2804 }
2805
2806 #[test]
2807 fn the_briefing_names_solo_task_add() {
2808 let brief = briefing(Path::new("/repo"), "en", false);
2809 assert!(brief.contains("magi task add --solo"));
2810 assert!(brief.contains("/repo"));
2811 assert!(!brief.contains("Hold this conversation in"));
2812 }
2813
2814 #[test]
2820 fn the_briefing_explains_targeting_a_different_repository_by_name() {
2821 let brief = briefing(Path::new("/repo"), "en", false);
2822 assert!(brief.contains("--repo does not have to be a full path"));
2823 assert!(brief.contains("owner/repo"));
2824 assert!(brief.contains("magi repos"));
2825 assert!(brief.contains("ask the operator"));
2826 }
2827
2828 #[test]
2829 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
2830 let brief = briefing(Path::new("/repo"), "en", false);
2831 assert!(brief.contains("--attach <path>"), "{brief}");
2832 assert!(brief.contains("deleting this conversation"), "{brief}");
2833 }
2834
2835 #[test]
2836 fn the_briefing_names_the_language_when_it_is_not_english() {
2837 let brief = briefing(Path::new("/repo"), "Japanese", false);
2838 assert!(brief.contains("Hold this conversation in Japanese"));
2839 }
2840
2841 #[test]
2842 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
2843 let read_only = briefing(Path::new("/repo"), "en", false);
2844 assert!(read_only.contains("Do not write files"));
2845 assert!(!read_only.contains("allow_write"));
2846
2847 let writable = briefing(Path::new("/repo"), "en", true);
2848 assert!(!writable.contains("Do not write files"));
2849 assert!(writable.contains("allow_write = true"));
2850 assert!(writable.contains("magi task add --solo"));
2853 assert!(writable.contains("say plainly what you"));
2854 }
2855}