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 The current state of the code is whatever origin/main holds, not \
1341 whatever a working tree shows: a primary checkout often lags \
1342 upstream, sits on a detached HEAD and carries uncommitted changes. \
1343 Before answering about code, run `git fetch origin` in that \
1344 repository if it is cheap, then read through \
1345 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
1346 If the working tree differs, say so; if the fetch fails, say that \
1347 too, so the operator knows the answer may be stale.\n\n\
1348 If the operator attached an image (a screenshot, say) that the task \
1349 is about, pass it with `--attach <path>`, using the absolute path \
1350 the turn's attachment note gives; repeat the flag for several. \
1351 `magi task add --solo --attach <path> <instruction>` copies the \
1352 file into the task, so the implementer receives it. Do not paste the \
1353 path into <instruction> instead: deleting this conversation deletes \
1354 its attachments, and then that path reaches no one.\n",
1355 repo = repo.display(),
1356 );
1357 out.push_str(&language_note(language));
1358 out
1359}
1360
1361fn language_note(language: &str) -> String {
1364 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
1365 String::new()
1366 } else {
1367 format!("\nHold this conversation in {language}.\n")
1368 }
1369}
1370
1371pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
1378 let mut tasks: Vec<Task> = queue
1379 .list()
1380 .into_iter()
1381 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
1382 .collect();
1383 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
1384 tasks
1385}
1386
1387fn read_path(path: &Path) -> Result<Talk> {
1388 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1389 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
1390}
1391
1392const PUT_RETRIES: u32 = 5;
1395
1396fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1408 let mut last_err = None;
1409 for attempt in 0..PUT_RETRIES {
1410 if attempt > 0 {
1411 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
1412 }
1413 match try_write_atomic(tmp, path, body) {
1414 Ok(()) => return Ok(()),
1415 Err(e) => last_err = Some(e),
1416 }
1417 }
1418 Err(last_err.expect("the loop above always runs at least once"))
1419}
1420
1421fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1422 #[cfg(test)]
1423 if failpoint::take_forced_put_failure() {
1424 bail!("simulated write failure (test)");
1425 }
1426 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
1427 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
1428 Ok(())
1429}
1430
1431fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
1436 let dir = store.artifacts_of(id);
1437 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1438 let path = dir.join(format!("{stem}-lost.txt"));
1439 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
1440 Ok(path)
1441}
1442
1443#[cfg(test)]
1450mod failpoint {
1451 use std::cell::Cell;
1452
1453 thread_local! {
1454 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
1455 }
1456
1457 pub(super) fn force_put_failures(count: u32) {
1460 FORCE_PUT_FAILURES.with(|c| c.set(count));
1461 }
1462
1463 pub(super) fn take_forced_put_failure() -> bool {
1466 FORCE_PUT_FAILURES.with(|c| {
1467 let n = c.get();
1468 if n == 0 {
1469 false
1470 } else {
1471 c.set(n - 1);
1472 true
1473 }
1474 })
1475 }
1476}
1477
1478fn short(id: &str) -> &str {
1479 id.split('-').next_back().unwrap_or(id)
1480}
1481
1482fn new_id() -> String {
1483 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1484 let seed = crate::rng::entropy();
1485 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1486}
1487
1488fn attachment_ext(mime: &str) -> Option<&'static str> {
1493 match mime {
1494 "image/png" => Some("png"),
1495 "image/jpeg" => Some("jpg"),
1496 "image/gif" => Some("gif"),
1497 "image/webp" => Some("webp"),
1498 _ => None,
1499 }
1500}
1501
1502pub fn valid_attachment_id(id: &str) -> bool {
1507 id.len() == 32
1508 && id
1509 .bytes()
1510 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
1511}
1512
1513fn new_attachment_id() -> String {
1517 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
1518 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
1519}
1520
1521#[cfg(test)]
1522mod tests {
1523 #[test]
1524 fn the_briefing_points_at_origin_main_not_the_working_tree() {
1525 let b = briefing(Path::new("/r"), "en", false);
1526 assert!(b.contains("origin/main"));
1527 assert!(b.contains("git show origin/main:"));
1528 }
1529 use std::collections::BTreeMap;
1530
1531 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
1532 use crate::queue::{Queue, Source, Task};
1533
1534 use super::*;
1535
1536 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
1537 AgentSpec {
1538 id: id.to_owned(),
1539 kind: AgentKind::Command,
1540 model: model.map(str::to_owned),
1541 command: Vec::new(),
1542 extra_args: Vec::new(),
1543 env: BTreeMap::new(),
1544 prompt_delivery: None,
1545 }
1546 }
1547
1548 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
1549 Talk {
1550 schema: SCHEMA,
1551 id: "20260904-014455-ab12".to_owned(),
1552 repo: PathBuf::from("."),
1553 agent: agent.to_owned(),
1554 status: TalkStatus::Open,
1555 turns,
1556 pending: String::new(),
1557 pending_attachments: Vec::new(),
1558 fallback: false,
1559 created_at: Timestamp::now(),
1560 updated_at: Timestamp::now(),
1561 seat: SeatState::new(SEAT, agent, 1),
1562 }
1563 }
1564
1565 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
1566 Turn {
1567 who: Who::Agent,
1568 body: body.to_owned(),
1569 at: Timestamp::now(),
1570 attachments: Vec::new(),
1571 usage: usage.map(|(t, a, m)| TurnUsage {
1572 context_tokens: t,
1573 agent: a.to_owned(),
1574 model: m.map(str::to_owned),
1575 }),
1576 }
1577 }
1578
1579 fn ctx_config(windows: &[(&str, u64)]) -> Config {
1580 Config {
1581 agents: vec![
1582 ctx_agent("small", Some("small-model")),
1583 ctx_agent("big", Some("big-model")),
1584 ctx_agent("plain", None),
1585 ],
1586 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
1587 ..Config::default()
1588 }
1589 }
1590
1591 #[test]
1592 fn context_usage_computes_percent_and_warns_at_eighty() {
1593 let cfg = ctx_config(&[("small-model", 1000)]);
1594 let at = |tokens| {
1595 let t = ctx_talk(
1596 "small",
1597 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
1598 );
1599 context_usage(&t, Some(&cfg))
1600 };
1601 let u = at(799);
1602 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
1603 let u = at(800);
1604 assert_eq!((u.percent, u.warn), (Some(80), true));
1605 let u = at(1500);
1606 assert_eq!((u.percent, u.warn), (Some(150), true));
1607 assert!(!u.since_switch);
1608 }
1609
1610 #[test]
1611 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
1612 let cfg = ctx_config(&[("small-model", 1000)]);
1613 let t = ctx_talk(
1614 "small",
1615 vec![
1616 reply("old", Some((900, "small", Some("small-model")))),
1617 reply("new", None),
1618 ],
1619 );
1620 let u = context_usage(&t, Some(&cfg));
1621 assert_eq!((u.tokens, u.percent, u.warn), (None, None, false));
1622 let t = ctx_talk(
1624 "small",
1625 vec![
1626 reply("old", Some((900, "small", Some("small-model")))),
1627 reply("magi: could not run agent", None),
1628 ],
1629 );
1630 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
1631 assert_eq!(
1632 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
1633 None
1634 );
1635 }
1636
1637 #[test]
1638 fn context_usage_without_a_window_shows_tokens_only() {
1639 let cfg = ctx_config(&[]);
1640 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
1642 let u = context_usage(&t, Some(&cfg));
1643 assert_eq!(
1644 (u.tokens, u.window, u.percent, u.warn),
1645 (Some(5000), None, None, false)
1646 );
1647 let t = ctx_talk(
1648 "small",
1649 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
1650 );
1651 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
1652 assert_eq!(context_usage(&t, None).window, None);
1654 }
1655
1656 #[test]
1657 fn context_usage_switching_model_changes_the_denominator() {
1658 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
1659 let used = reply("hi", Some((900, "small", Some("small-model"))));
1660 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
1661 assert_eq!(
1662 (before.percent, before.warn, before.since_switch),
1663 (Some(90), true, false)
1664 );
1665 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
1668 assert_eq!(after.window, Some(10_000));
1669 assert_eq!(
1670 (after.percent, after.warn, after.since_switch),
1671 (Some(9), false, true)
1672 );
1673 assert_eq!(after.model.as_deref(), Some("big-model"));
1674 }
1675
1676 #[test]
1677 fn a_turn_recorded_before_usage_existed_still_reads() {
1678 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
1679 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
1680 assert!(turn.usage.is_none());
1681 let json = serde_json::to_string(&turn).expect("serialize");
1682 assert!(
1683 !json.contains("usage"),
1684 "absent usage is not written: {json}"
1685 );
1686 }
1687
1688 fn store() -> (tempfile::TempDir, Talks) {
1690 let tmp = tempfile::tempdir().expect("tempdir");
1691 let talks = Talks::at(tmp.path().join("talks"));
1692 (tmp, talks)
1693 }
1694
1695 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
1699 let path = dir.join("mock-talk-agent.sh");
1700 std::fs::write(&path, script).expect("write mock");
1701 AgentSpec {
1702 id: "mock".to_owned(),
1703 kind: AgentKind::Command,
1704 model: None,
1705 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
1706 extra_args: Vec::new(),
1707 env,
1708 prompt_delivery: None,
1709 }
1710 }
1711
1712 fn config(spec: AgentSpec) -> Config {
1713 Config {
1714 agents: vec![spec],
1715 graph: Graph {
1716 language: "en".to_owned(),
1717 ..Graph::default()
1718 },
1719 ..Config::default()
1720 }
1721 }
1722
1723 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
1725
1726 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
1728
1729 const ECHO: &str = "#!/bin/sh\ncat\n";
1732
1733 fn env(reply: &str) -> BTreeMap<String, String> {
1734 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
1735 }
1736
1737 #[test]
1738 fn the_frozen_json_field_names_round_trip_through_disk() {
1739 let (tmp, talks) = store();
1740 let mut talk = Talk {
1741 schema: SCHEMA,
1742 id: "20260904-014455-ab12".to_owned(),
1743 repo: tmp.path().to_owned(),
1744 agent: "sonnet".to_owned(),
1745 status: TalkStatus::Open,
1746 turns: Vec::new(),
1747 pending: String::new(),
1748 pending_attachments: Vec::new(),
1749 fallback: false,
1750 created_at: Timestamp::now(),
1751 updated_at: Timestamp::now(),
1752 seat: SeatState::new(SEAT, "sonnet", 7),
1753 };
1754 talks.put(&mut talk).expect("put");
1755
1756 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
1757 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
1758 for field in [
1759 "schema",
1760 "id",
1761 "repo",
1762 "agent",
1763 "status",
1764 "turns",
1765 "created_at",
1766 "updated_at",
1767 ] {
1768 assert!(v.get(field).is_some(), "missing field `{field}`");
1769 }
1770 assert_eq!(v["schema"], 1);
1771 assert_eq!(v["status"], "open");
1772
1773 let back = talks.get(&talk.id).expect("get");
1774 assert_eq!(back.id, talk.id);
1775 assert_eq!(back.status, TalkStatus::Open);
1776 }
1777
1778 #[test]
1779 fn opening_a_talk_takes_no_agent_turn() {
1780 let (tmp, talks) = store();
1781 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1785 let cfg = config(spec);
1786
1787 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1788 assert_eq!(talk.status, TalkStatus::Open);
1789 assert!(talk.turns.is_empty(), "nothing has been said yet");
1790
1791 let on_disk = talks.get(&talk.id).expect("get");
1792 assert_eq!(on_disk.turns.len(), 0);
1793 }
1794
1795 #[test]
1803 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
1804 let (tmp, talks) = store();
1805 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1806 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
1807 chatter_spec.id = "chatter-mock".to_owned();
1808
1809 let mut cfg = Config {
1810 agents: vec![first_spec.clone(), chatter_spec.clone()],
1811 graph: Graph {
1812 language: "en".to_owned(),
1813 ..Graph::default()
1814 },
1815 ..Config::default()
1816 };
1817 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
1818
1819 let talk =
1820 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
1821 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
1822
1823 cfg.roles.chatter = None;
1824 let fallback =
1825 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
1826 assert_eq!(
1827 fallback.agent, first_spec.id,
1828 "unset chatter must fall back to agent::pick's own default order"
1829 );
1830 }
1831
1832 #[test]
1835 fn a_talk_recorded_without_attachments_still_reads() {
1836 let (tmp, talks) = store();
1837 let path = talks.path_of("20260904-014455-ab12");
1838 std::fs::create_dir_all(talks.root()).expect("talks dir");
1839 std::fs::write(
1840 &path,
1841 serde_json::json!({
1842 "schema": 1,
1843 "id": "20260904-014455-ab12",
1844 "repo": tmp.path(),
1845 "agent": "sonnet",
1846 "status": "open",
1847 "turns": [
1848 { "who": "operator", "body": "still there?",
1849 "at": Timestamp::now().to_string() },
1850 ],
1851 "created_at": Timestamp::now().to_string(),
1852 "updated_at": Timestamp::now().to_string(),
1853 "seat": SeatState::new(SEAT, "sonnet", 7),
1854 })
1855 .to_string(),
1856 )
1857 .expect("write pre-attachments talk");
1858
1859 let talk = talks.get("20260904-014455-ab12").expect("must still read");
1860 assert!(talk.turns[0].attachments.is_empty());
1861 }
1862
1863 #[test]
1864 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
1865 let (tmp, talks) = store();
1866 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1867 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1868
1869 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
1870 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
1871 let saved = talks.get(&talk.id).expect("reload queued talk");
1872 assert_eq!(saved.pending, "first\n\nsecond");
1873 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
1874
1875 let drained = drain(&mut talk, &talks).expect("drain");
1876 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
1877 let saved = talks.get(&talk.id).expect("reload drained talk");
1878 assert!(saved.pending.is_empty());
1879 assert_eq!(saved.turns.len(), 1);
1880 assert_eq!(saved.turns[0].body, "first\n\nsecond");
1881 }
1882
1883 #[test]
1884 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
1885 let (tmp, talks) = store();
1886 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
1887 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1888 let attachment = Attachment {
1889 id: "a".repeat(32),
1890 name: "shot.png".to_owned(),
1891 mime: "image/png".to_owned(),
1892 bytes: 3,
1893 };
1894
1895 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
1896 assert!(
1897 edit_pending_text(
1898 &mut talk,
1899 &talks,
1900 "corrected",
1901 "first",
1902 std::slice::from_ref(&attachment.id),
1903 )
1904 .expect("edit")
1905 );
1906 let saved = talks.get(&talk.id).expect("reload edited draft");
1907 assert_eq!(saved.pending, "corrected");
1908 assert_eq!(saved.pending_attachments, vec![attachment]);
1909
1910 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
1911 assert!(
1912 !edit_pending_text(
1913 &mut talk,
1914 &talks,
1915 "stale edit",
1916 "corrected",
1917 &["a".repeat(32)],
1918 )
1919 .expect("stale edit is a conflict")
1920 );
1921 assert_eq!(
1922 talks.get(&talk.id).expect("reload after conflict").pending,
1923 "corrected\n\nlater"
1924 );
1925 assert!(
1926 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
1927 .expect("stale clear is a conflict")
1928 );
1929 assert_eq!(
1930 talks
1931 .get(&talk.id)
1932 .expect("reload after stale clear")
1933 .pending,
1934 "corrected\n\nlater"
1935 );
1936 }
1937
1938 #[tokio::test]
1939 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
1940 let (tmp, talks) = store();
1941 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
1942 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
1943 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
1944 let id = running.id.clone();
1945 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
1946
1947 let response_talks = talks.clone();
1948 let response_cfg = cfg.clone();
1949 let reply = tokio::spawn(async move {
1950 respond(&mut running, &response_talks, &response_cfg, &first).await
1951 });
1952 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
1953
1954 let mut queued = talks.get(&id).expect("queued handle");
1955 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
1956 reply.await.expect("join").expect("reply");
1957
1958 let saved = talks.get(&id).expect("reload");
1959 assert_eq!(saved.pending, "next");
1960 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
1961 }
1962
1963 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
1966 let calls = dir.join(format!("{id}.calls"));
1967 let script = format!(
1968 "#!/bin/sh\necho x >> '{}'\n{body}\n",
1969 calls.to_string_lossy()
1970 );
1971 let path = dir.join(format!("mock-{id}.sh"));
1972 std::fs::write(&path, script).expect("write mock");
1973 AgentSpec {
1974 id: id.to_owned(),
1975 kind: AgentKind::Command,
1976 model: None,
1977 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
1978 extra_args: Vec::new(),
1979 env: BTreeMap::new(),
1980 prompt_delivery: None,
1981 }
1982 }
1983
1984 fn calls(dir: &Path, id: &str) -> usize {
1985 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
1986 }
1987
1988 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
1989 let mut cfg = config(specs[0].clone());
1990 cfg.agents = specs;
1991 cfg.roles.chatter = Some(AgentChoice::Chain(
1992 ids.iter().map(|s| (*s).to_owned()).collect(),
1993 ));
1994 cfg
1995 }
1996
1997 #[tokio::test]
1998 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
1999 let (tmp, talks) = store();
2000 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2001 let b = counting_agent(tmp.path(), "b", "cat");
2002 let cfg = chain_config(vec![a, b], &["a", "b"]);
2003 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2004 assert_eq!(talk.agent, "a");
2005
2006 say(&mut talk, &talks, &cfg, "hello there", Vec::new())
2007 .await
2008 .expect("turn");
2009 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
2010 assert_eq!(calls(tmp.path(), "b"), 1);
2011 assert_eq!(talk.agent, "b", "the switch persists");
2012 assert!(talks.get(&talk.id).unwrap().agent == "b");
2013 let reply = talk.turns.last().unwrap();
2014 assert!(reply.body.contains("hello there"));
2015 assert!(
2016 reply.body.contains("magi task add --solo"),
2017 "a fresh seat gets the full briefing"
2018 );
2019 assert!(
2020 talk.turns
2021 .iter()
2022 .any(|t| t.body.contains("agent changed from a to b")),
2023 "the switch is noted"
2024 );
2025 }
2026
2027 #[tokio::test]
2028 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
2029 let (tmp, talks) = store();
2030 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2031 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
2032 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
2033 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2034
2035 let err = say(&mut talk, &talks, &cfg, "hi", Vec::new())
2036 .await
2037 .expect_err("every agent failed");
2038 assert!(err.to_string().contains("`a`"), "{err:#}");
2039 assert_eq!(calls(tmp.path(), "a"), 1);
2040 assert_eq!(calls(tmp.path(), "b"), 1);
2041 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
2042 }
2043
2044 #[test]
2045 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
2046 let (tmp, talks) = store();
2047 let b = counting_agent(tmp.path(), "b", "cat");
2048 let cfg = chain_config(vec![b], &["ghost", "b"]);
2049 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2050 assert_eq!(talk.agent, "b");
2051 }
2052
2053 #[tokio::test]
2054 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
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 cfg = chain_config(vec![a, b], &["a", "b"]);
2059 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
2060 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2061 .await
2062 .expect_err("a alone, and it fails");
2063 assert_eq!(calls(tmp.path(), "b"), 0);
2064 assert_eq!(talk.agent, "a");
2065 }
2066
2067 #[tokio::test]
2068 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
2069 let (tmp, talks) = store();
2070 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2071 let b = counting_agent(tmp.path(), "b", "cat");
2072 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
2073 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
2074 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
2075 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2076 .await
2077 .expect_err("c alone, and it fails");
2078 assert_eq!(calls(tmp.path(), "b"), 0);
2079 }
2080
2081 #[tokio::test]
2082 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
2083 let (tmp, talks) = store();
2084 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2085 let cfg = config(spec);
2086 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2087
2088 say(
2089 &mut talk,
2090 &talks,
2091 &cfg,
2092 "what does the queue module do?",
2093 Vec::new(),
2094 )
2095 .await
2096 .expect("first turn");
2097 let first_prompt = &talk.turns[1].body;
2098 assert!(first_prompt.contains("magi task add --solo"));
2099 assert!(first_prompt.contains("what does the queue module do?"));
2100
2101 say(&mut talk, &talks, &cfg, "and how is it locked?", Vec::new())
2102 .await
2103 .expect("second turn");
2104 let second_prompt = &talk.turns[3].body;
2105 assert!(
2106 !second_prompt.contains("magi task add --solo"),
2107 "the briefing is sent once, not on every turn: {second_prompt}"
2108 );
2109 assert!(second_prompt.contains("and how is it locked?"));
2110 }
2111
2112 #[tokio::test]
2113 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
2114 let (tmp, talks) = store();
2115 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2116 let mut b = a.clone();
2117 b.id = "other".to_owned();
2118 let mut cfg = config(a.clone());
2119 cfg.agents.push(b.clone());
2120 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
2121 say(&mut talk, &talks, &cfg, "remember the walrus", Vec::new())
2122 .await
2123 .expect("first turn");
2124 let old_session = talk.seat.claude_session.clone();
2125 assert_eq!(talk.seat.turns, 1);
2126
2127 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
2128 assert_eq!(talk.agent, "other");
2129 assert_eq!(talk.seat.turns, 0);
2130 assert_eq!(talk.seat.agent, "other");
2131 assert_ne!(talk.seat.claude_session, old_session);
2132 let note = talk.turns.last().expect("note");
2133 assert_eq!(note.who, Who::Agent);
2134 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2135 assert!(note.body.contains("changed from"), "{}", note.body);
2136 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
2137
2138 let before = talk.turns.len();
2139 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
2140 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
2141
2142 say(&mut talk, &talks, &cfg, "what did I say?", Vec::new())
2143 .await
2144 .expect("turn after switch");
2145 let prompt = &talk.turns.last().expect("reply").body;
2146 assert!(prompt.contains("remember the walrus"), "{prompt}");
2147 assert!(prompt.contains("## magi"), "{prompt}");
2148 assert!(prompt.contains("what did I say?"), "{prompt}");
2149 }
2150
2151 #[tokio::test]
2152 async fn say_appends_the_operator_turn_then_the_agent_turn() {
2153 let (tmp, talks) = store();
2154 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2155 let cfg = config(spec);
2156 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2157
2158 say(
2159 &mut talk,
2160 &talks,
2161 &cfg,
2162 "can I rename this function?",
2163 Vec::new(),
2164 )
2165 .await
2166 .expect("say");
2167
2168 assert_eq!(talk.turns.len(), 2);
2169 assert_eq!(talk.turns[0].who, Who::Operator);
2170 assert_eq!(talk.turns[0].body, "can I rename this function?");
2171 assert_eq!(talk.turns[1].who, Who::Agent);
2172 assert_eq!(talk.turns[1].body, "go ahead");
2173 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
2174 }
2175
2176 #[tokio::test]
2177 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
2178 let (tmp, talks) = store();
2179 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2180 let cfg = config(spec);
2181 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2182
2183 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2184 .await
2185 .expect_err("a turn with no answer is an error");
2186 assert!(err.to_string().contains("no answer"), "{err}");
2187
2188 let on_disk = talks.get(&talk.id).expect("get");
2189 assert_eq!(on_disk.turns.len(), 2);
2190 assert_eq!(on_disk.turns[0].body, "check the tests");
2191 let note = &on_disk.turns[1];
2192 assert_eq!(note.who, Who::Agent);
2193 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2194 assert!(note.body.contains("your message is saved"));
2195 }
2196
2197 #[tokio::test]
2203 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
2204 let (tmp, talks) = store();
2205 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2206 let cfg = config(spec);
2207 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2208
2209 let text =
2210 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
2211 failpoint::force_put_failures(PUT_RETRIES - 1);
2214 respond(&mut talk, &talks, &cfg, &text)
2215 .await
2216 .expect("respond must survive a write failure its own retries can outlast");
2217
2218 assert_eq!(talk.turns.len(), 2);
2219 assert_eq!(talk.turns[1].who, Who::Agent);
2220 assert_eq!(talk.turns[1].body, "go ahead");
2221 let on_disk = talks.get(&talk.id).expect("get");
2222 assert_eq!(
2223 on_disk.turns, talk.turns,
2224 "the reply must reach disk despite the early write failures"
2225 );
2226 }
2227
2228 #[tokio::test]
2234 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
2235 let (tmp, talks) = store();
2236 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2237 let cfg = config(spec);
2238 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2239
2240 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2241 failpoint::force_put_failures(PUT_RETRIES);
2246 let err = respond(&mut talk, &talks, &cfg, &text)
2247 .await
2248 .expect_err("a reply that cannot be saved must be reported, not swallowed");
2249 assert!(err.to_string().contains("could not be saved"), "{err}");
2250
2251 let on_disk = talks.get(&talk.id).expect("get");
2252 assert_eq!(
2253 on_disk.turns.len(),
2254 2,
2255 "the operator turn plus a visible note"
2256 );
2257 assert_eq!(on_disk.turns[0].body, "check the tests");
2258 let note = &on_disk.turns[1];
2259 assert_eq!(note.who, Who::Agent);
2260 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2261 assert!(
2262 note.body.contains("could not be saved"),
2263 "the operator must be told the reply is missing, not left staring \
2264 at a gap with no explanation: {}",
2265 note.body
2266 );
2267 assert_eq!(
2268 talk.turns, on_disk.turns,
2269 "the in-memory talk must match what actually landed on disk"
2270 );
2271
2272 let artifacts = talks.artifacts_of(&talk.id);
2275 let stash = std::fs::read_dir(&artifacts)
2276 .expect("artifacts dir")
2277 .filter_map(|e| e.ok())
2278 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
2279 .expect("a stash file for the lost reply");
2280 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
2281 assert_eq!(stashed, "go ahead");
2282
2283 assert_eq!(
2291 on_disk.seat.turns, 1,
2292 "the note's write must carry the turn the CLI actually took"
2293 );
2294 assert_eq!(
2295 on_disk.seat.claude_session, talk.seat.claude_session,
2296 "the session id handed to the CLI must survive the failed reply"
2297 );
2298 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
2299 assert!(
2300 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
2301 "the next turn must resume, not open the same session id twice"
2302 );
2303 }
2304
2305 #[tokio::test]
2310 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
2311 let (tmp, talks) = store();
2312 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2313 let cfg = config(spec);
2314 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2315
2316 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2317 failpoint::force_put_failures(PUT_RETRIES * 2);
2320 let err = respond(&mut talk, &talks, &cfg, &text)
2321 .await
2322 .expect_err("neither the reply nor the note could be saved");
2323 assert!(err.to_string().contains("could not be saved"), "{err}");
2324
2325 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
2326 let on_disk = talks.get(&talk.id).expect("get");
2327 assert_eq!(on_disk.turns.len(), 1);
2328
2329 assert_eq!(
2339 on_disk.seat.turns, 0,
2340 "an unwritable file cannot record the turn the CLI took"
2341 );
2342 assert_eq!(
2343 talk.seat.turns, 1,
2344 "the in-memory seat still reports the turn the CLI actually took"
2345 );
2346 assert_eq!(
2347 on_disk.seat.claude_session, talk.seat.claude_session,
2348 "the session id was minted at `begin` and never changes here"
2349 );
2350 }
2351
2352 #[tokio::test]
2356 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
2357 let (tmp, talks) = store();
2358 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2359 let cfg = config(spec);
2360 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2361
2362 let att = talks
2363 .put_attachment(
2364 &talk.id,
2365 "image/png",
2366 "screenshot.png",
2367 b"pretend-png-bytes",
2368 )
2369 .expect("put attachment");
2370
2371 say(&mut talk, &talks, &cfg, "", vec![att.clone()])
2372 .await
2373 .expect("an empty body with an attachment is still a turn");
2374
2375 let operator_turn = &talk.turns[0];
2376 assert_eq!(operator_turn.who, Who::Operator);
2377 assert_eq!(operator_turn.body, "");
2378 assert_eq!(operator_turn.attachments, vec![att.clone()]);
2379
2380 let prompt = &talk.turns[1].body;
2381 let expected_path = talks
2382 .attachments_dir(&talk.id)
2383 .join(format!("{}.png", att.id));
2384 assert!(
2385 prompt.contains(&expected_path.display().to_string()),
2386 "the agent must be told the attachment's absolute path: {prompt}"
2387 );
2388 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
2389 }
2390
2391 #[test]
2400 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
2401 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
2402 let att = Attachment {
2403 id: "0".repeat(32),
2404 name: "shot.png".to_owned(),
2405 mime: "image/png".to_owned(),
2406 bytes: 3,
2407 };
2408 let path = talks
2409 .attachment_path("some-talk-id", &att)
2410 .expect("a supported mime always yields a path");
2411 assert!(
2412 path.is_absolute(),
2413 "must be absolute even off a relative store root: {}",
2414 path.display()
2415 );
2416 }
2417
2418 #[tokio::test]
2419 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
2420 let (tmp, talks) = store();
2425 let slow = mock_agent(
2426 tmp.path(),
2427 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
2428 BTreeMap::new(),
2429 );
2430 let mut cfg = config(slow);
2431 cfg.graph.timeout_talk = 1;
2432 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2433
2434 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2435 .await
2436 .expect_err("a turn that never answers is an error");
2437 assert!(
2438 err.to_string().contains("did not answer within 1s"),
2439 "{err}"
2440 );
2441
2442 let on_disk = talks.get(&talk.id).expect("get");
2443 let note = on_disk.turns.last().expect("a note turn was recorded");
2444 assert!(
2445 note.body.contains("did not answer within 1s"),
2446 "the transcript must show the configured timeout: {}",
2447 note.body
2448 );
2449 }
2450
2451 #[test]
2452 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
2453 let (tmp, talks) = store();
2454 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2455 let cfg = config(spec);
2456 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2457
2458 close(&mut talk, &talks).expect("close");
2459 assert_eq!(talk.status, TalkStatus::Closed);
2460 close(&mut talk, &talks).expect("closing twice is not an error");
2461
2462 let err =
2463 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
2464 assert!(err.to_string().contains("closed"));
2465 let _ = &cfg; }
2467
2468 #[tokio::test]
2469 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2470 let (tmp, talks) = store();
2471 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2472 let cfg = config(spec);
2473 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2476
2477 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
2481 close(&mut closed_elsewhere, &talks).expect("close");
2482 assert_eq!(
2483 talks.get(&in_flight.id).expect("reread").status,
2484 TalkStatus::Closed,
2485 "the close landed on disk before the turn finished"
2486 );
2487
2488 assert_eq!(in_flight.status, TalkStatus::Open);
2492 respond(&mut in_flight, &talks, &cfg, "one more question")
2493 .await
2494 .expect("the turn itself still completes");
2495
2496 let on_disk = talks.get(&in_flight.id).expect("reread");
2497 assert_eq!(
2498 on_disk.status,
2499 TalkStatus::Closed,
2500 "a close must stick even when a turn that started before it finishes after it"
2501 );
2502 assert!(
2505 on_disk.turns.iter().any(|t| t.body == "here you go"),
2506 "the in-flight turn's own reply is still recorded: {:?}",
2507 on_disk.turns
2508 );
2509 }
2510
2511 #[test]
2512 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
2513 let (tmp, talks) = store();
2514 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2515 let cfg = config(spec);
2516 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2519
2520 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
2523 close(&mut closed_elsewhere, &talks).expect("close");
2524 assert_eq!(
2525 talks.get(&stale.id).expect("reread").status,
2526 TalkStatus::Closed,
2527 "the close landed on disk before record was called"
2528 );
2529
2530 assert_eq!(stale.status, TalkStatus::Open);
2534 let err = record(&mut stale, &talks, "still there?", Vec::new())
2535 .expect_err("a close that landed first must be honored, not overwritten");
2536 assert!(err.to_string().contains("closed"));
2537
2538 let on_disk = talks.get(&stale.id).expect("reread");
2539 assert_eq!(
2540 on_disk.status,
2541 TalkStatus::Closed,
2542 "record must not resurrect a conversation closed while its snapshot was stale"
2543 );
2544 assert!(
2545 on_disk.turns.is_empty(),
2546 "the rejected turn must not have been appended: {:?}",
2547 on_disk.turns
2548 );
2549 let _ = &cfg; }
2551
2552 #[test]
2553 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
2554 let (tmp, talks) = store();
2555 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2556 let cfg = config(spec);
2557 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2558
2559 let held = talks.guard();
2563
2564 let talks2 = talks.clone();
2565 let id = talk.id.clone();
2566 let closing = std::thread::spawn(move || {
2567 let mut talk = talks2.get(&id).expect("get");
2568 close(&mut talk, &talks2).expect("close");
2569 });
2570
2571 std::thread::sleep(Duration::from_millis(50));
2572 assert!(
2573 !closing.is_finished(),
2574 "close must wait for the guard, not read and write while it is held - \
2575 a re-read alone narrows this window without closing it"
2576 );
2577
2578 drop(held);
2579 closing.join().expect("close thread panicked");
2580
2581 assert_eq!(
2582 talks.get(&talk.id).expect("reread").status,
2583 TalkStatus::Closed,
2584 "once the guard is free, close still lands"
2585 );
2586 let _ = &cfg; }
2588
2589 #[test]
2590 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
2591 let (tmp, talks) = store();
2592 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2593 let cfg = config(spec);
2594 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2595
2596 close(&mut talk, &talks).expect("close");
2597 assert_eq!(talk.status, TalkStatus::Closed);
2598
2599 reopen(&mut talk, &talks).expect("reopen");
2600 assert_eq!(talk.status, TalkStatus::Open);
2601 assert_eq!(
2602 talks.get(&talk.id).expect("reread").status,
2603 TalkStatus::Open
2604 );
2605
2606 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
2608 assert_eq!(talk.status, TalkStatus::Open);
2609
2610 record(&mut talk, &talks, "one more thing", Vec::new())
2611 .expect("a reopened talk takes turns again");
2612 let _ = &cfg; }
2614
2615 #[test]
2616 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
2617 let (tmp, talks) = store();
2618 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2619 let cfg = config(spec);
2620 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2621
2622 let artifacts = talks.artifacts_of(&talk.id);
2623 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
2624 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
2625
2626 talks.remove(&talk.id).expect("remove");
2627 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
2628 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
2629 assert!(
2630 talks.get(&talk.id).is_err(),
2631 "a removed talk cannot be read back"
2632 );
2633
2634 let err = talks
2635 .remove("nonexistent-id")
2636 .expect_err("unknown id refused");
2637 assert!(err.to_string().contains("no talk matches"), "{err}");
2638 let _ = &cfg; }
2640
2641 #[tokio::test]
2642 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2643 let (tmp, talks) = store();
2644 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2645 let cfg = config(spec);
2646 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2649
2650 talks.remove(&in_flight.id).expect("remove");
2651 assert!(
2652 talks.get(&in_flight.id).is_err(),
2653 "the delete landed on disk before the turn finished"
2654 );
2655
2656 respond(&mut in_flight, &talks, &cfg, "one more question")
2659 .await
2660 .expect("the turn itself still completes rather than erroring");
2661
2662 assert!(
2663 talks.get(&in_flight.id).is_err(),
2664 "a delete must stick even when a turn that started before it finishes after it"
2665 );
2666 }
2667
2668 #[test]
2669 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
2670 let (tmp, talks) = store();
2671 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2672 let cfg = config(spec);
2673 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2676
2677 talks.remove(&stale.id).expect("remove");
2678
2679 let err = record(&mut stale, &talks, "still there?", Vec::new())
2683 .expect_err("a delete that landed first must be honored, not overwritten");
2684 assert!(err.to_string().contains("deleted"), "{err}");
2685
2686 assert!(
2687 talks.get(&stale.id).is_err(),
2688 "record must not resurrect a conversation deleted while its snapshot was stale"
2689 );
2690 let _ = &cfg; }
2692
2693 #[test]
2694 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
2695 let (tmp, talks) = store();
2696 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2697 let cfg = config(spec);
2698 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2701
2702 talks.remove(&stale.id).expect("remove");
2703
2704 let err = close(&mut stale, &talks)
2708 .expect_err("a delete that landed first must be honored, not overwritten");
2709 assert!(err.to_string().contains("deleted"), "{err}");
2710
2711 assert!(
2712 talks.get(&stale.id).is_err(),
2713 "close must not resurrect a conversation deleted while its snapshot was stale"
2714 );
2715 let _ = &cfg; }
2717
2718 #[test]
2719 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
2720 let (tmp, talks) = store();
2721 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2722 let cfg = config(spec);
2723 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2726 close(&mut stale, &talks).expect("close");
2727
2728 talks.remove(&stale.id).expect("remove");
2729
2730 let err = reopen(&mut stale, &talks)
2734 .expect_err("a delete that landed first must be honored, not overwritten");
2735 assert!(err.to_string().contains("deleted"), "{err}");
2736
2737 assert!(
2738 talks.get(&stale.id).is_err(),
2739 "reopen must not resurrect a conversation deleted while its snapshot was stale"
2740 );
2741 let _ = &cfg; }
2743
2744 #[test]
2745 fn list_puts_open_talks_before_closed_ones() {
2746 let (tmp, talks) = store();
2747 let make = |id: &str, status: TalkStatus| {
2748 let mut t = Talk {
2749 schema: SCHEMA,
2750 id: id.to_owned(),
2751 repo: tmp.path().to_owned(),
2752 agent: "mock".to_owned(),
2753 status,
2754 turns: Vec::new(),
2755 pending: String::new(),
2756 pending_attachments: Vec::new(),
2757 fallback: false,
2758 created_at: Timestamp::now(),
2759 updated_at: Timestamp::now(),
2760 seat: SeatState::new(SEAT, "mock", 7),
2761 };
2762 talks.put(&mut t).expect("put");
2763 };
2764 make("20260901-000000-0001", TalkStatus::Open);
2765 make("20260902-000000-0002", TalkStatus::Open);
2766 make("20260903-000000-0003", TalkStatus::Closed);
2767
2768 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
2769 assert_eq!(
2770 ids,
2771 [
2772 "20260902-000000-0002",
2773 "20260901-000000-0001",
2774 "20260903-000000-0003"
2775 ]
2776 );
2777 assert_eq!(talks.count_open(), 2);
2778 }
2779
2780 #[test]
2781 fn tasks_of_finds_only_this_talks_own_tasks() {
2782 let dir = tempfile::tempdir().expect("tempdir");
2783 let queue = Queue::at(dir.path().join("queue"));
2784
2785 let mut mine = Task::new(
2786 "rework the loader".to_owned(),
2787 "rework the loader".to_owned(),
2788 PathBuf::from("/repo"),
2789 Source::Agent {
2790 run: "20260904-014455-ab12".to_owned(),
2791 node: "chat".to_owned(),
2792 },
2793 );
2794 queue.put(&mut mine).expect("put mine");
2795
2796 let mut theirs = Task::new(
2797 "unrelated".to_owned(),
2798 "unrelated".to_owned(),
2799 PathBuf::from("/repo"),
2800 Source::Agent {
2801 run: "20260904-090000-zz99".to_owned(),
2802 node: "implement".to_owned(),
2803 },
2804 );
2805 queue.put(&mut theirs).expect("put theirs");
2806
2807 let mut human = Task::new(
2808 "typed by hand".to_owned(),
2809 "typed by hand".to_owned(),
2810 PathBuf::from("/repo"),
2811 Source::Human,
2812 );
2813 queue.put(&mut human).expect("put human");
2814
2815 let found = tasks_of(&queue, "20260904-014455-ab12");
2816 assert_eq!(found.len(), 1);
2817 assert_eq!(found[0].id, mine.id);
2818 }
2819
2820 #[test]
2821 fn the_briefing_names_solo_task_add() {
2822 let brief = briefing(Path::new("/repo"), "en", false);
2823 assert!(brief.contains("magi task add --solo"));
2824 assert!(brief.contains("/repo"));
2825 assert!(!brief.contains("Hold this conversation in"));
2826 }
2827
2828 #[test]
2834 fn the_briefing_explains_targeting_a_different_repository_by_name() {
2835 let brief = briefing(Path::new("/repo"), "en", false);
2836 assert!(brief.contains("--repo does not have to be a full path"));
2837 assert!(brief.contains("owner/repo"));
2838 assert!(brief.contains("magi repos"));
2839 assert!(brief.contains("ask the operator"));
2840 }
2841
2842 #[test]
2843 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
2844 let brief = briefing(Path::new("/repo"), "en", false);
2845 assert!(brief.contains("--attach <path>"), "{brief}");
2846 assert!(brief.contains("deleting this conversation"), "{brief}");
2847 }
2848
2849 #[test]
2850 fn the_briefing_names_the_language_when_it_is_not_english() {
2851 let brief = briefing(Path::new("/repo"), "Japanese", false);
2852 assert!(brief.contains("Hold this conversation in Japanese"));
2853 }
2854
2855 #[test]
2856 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
2857 let read_only = briefing(Path::new("/repo"), "en", false);
2858 assert!(read_only.contains("Do not write files"));
2859 assert!(!read_only.contains("allow_write"));
2860
2861 let writable = briefing(Path::new("/repo"), "en", true);
2862 assert!(!writable.contains("Do not write files"));
2863 assert!(writable.contains("allow_write = true"));
2864 assert!(writable.contains("magi task add --solo"));
2867 assert!(writable.contains("say plainly what you"));
2868 }
2869}