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 = fresh_token();
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 fresh_token() -> String {
693 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
694 let n = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
695 let seed = crate::rng::entropy() ^ n.wrapping_mul(0x9E37_79B9_7F4A_7C15);
696 crate::rng::SplitMix64::new(seed).uuid_v4()
697}
698
699fn read_turn(path: &Path) -> Option<TurnRecord> {
700 serde_json::from_str(&std::fs::read_to_string(path).ok()?).ok()
701}
702
703const TAKEOVER_LOCK_TTL: Duration = Duration::from_secs(10);
705
706const TICKET_BUCKET: Duration = Duration::from_secs(60);
709
710fn create_exclusive(path: &Path, body: &str) -> Result<bool> {
712 use std::io::Write as _;
713 match std::fs::OpenOptions::new()
714 .write(true)
715 .create_new(true)
716 .open(path)
717 {
718 Ok(mut f) => {
719 if let Err(e) = f.write_all(body.as_bytes()) {
720 drop(f);
721 let _ = std::fs::remove_file(path);
722 return Err(e).with_context(|| format!("write {}", path.display()));
723 }
724 Ok(true)
725 }
726 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
727 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
728 }
729}
730
731fn create_turn(path: &Path, token: &str, now: Timestamp) -> Result<bool> {
732 let record = TurnRecord {
733 token: token.to_owned(),
734 pid: std::process::id(),
735 beat_at: now,
736 };
737 let body = serde_json::to_string(&record).context("serialize turn lease")?;
738 let tmp = path.with_extension(format!("turn.{token}.new"));
741 std::fs::write(&tmp, body).with_context(|| format!("write {}", tmp.display()))?;
742 let linked = std::fs::hard_link(&tmp, path);
743 let _ = std::fs::remove_file(&tmp);
744 match linked {
745 Ok(()) => Ok(true),
746 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
747 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
748 }
749}
750
751struct TurnLock {
755 path: PathBuf,
756 token: String,
757}
758
759impl TurnLock {
760 fn token() -> String {
761 fresh_token()
762 }
763
764 fn publish(path: &Path, token: &str) -> Result<bool> {
767 let tmp = path.with_extension(format!("lock.{token}.new"));
768 std::fs::write(&tmp, token).with_context(|| format!("write {}", tmp.display()))?;
769 let linked = std::fs::hard_link(&tmp, path);
770 let _ = std::fs::remove_file(&tmp);
771 match linked {
772 Ok(()) => Ok(true),
773 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
774 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
775 }
776 }
777
778 fn take(lease: &Path) -> Result<Option<Self>> {
779 let path = lease.with_extension("turn.lock");
780 let token = Self::token();
781 if Self::publish(&path, &token)? {
782 return Ok(Some(Self { path, token }));
783 }
784 let Ok(seen) = std::fs::read_to_string(&path) else {
785 return Ok(None);
786 };
787 let key = if !seen.is_empty() && seen.chars().all(|c| c.is_ascii_alphanumeric() || c == '-')
790 {
791 seen.as_str()
792 } else {
793 "invalid"
794 };
795 let aged = std::fs::metadata(&path)
796 .and_then(|m| m.modified())
797 .ok()
798 .and_then(|t| t.elapsed().ok())
799 .is_some_and(|age| age > TAKEOVER_LOCK_TTL);
800 if !aged {
801 return Ok(None);
802 }
803 let bucket = std::time::SystemTime::now()
812 .duration_since(std::time::UNIX_EPOCH)
813 .map_or(0, |d| d.as_secs() / TICKET_BUCKET.as_secs());
814 let ticket = path.with_extension(format!("lock.{key}.{bucket}.break"));
815 if !create_exclusive(&ticket, "")? {
816 return Ok(None);
817 }
818 Self::sweep_tickets(&path);
819 if std::fs::read_to_string(&path).ok().as_deref() != Some(seen.as_str()) {
822 return Ok(None);
823 }
824 let _ = std::fs::remove_file(&path);
825 if Self::publish(&path, &token)? {
826 return Ok(Some(Self { path, token }));
827 }
828 Ok(None)
829 }
830
831 fn sweep_tickets(path: &Path) {
834 let (Some(dir), Some(name)) = (path.parent(), path.file_name().and_then(|n| n.to_str()))
835 else {
836 return;
837 };
838 let prefix = format!("{name}.");
839 let Ok(entries) = std::fs::read_dir(dir) else {
840 return;
841 };
842 for entry in entries.flatten() {
843 let file = entry.file_name();
844 let Some(file) = file.to_str() else { continue };
845 if !(file.starts_with(&prefix) && file.ends_with(".break")) {
846 continue;
847 }
848 let old = entry
849 .metadata()
850 .and_then(|m| m.modified())
851 .ok()
852 .and_then(|t| t.elapsed().ok())
853 .is_some_and(|age| age > TICKET_BUCKET * 60);
854 if old {
855 let _ = std::fs::remove_file(entry.path());
856 }
857 }
858 }
859
860 fn take_patiently(lease: &Path) -> Option<Self> {
862 for _ in 0..50 {
863 match Self::take(lease) {
864 Ok(Some(lock)) => return Some(lock),
865 Ok(None) => std::thread::sleep(Duration::from_millis(10)),
866 Err(_) => return None,
867 }
868 }
869 None
870 }
871}
872
873impl Drop for TurnLock {
874 fn drop(&mut self) {
875 if std::fs::read_to_string(&self.path).is_ok_and(|t| t == self.token) {
877 let _ = std::fs::remove_file(&self.path);
878 }
879 }
880}
881
882pub const TURN_BEAT: Duration = Duration::from_secs(20);
885
886#[derive(Debug)]
891pub struct TurnLease {
892 path: PathBuf,
893 token: String,
894}
895
896impl TurnLease {
897 pub fn beat(&self) -> Result<bool> {
901 let _lock = TurnLock::take_patiently(&self.path)
902 .with_context(|| format!("lock {} to renew it", self.path.display()))?;
903 let Some(mut record) = read_turn(&self.path).filter(|r| r.token == self.token) else {
904 return Ok(false);
905 };
906 record.beat_at = Timestamp::now();
907 let body = serde_json::to_string(&record).context("serialize turn lease")?;
908 let tmp = self.path.with_extension(format!("turn.{}.tmp", self.token));
909 write_atomic(&tmp, &self.path, &body)?;
910 Ok(true)
911 }
912
913 pub async fn beating<T>(&self, fut: impl std::future::Future<Output = T>) -> Result<T> {
918 self.beating_every(TURN_BEAT, fut).await
919 }
920
921 async fn beating_every<T>(
922 &self,
923 period: Duration,
924 fut: impl std::future::Future<Output = T>,
925 ) -> Result<T> {
926 tokio::pin!(fut);
927 loop {
928 match tokio::time::timeout(period, &mut fut).await {
929 Ok(out) => return Ok(out),
930 Err(_) => match self.beat() {
931 Ok(true) => {}
932 Ok(false) => bail!(
933 "the turn lease {} was taken over; this turn is stopped",
934 self.path.display()
935 ),
936 Err(e) => tracing::warn!("{e:#}"),
937 },
938 }
939 }
940 }
941}
942
943impl Drop for TurnLease {
944 fn drop(&mut self) {
945 if let Some(_lock) = TurnLock::take_patiently(&self.path) {
948 if read_turn(&self.path).is_some_and(|r| r.token == self.token) {
949 let _ = std::fs::remove_file(&self.path);
950 }
951 }
952 }
953}
954
955pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
965 let repo = repo.canonicalize().unwrap_or(repo);
968 let spec = match agent {
971 Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
972 None => agent::pick_chain(
973 &cfg.agents,
974 cfg.roles.chatter.as_ref(),
975 &agent::installed,
976 "chatter",
977 )?
978 .remove(0),
979 };
980
981 let now = Timestamp::now();
982 let mut talk = Talk {
983 schema: SCHEMA,
984 id: new_id(),
985 repo,
986 agent: spec.id.clone(),
987 status: TalkStatus::Open,
988 turns: Vec::new(),
989 pending: String::new(),
990 pending_attachments: Vec::new(),
991 fallback: agent.is_none(),
992 created_at: now,
993 updated_at: now,
994 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
995 };
996 store.put(&mut talk)?;
997 Ok(talk)
998}
999
1000pub fn record(
1007 talk: &mut Talk,
1008 store: &Talks,
1009 text: &str,
1010 attachments: Vec<Attachment>,
1011) -> Result<String> {
1012 let _guard = store.guard();
1020 let Ok(fresh) = store.get(&talk.id) else {
1025 bail!("talk {} was deleted", talk.short());
1026 };
1027 talk.status = fresh.status;
1028 talk.pending = fresh.pending;
1031 talk.pending_attachments = fresh.pending_attachments;
1032 if !talk.status.open() {
1033 bail!(
1034 "talk {} is {} and takes no more turns",
1035 talk.short(),
1036 talk.status.as_str()
1037 );
1038 }
1039 let text = text.trim();
1040 if text.is_empty() && attachments.is_empty() {
1041 bail!("nothing to say");
1042 }
1043 talk.turns.push(Turn {
1044 who: Who::Operator,
1045 body: text.to_owned(),
1046 at: Timestamp::now(),
1047 attachments,
1048 usage: None,
1049 });
1050 store.put(talk)?;
1051 Ok(text.to_owned())
1052}
1053
1054pub fn queue(
1056 talk: &mut Talk,
1057 store: &Talks,
1058 text: &str,
1059 attachments: Vec<Attachment>,
1060) -> Result<()> {
1061 let text = text.trim();
1062 if text.is_empty() && attachments.is_empty() {
1063 bail!("nothing to say");
1064 }
1065 let _guard = store.guard();
1066 let mut fresh = store
1067 .get(&talk.id)
1068 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1069 if !fresh.status.open() {
1070 bail!(
1071 "talk {} is {} and takes no more turns",
1072 fresh.short(),
1073 fresh.status.as_str()
1074 );
1075 }
1076 if !text.is_empty() {
1077 if fresh.pending.is_empty() {
1078 fresh.pending = text.to_owned();
1079 } else {
1080 fresh.pending.push_str("\n\n");
1081 fresh.pending.push_str(text);
1082 }
1083 }
1084 fresh.pending_attachments.extend(attachments);
1085 store.put(&mut fresh)?;
1086 *talk = fresh;
1087 Ok(())
1088}
1089
1090pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
1092 let _guard = store.guard();
1093 let mut fresh = store
1094 .get(&talk.id)
1095 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1096 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
1097 *talk = fresh;
1098 return Ok(None);
1099 }
1100 let text = std::mem::take(&mut fresh.pending);
1101 let attachments = std::mem::take(&mut fresh.pending_attachments);
1102 fresh.turns.push(Turn {
1103 who: Who::Operator,
1104 body: text.clone(),
1105 at: Timestamp::now(),
1106 attachments,
1107 usage: None,
1108 });
1109 store.put(&mut fresh)?;
1110 *talk = fresh;
1111 Ok(Some(text))
1112}
1113
1114pub async fn say(
1117 talk: &mut Talk,
1118 store: &Talks,
1119 cfg: &Config,
1120 text: &str,
1121 attachments: Vec<Attachment>,
1122) -> Result<()> {
1123 let text = record(talk, store, text, attachments)?;
1124 turn(talk, store, cfg, &text).await
1125}
1126
1127pub async fn respond(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1129 turn(talk, store, cfg, text).await
1130}
1131
1132pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
1150 let _guard = store.guard();
1151 let mut fresh = store
1152 .get(&talk.id)
1153 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1154 fresh.status = TalkStatus::Closed;
1155 fresh.pending.clear();
1157 fresh.pending_attachments.clear();
1158 store.put(&mut fresh)?;
1159 *talk = fresh;
1160 Ok(())
1161}
1162
1163pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
1174 let _guard = store.guard();
1175 let mut fresh = store
1176 .get(&talk.id)
1177 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1178 fresh.status = TalkStatus::Open;
1179 store.put(&mut fresh)?;
1180 *talk = fresh;
1181 Ok(())
1182}
1183
1184pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
1197 let _guard = store.guard();
1198 let mut fresh = store
1199 .get(&talk.id)
1200 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1201 if fresh.agent == spec.id {
1202 *talk = fresh;
1203 return Ok(false);
1204 }
1205 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
1206 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1207 fresh.fallback = false;
1209 fresh.turns.push(Turn {
1210 who: Who::Agent,
1211 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
1212 at: Timestamp::now(),
1213 attachments: Vec::new(),
1214 usage: None,
1215 });
1216 store.put(&mut fresh)?;
1217 *talk = fresh;
1218 Ok(true)
1219}
1220
1221pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
1223 let _guard = store.guard();
1224 let mut fresh = store
1225 .get(&talk.id)
1226 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1227 fresh.pending.clear();
1228 fresh.pending_attachments.clear();
1229 store.put(&mut fresh)?;
1230 *talk = fresh;
1231 Ok(())
1232}
1233
1234pub fn clear_pending_if_matches(
1236 talk: &mut Talk,
1237 store: &Talks,
1238 expected_text: &str,
1239 expected_attachments: &[String],
1240) -> Result<bool> {
1241 let _guard = store.guard();
1242 let mut fresh = store
1243 .get(&talk.id)
1244 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1245 if !pending_matches(&fresh, expected_text, expected_attachments) {
1246 *talk = fresh;
1247 return Ok(false);
1248 }
1249 fresh.pending.clear();
1250 fresh.pending_attachments.clear();
1251 store.put(&mut fresh)?;
1252 *talk = fresh;
1253 Ok(true)
1254}
1255
1256pub fn edit_pending_text(
1260 talk: &mut Talk,
1261 store: &Talks,
1262 text: &str,
1263 expected_text: &str,
1264 expected_attachments: &[String],
1265) -> Result<bool> {
1266 let _guard = store.guard();
1267 let mut fresh = store
1268 .get(&talk.id)
1269 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1270 if !pending_matches(&fresh, expected_text, expected_attachments) {
1271 *talk = fresh;
1272 return Ok(false);
1273 }
1274 fresh.pending = text.trim().to_owned();
1275 store.put(&mut fresh)?;
1276 *talk = fresh;
1277 Ok(true)
1278}
1279
1280fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
1281 talk.pending == expected_text
1282 && talk
1283 .pending_attachments
1284 .iter()
1285 .map(|attachment| &attachment.id)
1286 .eq(expected_attachments.iter())
1287}
1288
1289async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1296 let spec = cfg
1297 .agents
1298 .iter()
1299 .find(|a| a.id == talk.agent)
1300 .with_context(|| {
1301 format!(
1302 "talk {} was opened with agent `{}`, which is no longer in \
1303 the roster; restore it in magi.toml or start a new \
1304 conversation",
1305 talk.short(),
1306 talk.agent
1307 )
1308 })?;
1309
1310 let last_note = attachment_note(
1314 store,
1315 &talk.id,
1316 talk.turns
1317 .last()
1318 .map_or(&[][..], |t| t.attachments.as_slice()),
1319 );
1320
1321 let attachment_paths: Vec<PathBuf> = talk
1327 .turns
1328 .iter()
1329 .flat_map(|t| t.attachments.iter())
1330 .filter_map(|a| store.attachment_path(&talk.id, a))
1331 .collect();
1332
1333 let questions = crate::ask::Questions::open();
1341 let consulted = crate::consult::pending_consults(&questions, &talk.id);
1342 let consult_roots: Vec<PathBuf> = if consulted {
1343 vec![questions.root().to_path_buf()]
1344 } else {
1345 Vec::new()
1346 };
1347
1348 let artifacts = store.artifacts_of(&talk.id);
1349 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
1352 let stem = format!("turn-{}", operator_turns.max(1));
1353 let cache_dir = cfg.cache_dir();
1356
1357 let mut chain = vec![spec.clone()];
1360 if let Some(choice) = cfg.roles.chatter.as_ref()
1361 && talk.fallback
1362 {
1363 for id in choice.ids() {
1364 if id == talk.agent || chain.iter().any(|s| s.id == id) {
1365 continue;
1366 }
1367 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1368 Ok(s) => chain.push(s),
1369 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1370 }
1371 }
1372 }
1373
1374 let mut outcome = None;
1375 let mut fell_back_from: Option<String> = None;
1376 let mut first_try: Option<(String, SeatState)> = None;
1379 for (n, spec) in chain.iter().enumerate() {
1380 if n > 0 {
1381 if first_try.is_none() {
1382 first_try = Some((talk.agent.clone(), talk.seat.clone()));
1383 }
1384 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1385 fell_back_from.get_or_insert_with(|| talk.agent.clone());
1388 talk.agent = spec.id.clone();
1389 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1390 }
1391 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1392 let first_ever = talk.turns.len() <= 1;
1393 let body = if talk.seat.turns == 0 && first_ever {
1394 format!(
1395 "{}\n\n# Operator\n\n{text}{last_note}",
1396 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write)
1397 )
1398 } else if talk.seat.turns == 0 {
1399 format!(
1402 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1403 briefing(&talk.repo, &cfg.graph.language, cfg.talk.allow_write),
1404 transcript(talk, store)
1405 )
1406 } else if resuming {
1407 format!("{text}{last_note}")
1408 } else {
1409 format!("{}\n\n{text}{last_note}", transcript(talk, store))
1410 };
1411 let attempt_stem = if n == 0 {
1412 stem.clone()
1413 } else {
1414 format!("{stem}-{}", spec.id)
1415 };
1416 let inv = Invocation {
1417 cwd: &talk.repo,
1418 prompt: &body,
1419 timeout: turn_timeout(cfg),
1420 allow_write: cfg.talk.allow_write || consulted,
1425 sessions: cfg.graph.sessions,
1426 artifacts: &artifacts,
1427 stem: &attempt_stem,
1428 run: &talk.id,
1431 node: crate::queue::CHAT_NODE,
1432 cache_dir: cache_dir.as_deref(),
1433 attachments: &attachment_paths,
1434 writable: &consult_roots,
1435 };
1436 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1437 let advance = agent::chain_advances(&result);
1438 if n == 0 || !advance {
1439 outcome = Some(result);
1440 } else {
1441 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1443 }
1444 if !advance {
1445 break;
1446 }
1447 }
1448 if outcome.as_ref().is_some_and(agent::chain_advances) {
1449 if let Some((id, seat)) = first_try {
1452 talk.agent = id;
1453 talk.seat = seat;
1454 fell_back_from = None;
1455 }
1456 }
1457 let outcome = outcome.expect("a chain holds at least one agent");
1458 let note = |why: String| Turn {
1459 who: Who::Agent,
1460 body: format!("{MAGI_NOTE}{why}"),
1461 at: Timestamp::now(),
1462 attachments: Vec::new(),
1463 usage: None,
1464 };
1465 let (reply, failure) = match outcome {
1466 Err(e) => (
1467 note(format!("could not run agent `{}`: {e}", talk.agent)),
1468 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1469 ),
1470 Ok(out) if out.quota_exhausted() => {
1471 let reset = out
1472 .quota
1473 .as_ref()
1474 .and_then(|q| q.reset.clone())
1475 .map_or_else(String::new, |r| format!(" (resets {r})"));
1476 let why = format!(
1477 "agent `{}` is out of quota{reset}; your message is saved, so \
1478 say it again when the window reopens",
1479 talk.agent
1480 );
1481 (note(why.clone()), Some(why))
1482 }
1483 Ok(out) if out.timed_out => {
1484 let why = format!(
1485 "agent `{}` did not answer within {}s; your message is saved",
1486 talk.agent,
1487 turn_timeout(cfg).as_secs()
1488 );
1489 (note(why.clone()), Some(why))
1490 }
1491 Ok(out) if !out.usable() => {
1492 let why = format!(
1493 "agent `{}` produced no answer (exit {}); your message is saved",
1494 talk.agent,
1495 out.exit_code
1496 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1497 );
1498 (note(why.clone()), Some(why))
1499 }
1500 Ok(out) => (
1501 Turn {
1502 who: Who::Agent,
1503 body: out.text.trim().to_owned(),
1504 at: Timestamp::now(),
1505 attachments: Vec::new(),
1506 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1509 context_tokens,
1510 agent: talk.agent.clone(),
1511 model: cfg
1512 .agents
1513 .iter()
1514 .find(|a| a.id == talk.agent)
1515 .and_then(|a| a.model.clone()),
1516 }),
1517 },
1518 None,
1519 ),
1520 };
1521
1522 let _guard = store.guard();
1534 let Ok(fresh) = store.get(&talk.id) else {
1540 return Ok(());
1541 };
1542 talk.status = fresh.status;
1543 talk.pending = fresh.pending;
1547 talk.pending_attachments = fresh.pending_attachments;
1548 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1549 talk.turns.push(note(format!(
1552 "agent changed from {from} to {} (fallback)",
1553 talk.agent
1554 )));
1555 }
1556 talk.turns.push(reply);
1557 if let Err(put_err) = store.put(talk) {
1558 let lost = talk.turns.pop().expect("just pushed above");
1567 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1568 let why = match &stash {
1569 Ok(path) => format!(
1570 "agent `{}` answered, but the reply could not be saved to \
1571 this conversation ({put_err:#}); the raw text was kept at \
1572 {} - your message is saved, ask again",
1573 talk.agent,
1574 path.display()
1575 ),
1576 Err(stash_err) => format!(
1577 "agent `{}` answered, but the reply could not be saved to \
1578 this conversation ({put_err:#}), and it could not be kept \
1579 anywhere else either ({stash_err:#}); your message is \
1580 saved, ask again",
1581 talk.agent
1582 ),
1583 };
1584 talk.turns.push(note(why.clone()));
1585 return match store.put(talk) {
1592 Ok(()) => bail!("{why}"),
1593 Err(note_err) => {
1594 talk.turns.pop();
1614 Err(note_err).context(why)
1615 }
1616 };
1617 }
1618
1619 match failure {
1620 Some(why) => bail!("{why}"),
1621 None => Ok(()),
1622 }
1623}
1624
1625fn transcript(talk: &Talk, store: &Talks) -> String {
1628 let mut out = String::from(
1629 "This conversation cannot resume on the CLI's side, so here is \
1630 everything said so far; answer only the last message.\n",
1631 );
1632 for t in &talk.turns {
1633 let who = match t.who {
1634 Who::Operator => "operator",
1635 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1636 Who::Agent => "you",
1637 };
1638 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1639 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1640 }
1641 out
1642}
1643
1644fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1649 if attachments.is_empty() {
1650 return String::new();
1651 }
1652 let mut out = String::from(
1653 "\n\nThe operator attached the image(s) below to this message. Open \
1654 and look at each one before you answer.\n",
1655 );
1656 for att in attachments {
1657 if let Some(path) = store.attachment_path(talk_id, att) {
1658 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1659 }
1660 }
1661 out.push('\n');
1662 out
1663}
1664
1665pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
1687 let write_policy = if allow_write {
1688 "Write access is enabled for this conversation (`allow_write = \
1689 true`), so you may write files - but only a small, \
1690 already-decided edit the operator names outright in this \
1691 conversation, not an implementation. This is a permission on the \
1692 conversation as a whole, not a property of whichever repository \
1693 it happened to start in: if the operator names a different \
1694 repository for that small edit, the policy allows it there too. \
1695 Your own tool may still confine writes to the repository this \
1696 conversation started in regardless - if a write elsewhere is \
1697 refused, say so plainly rather than working around it. Once you \
1698 have made an edit, say plainly what you edited. Anything bigger, \
1699 or anything still open-ended, still goes through the queue below \
1700 rather than being done here."
1701 } else {
1702 "Do not write files. Implementing a change is not this \
1703 conversation's job; a separate, blind competition of agents does \
1704 that, and a repository this conversation has already edited would \
1705 make their diffs unjudgeable."
1706 };
1707 let mut out = format!(
1708 "You are magi's standing conversation partner for its operator, who \
1709 usually has this open on a phone. Keep replies short: no preamble, \
1710 no restating what they just said.\n\n\
1711 # Repository\n\n{repo}\n\n\
1712 You may look around: read files, run shell commands, search history, \
1713 run tests - whatever answers the question. {write_policy}\n\n\
1714 A short, command-shaped message (\"list\", \"info <id>\", \"show \
1715 3cbf\") is almost always the operator asking you to look something \
1716 up, not an instruction to file - answer it yourself with `magi \
1717 list`, `magi show <id>`, `magi task list`, or the like, the same way \
1718 you would answer any other question in this conversation.\n\n\
1719 # When the operator wants something done\n\n\
1720 Run:\n\n\
1721 magi task add --solo --repo {repo} <instruction>\n\n\
1722 and tell the operator the task id it prints, so they can follow it \
1723 from the Queue. If it refuses with a duplicate warning (the \
1724 instruction names a branch, commit or pull request that an \
1725 unfinished task, run or PR already owns), do not repeat it with \
1726 --force yourself: tell the operator what it matched and let them \
1727 decide. Write <instruction> so that an implementer who has \
1728 never seen this conversation can act on it alone - it is everything \
1729 they get. Use --solo: it runs the task through one implementer \
1730 straight into review instead of the usual multi-agent competition, \
1731 which is the right shape for a change this conversation has already \
1732 settled, rather than one still worth several independent takes.\n\n\
1733 If the operator asks for something in a different repository, \
1734 --repo does not have to be a full path: --repo owner/repo (or just \
1735 repo, when that is unambiguous) is resolved against local checkouts \
1736 the same way `magi repos` lists them. If the command fails because \
1737 nothing matches or more than one checkout shares that name, ask the \
1738 operator which repository they mean (or run `magi repos` yourself \
1739 to see the candidates) rather than guessing.\n\n\
1740 The current state of the code is whatever origin/main holds, not \
1741 whatever a working tree shows: a primary checkout often lags \
1742 upstream, sits on a detached HEAD and carries uncommitted changes. \
1743 Before answering about code, run `git fetch origin` in that \
1744 repository if it is cheap, then read through \
1745 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
1746 If the working tree differs, say so; if the fetch fails, say that \
1747 too, so the operator knows the answer may be stale.\n\n\
1748 If the operator attached an image (a screenshot, say) that the task \
1749 is about, pass it with `--attach <path>`, using the absolute path \
1750 the turn's attachment note gives; repeat the flag for several. \
1751 `magi task add --solo --attach <path> <instruction>` copies the \
1752 file into the task, so the implementer receives it. Do not paste the \
1753 path into <instruction> instead: deleting this conversation deletes \
1754 its attachments, and then that path reaches no one.\n",
1755 repo = repo.display(),
1756 );
1757 out.push_str(&language_note(language));
1758 out
1759}
1760
1761fn language_note(language: &str) -> String {
1764 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
1765 String::new()
1766 } else {
1767 format!("\nHold this conversation in {language}.\n")
1768 }
1769}
1770
1771pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
1778 let mut tasks: Vec<Task> = queue
1779 .list()
1780 .into_iter()
1781 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
1782 .collect();
1783 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
1784 tasks
1785}
1786
1787fn read_path(path: &Path) -> Result<Talk> {
1788 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1789 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
1790}
1791
1792const PUT_RETRIES: u32 = 5;
1795
1796fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1808 let mut last_err = None;
1809 for attempt in 0..PUT_RETRIES {
1810 if attempt > 0 {
1811 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
1812 }
1813 match try_write_atomic(tmp, path, body) {
1814 Ok(()) => return Ok(()),
1815 Err(e) => last_err = Some(e),
1816 }
1817 }
1818 Err(last_err.expect("the loop above always runs at least once"))
1819}
1820
1821fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1822 #[cfg(test)]
1823 if failpoint::take_forced_put_failure() {
1824 bail!("simulated write failure (test)");
1825 }
1826 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
1827 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
1828 Ok(())
1829}
1830
1831fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
1836 let dir = store.artifacts_of(id);
1837 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1838 let path = dir.join(format!("{stem}-lost.txt"));
1839 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
1840 Ok(path)
1841}
1842
1843#[cfg(test)]
1850mod failpoint {
1851 use std::cell::Cell;
1852
1853 thread_local! {
1854 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
1855 }
1856
1857 pub(super) fn force_put_failures(count: u32) {
1860 FORCE_PUT_FAILURES.with(|c| c.set(count));
1861 }
1862
1863 pub(super) fn take_forced_put_failure() -> bool {
1866 FORCE_PUT_FAILURES.with(|c| {
1867 let n = c.get();
1868 if n == 0 {
1869 false
1870 } else {
1871 c.set(n - 1);
1872 true
1873 }
1874 })
1875 }
1876}
1877
1878fn short(id: &str) -> &str {
1879 id.split('-').next_back().unwrap_or(id)
1880}
1881
1882fn new_id() -> String {
1883 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
1884 let seed = crate::rng::entropy();
1885 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
1886}
1887
1888fn attachment_ext(mime: &str) -> Option<&'static str> {
1893 match mime {
1894 "image/png" => Some("png"),
1895 "image/jpeg" => Some("jpg"),
1896 "image/gif" => Some("gif"),
1897 "image/webp" => Some("webp"),
1898 _ => None,
1899 }
1900}
1901
1902pub fn valid_attachment_id(id: &str) -> bool {
1907 id.len() == 32
1908 && id
1909 .bytes()
1910 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
1911}
1912
1913fn new_attachment_id() -> String {
1917 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
1918 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
1919}
1920
1921#[cfg(test)]
1922mod tests {
1923 #[test]
1924 fn the_briefing_points_at_origin_main_not_the_working_tree() {
1925 let b = briefing(Path::new("/r"), "en", false);
1926 assert!(b.contains("origin/main"));
1927 assert!(b.contains("git show origin/main:"));
1928 }
1929 use std::collections::BTreeMap;
1930
1931 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
1932 use crate::queue::{Queue, Source, Task};
1933
1934 use super::*;
1935
1936 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
1937 AgentSpec {
1938 id: id.to_owned(),
1939 kind: AgentKind::Command,
1940 model: model.map(str::to_owned),
1941 command: Vec::new(),
1942 extra_args: Vec::new(),
1943 env: BTreeMap::new(),
1944 prompt_delivery: None,
1945 }
1946 }
1947
1948 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
1949 Talk {
1950 schema: SCHEMA,
1951 id: "20260904-014455-ab12".to_owned(),
1952 repo: PathBuf::from("."),
1953 agent: agent.to_owned(),
1954 status: TalkStatus::Open,
1955 turns,
1956 pending: String::new(),
1957 pending_attachments: Vec::new(),
1958 fallback: false,
1959 created_at: Timestamp::now(),
1960 updated_at: Timestamp::now(),
1961 seat: SeatState::new(SEAT, agent, 1),
1962 }
1963 }
1964
1965 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
1966 Turn {
1967 who: Who::Agent,
1968 body: body.to_owned(),
1969 at: Timestamp::now(),
1970 attachments: Vec::new(),
1971 usage: usage.map(|(t, a, m)| TurnUsage {
1972 context_tokens: t,
1973 agent: a.to_owned(),
1974 model: m.map(str::to_owned),
1975 }),
1976 }
1977 }
1978
1979 fn ctx_config(windows: &[(&str, u64)]) -> Config {
1980 Config {
1981 agents: vec![
1982 ctx_agent("small", Some("small-model")),
1983 ctx_agent("big", Some("big-model")),
1984 ctx_agent("plain", None),
1985 ],
1986 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
1987 ..Config::default()
1988 }
1989 }
1990
1991 #[test]
1992 fn context_usage_computes_percent_and_warns_at_eighty() {
1993 let cfg = ctx_config(&[("small-model", 1000)]);
1994 let at = |tokens| {
1995 let t = ctx_talk(
1996 "small",
1997 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
1998 );
1999 context_usage(&t, Some(&cfg))
2000 };
2001 let u = at(799);
2002 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
2003 let u = at(800);
2004 assert_eq!((u.percent, u.warn), (Some(80), true));
2005 let u = at(1500);
2006 assert_eq!((u.percent, u.warn), (Some(150), true));
2007 assert!(!u.since_switch);
2008 }
2009
2010 #[test]
2011 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
2012 let cfg = ctx_config(&[("small-model", 1000)]);
2013 let t = ctx_talk(
2014 "small",
2015 vec![
2016 reply("old", Some((900, "small", Some("small-model")))),
2017 reply("new", None),
2018 ],
2019 );
2020 let u = context_usage(&t, Some(&cfg));
2021 assert!(u.estimated);
2023 assert_ne!(u.tokens, Some(900));
2024 assert!(u.tokens.is_some());
2025 let t = ctx_talk(
2027 "small",
2028 vec![
2029 reply("old", Some((900, "small", Some("small-model")))),
2030 reply("magi: could not run agent", None),
2031 ],
2032 );
2033 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
2034 assert_eq!(
2035 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
2036 None
2037 );
2038 }
2039
2040 #[test]
2041 fn estimate_counts_chars_both_sides_and_standing_prompt() {
2042 let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
2043 assert_eq!(estimate_context_tokens(&t, 0), Some(2)); let op = Turn {
2045 who: Who::Operator,
2046 ..reply("abcdefg", None)
2047 };
2048 t.turns.push(op);
2049 assert_eq!(estimate_context_tokens(&t, 0), Some(4));
2050 assert!(
2051 estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
2052 );
2053 let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
2055 assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
2056 let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
2058 assert_eq!(estimate_context_tokens(¬e, 1000), None);
2059 assert_eq!(
2060 estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
2061 None
2062 );
2063 }
2064
2065 #[test]
2066 fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
2067 let cfg = ctx_config(&[("small-model", 1000)]);
2068 let t = ctx_talk(
2069 "small",
2070 vec![reply(
2071 &"x".repeat(5000),
2072 Some((10, "small", Some("small-model"))),
2073 )],
2074 );
2075 let u = context_usage(&t, Some(&cfg));
2076 assert_eq!((u.tokens, u.estimated), (Some(10), false));
2077 let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
2078 let u = context_usage(&t, Some(&cfg));
2079 assert!(u.estimated && !u.since_switch);
2080 assert_eq!(u.window, Some(1000));
2081 assert!(u.warn && u.percent.unwrap() >= 80);
2082 let t = ctx_talk("small", vec![reply("hi", None)]);
2083 let u = context_usage(&t, Some(&cfg));
2084 assert!(u.estimated && u.percent.is_some());
2085 }
2086
2087 #[test]
2088 fn context_usage_without_a_window_shows_tokens_only() {
2089 let cfg = ctx_config(&[]);
2090 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
2092 let u = context_usage(&t, Some(&cfg));
2093 assert_eq!(
2094 (u.tokens, u.window, u.percent, u.warn),
2095 (Some(5000), None, None, false)
2096 );
2097 let t = ctx_talk(
2098 "small",
2099 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
2100 );
2101 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
2102 assert_eq!(context_usage(&t, None).window, None);
2104 }
2105
2106 #[test]
2107 fn context_usage_switching_model_changes_the_denominator() {
2108 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
2109 let used = reply("hi", Some((900, "small", Some("small-model"))));
2110 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
2111 assert_eq!(
2112 (before.percent, before.warn, before.since_switch),
2113 (Some(90), true, false)
2114 );
2115 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
2118 assert_eq!(after.window, Some(10_000));
2119 assert_eq!(
2120 (after.percent, after.warn, after.since_switch),
2121 (Some(9), false, true)
2122 );
2123 assert_eq!(after.model.as_deref(), Some("big-model"));
2124 }
2125
2126 #[test]
2127 fn a_turn_recorded_before_usage_existed_still_reads() {
2128 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
2129 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
2130 assert!(turn.usage.is_none());
2131 let json = serde_json::to_string(&turn).expect("serialize");
2132 assert!(
2133 !json.contains("usage"),
2134 "absent usage is not written: {json}"
2135 );
2136 }
2137
2138 fn store() -> (tempfile::TempDir, Talks) {
2140 let tmp = tempfile::tempdir().expect("tempdir");
2141 let talks = Talks::at(tmp.path().join("talks"));
2142 (tmp, talks)
2143 }
2144
2145 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
2149 let path = dir.join("mock-talk-agent.sh");
2150 std::fs::write(&path, script).expect("write mock");
2151 AgentSpec {
2152 id: "mock".to_owned(),
2153 kind: AgentKind::Command,
2154 model: None,
2155 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2156 extra_args: Vec::new(),
2157 env,
2158 prompt_delivery: None,
2159 }
2160 }
2161
2162 fn config(spec: AgentSpec) -> Config {
2163 Config {
2164 agents: vec![spec],
2165 graph: Graph {
2166 language: "en".to_owned(),
2167 ..Graph::default()
2168 },
2169 ..Config::default()
2170 }
2171 }
2172
2173 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
2175
2176 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
2178
2179 const ECHO: &str = "#!/bin/sh\ncat\n";
2182
2183 fn env(reply: &str) -> BTreeMap<String, String> {
2184 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
2185 }
2186
2187 #[test]
2188 fn the_frozen_json_field_names_round_trip_through_disk() {
2189 let (tmp, talks) = store();
2190 let mut talk = Talk {
2191 schema: SCHEMA,
2192 id: "20260904-014455-ab12".to_owned(),
2193 repo: tmp.path().to_owned(),
2194 agent: "sonnet".to_owned(),
2195 status: TalkStatus::Open,
2196 turns: Vec::new(),
2197 pending: String::new(),
2198 pending_attachments: Vec::new(),
2199 fallback: false,
2200 created_at: Timestamp::now(),
2201 updated_at: Timestamp::now(),
2202 seat: SeatState::new(SEAT, "sonnet", 7),
2203 };
2204 talks.put(&mut talk).expect("put");
2205
2206 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
2207 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
2208 for field in [
2209 "schema",
2210 "id",
2211 "repo",
2212 "agent",
2213 "status",
2214 "turns",
2215 "created_at",
2216 "updated_at",
2217 ] {
2218 assert!(v.get(field).is_some(), "missing field `{field}`");
2219 }
2220 assert_eq!(v["schema"], 1);
2221 assert_eq!(v["status"], "open");
2222
2223 let back = talks.get(&talk.id).expect("get");
2224 assert_eq!(back.id, talk.id);
2225 assert_eq!(back.status, TalkStatus::Open);
2226 }
2227
2228 #[test]
2229 fn opening_a_talk_takes_no_agent_turn() {
2230 let (tmp, talks) = store();
2231 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2235 let cfg = config(spec);
2236
2237 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2238 assert_eq!(talk.status, TalkStatus::Open);
2239 assert!(talk.turns.is_empty(), "nothing has been said yet");
2240
2241 let on_disk = talks.get(&talk.id).expect("get");
2242 assert_eq!(on_disk.turns.len(), 0);
2243 }
2244
2245 #[test]
2253 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
2254 let (tmp, talks) = store();
2255 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2256 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2257 chatter_spec.id = "chatter-mock".to_owned();
2258
2259 let mut cfg = Config {
2260 agents: vec![first_spec.clone(), chatter_spec.clone()],
2261 graph: Graph {
2262 language: "en".to_owned(),
2263 ..Graph::default()
2264 },
2265 ..Config::default()
2266 };
2267 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
2268
2269 let talk =
2270 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
2271 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
2272
2273 cfg.roles.chatter = None;
2274 let fallback =
2275 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
2276 assert_eq!(
2277 fallback.agent, first_spec.id,
2278 "unset chatter must fall back to agent::pick's own default order"
2279 );
2280 }
2281
2282 #[test]
2285 fn a_talk_recorded_without_attachments_still_reads() {
2286 let (tmp, talks) = store();
2287 let path = talks.path_of("20260904-014455-ab12");
2288 std::fs::create_dir_all(talks.root()).expect("talks dir");
2289 std::fs::write(
2290 &path,
2291 serde_json::json!({
2292 "schema": 1,
2293 "id": "20260904-014455-ab12",
2294 "repo": tmp.path(),
2295 "agent": "sonnet",
2296 "status": "open",
2297 "turns": [
2298 { "who": "operator", "body": "still there?",
2299 "at": Timestamp::now().to_string() },
2300 ],
2301 "created_at": Timestamp::now().to_string(),
2302 "updated_at": Timestamp::now().to_string(),
2303 "seat": SeatState::new(SEAT, "sonnet", 7),
2304 })
2305 .to_string(),
2306 )
2307 .expect("write pre-attachments talk");
2308
2309 let talk = talks.get("20260904-014455-ab12").expect("must still read");
2310 assert!(talk.turns[0].attachments.is_empty());
2311 }
2312
2313 fn lease_store() -> (tempfile::TempDir, Talks, Talks) {
2314 let tmp = tempfile::TempDir::new().expect("tmp");
2315 let root = tmp.path().join("talks");
2316 (tmp, Talks::at(root.clone()), Talks::at(root))
2317 }
2318
2319 #[test]
2320 fn two_starters_on_one_talk_one_wins_and_the_other_is_refused() {
2321 let (_tmp, a, b) = lease_store();
2322 let won = a.claim_turn("t1").expect("claim").expect("first wins");
2323 assert!(
2324 b.claim_turn("t1").expect("claim").is_none(),
2325 "second is refused"
2326 );
2327 assert!(b.turn_held("t1"));
2328 assert!(
2329 b.claim_turn("t2").expect("claim").is_some(),
2330 "other talks are free"
2331 );
2332 drop(won);
2333 }
2334
2335 #[test]
2336 fn a_stale_lease_is_taken_over_and_the_old_guard_cannot_release_it() {
2337 let (_tmp, a, b) = lease_store();
2338 let old = a.claim_turn("t1").expect("claim").expect("held");
2339 let later = Timestamp::now()
2340 .checked_add(jiff::SignedDuration::from_secs(
2341 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2342 ))
2343 .expect("later");
2344 let new = b
2345 .claim_turn_at("t1", later)
2346 .expect("claim")
2347 .expect("a stale lease is taken over");
2348 drop(old);
2349 assert!(a.turn_held("t1"), "the old guard left the new lease alone");
2350 assert!(new.beat().expect("beat"), "the new owner still beats");
2351 drop(new);
2352 assert!(!a.turn_held("t1"));
2353 }
2354
2355 #[tokio::test]
2356 async fn a_turn_whose_lease_was_taken_over_is_stopped() {
2357 let (_tmp, a, b) = lease_store();
2358 let old = a.claim_turn("t1").expect("claim").expect("held");
2359 let later = Timestamp::now()
2360 .checked_add(jiff::SignedDuration::from_secs(
2361 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2362 ))
2363 .expect("later");
2364 let _new = b
2365 .claim_turn_at("t1", later)
2366 .expect("claim")
2367 .expect("taken over");
2368 let out = old
2369 .beating_every(Duration::from_millis(10), std::future::pending::<()>())
2370 .await;
2371 assert!(out.is_err(), "the displaced turn must stop, not run on");
2372 }
2373
2374 #[tokio::test]
2375 async fn a_turn_that_finishes_is_returned_and_keeps_its_lease_beating() {
2376 let (_tmp, a, _b) = lease_store();
2377 let lease = a.claim_turn("t1").expect("claim").expect("held");
2378 let out = lease
2379 .beating_every(Duration::from_millis(5), async {
2380 tokio::time::sleep(Duration::from_millis(40)).await;
2381 7
2382 })
2383 .await
2384 .expect("still ours");
2385 assert_eq!(out, 7);
2386 assert!(a.turn_held("t1"));
2387 }
2388
2389 #[test]
2390 fn an_unreadable_lease_counts_as_stale() {
2391 let (_tmp, a, b) = lease_store();
2392 std::fs::create_dir_all(a.root()).expect("dir");
2393 std::fs::write(a.turn_path("t1"), "not json").expect("write");
2394 assert!(!a.turn_held("t1"));
2395 assert!(b.claim_turn("t1").expect("claim").is_some());
2396 }
2397
2398 #[test]
2399 fn a_lease_is_released_when_the_turn_ends_or_fails() {
2400 let (_tmp, a, b) = lease_store();
2401 let lease = a.claim_turn("t1").expect("claim").expect("held");
2402 let failed: Result<()> = (|| {
2403 let _held = &lease;
2404 bail!("turn failed")
2405 })();
2406 assert!(failed.is_err());
2407 assert!(
2408 b.claim_turn("t1").expect("claim").is_none(),
2409 "held mid-turn"
2410 );
2411 drop(lease);
2412 assert!(
2413 b.claim_turn("t1").expect("claim").is_some(),
2414 "free after the turn"
2415 );
2416 }
2417
2418 #[test]
2419 fn concurrent_takeovers_of_a_stale_lease_have_one_winner() {
2420 let (_tmp, a, _b) = lease_store();
2421 drop(a.claim_turn("t1").expect("claim").expect("held"));
2422 std::fs::write(
2423 a.turn_path("t1"),
2424 serde_json::to_string(&TurnRecord {
2425 token: "gone".into(),
2426 pid: 1,
2427 beat_at: Timestamp::from_second(1).expect("ts"),
2428 })
2429 .expect("json"),
2430 )
2431 .expect("write");
2432 let wins: Vec<_> = std::thread::scope(|sc| {
2433 let hs: Vec<_> = (0..8)
2434 .map(|_| {
2435 let s = a.clone();
2436 sc.spawn(move || s.claim_turn("t1").expect("claim"))
2437 })
2438 .collect();
2439 hs.into_iter().map(|h| h.join().expect("join")).collect()
2440 });
2441 assert_eq!(wins.iter().filter(|w| w.is_some()).count(), 1);
2442 }
2443
2444 fn age_file(path: &Path) {
2445 let f = std::fs::OpenOptions::new()
2446 .write(true)
2447 .open(path)
2448 .expect("open");
2449 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
2450 .expect("age");
2451 }
2452
2453 #[test]
2454 fn a_late_taker_of_a_broken_lock_cannot_disturb_its_replacement() {
2455 let (_tmp, a, _b) = lease_store();
2456 let lease = a.turn_path("t1");
2457 let lock = lease.with_extension("turn.lock");
2458 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2459 std::fs::write(&lock, "t1-dead").expect("dead lock");
2460 age_file(&lock);
2461 let b = TurnLock::take(&lease).expect("take").expect("b wins");
2463 let fresh = std::fs::read_to_string(&lock).expect("read");
2464 assert_eq!(fresh, b.token);
2465 let ticket = std::fs::read_dir(lock.parent().expect("dir"))
2468 .expect("dir")
2469 .flatten()
2470 .map(|e| e.path())
2471 .find(|p| p.to_string_lossy().ends_with(".break"))
2472 .expect("ticket");
2473 assert!(!create_exclusive(&ticket, "").expect("ticket"));
2474 assert!(TurnLock::take(&lease).expect("take").is_none());
2475 assert_eq!(std::fs::read_to_string(&lock).expect("read"), fresh);
2476 age_file(&lock);
2478 std::mem::forget(b);
2479 let c = TurnLock::take(&lease)
2480 .expect("take")
2481 .expect("next generation");
2482 assert_ne!(c.token, fresh);
2483 }
2484
2485 #[test]
2486 fn an_aged_empty_lock_is_broken() {
2487 let (_tmp, a, _b) = lease_store();
2488 let lease = a.turn_path("t1");
2489 let lock = lease.with_extension("turn.lock");
2490 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2491 std::fs::write(&lock, "").expect("empty lock");
2492 age_file(&lock);
2493 assert!(TurnLock::take(&lease).expect("take").is_some());
2494 }
2495
2496 #[test]
2497 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
2498 let (tmp, talks) = store();
2499 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2500 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2501
2502 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
2503 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
2504 let saved = talks.get(&talk.id).expect("reload queued talk");
2505 assert_eq!(saved.pending, "first\n\nsecond");
2506 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
2507
2508 let drained = drain(&mut talk, &talks).expect("drain");
2509 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
2510 let saved = talks.get(&talk.id).expect("reload drained talk");
2511 assert!(saved.pending.is_empty());
2512 assert_eq!(saved.turns.len(), 1);
2513 assert_eq!(saved.turns[0].body, "first\n\nsecond");
2514 }
2515
2516 #[test]
2517 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
2518 let (tmp, talks) = store();
2519 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2520 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2521 let attachment = Attachment {
2522 id: "a".repeat(32),
2523 name: "shot.png".to_owned(),
2524 mime: "image/png".to_owned(),
2525 bytes: 3,
2526 };
2527
2528 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
2529 assert!(
2530 edit_pending_text(
2531 &mut talk,
2532 &talks,
2533 "corrected",
2534 "first",
2535 std::slice::from_ref(&attachment.id),
2536 )
2537 .expect("edit")
2538 );
2539 let saved = talks.get(&talk.id).expect("reload edited draft");
2540 assert_eq!(saved.pending, "corrected");
2541 assert_eq!(saved.pending_attachments, vec![attachment]);
2542
2543 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
2544 assert!(
2545 !edit_pending_text(
2546 &mut talk,
2547 &talks,
2548 "stale edit",
2549 "corrected",
2550 &["a".repeat(32)],
2551 )
2552 .expect("stale edit is a conflict")
2553 );
2554 assert_eq!(
2555 talks.get(&talk.id).expect("reload after conflict").pending,
2556 "corrected\n\nlater"
2557 );
2558 assert!(
2559 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
2560 .expect("stale clear is a conflict")
2561 );
2562 assert_eq!(
2563 talks
2564 .get(&talk.id)
2565 .expect("reload after stale clear")
2566 .pending,
2567 "corrected\n\nlater"
2568 );
2569 }
2570
2571 #[tokio::test]
2572 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
2573 let (tmp, talks) = store();
2574 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
2575 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
2576 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2577 let id = running.id.clone();
2578 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
2579
2580 let response_talks = talks.clone();
2581 let response_cfg = cfg.clone();
2582 let reply = tokio::spawn(async move {
2583 respond(&mut running, &response_talks, &response_cfg, &first).await
2584 });
2585 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
2586
2587 let mut queued = talks.get(&id).expect("queued handle");
2588 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
2589 reply.await.expect("join").expect("reply");
2590
2591 let saved = talks.get(&id).expect("reload");
2592 assert_eq!(saved.pending, "next");
2593 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
2594 }
2595
2596 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
2599 let calls = dir.join(format!("{id}.calls"));
2600 let script = format!(
2601 "#!/bin/sh\necho x >> '{}'\n{body}\n",
2602 calls.to_string_lossy()
2603 );
2604 let path = dir.join(format!("mock-{id}.sh"));
2605 std::fs::write(&path, script).expect("write mock");
2606 AgentSpec {
2607 id: id.to_owned(),
2608 kind: AgentKind::Command,
2609 model: None,
2610 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2611 extra_args: Vec::new(),
2612 env: BTreeMap::new(),
2613 prompt_delivery: None,
2614 }
2615 }
2616
2617 fn calls(dir: &Path, id: &str) -> usize {
2618 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
2619 }
2620
2621 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
2622 let mut cfg = config(specs[0].clone());
2623 cfg.agents = specs;
2624 cfg.roles.chatter = Some(AgentChoice::Chain(
2625 ids.iter().map(|s| (*s).to_owned()).collect(),
2626 ));
2627 cfg
2628 }
2629
2630 #[tokio::test]
2631 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
2632 let (tmp, talks) = store();
2633 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2634 let b = counting_agent(tmp.path(), "b", "cat");
2635 let cfg = chain_config(vec![a, b], &["a", "b"]);
2636 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2637 assert_eq!(talk.agent, "a");
2638
2639 say(&mut talk, &talks, &cfg, "hello there", Vec::new())
2640 .await
2641 .expect("turn");
2642 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
2643 assert_eq!(calls(tmp.path(), "b"), 1);
2644 assert_eq!(talk.agent, "b", "the switch persists");
2645 assert!(talks.get(&talk.id).unwrap().agent == "b");
2646 let reply = talk.turns.last().unwrap();
2647 assert!(reply.body.contains("hello there"));
2648 assert!(
2649 reply.body.contains("magi task add --solo"),
2650 "a fresh seat gets the full briefing"
2651 );
2652 assert!(
2653 talk.turns
2654 .iter()
2655 .any(|t| t.body.contains("agent changed from a to b")),
2656 "the switch is noted"
2657 );
2658 }
2659
2660 #[tokio::test]
2661 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
2662 let (tmp, talks) = store();
2663 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2664 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
2665 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
2666 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2667
2668 let err = say(&mut talk, &talks, &cfg, "hi", Vec::new())
2669 .await
2670 .expect_err("every agent failed");
2671 assert!(err.to_string().contains("`a`"), "{err:#}");
2672 assert_eq!(calls(tmp.path(), "a"), 1);
2673 assert_eq!(calls(tmp.path(), "b"), 1);
2674 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
2675 }
2676
2677 #[test]
2678 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
2679 let (tmp, talks) = store();
2680 let b = counting_agent(tmp.path(), "b", "cat");
2681 let cfg = chain_config(vec![b], &["ghost", "b"]);
2682 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2683 assert_eq!(talk.agent, "b");
2684 }
2685
2686 #[tokio::test]
2687 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
2688 let (tmp, talks) = store();
2689 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2690 let b = counting_agent(tmp.path(), "b", "cat");
2691 let cfg = chain_config(vec![a, b], &["a", "b"]);
2692 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
2693 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2694 .await
2695 .expect_err("a alone, and it fails");
2696 assert_eq!(calls(tmp.path(), "b"), 0);
2697 assert_eq!(talk.agent, "a");
2698 }
2699
2700 #[tokio::test]
2701 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
2702 let (tmp, talks) = store();
2703 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2704 let b = counting_agent(tmp.path(), "b", "cat");
2705 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
2706 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
2707 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
2708 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2709 .await
2710 .expect_err("c alone, and it fails");
2711 assert_eq!(calls(tmp.path(), "b"), 0);
2712 }
2713
2714 #[tokio::test]
2715 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
2716 let (tmp, talks) = store();
2717 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2718 let cfg = config(spec);
2719 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2720
2721 say(
2722 &mut talk,
2723 &talks,
2724 &cfg,
2725 "what does the queue module do?",
2726 Vec::new(),
2727 )
2728 .await
2729 .expect("first turn");
2730 let first_prompt = &talk.turns[1].body;
2731 assert!(first_prompt.contains("magi task add --solo"));
2732 assert!(first_prompt.contains("what does the queue module do?"));
2733
2734 say(&mut talk, &talks, &cfg, "and how is it locked?", Vec::new())
2735 .await
2736 .expect("second turn");
2737 let second_prompt = &talk.turns[3].body;
2738 assert!(
2739 !second_prompt.contains("magi task add --solo"),
2740 "the briefing is sent once, not on every turn: {second_prompt}"
2741 );
2742 assert!(second_prompt.contains("and how is it locked?"));
2743 }
2744
2745 #[tokio::test]
2746 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
2747 let (tmp, talks) = store();
2748 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2749 let mut b = a.clone();
2750 b.id = "other".to_owned();
2751 let mut cfg = config(a.clone());
2752 cfg.agents.push(b.clone());
2753 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
2754 say(&mut talk, &talks, &cfg, "remember the walrus", Vec::new())
2755 .await
2756 .expect("first turn");
2757 let old_session = talk.seat.claude_session.clone();
2758 assert_eq!(talk.seat.turns, 1);
2759
2760 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
2761 assert_eq!(talk.agent, "other");
2762 assert_eq!(talk.seat.turns, 0);
2763 assert_eq!(talk.seat.agent, "other");
2764 assert_ne!(talk.seat.claude_session, old_session);
2765 let note = talk.turns.last().expect("note");
2766 assert_eq!(note.who, Who::Agent);
2767 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2768 assert!(note.body.contains("changed from"), "{}", note.body);
2769 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
2770
2771 let before = talk.turns.len();
2772 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
2773 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
2774
2775 say(&mut talk, &talks, &cfg, "what did I say?", Vec::new())
2776 .await
2777 .expect("turn after switch");
2778 let prompt = &talk.turns.last().expect("reply").body;
2779 assert!(prompt.contains("remember the walrus"), "{prompt}");
2780 assert!(prompt.contains("## magi"), "{prompt}");
2781 assert!(prompt.contains("what did I say?"), "{prompt}");
2782 }
2783
2784 #[tokio::test]
2785 async fn say_appends_the_operator_turn_then_the_agent_turn() {
2786 let (tmp, talks) = store();
2787 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2788 let cfg = config(spec);
2789 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2790
2791 say(
2792 &mut talk,
2793 &talks,
2794 &cfg,
2795 "can I rename this function?",
2796 Vec::new(),
2797 )
2798 .await
2799 .expect("say");
2800
2801 assert_eq!(talk.turns.len(), 2);
2802 assert_eq!(talk.turns[0].who, Who::Operator);
2803 assert_eq!(talk.turns[0].body, "can I rename this function?");
2804 assert_eq!(talk.turns[1].who, Who::Agent);
2805 assert_eq!(talk.turns[1].body, "go ahead");
2806 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
2807 }
2808
2809 #[tokio::test]
2810 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
2811 let (tmp, talks) = store();
2812 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2813 let cfg = config(spec);
2814 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2815
2816 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
2817 .await
2818 .expect_err("a turn with no answer is an error");
2819 assert!(err.to_string().contains("no answer"), "{err}");
2820
2821 let on_disk = talks.get(&talk.id).expect("get");
2822 assert_eq!(on_disk.turns.len(), 2);
2823 assert_eq!(on_disk.turns[0].body, "check the tests");
2824 let note = &on_disk.turns[1];
2825 assert_eq!(note.who, Who::Agent);
2826 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2827 assert!(note.body.contains("your message is saved"));
2828 }
2829
2830 #[tokio::test]
2836 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
2837 let (tmp, talks) = store();
2838 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2839 let cfg = config(spec);
2840 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2841
2842 let text =
2843 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
2844 failpoint::force_put_failures(PUT_RETRIES - 1);
2847 respond(&mut talk, &talks, &cfg, &text)
2848 .await
2849 .expect("respond must survive a write failure its own retries can outlast");
2850
2851 assert_eq!(talk.turns.len(), 2);
2852 assert_eq!(talk.turns[1].who, Who::Agent);
2853 assert_eq!(talk.turns[1].body, "go ahead");
2854 let on_disk = talks.get(&talk.id).expect("get");
2855 assert_eq!(
2856 on_disk.turns, talk.turns,
2857 "the reply must reach disk despite the early write failures"
2858 );
2859 }
2860
2861 #[tokio::test]
2867 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
2868 let (tmp, talks) = store();
2869 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2870 let cfg = config(spec);
2871 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2872
2873 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2874 failpoint::force_put_failures(PUT_RETRIES);
2879 let err = respond(&mut talk, &talks, &cfg, &text)
2880 .await
2881 .expect_err("a reply that cannot be saved must be reported, not swallowed");
2882 assert!(err.to_string().contains("could not be saved"), "{err}");
2883
2884 let on_disk = talks.get(&talk.id).expect("get");
2885 assert_eq!(
2886 on_disk.turns.len(),
2887 2,
2888 "the operator turn plus a visible note"
2889 );
2890 assert_eq!(on_disk.turns[0].body, "check the tests");
2891 let note = &on_disk.turns[1];
2892 assert_eq!(note.who, Who::Agent);
2893 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
2894 assert!(
2895 note.body.contains("could not be saved"),
2896 "the operator must be told the reply is missing, not left staring \
2897 at a gap with no explanation: {}",
2898 note.body
2899 );
2900 assert_eq!(
2901 talk.turns, on_disk.turns,
2902 "the in-memory talk must match what actually landed on disk"
2903 );
2904
2905 let artifacts = talks.artifacts_of(&talk.id);
2908 let stash = std::fs::read_dir(&artifacts)
2909 .expect("artifacts dir")
2910 .filter_map(|e| e.ok())
2911 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
2912 .expect("a stash file for the lost reply");
2913 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
2914 assert_eq!(stashed, "go ahead");
2915
2916 assert_eq!(
2924 on_disk.seat.turns, 1,
2925 "the note's write must carry the turn the CLI actually took"
2926 );
2927 assert_eq!(
2928 on_disk.seat.claude_session, talk.seat.claude_session,
2929 "the session id handed to the CLI must survive the failed reply"
2930 );
2931 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
2932 assert!(
2933 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
2934 "the next turn must resume, not open the same session id twice"
2935 );
2936 }
2937
2938 #[tokio::test]
2943 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
2944 let (tmp, talks) = store();
2945 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
2946 let cfg = config(spec);
2947 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2948
2949 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
2950 failpoint::force_put_failures(PUT_RETRIES * 2);
2953 let err = respond(&mut talk, &talks, &cfg, &text)
2954 .await
2955 .expect_err("neither the reply nor the note could be saved");
2956 assert!(err.to_string().contains("could not be saved"), "{err}");
2957
2958 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
2959 let on_disk = talks.get(&talk.id).expect("get");
2960 assert_eq!(on_disk.turns.len(), 1);
2961
2962 assert_eq!(
2972 on_disk.seat.turns, 0,
2973 "an unwritable file cannot record the turn the CLI took"
2974 );
2975 assert_eq!(
2976 talk.seat.turns, 1,
2977 "the in-memory seat still reports the turn the CLI actually took"
2978 );
2979 assert_eq!(
2980 on_disk.seat.claude_session, talk.seat.claude_session,
2981 "the session id was minted at `begin` and never changes here"
2982 );
2983 }
2984
2985 #[tokio::test]
2989 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
2990 let (tmp, talks) = store();
2991 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2992 let cfg = config(spec);
2993 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2994
2995 let att = talks
2996 .put_attachment(
2997 &talk.id,
2998 "image/png",
2999 "screenshot.png",
3000 b"pretend-png-bytes",
3001 )
3002 .expect("put attachment");
3003
3004 say(&mut talk, &talks, &cfg, "", vec![att.clone()])
3005 .await
3006 .expect("an empty body with an attachment is still a turn");
3007
3008 let operator_turn = &talk.turns[0];
3009 assert_eq!(operator_turn.who, Who::Operator);
3010 assert_eq!(operator_turn.body, "");
3011 assert_eq!(operator_turn.attachments, vec![att.clone()]);
3012
3013 let prompt = &talk.turns[1].body;
3014 let expected_path = talks
3015 .attachments_dir(&talk.id)
3016 .join(format!("{}.png", att.id));
3017 assert!(
3018 prompt.contains(&expected_path.display().to_string()),
3019 "the agent must be told the attachment's absolute path: {prompt}"
3020 );
3021 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
3022 }
3023
3024 #[test]
3033 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
3034 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
3035 let att = Attachment {
3036 id: "0".repeat(32),
3037 name: "shot.png".to_owned(),
3038 mime: "image/png".to_owned(),
3039 bytes: 3,
3040 };
3041 let path = talks
3042 .attachment_path("some-talk-id", &att)
3043 .expect("a supported mime always yields a path");
3044 assert!(
3045 path.is_absolute(),
3046 "must be absolute even off a relative store root: {}",
3047 path.display()
3048 );
3049 }
3050
3051 #[tokio::test]
3052 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
3053 let (tmp, talks) = store();
3058 let slow = mock_agent(
3059 tmp.path(),
3060 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
3061 BTreeMap::new(),
3062 );
3063 let mut cfg = config(slow);
3064 cfg.graph.timeout_talk = 1;
3065 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3066
3067 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
3068 .await
3069 .expect_err("a turn that never answers is an error");
3070 assert!(
3071 err.to_string().contains("did not answer within 1s"),
3072 "{err}"
3073 );
3074
3075 let on_disk = talks.get(&talk.id).expect("get");
3076 let note = on_disk.turns.last().expect("a note turn was recorded");
3077 assert!(
3078 note.body.contains("did not answer within 1s"),
3079 "the transcript must show the configured timeout: {}",
3080 note.body
3081 );
3082 }
3083
3084 #[test]
3085 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
3086 let (tmp, talks) = store();
3087 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3088 let cfg = config(spec);
3089 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3090
3091 close(&mut talk, &talks).expect("close");
3092 assert_eq!(talk.status, TalkStatus::Closed);
3093 close(&mut talk, &talks).expect("closing twice is not an error");
3094
3095 let err =
3096 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
3097 assert!(err.to_string().contains("closed"));
3098 let _ = &cfg; }
3100
3101 #[tokio::test]
3102 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3103 let (tmp, talks) = store();
3104 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3105 let cfg = config(spec);
3106 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3109
3110 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
3114 close(&mut closed_elsewhere, &talks).expect("close");
3115 assert_eq!(
3116 talks.get(&in_flight.id).expect("reread").status,
3117 TalkStatus::Closed,
3118 "the close landed on disk before the turn finished"
3119 );
3120
3121 assert_eq!(in_flight.status, TalkStatus::Open);
3125 respond(&mut in_flight, &talks, &cfg, "one more question")
3126 .await
3127 .expect("the turn itself still completes");
3128
3129 let on_disk = talks.get(&in_flight.id).expect("reread");
3130 assert_eq!(
3131 on_disk.status,
3132 TalkStatus::Closed,
3133 "a close must stick even when a turn that started before it finishes after it"
3134 );
3135 assert!(
3138 on_disk.turns.iter().any(|t| t.body == "here you go"),
3139 "the in-flight turn's own reply is still recorded: {:?}",
3140 on_disk.turns
3141 );
3142 }
3143
3144 #[test]
3145 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
3146 let (tmp, talks) = store();
3147 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3148 let cfg = config(spec);
3149 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3152
3153 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
3156 close(&mut closed_elsewhere, &talks).expect("close");
3157 assert_eq!(
3158 talks.get(&stale.id).expect("reread").status,
3159 TalkStatus::Closed,
3160 "the close landed on disk before record was called"
3161 );
3162
3163 assert_eq!(stale.status, TalkStatus::Open);
3167 let err = record(&mut stale, &talks, "still there?", Vec::new())
3168 .expect_err("a close that landed first must be honored, not overwritten");
3169 assert!(err.to_string().contains("closed"));
3170
3171 let on_disk = talks.get(&stale.id).expect("reread");
3172 assert_eq!(
3173 on_disk.status,
3174 TalkStatus::Closed,
3175 "record must not resurrect a conversation closed while its snapshot was stale"
3176 );
3177 assert!(
3178 on_disk.turns.is_empty(),
3179 "the rejected turn must not have been appended: {:?}",
3180 on_disk.turns
3181 );
3182 let _ = &cfg; }
3184
3185 #[test]
3186 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
3187 let (tmp, talks) = store();
3188 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3189 let cfg = config(spec);
3190 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3191
3192 let held = talks.guard();
3196
3197 let talks2 = talks.clone();
3198 let id = talk.id.clone();
3199 let closing = std::thread::spawn(move || {
3200 let mut talk = talks2.get(&id).expect("get");
3201 close(&mut talk, &talks2).expect("close");
3202 });
3203
3204 std::thread::sleep(Duration::from_millis(50));
3205 assert!(
3206 !closing.is_finished(),
3207 "close must wait for the guard, not read and write while it is held - \
3208 a re-read alone narrows this window without closing it"
3209 );
3210
3211 drop(held);
3212 closing.join().expect("close thread panicked");
3213
3214 assert_eq!(
3215 talks.get(&talk.id).expect("reread").status,
3216 TalkStatus::Closed,
3217 "once the guard is free, close still lands"
3218 );
3219 let _ = &cfg; }
3221
3222 #[test]
3223 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
3224 let (tmp, talks) = store();
3225 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3226 let cfg = config(spec);
3227 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3228
3229 close(&mut talk, &talks).expect("close");
3230 assert_eq!(talk.status, TalkStatus::Closed);
3231
3232 reopen(&mut talk, &talks).expect("reopen");
3233 assert_eq!(talk.status, TalkStatus::Open);
3234 assert_eq!(
3235 talks.get(&talk.id).expect("reread").status,
3236 TalkStatus::Open
3237 );
3238
3239 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
3241 assert_eq!(talk.status, TalkStatus::Open);
3242
3243 record(&mut talk, &talks, "one more thing", Vec::new())
3244 .expect("a reopened talk takes turns again");
3245 let _ = &cfg; }
3247
3248 #[test]
3249 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
3250 let (tmp, talks) = store();
3251 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3252 let cfg = config(spec);
3253 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3254
3255 let artifacts = talks.artifacts_of(&talk.id);
3256 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
3257 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
3258
3259 talks.remove(&talk.id).expect("remove");
3260 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
3261 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
3262 assert!(
3263 talks.get(&talk.id).is_err(),
3264 "a removed talk cannot be read back"
3265 );
3266
3267 let err = talks
3268 .remove("nonexistent-id")
3269 .expect_err("unknown id refused");
3270 assert!(err.to_string().contains("no talk matches"), "{err}");
3271 let _ = &cfg; }
3273
3274 #[tokio::test]
3275 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3276 let (tmp, talks) = store();
3277 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3278 let cfg = config(spec);
3279 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3282
3283 talks.remove(&in_flight.id).expect("remove");
3284 assert!(
3285 talks.get(&in_flight.id).is_err(),
3286 "the delete landed on disk before the turn finished"
3287 );
3288
3289 respond(&mut in_flight, &talks, &cfg, "one more question")
3292 .await
3293 .expect("the turn itself still completes rather than erroring");
3294
3295 assert!(
3296 talks.get(&in_flight.id).is_err(),
3297 "a delete must stick even when a turn that started before it finishes after it"
3298 );
3299 }
3300
3301 #[test]
3302 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
3303 let (tmp, talks) = store();
3304 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3305 let cfg = config(spec);
3306 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3309
3310 talks.remove(&stale.id).expect("remove");
3311
3312 let err = record(&mut stale, &talks, "still there?", Vec::new())
3316 .expect_err("a delete that landed first must be honored, not overwritten");
3317 assert!(err.to_string().contains("deleted"), "{err}");
3318
3319 assert!(
3320 talks.get(&stale.id).is_err(),
3321 "record must not resurrect a conversation deleted while its snapshot was stale"
3322 );
3323 let _ = &cfg; }
3325
3326 #[test]
3327 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
3328 let (tmp, talks) = store();
3329 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3330 let cfg = config(spec);
3331 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3334
3335 talks.remove(&stale.id).expect("remove");
3336
3337 let err = close(&mut stale, &talks)
3341 .expect_err("a delete that landed first must be honored, not overwritten");
3342 assert!(err.to_string().contains("deleted"), "{err}");
3343
3344 assert!(
3345 talks.get(&stale.id).is_err(),
3346 "close must not resurrect a conversation deleted while its snapshot was stale"
3347 );
3348 let _ = &cfg; }
3350
3351 #[test]
3352 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
3353 let (tmp, talks) = store();
3354 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3355 let cfg = config(spec);
3356 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3359 close(&mut stale, &talks).expect("close");
3360
3361 talks.remove(&stale.id).expect("remove");
3362
3363 let err = reopen(&mut stale, &talks)
3367 .expect_err("a delete that landed first must be honored, not overwritten");
3368 assert!(err.to_string().contains("deleted"), "{err}");
3369
3370 assert!(
3371 talks.get(&stale.id).is_err(),
3372 "reopen must not resurrect a conversation deleted while its snapshot was stale"
3373 );
3374 let _ = &cfg; }
3376
3377 #[test]
3378 fn list_puts_open_talks_before_closed_ones() {
3379 let (tmp, talks) = store();
3380 let make = |id: &str, status: TalkStatus| {
3381 let mut t = Talk {
3382 schema: SCHEMA,
3383 id: id.to_owned(),
3384 repo: tmp.path().to_owned(),
3385 agent: "mock".to_owned(),
3386 status,
3387 turns: Vec::new(),
3388 pending: String::new(),
3389 pending_attachments: Vec::new(),
3390 fallback: false,
3391 created_at: Timestamp::now(),
3392 updated_at: Timestamp::now(),
3393 seat: SeatState::new(SEAT, "mock", 7),
3394 };
3395 talks.put(&mut t).expect("put");
3396 };
3397 make("20260901-000000-0001", TalkStatus::Open);
3398 make("20260902-000000-0002", TalkStatus::Open);
3399 make("20260903-000000-0003", TalkStatus::Closed);
3400
3401 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
3402 assert_eq!(
3403 ids,
3404 [
3405 "20260902-000000-0002",
3406 "20260901-000000-0001",
3407 "20260903-000000-0003"
3408 ]
3409 );
3410 assert_eq!(talks.count_open(), 2);
3411 }
3412
3413 #[test]
3414 fn tasks_of_finds_only_this_talks_own_tasks() {
3415 let dir = tempfile::tempdir().expect("tempdir");
3416 let queue = Queue::at(dir.path().join("queue"));
3417
3418 let mut mine = Task::new(
3419 "rework the loader".to_owned(),
3420 "rework the loader".to_owned(),
3421 PathBuf::from("/repo"),
3422 Source::Agent {
3423 run: "20260904-014455-ab12".to_owned(),
3424 node: "chat".to_owned(),
3425 },
3426 );
3427 queue.put(&mut mine).expect("put mine");
3428
3429 let mut theirs = Task::new(
3430 "unrelated".to_owned(),
3431 "unrelated".to_owned(),
3432 PathBuf::from("/repo"),
3433 Source::Agent {
3434 run: "20260904-090000-zz99".to_owned(),
3435 node: "implement".to_owned(),
3436 },
3437 );
3438 queue.put(&mut theirs).expect("put theirs");
3439
3440 let mut human = Task::new(
3441 "typed by hand".to_owned(),
3442 "typed by hand".to_owned(),
3443 PathBuf::from("/repo"),
3444 Source::Human,
3445 );
3446 queue.put(&mut human).expect("put human");
3447
3448 let found = tasks_of(&queue, "20260904-014455-ab12");
3449 assert_eq!(found.len(), 1);
3450 assert_eq!(found[0].id, mine.id);
3451 }
3452
3453 #[test]
3454 fn the_briefing_names_solo_task_add() {
3455 let brief = briefing(Path::new("/repo"), "en", false);
3456 assert!(brief.contains("magi task add --solo"));
3457 assert!(brief.contains("/repo"));
3458 assert!(!brief.contains("Hold this conversation in"));
3459 }
3460
3461 #[test]
3467 fn the_briefing_explains_targeting_a_different_repository_by_name() {
3468 let brief = briefing(Path::new("/repo"), "en", false);
3469 assert!(brief.contains("--repo does not have to be a full path"));
3470 assert!(brief.contains("owner/repo"));
3471 assert!(brief.contains("magi repos"));
3472 assert!(brief.contains("ask the operator"));
3473 }
3474
3475 #[test]
3476 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
3477 let brief = briefing(Path::new("/repo"), "en", false);
3478 assert!(brief.contains("--attach <path>"), "{brief}");
3479 assert!(brief.contains("deleting this conversation"), "{brief}");
3480 }
3481
3482 #[test]
3483 fn the_briefing_names_the_language_when_it_is_not_english() {
3484 let brief = briefing(Path::new("/repo"), "Japanese", false);
3485 assert!(brief.contains("Hold this conversation in Japanese"));
3486 }
3487
3488 #[test]
3489 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
3490 let read_only = briefing(Path::new("/repo"), "en", false);
3491 assert!(read_only.contains("Do not write files"));
3492 assert!(!read_only.contains("allow_write"));
3493
3494 let writable = briefing(Path::new("/repo"), "en", true);
3495 assert!(!writable.contains("Do not write files"));
3496 assert!(writable.contains("allow_write = true"));
3497 assert!(writable.contains("magi task add --solo"));
3500 assert!(writable.contains("say plainly what you"));
3501 }
3502}