1use std::path::{Path, PathBuf};
47use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
48use std::time::Duration;
49
50use anyhow::{Context, Result, bail};
51use jiff::Timestamp;
52use serde::{Deserialize, Serialize};
53
54use crate::agent::{self, Invocation, SeatState};
55use crate::config::{AgentSpec, Config};
56use crate::queue::{Queue, Source, Task};
57
58pub const SCHEMA: u32 = 1;
60
61fn turn_timeout(cfg: &Config) -> Duration {
72 Duration::from_secs(cfg.graph.timeout_talk)
73}
74
75const SEAT: &str = "talk";
78
79const MAGI_NOTE: &str = "magi: ";
81
82#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
84#[serde(rename_all = "lowercase")]
85pub enum Who {
86 Operator,
88 Agent,
91}
92
93#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
101#[serde(deny_unknown_fields)]
102pub struct Attachment {
103 pub id: String,
105 pub name: String,
107 pub mime: String,
110 pub bytes: u64,
112}
113
114#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
116#[serde(deny_unknown_fields)]
117pub struct Turn {
118 pub who: Who,
120 pub body: String,
122 pub at: Timestamp,
124 #[serde(default)]
127 pub attachments: Vec<Attachment>,
128 #[serde(default, skip_serializing_if = "Option::is_none")]
133 pub usage: Option<TurnUsage>,
134}
135
136#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
144pub struct TurnUsage {
145 pub context_tokens: u64,
147 pub agent: String,
149 #[serde(default)]
151 pub model: Option<String>,
152}
153
154#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
159pub struct ContextUsage {
160 pub tokens: Option<u64>,
162 pub window: Option<u64>,
164 pub percent: Option<u64>,
166 pub warn: bool,
168 pub since_switch: bool,
172 pub model: Option<String>,
174 #[serde(default)]
177 pub estimated: bool,
178}
179
180const CONTEXT_WARN_PERCENT: u64 = 80;
182
183const STANDING_PROMPT_FALLBACK_CHARS: u64 = 7000;
186
187pub fn estimate_context_tokens(talk: &Talk, standing_chars: u64) -> Option<u64> {
200 let mut counted = false;
201 let mut chars = standing_chars;
202 for t in talk.turns.iter().filter(|t| !t.body.starts_with(MAGI_NOTE)) {
203 counted = true;
204 chars += t.body.chars().count() as u64;
205 }
206 counted.then(|| (chars * 2).div_ceil(7))
208}
209
210pub fn context_usage(talk: &Talk, cfg: Option<&Config>) -> ContextUsage {
223 let current = cfg.and_then(|c| c.agents.iter().find(|a| a.id == talk.agent));
224 let model = current.and_then(|a| a.model.clone());
225 let window = cfg
226 .zip(model.as_deref())
227 .and_then(|(c, m)| c.context_window(m))
228 .filter(|w| *w > 0);
229 let usage = talk
230 .turns
231 .iter()
232 .rev()
233 .find(|t| t.who == Who::Agent && !t.body.starts_with(MAGI_NOTE))
234 .and_then(|t| t.usage.as_ref());
235 let measured = usage.map(|u| u.context_tokens);
236 let tokens = measured.or_else(|| {
237 let standing = cfg.map_or(STANDING_PROMPT_FALLBACK_CHARS, |c| {
238 briefing(&talk.repo, &c.graph.language, c.talk.allow_write)
239 .chars()
240 .count() as u64
241 });
242 estimate_context_tokens(talk, standing)
243 });
244 let estimated = measured.is_none() && tokens.is_some();
245 let since_switch =
246 usage.is_some_and(|u| u.agent != talk.agent || (current.is_some() && u.model != model));
247 let (percent, warn) = match (tokens, window) {
248 (Some(t), Some(w)) => (
249 Some(t.saturating_mul(100) / w),
250 t.saturating_mul(100) >= w.saturating_mul(CONTEXT_WARN_PERCENT),
251 ),
252 _ => (None, false),
253 };
254 ContextUsage {
255 tokens,
256 window,
257 percent,
258 warn,
259 since_switch,
260 model,
261 estimated,
262 }
263}
264
265#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
269#[serde(rename_all = "lowercase")]
270pub enum TalkStatus {
271 Open,
274 Closed,
276}
277
278impl TalkStatus {
279 pub fn open(self) -> bool {
281 matches!(self, Self::Open)
282 }
283
284 pub fn as_str(self) -> &'static str {
286 match self {
287 Self::Open => "open",
288 Self::Closed => "closed",
289 }
290 }
291}
292
293#[derive(Debug, Clone, Serialize, Deserialize)]
295#[serde(deny_unknown_fields)]
296pub struct Talk {
297 pub schema: u32,
299 pub id: String,
301 pub repo: PathBuf,
303 pub agent: String,
305 pub status: TalkStatus,
307 pub turns: Vec<Turn>,
309 #[serde(default)]
312 pub pending: String,
313 #[serde(default)]
315 pub pending_attachments: Vec<Attachment>,
316 #[serde(default)]
321 pub fallback: bool,
322 pub created_at: Timestamp,
324 pub updated_at: Timestamp,
326 seat: SeatState,
331}
332
333impl Talk {
334 pub fn short(&self) -> &str {
336 short(&self.id)
337 }
338}
339
340#[derive(Debug, Clone)]
342pub struct Talks {
343 root: PathBuf,
344 lock: Arc<Mutex<()>>,
353}
354
355impl Talks {
356 pub fn open() -> Self {
358 Self::at(crate::run::home().join("talks"))
359 }
360
361 pub fn at(root: PathBuf) -> Self {
364 Self {
365 root,
366 lock: Arc::new(Mutex::new(())),
367 }
368 }
369
370 fn guard(&self) -> MutexGuard<'_, ()> {
379 self.lock.lock().unwrap_or_else(PoisonError::into_inner)
380 }
381
382 pub fn root(&self) -> &Path {
384 &self.root
385 }
386
387 pub fn path_of(&self, id: &str) -> PathBuf {
389 self.root.join(format!("{id}.json"))
390 }
391
392 pub fn artifacts_of(&self, id: &str) -> PathBuf {
395 self.root.join(format!("{id}.artifacts"))
396 }
397
398 pub fn attachments_dir(&self, id: &str) -> PathBuf {
402 self.artifacts_of(id).join("attachments")
403 }
404
405 pub fn put_attachment(
414 &self,
415 id: &str,
416 mime: &str,
417 name: &str,
418 data: &[u8],
419 ) -> Result<Attachment> {
420 let dir = self.attachments_dir(id);
421 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
422 let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
423 let att = Attachment {
424 id: new_attachment_id(),
425 name: name.to_owned(),
426 mime: mime.to_owned(),
427 bytes: data.len() as u64,
428 };
429 std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
430 .with_context(|| format!("write attachment {}", att.id))?;
431 std::fs::write(
432 dir.join(format!("{}.json", att.id)),
433 serde_json::to_string(&att).context("serialize attachment")?,
434 )
435 .with_context(|| format!("write attachment metadata {}", att.id))?;
436 Ok(att)
437 }
438
439 pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
448 if !valid_attachment_id(att_id) {
449 return Ok(None);
450 }
451 let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
452 if !meta_path.is_file() {
453 return Ok(None);
454 }
455 let att = serde_json::from_str(
456 &std::fs::read_to_string(&meta_path)
457 .with_context(|| format!("read {}", meta_path.display()))?,
458 )
459 .with_context(|| format!("parse {}", meta_path.display()))?;
460 Ok(Some(att))
461 }
462
463 pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
467 let Some(att) = self.attachment_meta(id, att_id)? else {
468 return Ok(None);
469 };
470 let ext = attachment_ext(&att.mime).with_context(|| {
471 format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
472 })?;
473 let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
474 let data =
475 std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
476 Ok(Some((att, data)))
477 }
478
479 fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
496 let ext = attachment_ext(&att.mime)?;
497 let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
498 std::path::absolute(&path).ok()
499 }
500
501 pub fn put(&self, t: &mut Talk) -> Result<()> {
510 std::fs::create_dir_all(&self.root)
511 .with_context(|| format!("create {}", self.root.display()))?;
512 t.updated_at = Timestamp::now();
513 let body = serde_json::to_string_pretty(t).context("serialize talk")?;
514 let path = self.path_of(&t.id);
515 let tmp = path.with_extension("json.tmp");
516 write_atomic(&tmp, &path, &body)
517 }
518
519 pub fn get(&self, id: &str) -> Result<Talk> {
521 let resolved = self.resolve_id(id)?;
522 read_path(&self.path_of(&resolved))
523 }
524
525 pub fn list(&self) -> Vec<Talk> {
528 self.list_counting_unreadable().0
529 }
530
531 pub fn list_counting_unreadable(&self) -> (Vec<Talk>, usize) {
534 let mut unreadable = 0;
535 let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
536 .into_iter()
537 .flatten()
538 .flatten()
539 .map(|e| e.path())
540 .filter(|p| p.extension().is_some_and(|x| x == "json"))
541 .filter_map(|p| {
542 let talk = read_path(&p).ok();
543 if talk.is_none() {
544 unreadable += 1;
545 }
546 talk
547 })
548 .collect();
549 all.sort_unstable_by(|a, b| {
550 let rank = |t: &Talk| u8::from(!t.status.open());
551 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
552 });
553 (all, unreadable)
554 }
555
556 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
558 if self.path_of(prefix).is_file() {
559 return Ok(prefix.to_owned());
560 }
561 let hits: Vec<String> = self
562 .list()
563 .into_iter()
564 .map(|t| t.id)
565 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
566 .collect();
567 match hits.len() {
568 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
569 0 => bail!("no talk matches `{prefix}`"),
570 _ => bail!(
571 "`{prefix}` matches {} talks: {}",
572 hits.len(),
573 hits.join(", ")
574 ),
575 }
576 }
577
578 pub fn revision(&self) -> u64 {
581 std::fs::read_dir(&self.root)
582 .into_iter()
583 .flatten()
584 .flatten()
585 .filter(|e| e.path().extension().is_none_or(|x| x != "turn"))
586 .filter_map(|e| e.metadata().ok())
587 .filter_map(|m| m.modified().ok())
588 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
589 .map(|d| d.as_millis() as u64)
590 .max()
591 .unwrap_or(0)
592 }
593
594 pub fn count_open(&self) -> usize {
596 self.list().iter().filter(|t| t.status.open()).count()
597 }
598
599 pub fn turn_path(&self, id: &str) -> PathBuf {
602 self.root.join(format!("{id}.turn"))
603 }
604
605 pub fn claim_turn(&self, id: &str) -> Result<Option<TurnLease>> {
614 self.claim_turn_at(id, Timestamp::now())
615 }
616
617 fn claim_turn_at(&self, id: &str, now: Timestamp) -> Result<Option<TurnLease>> {
618 std::fs::create_dir_all(&self.root)
619 .with_context(|| format!("create {}", self.root.display()))?;
620 let path = self.turn_path(id);
621 let token = crate::rng::SplitMix64::new(crate::rng::entropy()).uuid_v4();
622 if create_turn(&path, &token, now)? {
623 return Ok(Some(TurnLease { path, token }));
624 }
625 if read_turn(&path).is_some_and(|r| r.fresh(now)) {
626 return Ok(None);
627 }
628 let Some(_lock) = TurnLock::take(&path)? else {
633 return Ok(None);
634 };
635 if read_turn(&path).is_some_and(|r| r.fresh(now)) {
636 return Ok(None);
637 }
638 let _ = std::fs::remove_file(&path);
639 Ok(create_turn(&path, &token, now)?.then_some(TurnLease { path, token }))
640 }
641
642 pub fn turn_held(&self, id: &str) -> bool {
644 read_turn(&self.turn_path(id)).is_some_and(|r| r.fresh(Timestamp::now()))
645 }
646
647 pub fn remove(&self, id: &str) -> Result<()> {
660 let _guard = self.guard();
661 let resolved = self.resolve_id(id)?;
662 let path = self.path_of(&resolved);
663 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
664 let artifacts = self.artifacts_of(&resolved);
665 if artifacts.is_dir() {
666 std::fs::remove_dir_all(&artifacts)
667 .with_context(|| format!("remove {}", artifacts.display()))?;
668 }
669 let _ = std::fs::remove_file(self.turn_path(&resolved));
670 Ok(())
671 }
672}
673
674#[derive(Debug, Serialize, Deserialize)]
677struct TurnRecord {
678 token: String,
679 pid: u32,
680 beat_at: Timestamp,
681}
682
683impl TurnRecord {
684 fn fresh(&self, now: Timestamp) -> bool {
685 now.as_second() - self.beat_at.as_second() <= crate::ask::LEASE_TTL.as_secs() as i64
686 }
687}
688
689fn read_turn(path: &Path) -> Option<TurnRecord> {
690 serde_json::from_str(&std::fs::read_to_string(path).ok()?).ok()
691}
692
693const TAKEOVER_LOCK_TTL: Duration = Duration::from_secs(10);
695
696fn create_exclusive(path: &Path, body: &str) -> Result<bool> {
698 use std::io::Write as _;
699 match std::fs::OpenOptions::new()
700 .write(true)
701 .create_new(true)
702 .open(path)
703 {
704 Ok(mut f) => {
705 if let Err(e) = f.write_all(body.as_bytes()) {
706 drop(f);
707 let _ = std::fs::remove_file(path);
708 return Err(e).with_context(|| format!("write {}", path.display()));
709 }
710 Ok(true)
711 }
712 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
713 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
714 }
715}
716
717fn create_turn(path: &Path, token: &str, now: Timestamp) -> Result<bool> {
718 let record = TurnRecord {
719 token: token.to_owned(),
720 pid: std::process::id(),
721 beat_at: now,
722 };
723 let body = serde_json::to_string(&record).context("serialize turn lease")?;
724 let tmp = path.with_extension(format!("turn.{token}.new"));
727 std::fs::write(&tmp, body).with_context(|| format!("write {}", tmp.display()))?;
728 let linked = std::fs::hard_link(&tmp, path);
729 let _ = std::fs::remove_file(&tmp);
730 match linked {
731 Ok(()) => Ok(true),
732 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
733 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
734 }
735}
736
737struct TurnLock {
741 path: PathBuf,
742 token: String,
743}
744
745impl TurnLock {
746 fn token() -> String {
747 crate::rng::SplitMix64::new(crate::rng::entropy()).uuid_v4()
748 }
749
750 fn take(lease: &Path) -> Result<Option<Self>> {
751 let path = lease.with_extension("turn.lock");
752 let token = Self::token();
753 if create_exclusive(&path, &token)? {
754 return Ok(Some(Self { path, token }));
755 }
756 let seen = std::fs::read_to_string(&path).ok();
757 let aged = std::fs::metadata(&path)
758 .and_then(|m| m.modified())
759 .ok()
760 .and_then(|t| t.elapsed().ok())
761 .is_some_and(|age| age > TAKEOVER_LOCK_TTL);
762 if !aged {
763 return Ok(None);
764 }
765 let aside = path.with_extension(format!("lock.{token}.dead"));
770 if std::fs::rename(&path, &aside).is_err() {
771 return Ok(None);
772 }
773 let moved = std::fs::read_to_string(&aside).ok();
774 if moved != seen {
775 let _ = std::fs::hard_link(&aside, &path);
776 let _ = std::fs::remove_file(&aside);
777 return Ok(None);
778 }
779 let _ = std::fs::remove_file(&aside);
780 if create_exclusive(&path, &token)? {
781 return Ok(Some(Self { path, token }));
782 }
783 Ok(None)
784 }
785
786 fn take_patiently(lease: &Path) -> Option<Self> {
788 for _ in 0..50 {
789 match Self::take(lease) {
790 Ok(Some(lock)) => return Some(lock),
791 Ok(None) => std::thread::sleep(Duration::from_millis(10)),
792 Err(_) => return None,
793 }
794 }
795 None
796 }
797}
798
799impl Drop for TurnLock {
800 fn drop(&mut self) {
801 if std::fs::read_to_string(&self.path).is_ok_and(|t| t == self.token) {
803 let _ = std::fs::remove_file(&self.path);
804 }
805 }
806}
807
808pub const TURN_BEAT: Duration = Duration::from_secs(20);
811
812#[derive(Debug)]
817pub struct TurnLease {
818 path: PathBuf,
819 token: String,
820}
821
822impl TurnLease {
823 pub fn beat(&self) -> Result<bool> {
827 let _lock = TurnLock::take_patiently(&self.path)
828 .with_context(|| format!("lock {} to renew it", self.path.display()))?;
829 let Some(mut record) = read_turn(&self.path).filter(|r| r.token == self.token) else {
830 return Ok(false);
831 };
832 record.beat_at = Timestamp::now();
833 let body = serde_json::to_string(&record).context("serialize turn lease")?;
834 let tmp = self.path.with_extension(format!("turn.{}.tmp", self.token));
835 write_atomic(&tmp, &self.path, &body)?;
836 Ok(true)
837 }
838
839 pub async fn beating<T>(&self, fut: impl std::future::Future<Output = T>) -> Result<T> {
844 self.beating_every(TURN_BEAT, fut).await
845 }
846
847 async fn beating_every<T>(
848 &self,
849 period: Duration,
850 fut: impl std::future::Future<Output = T>,
851 ) -> Result<T> {
852 tokio::pin!(fut);
853 loop {
854 match tokio::time::timeout(period, &mut fut).await {
855 Ok(out) => return Ok(out),
856 Err(_) => match self.beat() {
857 Ok(true) => {}
858 Ok(false) => bail!(
859 "the turn lease {} was taken over; this turn is stopped",
860 self.path.display()
861 ),
862 Err(e) => tracing::warn!("{e:#}"),
863 },
864 }
865 }
866 }
867}
868
869impl Drop for TurnLease {
870 fn drop(&mut self) {
871 if let Some(_lock) = TurnLock::take_patiently(&self.path) {
874 if read_turn(&self.path).is_some_and(|r| r.token == self.token) {
875 let _ = std::fs::remove_file(&self.path);
876 }
877 }
878 }
879}
880
881pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
891 let repo = repo.canonicalize().unwrap_or(repo);
894 let spec = match agent {
897 Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
898 None => agent::pick_chain(
899 &cfg.agents,
900 cfg.roles.chatter.as_ref(),
901 &agent::installed,
902 "chatter",
903 )?
904 .remove(0),
905 };
906
907 let now = Timestamp::now();
908 let mut talk = Talk {
909 schema: SCHEMA,
910 id: new_id(),
911 repo,
912 agent: spec.id.clone(),
913 status: TalkStatus::Open,
914 turns: Vec::new(),
915 pending: String::new(),
916 pending_attachments: Vec::new(),
917 fallback: agent.is_none(),
918 created_at: now,
919 updated_at: now,
920 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
921 };
922 store.put(&mut talk)?;
923 Ok(talk)
924}
925
926pub fn record(
933 talk: &mut Talk,
934 store: &Talks,
935 text: &str,
936 attachments: Vec<Attachment>,
937) -> Result<String> {
938 let _guard = store.guard();
946 let Ok(fresh) = store.get(&talk.id) else {
951 bail!("talk {} was deleted", talk.short());
952 };
953 talk.status = fresh.status;
954 talk.pending = fresh.pending;
957 talk.pending_attachments = fresh.pending_attachments;
958 if !talk.status.open() {
959 bail!(
960 "talk {} is {} and takes no more turns",
961 talk.short(),
962 talk.status.as_str()
963 );
964 }
965 let text = text.trim();
966 if text.is_empty() && attachments.is_empty() {
967 bail!("nothing to say");
968 }
969 talk.turns.push(Turn {
970 who: Who::Operator,
971 body: text.to_owned(),
972 at: Timestamp::now(),
973 attachments,
974 usage: None,
975 });
976 store.put(talk)?;
977 Ok(text.to_owned())
978}
979
980pub fn queue(
982 talk: &mut Talk,
983 store: &Talks,
984 text: &str,
985 attachments: Vec<Attachment>,
986) -> Result<()> {
987 let text = text.trim();
988 if text.is_empty() && attachments.is_empty() {
989 bail!("nothing to say");
990 }
991 let _guard = store.guard();
992 let mut fresh = store
993 .get(&talk.id)
994 .with_context(|| format!("talk {} was deleted", talk.short()))?;
995 if !fresh.status.open() {
996 bail!(
997 "talk {} is {} and takes no more turns",
998 fresh.short(),
999 fresh.status.as_str()
1000 );
1001 }
1002 if !text.is_empty() {
1003 if fresh.pending.is_empty() {
1004 fresh.pending = text.to_owned();
1005 } else {
1006 fresh.pending.push_str("\n\n");
1007 fresh.pending.push_str(text);
1008 }
1009 }
1010 fresh.pending_attachments.extend(attachments);
1011 store.put(&mut fresh)?;
1012 *talk = fresh;
1013 Ok(())
1014}
1015
1016pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
1018 let _guard = store.guard();
1019 let mut fresh = store
1020 .get(&talk.id)
1021 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1022 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
1023 *talk = fresh;
1024 return Ok(None);
1025 }
1026 let text = std::mem::take(&mut fresh.pending);
1027 let attachments = std::mem::take(&mut fresh.pending_attachments);
1028 fresh.turns.push(Turn {
1029 who: Who::Operator,
1030 body: text.clone(),
1031 at: Timestamp::now(),
1032 attachments,
1033 usage: None,
1034 });
1035 store.put(&mut fresh)?;
1036 *talk = fresh;
1037 Ok(Some(text))
1038}
1039
1040pub async fn say(
1043 talk: &mut Talk,
1044 store: &Talks,
1045 cfg: &Config,
1046 text: &str,
1047 attachments: Vec<Attachment>,
1048) -> Result<()> {
1049 let text = record(talk, store, text, attachments)?;
1050 turn(talk, store, cfg, &text).await
1051}
1052
1053pub async fn respond(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1055 turn(talk, store, cfg, text).await
1056}
1057
1058pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
1076 let _guard = store.guard();
1077 let mut fresh = store
1078 .get(&talk.id)
1079 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1080 fresh.status = TalkStatus::Closed;
1081 fresh.pending.clear();
1083 fresh.pending_attachments.clear();
1084 store.put(&mut fresh)?;
1085 *talk = fresh;
1086 Ok(())
1087}
1088
1089pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
1100 let _guard = store.guard();
1101 let mut fresh = store
1102 .get(&talk.id)
1103 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1104 fresh.status = TalkStatus::Open;
1105 store.put(&mut fresh)?;
1106 *talk = fresh;
1107 Ok(())
1108}
1109
1110pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
1123 let _guard = store.guard();
1124 let mut fresh = store
1125 .get(&talk.id)
1126 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1127 if fresh.agent == spec.id {
1128 *talk = fresh;
1129 return Ok(false);
1130 }
1131 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
1132 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1133 fresh.fallback = false;
1135 fresh.turns.push(Turn {
1136 who: Who::Agent,
1137 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
1138 at: Timestamp::now(),
1139 attachments: Vec::new(),
1140 usage: None,
1141 });
1142 store.put(&mut fresh)?;
1143 *talk = fresh;
1144 Ok(true)
1145}
1146
1147pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
1149 let _guard = store.guard();
1150 let mut fresh = store
1151 .get(&talk.id)
1152 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1153 fresh.pending.clear();
1154 fresh.pending_attachments.clear();
1155 store.put(&mut fresh)?;
1156 *talk = fresh;
1157 Ok(())
1158}
1159
1160pub fn clear_pending_if_matches(
1162 talk: &mut Talk,
1163 store: &Talks,
1164 expected_text: &str,
1165 expected_attachments: &[String],
1166) -> Result<bool> {
1167 let _guard = store.guard();
1168 let mut fresh = store
1169 .get(&talk.id)
1170 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1171 if !pending_matches(&fresh, expected_text, expected_attachments) {
1172 *talk = fresh;
1173 return Ok(false);
1174 }
1175 fresh.pending.clear();
1176 fresh.pending_attachments.clear();
1177 store.put(&mut fresh)?;
1178 *talk = fresh;
1179 Ok(true)
1180}
1181
1182pub fn edit_pending_text(
1186 talk: &mut Talk,
1187 store: &Talks,
1188 text: &str,
1189 expected_text: &str,
1190 expected_attachments: &[String],
1191) -> Result<bool> {
1192 let _guard = store.guard();
1193 let mut fresh = store
1194 .get(&talk.id)
1195 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1196 if !pending_matches(&fresh, expected_text, expected_attachments) {
1197 *talk = fresh;
1198 return Ok(false);
1199 }
1200 fresh.pending = text.trim().to_owned();
1201 store.put(&mut fresh)?;
1202 *talk = fresh;
1203 Ok(true)
1204}
1205
1206fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
1207 talk.pending == expected_text
1208 && talk
1209 .pending_attachments
1210 .iter()
1211 .map(|attachment| &attachment.id)
1212 .eq(expected_attachments.iter())
1213}
1214
1215async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1222 let spec = cfg
1223 .agents
1224 .iter()
1225 .find(|a| a.id == talk.agent)
1226 .with_context(|| {
1227 format!(
1228 "talk {} was opened with agent `{}`, which is no longer in \
1229 the roster; restore it in magi.toml or start a new \
1230 conversation",
1231 talk.short(),
1232 talk.agent
1233 )
1234 })?;
1235
1236 let last_note = attachment_note(
1240 store,
1241 &talk.id,
1242 talk.turns
1243 .last()
1244 .map_or(&[][..], |t| t.attachments.as_slice()),
1245 );
1246
1247 let attachment_paths: Vec<PathBuf> = talk
1253 .turns
1254 .iter()
1255 .flat_map(|t| t.attachments.iter())
1256 .filter_map(|a| store.attachment_path(&talk.id, a))
1257 .collect();
1258
1259 let questions = crate::ask::Questions::open();
1267 let consulted = crate::consult::pending_consults(&questions, &talk.id);
1268 let consult_roots: Vec<PathBuf> = if consulted {
1269 vec![questions.root().to_path_buf()]
1270 } else {
1271 Vec::new()
1272 };
1273
1274 let artifacts = store.artifacts_of(&talk.id);
1275 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
1278 let stem = format!("turn-{}", operator_turns.max(1));
1279 let cache_dir = cfg.cache_dir();
1282
1283 let mut chain = vec![spec.clone()];
1286 if let Some(choice) = cfg.roles.chatter.as_ref()
1287 && talk.fallback
1288 {
1289 for id in choice.ids() {
1290 if id == talk.agent || chain.iter().any(|s| s.id == id) {
1291 continue;
1292 }
1293 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1294 Ok(s) => chain.push(s),
1295 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1296 }
1297 }
1298 }
1299
1300 let mut outcome = None;
1301 let mut fell_back_from: Option<String> = None;
1302 let mut first_try: Option<(String, SeatState)> = None;
1305 for (n, spec) in chain.iter().enumerate() {
1306 if n > 0 {
1307 if first_try.is_none() {
1308 first_try = Some((talk.agent.clone(), talk.seat.clone()));
1309 }
1310 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1311 fell_back_from.get_or_insert_with(|| talk.agent.clone());
1314 talk.agent = spec.id.clone();
1315 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1316 }
1317 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1318 let first_ever = talk.turns.len() <= 1;
1319 let body = if talk.seat.turns == 0 && first_ever {
1320 format!(
1321 "{}\n\n# Operator\n\n{text}{last_note}",
1322 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write)
1323 )
1324 } else if talk.seat.turns == 0 {
1325 format!(
1328 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1329 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write),
1330 transcript(talk, store)
1331 )
1332 } else if resuming {
1333 format!("{text}{last_note}")
1334 } else {
1335 format!("{}\n\n{text}{last_note}", transcript(talk, store))
1336 };
1337 let attempt_stem = if n == 0 {
1338 stem.clone()
1339 } else {
1340 format!("{stem}-{}", spec.id)
1341 };
1342 let inv = Invocation {
1343 cwd: &talk.repo,
1344 prompt: &body,
1345 timeout: turn_timeout(cfg),
1346 allow_write: cfg.talk.allow_write || consulted,
1351 sessions: cfg.graph.sessions,
1352 artifacts: &artifacts,
1353 stem: &attempt_stem,
1354 run: &talk.id,
1357 node: crate::queue::CHAT_NODE,
1358 cache_dir: cache_dir.as_deref(),
1359 attachments: &attachment_paths,
1360 writable: &consult_roots,
1361 };
1362 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1363 let advance = agent::chain_advances(&result);
1364 if n == 0 || !advance {
1365 outcome = Some(result);
1366 } else {
1367 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1369 }
1370 if !advance {
1371 break;
1372 }
1373 }
1374 if outcome.as_ref().is_some_and(agent::chain_advances) {
1375 if let Some((id, seat)) = first_try {
1378 talk.agent = id;
1379 talk.seat = seat;
1380 fell_back_from = None;
1381 }
1382 }
1383 let outcome = outcome.expect("a chain holds at least one agent");
1384 let note = |why: String| Turn {
1385 who: Who::Agent,
1386 body: format!("{MAGI_NOTE}{why}"),
1387 at: Timestamp::now(),
1388 attachments: Vec::new(),
1389 usage: None,
1390 };
1391 let (reply, failure) = match outcome {
1392 Err(e) => (
1393 note(format!("could not run agent `{}`: {e}", talk.agent)),
1394 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1395 ),
1396 Ok(out) if out.quota_exhausted() => {
1397 let reset = out
1398 .quota
1399 .as_ref()
1400 .and_then(|q| q.reset.clone())
1401 .map_or_else(String::new, |r| format!(" (resets {r})"));
1402 let why = format!(
1403 "agent `{}` is out of quota{reset}; your message is saved, so \
1404 say it again when the window reopens",
1405 talk.agent
1406 );
1407 (note(why.clone()), Some(why))
1408 }
1409 Ok(out) if out.timed_out => {
1410 let why = format!(
1411 "agent `{}` did not answer within {}s; your message is saved",
1412 talk.agent,
1413 turn_timeout(cfg).as_secs()
1414 );
1415 (note(why.clone()), Some(why))
1416 }
1417 Ok(out) if !out.usable() => {
1418 let why = format!(
1419 "agent `{}` produced no answer (exit {}); your message is saved",
1420 talk.agent,
1421 out.exit_code
1422 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1423 );
1424 (note(why.clone()), Some(why))
1425 }
1426 Ok(out) => (
1427 Turn {
1428 who: Who::Agent,
1429 body: out.text.trim().to_owned(),
1430 at: Timestamp::now(),
1431 attachments: Vec::new(),
1432 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1435 context_tokens,
1436 agent: talk.agent.clone(),
1437 model: cfg
1438 .agents
1439 .iter()
1440 .find(|a| a.id == talk.agent)
1441 .and_then(|a| a.model.clone()),
1442 }),
1443 },
1444 None,
1445 ),
1446 };
1447
1448 let _guard = store.guard();
1460 let Ok(fresh) = store.get(&talk.id) else {
1466 return Ok(());
1467 };
1468 talk.status = fresh.status;
1469 talk.pending = fresh.pending;
1473 talk.pending_attachments = fresh.pending_attachments;
1474 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1475 talk.turns.push(note(format!(
1478 "agent changed from {from} to {} (fallback)",
1479 talk.agent
1480 )));
1481 }
1482 talk.turns.push(reply);
1483 if let Err(put_err) = store.put(talk) {
1484 let lost = talk.turns.pop().expect("just pushed above");
1493 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1494 let why = match &stash {
1495 Ok(path) => format!(
1496 "agent `{}` answered, but the reply could not be saved to \
1497 this conversation ({put_err:#}); the raw text was kept at \
1498 {} - your message is saved, ask again",
1499 talk.agent,
1500 path.display()
1501 ),
1502 Err(stash_err) => format!(
1503 "agent `{}` answered, but the reply could not be saved to \
1504 this conversation ({put_err:#}), and it could not be kept \
1505 anywhere else either ({stash_err:#}); your message is \
1506 saved, ask again",
1507 talk.agent
1508 ),
1509 };
1510 talk.turns.push(note(why.clone()));
1511 return match store.put(talk) {
1518 Ok(()) => bail!("{why}"),
1519 Err(note_err) => {
1520 talk.turns.pop();
1540 Err(note_err).context(why)
1541 }
1542 };
1543 }
1544
1545 match failure {
1546 Some(why) => bail!("{why}"),
1547 None => Ok(()),
1548 }
1549}
1550
1551fn transcript(talk: &Talk, store: &Talks) -> String {
1554 let mut out = String::from(
1555 "This conversation cannot resume on the CLI's side, so here is \
1556 everything said so far; answer only the last message.\n",
1557 );
1558 for t in &talk.turns {
1559 let who = match t.who {
1560 Who::Operator => "operator",
1561 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1562 Who::Agent => "you",
1563 };
1564 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1565 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1566 }
1567 out
1568}
1569
1570fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1575 if attachments.is_empty() {
1576 return String::new();
1577 }
1578 let mut out = String::from(
1579 "\n\nThe operator attached the image(s) below to this message. Open \
1580 and look at each one before you answer.\n",
1581 );
1582 for att in attachments {
1583 if let Some(path) = store.attachment_path(talk_id, att) {
1584 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1585 }
1586 }
1587 out.push('\n');
1588 out
1589}
1590
1591pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
1613 let write_policy = if allow_write {
1614 "Write access is enabled for this conversation (`allow_write = \
1615 true`), so you may write files - but only a small, \
1616 already-decided edit the operator names outright in this \
1617 conversation, not an implementation. This is a permission on the \
1618 conversation as a whole, not a property of whichever repository \
1619 it happened to start in: if the operator names a different \
1620 repository for that small edit, the policy allows it there too. \
1621 Your own tool may still confine writes to the repository this \
1622 conversation started in regardless - if a write elsewhere is \
1623 refused, say so plainly rather than working around it. Once you \
1624 have made an edit, say plainly what you edited. Anything bigger, \
1625 or anything still open-ended, still goes through the queue below \
1626 rather than being done here."
1627 } else {
1628 "Do not write files. Implementing a change is not this \
1629 conversation's job; a separate, blind competition of agents does \
1630 that, and a repository this conversation has already edited would \
1631 make their diffs unjudgeable."
1632 };
1633 let mut out = format!(
1634 "You are magi's standing conversation partner for its operator, who \
1635 usually has this open on a phone. Keep replies short: no preamble, \
1636 no restating what they just said.\n\n\
1637 # Repository\n\n{repo}\n\n\
1638 You may look around: read files, run shell commands, search history, \
1639 run tests - whatever answers the question. {write_policy}\n\n\
1640 A short, command-shaped message (\"list\", \"info <id>\", \"show \
1641 3cbf\") is almost always the operator asking you to look something \
1642 up, not an instruction to file - answer it yourself with `magi \
1643 list`, `magi show <id>`, `magi task list`, or the like, the same way \
1644 you would answer any other question in this conversation.\n\n\
1645 # When the operator wants something done\n\n\
1646 Run:\n\n\
1647 magi task add --solo --repo {repo} <instruction>\n\n\
1648 and tell the operator the task id it prints, so they can follow it \
1649 from the Queue. If it refuses with a duplicate warning (the \
1650 instruction names a branch, commit or pull request that an \
1651 unfinished task, run or PR already owns), do not repeat it with \
1652 --force yourself: tell the operator what it matched and let them \
1653 decide. Write <instruction> so that an implementer who has \
1654 never seen this conversation can act on it alone - it is everything \
1655 they get. Use --solo: it runs the task through one implementer \
1656 straight into review instead of the usual multi-agent competition, \
1657 which is the right shape for a change this conversation has already \
1658 settled, rather than one still worth several independent takes.\n\n\
1659 If the operator asks for something in a different repository, \
1660 --repo does not have to be a full path: --repo owner/repo (or just \
1661 repo, when that is unambiguous) is resolved against local checkouts \
1662 the same way `magi repos` lists them. If the command fails because \
1663 nothing matches or more than one checkout shares that name, ask the \
1664 operator which repository they mean (or run `magi repos` yourself \
1665 to see the candidates) rather than guessing.\n\n\
1666 The current state of the code is whatever origin/main holds, not \
1667 whatever a working tree shows: a primary checkout often lags \
1668 upstream, sits on a detached HEAD and carries uncommitted changes. \
1669 Before answering about code, run `git fetch origin` in that \
1670 repository if it is cheap, then read through \
1671 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
1672 If the working tree differs, say so; if the fetch fails, say that \
1673 too, so the operator knows the answer may be stale.\n\n\
1674 If the operator attached an image (a screenshot, say) that the task \
1675 is about, pass it with `--attach <path>`, using the absolute path \
1676 the turn's attachment note gives; repeat the flag for several. \
1677 `magi task add --solo --attach <path> <instruction>` copies the \
1678 file into the task, so the implementer receives it. Do not paste the \
1679 path into <instruction> instead: deleting this conversation deletes \
1680 its attachments, and then that path reaches no one.\n",
1681 repo = repo.display(),
1682 );
1683 out.push_str(&language_note(language));
1684 out
1685}
1686
1687fn language_note(language: &str) -> String {
1690 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
1691 String::new()
1692 } else {
1693 format!("\nHold this conversation in {language}.\n")
1694 }
1695}
1696
1697pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
1704 let mut tasks: Vec<Task> = queue
1705 .list()
1706 .into_iter()
1707 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
1708 .collect();
1709 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
1710 tasks
1711}
1712
1713fn read_path(path: &Path) -> Result<Talk> {
1714 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1715 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
1716}
1717
1718const PUT_RETRIES: u32 = 5;
1721
1722fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1734 let mut last_err = None;
1735 for attempt in 0..PUT_RETRIES {
1736 if attempt > 0 {
1737 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
1738 }
1739 match try_write_atomic(tmp, path, body) {
1740 Ok(()) => return Ok(()),
1741 Err(e) => last_err = Some(e),
1742 }
1743 }
1744 Err(last_err.expect("the loop above always runs at least once"))
1745}
1746
1747fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1748 #[cfg(test)]
1749 if failpoint::take_forced_put_failure() {
1750 bail!("simulated write failure (test)");
1751 }
1752 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
1753 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
1754 Ok(())
1755}
1756
1757fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
1762 let dir = store.artifacts_of(id);
1763 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1764 let path = dir.join(format!("{stem}-lost.txt"));
1765 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
1766 Ok(path)
1767}
1768
1769#[cfg(test)]
1776mod failpoint {
1777 use std::cell::Cell;
1778
1779 thread_local! {
1780 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
1781 }
1782
1783 pub(super) fn force_put_failures(count: u32) {
1786 FORCE_PUT_FAILURES.with(|c| c.set(count));
1787 }
1788
1789 pub(super) fn take_forced_put_failure() -> bool {
1792 FORCE_PUT_FAILURES.with(|c| {
1793 let n = c.get();
1794 if n == 0 {
1795 false
1796 } else {
1797 c.set(n - 1);
1798 true
1799 }
1800 })
1801 }
1802}
1803
1804fn short(id: &str) -> &str {
1805 id.split('-').next_back().unwrap_or(id)
1806}
1807
1808fn new_id() -> String {
1809 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1810 let seed = crate::rng::entropy();
1811 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1812}
1813
1814fn attachment_ext(mime: &str) -> Option<&'static str> {
1819 match mime {
1820 "image/png" => Some("png"),
1821 "image/jpeg" => Some("jpg"),
1822 "image/gif" => Some("gif"),
1823 "image/webp" => Some("webp"),
1824 _ => None,
1825 }
1826}
1827
1828pub fn valid_attachment_id(id: &str) -> bool {
1833 id.len() == 32
1834 && id
1835 .bytes()
1836 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
1837}
1838
1839fn new_attachment_id() -> String {
1843 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
1844 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
1845}
1846
1847#[cfg(test)]
1848mod tests {
1849 #[test]
1850 fn the_briefing_points_at_origin_main_not_the_working_tree() {
1851 let b = briefing(Path::new("/r"), "en", false);
1852 assert!(b.contains("origin/main"));
1853 assert!(b.contains("git show origin/main:"));
1854 }
1855 use std::collections::BTreeMap;
1856
1857 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
1858 use crate::queue::{Queue, Source, Task};
1859
1860 use super::*;
1861
1862 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
1863 AgentSpec {
1864 id: id.to_owned(),
1865 kind: AgentKind::Command,
1866 model: model.map(str::to_owned),
1867 command: Vec::new(),
1868 extra_args: Vec::new(),
1869 env: BTreeMap::new(),
1870 prompt_delivery: None,
1871 }
1872 }
1873
1874 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
1875 Talk {
1876 schema: SCHEMA,
1877 id: "20260904-014455-ab12".to_owned(),
1878 repo: PathBuf::from("."),
1879 agent: agent.to_owned(),
1880 status: TalkStatus::Open,
1881 turns,
1882 pending: String::new(),
1883 pending_attachments: Vec::new(),
1884 fallback: false,
1885 created_at: Timestamp::now(),
1886 updated_at: Timestamp::now(),
1887 seat: SeatState::new(SEAT, agent, 1),
1888 }
1889 }
1890
1891 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
1892 Turn {
1893 who: Who::Agent,
1894 body: body.to_owned(),
1895 at: Timestamp::now(),
1896 attachments: Vec::new(),
1897 usage: usage.map(|(t, a, m)| TurnUsage {
1898 context_tokens: t,
1899 agent: a.to_owned(),
1900 model: m.map(str::to_owned),
1901 }),
1902 }
1903 }
1904
1905 fn ctx_config(windows: &[(&str, u64)]) -> Config {
1906 Config {
1907 agents: vec![
1908 ctx_agent("small", Some("small-model")),
1909 ctx_agent("big", Some("big-model")),
1910 ctx_agent("plain", None),
1911 ],
1912 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
1913 ..Config::default()
1914 }
1915 }
1916
1917 #[test]
1918 fn context_usage_computes_percent_and_warns_at_eighty() {
1919 let cfg = ctx_config(&[("small-model", 1000)]);
1920 let at = |tokens| {
1921 let t = ctx_talk(
1922 "small",
1923 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
1924 );
1925 context_usage(&t, Some(&cfg))
1926 };
1927 let u = at(799);
1928 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
1929 let u = at(800);
1930 assert_eq!((u.percent, u.warn), (Some(80), true));
1931 let u = at(1500);
1932 assert_eq!((u.percent, u.warn), (Some(150), true));
1933 assert!(!u.since_switch);
1934 }
1935
1936 #[test]
1937 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
1938 let cfg = ctx_config(&[("small-model", 1000)]);
1939 let t = ctx_talk(
1940 "small",
1941 vec![
1942 reply("old", Some((900, "small", Some("small-model")))),
1943 reply("new", None),
1944 ],
1945 );
1946 let u = context_usage(&t, Some(&cfg));
1947 assert!(u.estimated);
1949 assert_ne!(u.tokens, Some(900));
1950 assert!(u.tokens.is_some());
1951 let t = ctx_talk(
1953 "small",
1954 vec![
1955 reply("old", Some((900, "small", Some("small-model")))),
1956 reply("magi: could not run agent", None),
1957 ],
1958 );
1959 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
1960 assert_eq!(
1961 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
1962 None
1963 );
1964 }
1965
1966 #[test]
1967 fn estimate_counts_chars_both_sides_and_standing_prompt() {
1968 let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
1969 assert_eq!(estimate_context_tokens(&t, 0), Some(2)); let op = Turn {
1971 who: Who::Operator,
1972 ..reply("abcdefg", None)
1973 };
1974 t.turns.push(op);
1975 assert_eq!(estimate_context_tokens(&t, 0), Some(4));
1976 assert!(
1977 estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
1978 );
1979 let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
1981 assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
1982 let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
1984 assert_eq!(estimate_context_tokens(¬e, 1000), None);
1985 assert_eq!(
1986 estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
1987 None
1988 );
1989 }
1990
1991 #[test]
1992 fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
1993 let cfg = ctx_config(&[("small-model", 1000)]);
1994 let t = ctx_talk(
1995 "small",
1996 vec![reply(
1997 &"x".repeat(5000),
1998 Some((10, "small", Some("small-model"))),
1999 )],
2000 );
2001 let u = context_usage(&t, Some(&cfg));
2002 assert_eq!((u.tokens, u.estimated), (Some(10), false));
2003 let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
2004 let u = context_usage(&t, Some(&cfg));
2005 assert!(u.estimated && !u.since_switch);
2006 assert_eq!(u.window, Some(1000));
2007 assert!(u.warn && u.percent.unwrap() >= 80);
2008 let t = ctx_talk("small", vec![reply("hi", None)]);
2009 let u = context_usage(&t, Some(&cfg));
2010 assert!(u.estimated && u.percent.is_some());
2011 }
2012
2013 #[test]
2014 fn context_usage_without_a_window_shows_tokens_only() {
2015 let cfg = ctx_config(&[]);
2016 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
2018 let u = context_usage(&t, Some(&cfg));
2019 assert_eq!(
2020 (u.tokens, u.window, u.percent, u.warn),
2021 (Some(5000), None, None, false)
2022 );
2023 let t = ctx_talk(
2024 "small",
2025 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
2026 );
2027 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
2028 assert_eq!(context_usage(&t, None).window, None);
2030 }
2031
2032 #[test]
2033 fn context_usage_switching_model_changes_the_denominator() {
2034 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
2035 let used = reply("hi", Some((900, "small", Some("small-model"))));
2036 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
2037 assert_eq!(
2038 (before.percent, before.warn, before.since_switch),
2039 (Some(90), true, false)
2040 );
2041 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
2044 assert_eq!(after.window, Some(10_000));
2045 assert_eq!(
2046 (after.percent, after.warn, after.since_switch),
2047 (Some(9), false, true)
2048 );
2049 assert_eq!(after.model.as_deref(), Some("big-model"));
2050 }
2051
2052 #[test]
2053 fn a_turn_recorded_before_usage_existed_still_reads() {
2054 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
2055 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
2056 assert!(turn.usage.is_none());
2057 let json = serde_json::to_string(&turn).expect("serialize");
2058 assert!(
2059 !json.contains("usage"),
2060 "absent usage is not written: {json}"
2061 );
2062 }
2063
2064 fn store() -> (tempfile::TempDir, Talks) {
2066 let tmp = tempfile::tempdir().expect("tempdir");
2067 let talks = Talks::at(tmp.path().join("talks"));
2068 (tmp, talks)
2069 }
2070
2071 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
2075 let path = dir.join("mock-talk-agent.sh");
2076 std::fs::write(&path, script).expect("write mock");
2077 AgentSpec {
2078 id: "mock".to_owned(),
2079 kind: AgentKind::Command,
2080 model: None,
2081 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2082 extra_args: Vec::new(),
2083 env,
2084 prompt_delivery: None,
2085 }
2086 }
2087
2088 fn config(spec: AgentSpec) -> Config {
2089 Config {
2090 agents: vec![spec],
2091 graph: Graph {
2092 language: "en".to_owned(),
2093 ..Graph::default()
2094 },
2095 ..Config::default()
2096 }
2097 }
2098
2099 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
2101
2102 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
2104
2105 const ECHO: &str = "#!/bin/sh\ncat\n";
2108
2109 fn env(reply: &str) -> BTreeMap<String, String> {
2110 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
2111 }
2112
2113 #[test]
2114 fn the_frozen_json_field_names_round_trip_through_disk() {
2115 let (tmp, talks) = store();
2116 let mut talk = Talk {
2117 schema: SCHEMA,
2118 id: "20260904-014455-ab12".to_owned(),
2119 repo: tmp.path().to_owned(),
2120 agent: "sonnet".to_owned(),
2121 status: TalkStatus::Open,
2122 turns: Vec::new(),
2123 pending: String::new(),
2124 pending_attachments: Vec::new(),
2125 fallback: false,
2126 created_at: Timestamp::now(),
2127 updated_at: Timestamp::now(),
2128 seat: SeatState::new(SEAT, "sonnet", 7),
2129 };
2130 talks.put(&mut talk).expect("put");
2131
2132 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
2133 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
2134 for field in [
2135 "schema",
2136 "id",
2137 "repo",
2138 "agent",
2139 "status",
2140 "turns",
2141 "created_at",
2142 "updated_at",
2143 ] {
2144 assert!(v.get(field).is_some(), "missing field `{field}`");
2145 }
2146 assert_eq!(v["schema"], 1);
2147 assert_eq!(v["status"], "open");
2148
2149 let back = talks.get(&talk.id).expect("get");
2150 assert_eq!(back.id, talk.id);
2151 assert_eq!(back.status, TalkStatus::Open);
2152 }
2153
2154 #[test]
2155 fn opening_a_talk_takes_no_agent_turn() {
2156 let (tmp, talks) = store();
2157 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2161 let cfg = config(spec);
2162
2163 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2164 assert_eq!(talk.status, TalkStatus::Open);
2165 assert!(talk.turns.is_empty(), "nothing has been said yet");
2166
2167 let on_disk = talks.get(&talk.id).expect("get");
2168 assert_eq!(on_disk.turns.len(), 0);
2169 }
2170
2171 #[test]
2179 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
2180 let (tmp, talks) = store();
2181 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2182 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2183 chatter_spec.id = "chatter-mock".to_owned();
2184
2185 let mut cfg = Config {
2186 agents: vec![first_spec.clone(), chatter_spec.clone()],
2187 graph: Graph {
2188 language: "en".to_owned(),
2189 ..Graph::default()
2190 },
2191 ..Config::default()
2192 };
2193 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
2194
2195 let talk =
2196 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
2197 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
2198
2199 cfg.roles.chatter = None;
2200 let fallback =
2201 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
2202 assert_eq!(
2203 fallback.agent, first_spec.id,
2204 "unset chatter must fall back to agent::pick's own default order"
2205 );
2206 }
2207
2208 #[test]
2211 fn a_talk_recorded_without_attachments_still_reads() {
2212 let (tmp, talks) = store();
2213 let path = talks.path_of("20260904-014455-ab12");
2214 std::fs::create_dir_all(talks.root()).expect("talks dir");
2215 std::fs::write(
2216 &path,
2217 serde_json::json!({
2218 "schema": 1,
2219 "id": "20260904-014455-ab12",
2220 "repo": tmp.path(),
2221 "agent": "sonnet",
2222 "status": "open",
2223 "turns": [
2224 { "who": "operator", "body": "still there?",
2225 "at": Timestamp::now().to_string() },
2226 ],
2227 "created_at": Timestamp::now().to_string(),
2228 "updated_at": Timestamp::now().to_string(),
2229 "seat": SeatState::new(SEAT, "sonnet", 7),
2230 })
2231 .to_string(),
2232 )
2233 .expect("write pre-attachments talk");
2234
2235 let talk = talks.get("20260904-014455-ab12").expect("must still read");
2236 assert!(talk.turns[0].attachments.is_empty());
2237 }
2238
2239 fn lease_store() -> (tempfile::TempDir, Talks, Talks) {
2240 let tmp = tempfile::TempDir::new().expect("tmp");
2241 let root = tmp.path().join("talks");
2242 (tmp, Talks::at(root.clone()), Talks::at(root))
2243 }
2244
2245 #[test]
2246 fn two_starters_on_one_talk_one_wins_and_the_other_is_refused() {
2247 let (_tmp, a, b) = lease_store();
2248 let won = a.claim_turn("t1").expect("claim").expect("first wins");
2249 assert!(
2250 b.claim_turn("t1").expect("claim").is_none(),
2251 "second is refused"
2252 );
2253 assert!(b.turn_held("t1"));
2254 assert!(
2255 b.claim_turn("t2").expect("claim").is_some(),
2256 "other talks are free"
2257 );
2258 drop(won);
2259 }
2260
2261 #[test]
2262 fn a_stale_lease_is_taken_over_and_the_old_guard_cannot_release_it() {
2263 let (_tmp, a, b) = lease_store();
2264 let old = a.claim_turn("t1").expect("claim").expect("held");
2265 let later = Timestamp::now()
2266 .checked_add(jiff::SignedDuration::from_secs(
2267 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2268 ))
2269 .expect("later");
2270 let new = b
2271 .claim_turn_at("t1", later)
2272 .expect("claim")
2273 .expect("a stale lease is taken over");
2274 drop(old);
2275 assert!(a.turn_held("t1"), "the old guard left the new lease alone");
2276 assert!(new.beat().expect("beat"), "the new owner still beats");
2277 drop(new);
2278 assert!(!a.turn_held("t1"));
2279 }
2280
2281 #[tokio::test]
2282 async fn a_turn_whose_lease_was_taken_over_is_stopped() {
2283 let (_tmp, a, b) = lease_store();
2284 let old = a.claim_turn("t1").expect("claim").expect("held");
2285 let later = Timestamp::now()
2286 .checked_add(jiff::SignedDuration::from_secs(
2287 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2288 ))
2289 .expect("later");
2290 let _new = b
2291 .claim_turn_at("t1", later)
2292 .expect("claim")
2293 .expect("taken over");
2294 let out = old
2295 .beating_every(Duration::from_millis(10), std::future::pending::<()>())
2296 .await;
2297 assert!(out.is_err(), "the displaced turn must stop, not run on");
2298 }
2299
2300 #[tokio::test]
2301 async fn a_turn_that_finishes_is_returned_and_keeps_its_lease_beating() {
2302 let (_tmp, a, _b) = lease_store();
2303 let lease = a.claim_turn("t1").expect("claim").expect("held");
2304 let out = lease
2305 .beating_every(Duration::from_millis(5), async {
2306 tokio::time::sleep(Duration::from_millis(40)).await;
2307 7
2308 })
2309 .await
2310 .expect("still ours");
2311 assert_eq!(out, 7);
2312 assert!(a.turn_held("t1"));
2313 }
2314
2315 #[test]
2316 fn an_unreadable_lease_counts_as_stale() {
2317 let (_tmp, a, b) = lease_store();
2318 std::fs::create_dir_all(a.root()).expect("dir");
2319 std::fs::write(a.turn_path("t1"), "not json").expect("write");
2320 assert!(!a.turn_held("t1"));
2321 assert!(b.claim_turn("t1").expect("claim").is_some());
2322 }
2323
2324 #[test]
2325 fn a_lease_is_released_when_the_turn_ends_or_fails() {
2326 let (_tmp, a, b) = lease_store();
2327 let lease = a.claim_turn("t1").expect("claim").expect("held");
2328 let failed: Result<()> = (|| {
2329 let _held = &lease;
2330 bail!("turn failed")
2331 })();
2332 assert!(failed.is_err());
2333 assert!(
2334 b.claim_turn("t1").expect("claim").is_none(),
2335 "held mid-turn"
2336 );
2337 drop(lease);
2338 assert!(
2339 b.claim_turn("t1").expect("claim").is_some(),
2340 "free after the turn"
2341 );
2342 }
2343
2344 #[test]
2345 fn concurrent_takeovers_of_a_stale_lease_have_one_winner() {
2346 let (_tmp, a, _b) = lease_store();
2347 drop(a.claim_turn("t1").expect("claim").expect("held"));
2348 std::fs::write(
2349 a.turn_path("t1"),
2350 serde_json::to_string(&TurnRecord {
2351 token: "gone".into(),
2352 pid: 1,
2353 beat_at: Timestamp::from_second(1).expect("ts"),
2354 })
2355 .expect("json"),
2356 )
2357 .expect("write");
2358 let wins: Vec<_> = std::thread::scope(|sc| {
2359 let hs: Vec<_> = (0..8)
2360 .map(|_| {
2361 let s = a.clone();
2362 sc.spawn(move || s.claim_turn("t1").expect("claim"))
2363 })
2364 .collect();
2365 hs.into_iter().map(|h| h.join().expect("join")).collect()
2366 });
2367 assert_eq!(wins.iter().filter(|w| w.is_some()).count(), 1);
2368 }
2369
2370 #[test]
2371 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
2372 let (tmp, talks) = store();
2373 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2374 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2375
2376 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
2377 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
2378 let saved = talks.get(&talk.id).expect("reload queued talk");
2379 assert_eq!(saved.pending, "first\n\nsecond");
2380 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
2381
2382 let drained = drain(&mut talk, &talks).expect("drain");
2383 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
2384 let saved = talks.get(&talk.id).expect("reload drained talk");
2385 assert!(saved.pending.is_empty());
2386 assert_eq!(saved.turns.len(), 1);
2387 assert_eq!(saved.turns[0].body, "first\n\nsecond");
2388 }
2389
2390 #[test]
2391 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
2392 let (tmp, talks) = store();
2393 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2394 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2395 let attachment = Attachment {
2396 id: "a".repeat(32),
2397 name: "shot.png".to_owned(),
2398 mime: "image/png".to_owned(),
2399 bytes: 3,
2400 };
2401
2402 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
2403 assert!(
2404 edit_pending_text(
2405 &mut talk,
2406 &talks,
2407 "corrected",
2408 "first",
2409 std::slice::from_ref(&attachment.id),
2410 )
2411 .expect("edit")
2412 );
2413 let saved = talks.get(&talk.id).expect("reload edited draft");
2414 assert_eq!(saved.pending, "corrected");
2415 assert_eq!(saved.pending_attachments, vec![attachment]);
2416
2417 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
2418 assert!(
2419 !edit_pending_text(
2420 &mut talk,
2421 &talks,
2422 "stale edit",
2423 "corrected",
2424 &["a".repeat(32)],
2425 )
2426 .expect("stale edit is a conflict")
2427 );
2428 assert_eq!(
2429 talks.get(&talk.id).expect("reload after conflict").pending,
2430 "corrected\n\nlater"
2431 );
2432 assert!(
2433 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
2434 .expect("stale clear is a conflict")
2435 );
2436 assert_eq!(
2437 talks
2438 .get(&talk.id)
2439 .expect("reload after stale clear")
2440 .pending,
2441 "corrected\n\nlater"
2442 );
2443 }
2444
2445 #[tokio::test]
2446 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
2447 let (tmp, talks) = store();
2448 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
2449 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
2450 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2451 let id = running.id.clone();
2452 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
2453
2454 let response_talks = talks.clone();
2455 let response_cfg = cfg.clone();
2456 let reply = tokio::spawn(async move {
2457 respond(&mut running, &response_talks, &response_cfg, &first).await
2458 });
2459 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
2460
2461 let mut queued = talks.get(&id).expect("queued handle");
2462 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
2463 reply.await.expect("join").expect("reply");
2464
2465 let saved = talks.get(&id).expect("reload");
2466 assert_eq!(saved.pending, "next");
2467 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
2468 }
2469
2470 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
2473 let calls = dir.join(format!("{id}.calls"));
2474 let script = format!(
2475 "#!/bin/sh\necho x >> '{}'\n{body}\n",
2476 calls.to_string_lossy()
2477 );
2478 let path = dir.join(format!("mock-{id}.sh"));
2479 std::fs::write(&path, script).expect("write mock");
2480 AgentSpec {
2481 id: id.to_owned(),
2482 kind: AgentKind::Command,
2483 model: None,
2484 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2485 extra_args: Vec::new(),
2486 env: BTreeMap::new(),
2487 prompt_delivery: None,
2488 }
2489 }
2490
2491 fn calls(dir: &Path, id: &str) -> usize {
2492 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
2493 }
2494
2495 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
2496 let mut cfg = config(specs[0].clone());
2497 cfg.agents = specs;
2498 cfg.roles.chatter = Some(AgentChoice::Chain(
2499 ids.iter().map(|s| (*s).to_owned()).collect(),
2500 ));
2501 cfg
2502 }
2503
2504 #[tokio::test]
2505 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
2506 let (tmp, talks) = store();
2507 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2508 let b = counting_agent(tmp.path(), "b", "cat");
2509 let cfg = chain_config(vec![a, b], &["a", "b"]);
2510 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2511 assert_eq!(talk.agent, "a");
2512
2513 say(&mut talk, &talks, &cfg, "hello there", Vec::new())
2514 .await
2515 .expect("turn");
2516 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
2517 assert_eq!(calls(tmp.path(), "b"), 1);
2518 assert_eq!(talk.agent, "b", "the switch persists");
2519 assert!(talks.get(&talk.id).unwrap().agent == "b");
2520 let reply = talk.turns.last().unwrap();
2521 assert!(reply.body.contains("hello there"));
2522 assert!(
2523 reply.body.contains("magi task add --solo"),
2524 "a fresh seat gets the full briefing"
2525 );
2526 assert!(
2527 talk.turns
2528 .iter()
2529 .any(|t| t.body.contains("agent changed from a to b")),
2530 "the switch is noted"
2531 );
2532 }
2533
2534 #[tokio::test]
2535 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
2536 let (tmp, talks) = store();
2537 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2538 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
2539 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
2540 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2541
2542 let err = say(&mut talk, &talks, &cfg, "hi", Vec::new())
2543 .await
2544 .expect_err("every agent failed");
2545 assert!(err.to_string().contains("`a`"), "{err:#}");
2546 assert_eq!(calls(tmp.path(), "a"), 1);
2547 assert_eq!(calls(tmp.path(), "b"), 1);
2548 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
2549 }
2550
2551 #[test]
2552 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
2553 let (tmp, talks) = store();
2554 let b = counting_agent(tmp.path(), "b", "cat");
2555 let cfg = chain_config(vec![b], &["ghost", "b"]);
2556 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2557 assert_eq!(talk.agent, "b");
2558 }
2559
2560 #[tokio::test]
2561 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
2562 let (tmp, talks) = store();
2563 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2564 let b = counting_agent(tmp.path(), "b", "cat");
2565 let cfg = chain_config(vec![a, b], &["a", "b"]);
2566 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
2567 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2568 .await
2569 .expect_err("a alone, and it fails");
2570 assert_eq!(calls(tmp.path(), "b"), 0);
2571 assert_eq!(talk.agent, "a");
2572 }
2573
2574 #[tokio::test]
2575 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
2576 let (tmp, talks) = store();
2577 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2578 let b = counting_agent(tmp.path(), "b", "cat");
2579 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
2580 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
2581 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
2582 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2583 .await
2584 .expect_err("c alone, and it fails");
2585 assert_eq!(calls(tmp.path(), "b"), 0);
2586 }
2587
2588 #[tokio::test]
2589 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
2590 let (tmp, talks) = store();
2591 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2592 let cfg = config(spec);
2593 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2594
2595 say(
2596 &mut talk,
2597 &talks,
2598 &cfg,
2599 "what does the queue module do?",
2600 Vec::new(),
2601 )
2602 .await
2603 .expect("first turn");
2604 let first_prompt = &talk.turns[1].body;
2605 assert!(first_prompt.contains("magi task add --solo"));
2606 assert!(first_prompt.contains("what does the queue module do?"));
2607
2608 say(&mut talk, &talks, &cfg, "and how is it locked?", Vec::new())
2609 .await
2610 .expect("second turn");
2611 let second_prompt = &talk.turns[3].body;
2612 assert!(
2613 !second_prompt.contains("magi task add --solo"),
2614 "the briefing is sent once, not on every turn: {second_prompt}"
2615 );
2616 assert!(second_prompt.contains("and how is it locked?"));
2617 }
2618
2619 #[tokio::test]
2620 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
2621 let (tmp, talks) = store();
2622 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2623 let mut b = a.clone();
2624 b.id = "other".to_owned();
2625 let mut cfg = config(a.clone());
2626 cfg.agents.push(b.clone());
2627 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
2628 say(&mut talk, &talks, &cfg, "remember the walrus", Vec::new())
2629 .await
2630 .expect("first turn");
2631 let old_session = talk.seat.claude_session.clone();
2632 assert_eq!(talk.seat.turns, 1);
2633
2634 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
2635 assert_eq!(talk.agent, "other");
2636 assert_eq!(talk.seat.turns, 0);
2637 assert_eq!(talk.seat.agent, "other");
2638 assert_ne!(talk.seat.claude_session, old_session);
2639 let note = talk.turns.last().expect("note");
2640 assert_eq!(note.who, Who::Agent);
2641 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2642 assert!(note.body.contains("changed from"), "{}", note.body);
2643 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
2644
2645 let before = talk.turns.len();
2646 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
2647 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
2648
2649 say(&mut talk, &talks, &cfg, "what did I say?", Vec::new())
2650 .await
2651 .expect("turn after switch");
2652 let prompt = &talk.turns.last().expect("reply").body;
2653 assert!(prompt.contains("remember the walrus"), "{prompt}");
2654 assert!(prompt.contains("## magi"), "{prompt}");
2655 assert!(prompt.contains("what did I say?"), "{prompt}");
2656 }
2657
2658 #[tokio::test]
2659 async fn say_appends_the_operator_turn_then_the_agent_turn() {
2660 let (tmp, talks) = store();
2661 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2662 let cfg = config(spec);
2663 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2664
2665 say(
2666 &mut talk,
2667 &talks,
2668 &cfg,
2669 "can I rename this function?",
2670 Vec::new(),
2671 )
2672 .await
2673 .expect("say");
2674
2675 assert_eq!(talk.turns.len(), 2);
2676 assert_eq!(talk.turns[0].who, Who::Operator);
2677 assert_eq!(talk.turns[0].body, "can I rename this function?");
2678 assert_eq!(talk.turns[1].who, Who::Agent);
2679 assert_eq!(talk.turns[1].body, "go ahead");
2680 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
2681 }
2682
2683 #[tokio::test]
2684 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
2685 let (tmp, talks) = store();
2686 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2687 let cfg = config(spec);
2688 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2689
2690 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2691 .await
2692 .expect_err("a turn with no answer is an error");
2693 assert!(err.to_string().contains("no answer"), "{err}");
2694
2695 let on_disk = talks.get(&talk.id).expect("get");
2696 assert_eq!(on_disk.turns.len(), 2);
2697 assert_eq!(on_disk.turns[0].body, "check the tests");
2698 let note = &on_disk.turns[1];
2699 assert_eq!(note.who, Who::Agent);
2700 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2701 assert!(note.body.contains("your message is saved"));
2702 }
2703
2704 #[tokio::test]
2710 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
2711 let (tmp, talks) = store();
2712 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2713 let cfg = config(spec);
2714 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2715
2716 let text =
2717 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
2718 failpoint::force_put_failures(PUT_RETRIES - 1);
2721 respond(&mut talk, &talks, &cfg, &text)
2722 .await
2723 .expect("respond must survive a write failure its own retries can outlast");
2724
2725 assert_eq!(talk.turns.len(), 2);
2726 assert_eq!(talk.turns[1].who, Who::Agent);
2727 assert_eq!(talk.turns[1].body, "go ahead");
2728 let on_disk = talks.get(&talk.id).expect("get");
2729 assert_eq!(
2730 on_disk.turns, talk.turns,
2731 "the reply must reach disk despite the early write failures"
2732 );
2733 }
2734
2735 #[tokio::test]
2741 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
2742 let (tmp, talks) = store();
2743 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2744 let cfg = config(spec);
2745 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2746
2747 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2748 failpoint::force_put_failures(PUT_RETRIES);
2753 let err = respond(&mut talk, &talks, &cfg, &text)
2754 .await
2755 .expect_err("a reply that cannot be saved must be reported, not swallowed");
2756 assert!(err.to_string().contains("could not be saved"), "{err}");
2757
2758 let on_disk = talks.get(&talk.id).expect("get");
2759 assert_eq!(
2760 on_disk.turns.len(),
2761 2,
2762 "the operator turn plus a visible note"
2763 );
2764 assert_eq!(on_disk.turns[0].body, "check the tests");
2765 let note = &on_disk.turns[1];
2766 assert_eq!(note.who, Who::Agent);
2767 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2768 assert!(
2769 note.body.contains("could not be saved"),
2770 "the operator must be told the reply is missing, not left staring \
2771 at a gap with no explanation: {}",
2772 note.body
2773 );
2774 assert_eq!(
2775 talk.turns, on_disk.turns,
2776 "the in-memory talk must match what actually landed on disk"
2777 );
2778
2779 let artifacts = talks.artifacts_of(&talk.id);
2782 let stash = std::fs::read_dir(&artifacts)
2783 .expect("artifacts dir")
2784 .filter_map(|e| e.ok())
2785 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
2786 .expect("a stash file for the lost reply");
2787 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
2788 assert_eq!(stashed, "go ahead");
2789
2790 assert_eq!(
2798 on_disk.seat.turns, 1,
2799 "the note's write must carry the turn the CLI actually took"
2800 );
2801 assert_eq!(
2802 on_disk.seat.claude_session, talk.seat.claude_session,
2803 "the session id handed to the CLI must survive the failed reply"
2804 );
2805 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
2806 assert!(
2807 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
2808 "the next turn must resume, not open the same session id twice"
2809 );
2810 }
2811
2812 #[tokio::test]
2817 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
2818 let (tmp, talks) = store();
2819 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2820 let cfg = config(spec);
2821 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2822
2823 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2824 failpoint::force_put_failures(PUT_RETRIES * 2);
2827 let err = respond(&mut talk, &talks, &cfg, &text)
2828 .await
2829 .expect_err("neither the reply nor the note could be saved");
2830 assert!(err.to_string().contains("could not be saved"), "{err}");
2831
2832 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
2833 let on_disk = talks.get(&talk.id).expect("get");
2834 assert_eq!(on_disk.turns.len(), 1);
2835
2836 assert_eq!(
2846 on_disk.seat.turns, 0,
2847 "an unwritable file cannot record the turn the CLI took"
2848 );
2849 assert_eq!(
2850 talk.seat.turns, 1,
2851 "the in-memory seat still reports the turn the CLI actually took"
2852 );
2853 assert_eq!(
2854 on_disk.seat.claude_session, talk.seat.claude_session,
2855 "the session id was minted at `begin` and never changes here"
2856 );
2857 }
2858
2859 #[tokio::test]
2863 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
2864 let (tmp, talks) = store();
2865 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2866 let cfg = config(spec);
2867 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2868
2869 let att = talks
2870 .put_attachment(
2871 &talk.id,
2872 "image/png",
2873 "screenshot.png",
2874 b"pretend-png-bytes",
2875 )
2876 .expect("put attachment");
2877
2878 say(&mut talk, &talks, &cfg, "", vec![att.clone()])
2879 .await
2880 .expect("an empty body with an attachment is still a turn");
2881
2882 let operator_turn = &talk.turns[0];
2883 assert_eq!(operator_turn.who, Who::Operator);
2884 assert_eq!(operator_turn.body, "");
2885 assert_eq!(operator_turn.attachments, vec![att.clone()]);
2886
2887 let prompt = &talk.turns[1].body;
2888 let expected_path = talks
2889 .attachments_dir(&talk.id)
2890 .join(format!("{}.png", att.id));
2891 assert!(
2892 prompt.contains(&expected_path.display().to_string()),
2893 "the agent must be told the attachment's absolute path: {prompt}"
2894 );
2895 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
2896 }
2897
2898 #[test]
2907 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
2908 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
2909 let att = Attachment {
2910 id: "0".repeat(32),
2911 name: "shot.png".to_owned(),
2912 mime: "image/png".to_owned(),
2913 bytes: 3,
2914 };
2915 let path = talks
2916 .attachment_path("some-talk-id", &att)
2917 .expect("a supported mime always yields a path");
2918 assert!(
2919 path.is_absolute(),
2920 "must be absolute even off a relative store root: {}",
2921 path.display()
2922 );
2923 }
2924
2925 #[tokio::test]
2926 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
2927 let (tmp, talks) = store();
2932 let slow = mock_agent(
2933 tmp.path(),
2934 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
2935 BTreeMap::new(),
2936 );
2937 let mut cfg = config(slow);
2938 cfg.graph.timeout_talk = 1;
2939 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2940
2941 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2942 .await
2943 .expect_err("a turn that never answers is an error");
2944 assert!(
2945 err.to_string().contains("did not answer within 1s"),
2946 "{err}"
2947 );
2948
2949 let on_disk = talks.get(&talk.id).expect("get");
2950 let note = on_disk.turns.last().expect("a note turn was recorded");
2951 assert!(
2952 note.body.contains("did not answer within 1s"),
2953 "the transcript must show the configured timeout: {}",
2954 note.body
2955 );
2956 }
2957
2958 #[test]
2959 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
2960 let (tmp, talks) = store();
2961 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
2962 let cfg = config(spec);
2963 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2964
2965 close(&mut talk, &talks).expect("close");
2966 assert_eq!(talk.status, TalkStatus::Closed);
2967 close(&mut talk, &talks).expect("closing twice is not an error");
2968
2969 let err =
2970 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
2971 assert!(err.to_string().contains("closed"));
2972 let _ = &cfg; }
2974
2975 #[tokio::test]
2976 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
2977 let (tmp, talks) = store();
2978 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
2979 let cfg = config(spec);
2980 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2983
2984 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
2988 close(&mut closed_elsewhere, &talks).expect("close");
2989 assert_eq!(
2990 talks.get(&in_flight.id).expect("reread").status,
2991 TalkStatus::Closed,
2992 "the close landed on disk before the turn finished"
2993 );
2994
2995 assert_eq!(in_flight.status, TalkStatus::Open);
2999 respond(&mut in_flight, &talks, &cfg, "one more question")
3000 .await
3001 .expect("the turn itself still completes");
3002
3003 let on_disk = talks.get(&in_flight.id).expect("reread");
3004 assert_eq!(
3005 on_disk.status,
3006 TalkStatus::Closed,
3007 "a close must stick even when a turn that started before it finishes after it"
3008 );
3009 assert!(
3012 on_disk.turns.iter().any(|t| t.body == "here you go"),
3013 "the in-flight turn's own reply is still recorded: {:?}",
3014 on_disk.turns
3015 );
3016 }
3017
3018 #[test]
3019 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
3020 let (tmp, talks) = store();
3021 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3022 let cfg = config(spec);
3023 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3026
3027 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
3030 close(&mut closed_elsewhere, &talks).expect("close");
3031 assert_eq!(
3032 talks.get(&stale.id).expect("reread").status,
3033 TalkStatus::Closed,
3034 "the close landed on disk before record was called"
3035 );
3036
3037 assert_eq!(stale.status, TalkStatus::Open);
3041 let err = record(&mut stale, &talks, "still there?", Vec::new())
3042 .expect_err("a close that landed first must be honored, not overwritten");
3043 assert!(err.to_string().contains("closed"));
3044
3045 let on_disk = talks.get(&stale.id).expect("reread");
3046 assert_eq!(
3047 on_disk.status,
3048 TalkStatus::Closed,
3049 "record must not resurrect a conversation closed while its snapshot was stale"
3050 );
3051 assert!(
3052 on_disk.turns.is_empty(),
3053 "the rejected turn must not have been appended: {:?}",
3054 on_disk.turns
3055 );
3056 let _ = &cfg; }
3058
3059 #[test]
3060 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
3061 let (tmp, talks) = store();
3062 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3063 let cfg = config(spec);
3064 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3065
3066 let held = talks.guard();
3070
3071 let talks2 = talks.clone();
3072 let id = talk.id.clone();
3073 let closing = std::thread::spawn(move || {
3074 let mut talk = talks2.get(&id).expect("get");
3075 close(&mut talk, &talks2).expect("close");
3076 });
3077
3078 std::thread::sleep(Duration::from_millis(50));
3079 assert!(
3080 !closing.is_finished(),
3081 "close must wait for the guard, not read and write while it is held - \
3082 a re-read alone narrows this window without closing it"
3083 );
3084
3085 drop(held);
3086 closing.join().expect("close thread panicked");
3087
3088 assert_eq!(
3089 talks.get(&talk.id).expect("reread").status,
3090 TalkStatus::Closed,
3091 "once the guard is free, close still lands"
3092 );
3093 let _ = &cfg; }
3095
3096 #[test]
3097 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
3098 let (tmp, talks) = store();
3099 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3100 let cfg = config(spec);
3101 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3102
3103 close(&mut talk, &talks).expect("close");
3104 assert_eq!(talk.status, TalkStatus::Closed);
3105
3106 reopen(&mut talk, &talks).expect("reopen");
3107 assert_eq!(talk.status, TalkStatus::Open);
3108 assert_eq!(
3109 talks.get(&talk.id).expect("reread").status,
3110 TalkStatus::Open
3111 );
3112
3113 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
3115 assert_eq!(talk.status, TalkStatus::Open);
3116
3117 record(&mut talk, &talks, "one more thing", Vec::new())
3118 .expect("a reopened talk takes turns again");
3119 let _ = &cfg; }
3121
3122 #[test]
3123 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
3124 let (tmp, talks) = store();
3125 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3126 let cfg = config(spec);
3127 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3128
3129 let artifacts = talks.artifacts_of(&talk.id);
3130 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
3131 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
3132
3133 talks.remove(&talk.id).expect("remove");
3134 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
3135 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
3136 assert!(
3137 talks.get(&talk.id).is_err(),
3138 "a removed talk cannot be read back"
3139 );
3140
3141 let err = talks
3142 .remove("nonexistent-id")
3143 .expect_err("unknown id refused");
3144 assert!(err.to_string().contains("no talk matches"), "{err}");
3145 let _ = &cfg; }
3147
3148 #[tokio::test]
3149 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3150 let (tmp, talks) = store();
3151 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3152 let cfg = config(spec);
3153 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3156
3157 talks.remove(&in_flight.id).expect("remove");
3158 assert!(
3159 talks.get(&in_flight.id).is_err(),
3160 "the delete landed on disk before the turn finished"
3161 );
3162
3163 respond(&mut in_flight, &talks, &cfg, "one more question")
3166 .await
3167 .expect("the turn itself still completes rather than erroring");
3168
3169 assert!(
3170 talks.get(&in_flight.id).is_err(),
3171 "a delete must stick even when a turn that started before it finishes after it"
3172 );
3173 }
3174
3175 #[test]
3176 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
3177 let (tmp, talks) = store();
3178 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3179 let cfg = config(spec);
3180 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3183
3184 talks.remove(&stale.id).expect("remove");
3185
3186 let err = record(&mut stale, &talks, "still there?", Vec::new())
3190 .expect_err("a delete that landed first must be honored, not overwritten");
3191 assert!(err.to_string().contains("deleted"), "{err}");
3192
3193 assert!(
3194 talks.get(&stale.id).is_err(),
3195 "record must not resurrect a conversation deleted while its snapshot was stale"
3196 );
3197 let _ = &cfg; }
3199
3200 #[test]
3201 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
3202 let (tmp, talks) = store();
3203 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3204 let cfg = config(spec);
3205 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3208
3209 talks.remove(&stale.id).expect("remove");
3210
3211 let err = close(&mut stale, &talks)
3215 .expect_err("a delete that landed first must be honored, not overwritten");
3216 assert!(err.to_string().contains("deleted"), "{err}");
3217
3218 assert!(
3219 talks.get(&stale.id).is_err(),
3220 "close must not resurrect a conversation deleted while its snapshot was stale"
3221 );
3222 let _ = &cfg; }
3224
3225 #[test]
3226 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
3227 let (tmp, talks) = store();
3228 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3229 let cfg = config(spec);
3230 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3233 close(&mut stale, &talks).expect("close");
3234
3235 talks.remove(&stale.id).expect("remove");
3236
3237 let err = reopen(&mut stale, &talks)
3241 .expect_err("a delete that landed first must be honored, not overwritten");
3242 assert!(err.to_string().contains("deleted"), "{err}");
3243
3244 assert!(
3245 talks.get(&stale.id).is_err(),
3246 "reopen must not resurrect a conversation deleted while its snapshot was stale"
3247 );
3248 let _ = &cfg; }
3250
3251 #[test]
3252 fn list_puts_open_talks_before_closed_ones() {
3253 let (tmp, talks) = store();
3254 let make = |id: &str, status: TalkStatus| {
3255 let mut t = Talk {
3256 schema: SCHEMA,
3257 id: id.to_owned(),
3258 repo: tmp.path().to_owned(),
3259 agent: "mock".to_owned(),
3260 status,
3261 turns: Vec::new(),
3262 pending: String::new(),
3263 pending_attachments: Vec::new(),
3264 fallback: false,
3265 created_at: Timestamp::now(),
3266 updated_at: Timestamp::now(),
3267 seat: SeatState::new(SEAT, "mock", 7),
3268 };
3269 talks.put(&mut t).expect("put");
3270 };
3271 make("20260901-000000-0001", TalkStatus::Open);
3272 make("20260902-000000-0002", TalkStatus::Open);
3273 make("20260903-000000-0003", TalkStatus::Closed);
3274
3275 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
3276 assert_eq!(
3277 ids,
3278 [
3279 "20260902-000000-0002",
3280 "20260901-000000-0001",
3281 "20260903-000000-0003"
3282 ]
3283 );
3284 assert_eq!(talks.count_open(), 2);
3285 }
3286
3287 #[test]
3288 fn tasks_of_finds_only_this_talks_own_tasks() {
3289 let dir = tempfile::tempdir().expect("tempdir");
3290 let queue = Queue::at(dir.path().join("queue"));
3291
3292 let mut mine = Task::new(
3293 "rework the loader".to_owned(),
3294 "rework the loader".to_owned(),
3295 PathBuf::from("/repo"),
3296 Source::Agent {
3297 run: "20260904-014455-ab12".to_owned(),
3298 node: "chat".to_owned(),
3299 },
3300 );
3301 queue.put(&mut mine).expect("put mine");
3302
3303 let mut theirs = Task::new(
3304 "unrelated".to_owned(),
3305 "unrelated".to_owned(),
3306 PathBuf::from("/repo"),
3307 Source::Agent {
3308 run: "20260904-090000-zz99".to_owned(),
3309 node: "implement".to_owned(),
3310 },
3311 );
3312 queue.put(&mut theirs).expect("put theirs");
3313
3314 let mut human = Task::new(
3315 "typed by hand".to_owned(),
3316 "typed by hand".to_owned(),
3317 PathBuf::from("/repo"),
3318 Source::Human,
3319 );
3320 queue.put(&mut human).expect("put human");
3321
3322 let found = tasks_of(&queue, "20260904-014455-ab12");
3323 assert_eq!(found.len(), 1);
3324 assert_eq!(found[0].id, mine.id);
3325 }
3326
3327 #[test]
3328 fn the_briefing_names_solo_task_add() {
3329 let brief = briefing(Path::new("/repo"), "en", false);
3330 assert!(brief.contains("magi task add --solo"));
3331 assert!(brief.contains("/repo"));
3332 assert!(!brief.contains("Hold this conversation in"));
3333 }
3334
3335 #[test]
3341 fn the_briefing_explains_targeting_a_different_repository_by_name() {
3342 let brief = briefing(Path::new("/repo"), "en", false);
3343 assert!(brief.contains("--repo does not have to be a full path"));
3344 assert!(brief.contains("owner/repo"));
3345 assert!(brief.contains("magi repos"));
3346 assert!(brief.contains("ask the operator"));
3347 }
3348
3349 #[test]
3350 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
3351 let brief = briefing(Path::new("/repo"), "en", false);
3352 assert!(brief.contains("--attach <path>"), "{brief}");
3353 assert!(brief.contains("deleting this conversation"), "{brief}");
3354 }
3355
3356 #[test]
3357 fn the_briefing_names_the_language_when_it_is_not_english() {
3358 let brief = briefing(Path::new("/repo"), "Japanese", false);
3359 assert!(brief.contains("Hold this conversation in Japanese"));
3360 }
3361
3362 #[test]
3363 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
3364 let read_only = briefing(Path::new("/repo"), "en", false);
3365 assert!(read_only.contains("Do not write files"));
3366 assert!(!read_only.contains("allow_write"));
3367
3368 let writable = briefing(Path::new("/repo"), "en", true);
3369 assert!(!writable.contains("Do not write files"));
3370 assert!(writable.contains("allow_write = true"));
3371 assert!(writable.contains("magi task add --solo"));
3374 assert!(writable.contains("say plainly what you"));
3375 }
3376}