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_with(
239 &talk.repo,
240 &c.graph.language,
241 c.talk.allow_write,
242 crate::persona::active(&c.talk.personas, &talk.persona).as_ref(),
243 )
244 .chars()
245 .count() as u64
246 });
247 estimate_context_tokens(talk, standing)
248 });
249 let estimated = measured.is_none() && tokens.is_some();
250 let since_switch =
251 usage.is_some_and(|u| u.agent != talk.agent || (current.is_some() && u.model != model));
252 let (percent, warn) = match (tokens, window) {
253 (Some(t), Some(w)) => (
254 Some(t.saturating_mul(100) / w),
255 t.saturating_mul(100) >= w.saturating_mul(CONTEXT_WARN_PERCENT),
256 ),
257 _ => (None, false),
258 };
259 ContextUsage {
260 tokens,
261 window,
262 percent,
263 warn,
264 since_switch,
265 model,
266 estimated,
267 }
268}
269
270#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
274#[serde(rename_all = "lowercase")]
275pub enum TalkStatus {
276 Open,
279 Closed,
281}
282
283impl TalkStatus {
284 pub fn open(self) -> bool {
286 matches!(self, Self::Open)
287 }
288
289 pub fn as_str(self) -> &'static str {
291 match self {
292 Self::Open => "open",
293 Self::Closed => "closed",
294 }
295 }
296}
297
298#[derive(Debug, Clone, Serialize, Deserialize)]
300#[serde(deny_unknown_fields)]
301pub struct Talk {
302 pub schema: u32,
304 pub id: String,
306 pub repo: PathBuf,
308 pub agent: String,
310 pub status: TalkStatus,
312 pub turns: Vec<Turn>,
314 #[serde(default)]
317 pub pending: String,
318 #[serde(default)]
320 pub pending_attachments: Vec<Attachment>,
321 #[serde(default)]
326 pub fallback: bool,
327 #[serde(default)]
330 pub persona: String,
331 #[serde(default)]
335 pub persona_dirty: bool,
336 pub created_at: Timestamp,
338 pub updated_at: Timestamp,
340 seat: SeatState,
345}
346
347impl Talk {
348 pub fn short(&self) -> &str {
350 short(&self.id)
351 }
352}
353
354#[derive(Debug, Clone)]
356pub struct Talks {
357 root: PathBuf,
358 lock: Arc<Mutex<()>>,
367}
368
369impl Talks {
370 pub fn open() -> Self {
372 Self::at(crate::run::home().join("talks"))
373 }
374
375 pub fn at(root: PathBuf) -> Self {
378 Self {
379 root,
380 lock: Arc::new(Mutex::new(())),
381 }
382 }
383
384 fn guard(&self) -> MutexGuard<'_, ()> {
393 self.lock.lock().unwrap_or_else(PoisonError::into_inner)
394 }
395
396 pub fn root(&self) -> &Path {
398 &self.root
399 }
400
401 pub fn path_of(&self, id: &str) -> PathBuf {
403 self.root.join(format!("{id}.json"))
404 }
405
406 pub fn artifacts_of(&self, id: &str) -> PathBuf {
409 self.root.join(format!("{id}.artifacts"))
410 }
411
412 pub fn attachments_dir(&self, id: &str) -> PathBuf {
416 self.artifacts_of(id).join("attachments")
417 }
418
419 pub fn put_attachment(
428 &self,
429 id: &str,
430 mime: &str,
431 name: &str,
432 data: &[u8],
433 ) -> Result<Attachment> {
434 let dir = self.attachments_dir(id);
435 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
436 let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
437 let att = Attachment {
438 id: new_attachment_id(),
439 name: name.to_owned(),
440 mime: mime.to_owned(),
441 bytes: data.len() as u64,
442 };
443 std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
444 .with_context(|| format!("write attachment {}", att.id))?;
445 std::fs::write(
446 dir.join(format!("{}.json", att.id)),
447 serde_json::to_string(&att).context("serialize attachment")?,
448 )
449 .with_context(|| format!("write attachment metadata {}", att.id))?;
450 Ok(att)
451 }
452
453 pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
462 if !valid_attachment_id(att_id) {
463 return Ok(None);
464 }
465 let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
466 if !meta_path.is_file() {
467 return Ok(None);
468 }
469 let att = serde_json::from_str(
470 &std::fs::read_to_string(&meta_path)
471 .with_context(|| format!("read {}", meta_path.display()))?,
472 )
473 .with_context(|| format!("parse {}", meta_path.display()))?;
474 Ok(Some(att))
475 }
476
477 pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
481 let Some(att) = self.attachment_meta(id, att_id)? else {
482 return Ok(None);
483 };
484 let ext = attachment_ext(&att.mime).with_context(|| {
485 format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
486 })?;
487 let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
488 let data =
489 std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
490 Ok(Some((att, data)))
491 }
492
493 fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
510 let ext = attachment_ext(&att.mime)?;
511 let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
512 std::path::absolute(&path).ok()
513 }
514
515 pub fn put(&self, t: &mut Talk) -> Result<()> {
524 std::fs::create_dir_all(&self.root)
525 .with_context(|| format!("create {}", self.root.display()))?;
526 t.updated_at = Timestamp::now();
527 let body = serde_json::to_string_pretty(t).context("serialize talk")?;
528 let path = self.path_of(&t.id);
529 let tmp = path.with_extension("json.tmp");
530 write_atomic(&tmp, &path, &body)
531 }
532
533 pub fn get(&self, id: &str) -> Result<Talk> {
535 let resolved = self.resolve_id(id)?;
536 read_path(&self.path_of(&resolved))
537 }
538
539 pub fn list(&self) -> Vec<Talk> {
542 self.list_counting_unreadable().0
543 }
544
545 pub fn list_counting_unreadable(&self) -> (Vec<Talk>, usize) {
548 let mut unreadable = 0;
549 let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
550 .into_iter()
551 .flatten()
552 .flatten()
553 .map(|e| e.path())
554 .filter(|p| p.extension().is_some_and(|x| x == "json"))
555 .filter_map(|p| {
556 let talk = read_path(&p).ok();
557 if talk.is_none() {
558 unreadable += 1;
559 }
560 talk
561 })
562 .collect();
563 all.sort_unstable_by(|a, b| {
564 let rank = |t: &Talk| u8::from(!t.status.open());
565 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
566 });
567 (all, unreadable)
568 }
569
570 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
572 if self.path_of(prefix).is_file() {
573 return Ok(prefix.to_owned());
574 }
575 let hits: Vec<String> = self
576 .list()
577 .into_iter()
578 .map(|t| t.id)
579 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
580 .collect();
581 match hits.len() {
582 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
583 0 => bail!("no talk matches `{prefix}`"),
584 _ => bail!(
585 "`{prefix}` matches {} talks: {}",
586 hits.len(),
587 hits.join(", ")
588 ),
589 }
590 }
591
592 pub fn revision(&self) -> u64 {
595 std::fs::read_dir(&self.root)
596 .into_iter()
597 .flatten()
598 .flatten()
599 .filter(|e| e.path().extension().is_none_or(|x| x != "turn"))
600 .filter_map(|e| e.metadata().ok())
601 .filter_map(|m| m.modified().ok())
602 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
603 .map(|d| d.as_millis() as u64)
604 .max()
605 .unwrap_or(0)
606 }
607
608 pub fn count_open(&self) -> usize {
610 self.list().iter().filter(|t| t.status.open()).count()
611 }
612
613 pub fn turn_path(&self, id: &str) -> PathBuf {
616 self.root.join(format!("{id}.turn"))
617 }
618
619 pub fn claim_turn(&self, id: &str) -> Result<Option<TurnLease>> {
628 self.claim_turn_at(id, Timestamp::now())
629 }
630
631 fn claim_turn_at(&self, id: &str, now: Timestamp) -> Result<Option<TurnLease>> {
632 std::fs::create_dir_all(&self.root)
633 .with_context(|| format!("create {}", self.root.display()))?;
634 let path = self.turn_path(id);
635 let token = fresh_token();
636 if create_turn(&path, &token, now)? {
637 return Ok(Some(TurnLease { path, token }));
638 }
639 if read_turn(&path).is_some_and(|r| r.fresh(now)) {
640 return Ok(None);
641 }
642 let Some(_lock) = TurnLock::take(&path)? else {
647 return Ok(None);
648 };
649 if read_turn(&path).is_some_and(|r| r.fresh(now)) {
650 return Ok(None);
651 }
652 let _ = std::fs::remove_file(&path);
653 Ok(create_turn(&path, &token, now)?.then_some(TurnLease { path, token }))
654 }
655
656 pub fn turn_held(&self, id: &str) -> bool {
658 read_turn(&self.turn_path(id)).is_some_and(|r| r.fresh(Timestamp::now()))
659 }
660
661 pub fn remove(&self, id: &str) -> Result<()> {
674 let _guard = self.guard();
675 let resolved = self.resolve_id(id)?;
676 let path = self.path_of(&resolved);
677 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
678 let artifacts = self.artifacts_of(&resolved);
679 if artifacts.is_dir() {
680 std::fs::remove_dir_all(&artifacts)
681 .with_context(|| format!("remove {}", artifacts.display()))?;
682 }
683 let _ = std::fs::remove_file(self.turn_path(&resolved));
684 Ok(())
685 }
686}
687
688#[derive(Debug, Serialize, Deserialize)]
691struct TurnRecord {
692 token: String,
693 pid: u32,
694 beat_at: Timestamp,
695}
696
697impl TurnRecord {
698 fn fresh(&self, now: Timestamp) -> bool {
699 now.as_second() - self.beat_at.as_second() <= crate::ask::LEASE_TTL.as_secs() as i64
700 }
701}
702
703fn fresh_token() -> String {
707 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
708 let n = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
709 let seed = crate::rng::entropy() ^ n.wrapping_mul(0x9E37_79B9_7F4A_7C15);
710 crate::rng::SplitMix64::new(seed).uuid_v4()
711}
712
713fn read_turn(path: &Path) -> Option<TurnRecord> {
714 serde_json::from_str(&std::fs::read_to_string(path).ok()?).ok()
715}
716
717const TAKEOVER_LOCK_TTL: Duration = Duration::from_secs(10);
719
720const TICKET_TTL: Duration = Duration::from_secs(10);
723
724const TICKET_GENERATIONS: u32 = 16;
726
727const TICKET_SWEEP_AGE: Duration = Duration::from_secs(3600);
729
730fn create_exclusive(path: &Path, body: &str) -> Result<bool> {
732 use std::io::Write as _;
733 match std::fs::OpenOptions::new()
734 .write(true)
735 .create_new(true)
736 .open(path)
737 {
738 Ok(mut f) => {
739 if let Err(e) = f.write_all(body.as_bytes()) {
740 drop(f);
741 let _ = std::fs::remove_file(path);
742 return Err(e).with_context(|| format!("write {}", path.display()));
743 }
744 Ok(true)
745 }
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
751fn create_turn(path: &Path, token: &str, now: Timestamp) -> Result<bool> {
752 let record = TurnRecord {
753 token: token.to_owned(),
754 pid: std::process::id(),
755 beat_at: now,
756 };
757 let body = serde_json::to_string(&record).context("serialize turn lease")?;
758 let tmp = path.with_extension(format!("turn.{token}.new"));
761 std::fs::write(&tmp, body).with_context(|| format!("write {}", tmp.display()))?;
762 let linked = std::fs::hard_link(&tmp, path);
763 let _ = std::fs::remove_file(&tmp);
764 match linked {
765 Ok(()) => Ok(true),
766 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
767 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
768 }
769}
770
771struct TurnLock {
775 path: PathBuf,
776 token: String,
777}
778
779impl TurnLock {
780 fn token() -> String {
781 fresh_token()
782 }
783
784 fn publish(path: &Path, token: &str) -> Result<bool> {
787 let tmp = path.with_extension(format!("lock.{token}.new"));
788 std::fs::write(&tmp, token).with_context(|| format!("write {}", tmp.display()))?;
789 let linked = std::fs::hard_link(&tmp, path);
790 let _ = std::fs::remove_file(&tmp);
791 match linked {
792 Ok(()) => Ok(true),
793 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
794 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
795 }
796 }
797
798 fn take(lease: &Path) -> Result<Option<Self>> {
799 let path = lease.with_extension("turn.lock");
800 let token = Self::token();
801 if Self::publish(&path, &token)? {
802 return Ok(Some(Self { path, token }));
803 }
804 let Ok(seen) = std::fs::read_to_string(&path) else {
805 return Ok(None);
806 };
807 let key = if !seen.is_empty() && seen.chars().all(|c| c.is_ascii_alphanumeric() || c == '-')
810 {
811 seen.as_str()
812 } else {
813 "invalid"
814 };
815 let aged = std::fs::metadata(&path)
816 .and_then(|m| m.modified())
817 .ok()
818 .and_then(|t| t.elapsed().ok())
819 .is_some_and(|age| age > TAKEOVER_LOCK_TTL);
820 if !aged {
821 return Ok(None);
822 }
823 let mut won = false;
835 for n in 0..TICKET_GENERATIONS {
836 let ticket = path.with_extension(format!("lock.{key}.break.{n}"));
837 if create_exclusive(&ticket, "")? {
838 won = true;
839 break;
840 }
841 let stale = std::fs::metadata(&ticket)
842 .and_then(|m| m.modified())
843 .ok()
844 .and_then(|t| t.elapsed().ok())
845 .is_some_and(|age| age > TICKET_TTL);
846 if !stale {
847 return Ok(None);
848 }
849 }
850 if !won {
851 Self::sweep_tickets(&path);
856 return Ok(None);
857 }
858 Self::sweep_tickets(&path);
859 if std::fs::read_to_string(&path).ok().as_deref() != Some(seen.as_str()) {
862 return Ok(None);
863 }
864 let _ = std::fs::remove_file(&path);
865 if Self::publish(&path, &token)? {
866 return Ok(Some(Self { path, token }));
867 }
868 Ok(None)
869 }
870
871 fn sweep_tickets(path: &Path) {
874 let (Some(dir), Some(name)) = (path.parent(), path.file_name().and_then(|n| n.to_str()))
875 else {
876 return;
877 };
878 let prefix = format!("{name}.");
879 let Ok(entries) = std::fs::read_dir(dir) else {
880 return;
881 };
882 for entry in entries.flatten() {
883 let file = entry.file_name();
884 let Some(file) = file.to_str() else { continue };
885 if !(file.starts_with(&prefix) && file.contains(".break.")) {
886 continue;
887 }
888 let old = entry
889 .metadata()
890 .and_then(|m| m.modified())
891 .ok()
892 .and_then(|t| t.elapsed().ok())
893 .is_some_and(|age| age > TICKET_SWEEP_AGE);
894 if old {
895 let _ = std::fs::remove_file(entry.path());
896 }
897 }
898 }
899
900 fn take_patiently(lease: &Path) -> Option<Self> {
902 for _ in 0..50 {
903 match Self::take(lease) {
904 Ok(Some(lock)) => return Some(lock),
905 Ok(None) => std::thread::sleep(Duration::from_millis(10)),
906 Err(_) => return None,
907 }
908 }
909 None
910 }
911}
912
913impl Drop for TurnLock {
914 fn drop(&mut self) {
915 if std::fs::read_to_string(&self.path).is_ok_and(|t| t == self.token) {
917 let _ = std::fs::remove_file(&self.path);
918 }
919 }
920}
921
922pub const TURN_BEAT: Duration = Duration::from_secs(20);
925
926#[derive(Debug)]
931pub struct TurnLease {
932 path: PathBuf,
933 token: String,
934}
935
936impl TurnLease {
937 pub fn beat(&self) -> Result<bool> {
941 let _lock = TurnLock::take_patiently(&self.path)
942 .with_context(|| format!("lock {} to renew it", self.path.display()))?;
943 let Some(mut record) = read_turn(&self.path).filter(|r| r.token == self.token) else {
944 return Ok(false);
945 };
946 record.beat_at = Timestamp::now();
947 let body = serde_json::to_string(&record).context("serialize turn lease")?;
948 let tmp = self.path.with_extension(format!("turn.{}.tmp", self.token));
949 write_atomic(&tmp, &self.path, &body)?;
950 Ok(true)
951 }
952
953 pub async fn beating<T>(&self, fut: impl std::future::Future<Output = T>) -> Result<T> {
958 self.beating_every(TURN_BEAT, fut).await
959 }
960
961 async fn beating_every<T>(
962 &self,
963 period: Duration,
964 fut: impl std::future::Future<Output = T>,
965 ) -> Result<T> {
966 tokio::pin!(fut);
967 loop {
968 match tokio::time::timeout(period, &mut fut).await {
969 Ok(out) => return Ok(out),
970 Err(_) => match self.beat() {
971 Ok(true) => {}
972 Ok(false) => bail!(
973 "the turn lease {} was taken over; this turn is stopped",
974 self.path.display()
975 ),
976 Err(e) => tracing::warn!("{e:#}"),
977 },
978 }
979 }
980 }
981}
982
983impl Drop for TurnLease {
984 fn drop(&mut self) {
985 if let Some(_lock) = TurnLock::take_patiently(&self.path) {
988 if read_turn(&self.path).is_some_and(|r| r.token == self.token) {
989 let _ = std::fs::remove_file(&self.path);
990 }
991 }
992 }
993}
994
995pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
1005 let repo = repo.canonicalize().unwrap_or(repo);
1008 let spec = match agent {
1011 Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
1012 None => agent::pick_chain(
1013 &cfg.agents,
1014 cfg.roles.chatter.as_ref(),
1015 &agent::installed,
1016 "chatter",
1017 )?
1018 .remove(0),
1019 };
1020
1021 let now = Timestamp::now();
1022 let mut talk = Talk {
1023 schema: SCHEMA,
1024 id: new_id(),
1025 repo,
1026 agent: spec.id.clone(),
1027 status: TalkStatus::Open,
1028 turns: Vec::new(),
1029 pending: String::new(),
1030 pending_attachments: Vec::new(),
1031 fallback: agent.is_none(),
1032 persona: String::new(),
1033 persona_dirty: false,
1034 created_at: now,
1035 updated_at: now,
1036 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
1037 };
1038 store.put(&mut talk)?;
1039 Ok(talk)
1040}
1041
1042pub fn record(
1049 talk: &mut Talk,
1050 store: &Talks,
1051 text: &str,
1052 attachments: Vec<Attachment>,
1053) -> Result<String> {
1054 let _guard = store.guard();
1062 let Ok(fresh) = store.get(&talk.id) else {
1067 bail!("talk {} was deleted", talk.short());
1068 };
1069 talk.status = fresh.status;
1070 talk.pending = fresh.pending;
1073 talk.pending_attachments = fresh.pending_attachments;
1074 if !talk.status.open() {
1075 bail!(
1076 "talk {} is {} and takes no more turns",
1077 talk.short(),
1078 talk.status.as_str()
1079 );
1080 }
1081 let text = text.trim();
1082 if text.is_empty() && attachments.is_empty() {
1083 bail!("nothing to say");
1084 }
1085 talk.turns.push(Turn {
1086 who: Who::Operator,
1087 body: text.to_owned(),
1088 at: Timestamp::now(),
1089 attachments,
1090 usage: None,
1091 });
1092 store.put(talk)?;
1093 Ok(text.to_owned())
1094}
1095
1096pub fn queue(
1098 talk: &mut Talk,
1099 store: &Talks,
1100 text: &str,
1101 attachments: Vec<Attachment>,
1102) -> Result<()> {
1103 let text = text.trim();
1104 if text.is_empty() && attachments.is_empty() {
1105 bail!("nothing to say");
1106 }
1107 let _guard = store.guard();
1108 let mut fresh = store
1109 .get(&talk.id)
1110 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1111 if !fresh.status.open() {
1112 bail!(
1113 "talk {} is {} and takes no more turns",
1114 fresh.short(),
1115 fresh.status.as_str()
1116 );
1117 }
1118 if !text.is_empty() {
1119 if fresh.pending.is_empty() {
1120 fresh.pending = text.to_owned();
1121 } else {
1122 fresh.pending.push_str("\n\n");
1123 fresh.pending.push_str(text);
1124 }
1125 }
1126 fresh.pending_attachments.extend(attachments);
1127 store.put(&mut fresh)?;
1128 *talk = fresh;
1129 Ok(())
1130}
1131
1132pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
1134 let _guard = store.guard();
1135 let mut fresh = store
1136 .get(&talk.id)
1137 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1138 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
1139 *talk = fresh;
1140 return Ok(None);
1141 }
1142 let text = std::mem::take(&mut fresh.pending);
1143 let attachments = std::mem::take(&mut fresh.pending_attachments);
1144 fresh.turns.push(Turn {
1145 who: Who::Operator,
1146 body: text.clone(),
1147 at: Timestamp::now(),
1148 attachments,
1149 usage: None,
1150 });
1151 store.put(&mut fresh)?;
1152 *talk = fresh;
1153 Ok(Some(text))
1154}
1155
1156pub async fn say(
1159 talk: &mut Talk,
1160 store: &Talks,
1161 cfg: &Config,
1162 text: &str,
1163 attachments: Vec<Attachment>,
1164) -> Result<()> {
1165 let text = record(talk, store, text, attachments)?;
1166 turn(talk, store, cfg, &text).await
1167}
1168
1169pub async fn respond(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1171 turn(talk, store, cfg, text).await
1172}
1173
1174pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
1192 let _guard = store.guard();
1193 let mut fresh = store
1194 .get(&talk.id)
1195 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1196 fresh.status = TalkStatus::Closed;
1197 fresh.pending.clear();
1199 fresh.pending_attachments.clear();
1200 store.put(&mut fresh)?;
1201 *talk = fresh;
1202 Ok(())
1203}
1204
1205pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
1216 let _guard = store.guard();
1217 let mut fresh = store
1218 .get(&talk.id)
1219 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1220 fresh.status = TalkStatus::Open;
1221 store.put(&mut fresh)?;
1222 *talk = fresh;
1223 Ok(())
1224}
1225
1226pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
1239 let _guard = store.guard();
1240 let mut fresh = store
1241 .get(&talk.id)
1242 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1243 if fresh.agent == spec.id {
1244 *talk = fresh;
1245 return Ok(false);
1246 }
1247 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
1248 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1249 fresh.fallback = false;
1251 fresh.turns.push(Turn {
1252 who: Who::Agent,
1253 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
1254 at: Timestamp::now(),
1255 attachments: Vec::new(),
1256 usage: None,
1257 });
1258 store.put(&mut fresh)?;
1259 *talk = fresh;
1260 Ok(true)
1261}
1262
1263pub fn switch_persona(talk: &mut Talk, store: &Talks, id: &str) -> Result<bool> {
1271 let id = if id.trim() == crate::persona::DEFAULT_ID {
1272 ""
1273 } else {
1274 id.trim()
1275 };
1276 let _guard = store.guard();
1277 let mut fresh = store
1278 .get(&talk.id)
1279 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1280 if fresh.persona == id {
1281 *talk = fresh;
1282 return Ok(false);
1283 }
1284 let from = std::mem::replace(&mut fresh.persona, id.to_owned());
1285 fresh.persona_dirty = true;
1286 let label = |p: &str| {
1287 if p.is_empty() {
1288 crate::persona::DEFAULT_ID.to_owned()
1289 } else {
1290 p.to_owned()
1291 }
1292 };
1293 fresh.turns.push(Turn {
1294 who: Who::Agent,
1295 body: format!(
1296 "{MAGI_NOTE}persona changed from {} to {}",
1297 label(&from),
1298 label(id)
1299 ),
1300 at: Timestamp::now(),
1301 attachments: Vec::new(),
1302 usage: None,
1303 });
1304 store.put(&mut fresh)?;
1305 *talk = fresh;
1306 Ok(true)
1307}
1308
1309pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
1311 let _guard = store.guard();
1312 let mut fresh = store
1313 .get(&talk.id)
1314 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1315 fresh.pending.clear();
1316 fresh.pending_attachments.clear();
1317 store.put(&mut fresh)?;
1318 *talk = fresh;
1319 Ok(())
1320}
1321
1322pub fn clear_pending_if_matches(
1324 talk: &mut Talk,
1325 store: &Talks,
1326 expected_text: &str,
1327 expected_attachments: &[String],
1328) -> Result<bool> {
1329 let _guard = store.guard();
1330 let mut fresh = store
1331 .get(&talk.id)
1332 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1333 if !pending_matches(&fresh, expected_text, expected_attachments) {
1334 *talk = fresh;
1335 return Ok(false);
1336 }
1337 fresh.pending.clear();
1338 fresh.pending_attachments.clear();
1339 store.put(&mut fresh)?;
1340 *talk = fresh;
1341 Ok(true)
1342}
1343
1344pub fn edit_pending_text(
1348 talk: &mut Talk,
1349 store: &Talks,
1350 text: &str,
1351 expected_text: &str,
1352 expected_attachments: &[String],
1353) -> Result<bool> {
1354 let _guard = store.guard();
1355 let mut fresh = store
1356 .get(&talk.id)
1357 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1358 if !pending_matches(&fresh, expected_text, expected_attachments) {
1359 *talk = fresh;
1360 return Ok(false);
1361 }
1362 fresh.pending = text.trim().to_owned();
1363 store.put(&mut fresh)?;
1364 *talk = fresh;
1365 Ok(true)
1366}
1367
1368fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
1369 talk.pending == expected_text
1370 && talk
1371 .pending_attachments
1372 .iter()
1373 .map(|attachment| &attachment.id)
1374 .eq(expected_attachments.iter())
1375}
1376
1377async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1384 let spec = cfg
1385 .agents
1386 .iter()
1387 .find(|a| a.id == talk.agent)
1388 .with_context(|| {
1389 format!(
1390 "talk {} was opened with agent `{}`, which is no longer in \
1391 the roster; restore it in magi.toml or start a new \
1392 conversation",
1393 talk.short(),
1394 talk.agent
1395 )
1396 })?;
1397
1398 let last_note = attachment_note(
1402 store,
1403 &talk.id,
1404 talk.turns
1405 .last()
1406 .map_or(&[][..], |t| t.attachments.as_slice()),
1407 );
1408
1409 let attachment_paths: Vec<PathBuf> = talk
1415 .turns
1416 .iter()
1417 .flat_map(|t| t.attachments.iter())
1418 .filter_map(|a| store.attachment_path(&talk.id, a))
1419 .collect();
1420
1421 let questions = crate::ask::Questions::open();
1429 let consulted = crate::consult::pending_consults(&questions, &talk.id);
1430 let consult_roots: Vec<PathBuf> = if consulted {
1431 vec![questions.root().to_path_buf()]
1432 } else {
1433 Vec::new()
1434 };
1435
1436 let artifacts = store.artifacts_of(&talk.id);
1437 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
1440 let stem = format!("turn-{}", operator_turns.max(1));
1441 let cache_dir = cfg.cache_dir();
1444
1445 let mut chain = vec![spec.clone()];
1448 if let Some(choice) = cfg.roles.chatter.as_ref()
1449 && talk.fallback
1450 {
1451 for id in choice.ids() {
1452 if id == talk.agent || chain.iter().any(|s| s.id == id) {
1453 continue;
1454 }
1455 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1456 Ok(s) => chain.push(s),
1457 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1458 }
1459 }
1460 }
1461
1462 let persona = crate::persona::active(&cfg.talk.personas, &talk.persona);
1463 let persona_update = if talk.persona_dirty {
1464 format!("{}\n\n", crate::persona::update_block(persona.as_ref()))
1465 } else {
1466 String::new()
1467 };
1468
1469 let mut outcome = None;
1470 let mut fell_back_from: Option<String> = None;
1471 let mut first_try: Option<(String, SeatState)> = None;
1474 for (n, spec) in chain.iter().enumerate() {
1475 if n > 0 {
1476 if first_try.is_none() {
1477 first_try = Some((talk.agent.clone(), talk.seat.clone()));
1478 }
1479 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1480 fell_back_from.get_or_insert_with(|| talk.agent.clone());
1483 talk.agent = spec.id.clone();
1484 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1485 }
1486 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1487 let first_ever = talk.turns.len() <= 1;
1488 let body = if talk.seat.turns == 0 && first_ever {
1489 format!(
1490 "{}\n\n# Operator\n\n{text}{last_note}",
1491 briefing_with(
1492 &talk.repo,
1493 &cfg.graph.language,
1494 cfg.talk.allow_write,
1495 persona.as_ref()
1496 )
1497 )
1498 } else if talk.seat.turns == 0 {
1499 format!(
1502 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1503 briefing_with(
1504 &talk.repo,
1505 &cfg.graph.language,
1506 cfg.talk.allow_write,
1507 persona.as_ref()
1508 ),
1509 transcript(talk, store)
1510 )
1511 } else if resuming {
1512 format!("{persona_update}{text}{last_note}")
1513 } else {
1514 let standing = match (&persona, talk.persona_dirty) {
1518 (Some(p), false) => format!("{}\n", crate::persona::section(p)),
1519 _ => persona_update.clone(),
1520 };
1521 format!("{}\n\n{standing}{text}{last_note}", transcript(talk, store))
1522 };
1523 let attempt_stem = if n == 0 {
1524 stem.clone()
1525 } else {
1526 format!("{stem}-{}", spec.id)
1527 };
1528 let inv = Invocation {
1529 cwd: &talk.repo,
1530 prompt: &body,
1531 timeout: turn_timeout(cfg),
1532 allow_write: cfg.talk.allow_write || consulted,
1537 sessions: cfg.graph.sessions,
1538 artifacts: &artifacts,
1539 stem: &attempt_stem,
1540 run: &talk.id,
1543 node: crate::queue::CHAT_NODE,
1544 cache_dir: cache_dir.as_deref(),
1545 attachments: &attachment_paths,
1546 writable: &consult_roots,
1547 };
1548 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1549 let advance = agent::chain_advances(&result);
1550 if n == 0 || !advance {
1551 outcome = Some(result);
1552 } else {
1553 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1555 }
1556 if !advance {
1557 break;
1558 }
1559 }
1560 if outcome.as_ref().is_some_and(agent::chain_advances) {
1561 if let Some((id, seat)) = first_try {
1564 talk.agent = id;
1565 talk.seat = seat;
1566 fell_back_from = None;
1567 }
1568 }
1569 let outcome = outcome.expect("a chain holds at least one agent");
1570 let note = |why: String| Turn {
1571 who: Who::Agent,
1572 body: format!("{MAGI_NOTE}{why}"),
1573 at: Timestamp::now(),
1574 attachments: Vec::new(),
1575 usage: None,
1576 };
1577 let (reply, failure) = match outcome {
1578 Err(e) => (
1579 note(format!("could not run agent `{}`: {e}", talk.agent)),
1580 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1581 ),
1582 Ok(out) if out.quota_exhausted() => {
1583 let reset = out
1584 .quota
1585 .as_ref()
1586 .and_then(|q| q.reset.clone())
1587 .map_or_else(String::new, |r| format!(" (resets {r})"));
1588 let why = format!(
1589 "agent `{}` is out of quota{reset}; your message is saved, so \
1590 say it again when the window reopens",
1591 talk.agent
1592 );
1593 (note(why.clone()), Some(why))
1594 }
1595 Ok(out) if out.timed_out => {
1596 let why = format!(
1597 "agent `{}` did not answer within {}s; your message is saved",
1598 talk.agent,
1599 turn_timeout(cfg).as_secs()
1600 );
1601 (note(why.clone()), Some(why))
1602 }
1603 Ok(out) if !out.usable() => {
1604 let why = format!(
1605 "agent `{}` produced no answer (exit {}); your message is saved",
1606 talk.agent,
1607 out.exit_code
1608 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1609 );
1610 (note(why.clone()), Some(why))
1611 }
1612 Ok(out) => (
1613 Turn {
1614 who: Who::Agent,
1615 body: out.text.trim().to_owned(),
1616 at: Timestamp::now(),
1617 attachments: Vec::new(),
1618 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1621 context_tokens,
1622 agent: talk.agent.clone(),
1623 model: cfg
1624 .agents
1625 .iter()
1626 .find(|a| a.id == talk.agent)
1627 .and_then(|a| a.model.clone()),
1628 }),
1629 },
1630 None,
1631 ),
1632 };
1633
1634 let _guard = store.guard();
1646 let Ok(fresh) = store.get(&talk.id) else {
1652 return Ok(());
1653 };
1654 talk.status = fresh.status;
1655 talk.pending = fresh.pending;
1659 talk.pending_attachments = fresh.pending_attachments;
1660 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1661 talk.turns.push(note(format!(
1664 "agent changed from {from} to {} (fallback)",
1665 talk.agent
1666 )));
1667 }
1668 if failure.is_none() {
1671 talk.persona_dirty = false;
1672 }
1673 talk.turns.push(reply);
1674 if let Err(put_err) = store.put(talk) {
1675 let lost = talk.turns.pop().expect("just pushed above");
1684 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1685 let why = match &stash {
1686 Ok(path) => format!(
1687 "agent `{}` answered, but the reply could not be saved to \
1688 this conversation ({put_err:#}); the raw text was kept at \
1689 {} - your message is saved, ask again",
1690 talk.agent,
1691 path.display()
1692 ),
1693 Err(stash_err) => format!(
1694 "agent `{}` answered, but the reply could not be saved to \
1695 this conversation ({put_err:#}), and it could not be kept \
1696 anywhere else either ({stash_err:#}); your message is \
1697 saved, ask again",
1698 talk.agent
1699 ),
1700 };
1701 talk.turns.push(note(why.clone()));
1702 return match store.put(talk) {
1709 Ok(()) => bail!("{why}"),
1710 Err(note_err) => {
1711 talk.turns.pop();
1731 Err(note_err).context(why)
1732 }
1733 };
1734 }
1735
1736 match failure {
1737 Some(why) => bail!("{why}"),
1738 None => Ok(()),
1739 }
1740}
1741
1742fn transcript(talk: &Talk, store: &Talks) -> String {
1745 let mut out = String::from(
1746 "This conversation cannot resume on the CLI's side, so here is \
1747 everything said so far; answer only the last message.\n",
1748 );
1749 for t in &talk.turns {
1750 let who = match t.who {
1751 Who::Operator => "operator",
1752 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1753 Who::Agent => "you",
1754 };
1755 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1756 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1757 }
1758 out
1759}
1760
1761fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1766 if attachments.is_empty() {
1767 return String::new();
1768 }
1769 let mut out = String::from(
1770 "\n\nThe operator attached the image(s) below to this message. Open \
1771 and look at each one before you answer.\n",
1772 );
1773 for att in attachments {
1774 if let Some(path) = store.attachment_path(talk_id, att) {
1775 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1776 }
1777 }
1778 out.push('\n');
1779 out
1780}
1781
1782pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
1804 briefing_with(repo, language, allow_write, None)
1805}
1806
1807pub fn briefing_with(
1810 repo: &Path,
1811 language: &str,
1812 allow_write: bool,
1813 persona: Option<&crate::persona::Persona>,
1814) -> String {
1815 let write_policy = if allow_write {
1816 "Write access is enabled for this conversation (`allow_write = \
1817 true`), so you may write files - but only a small, \
1818 already-decided edit the operator names outright in this \
1819 conversation, not an implementation. This is a permission on the \
1820 conversation as a whole, not a property of whichever repository \
1821 it happened to start in: if the operator names a different \
1822 repository for that small edit, the policy allows it there too. \
1823 Your own tool may still confine writes to the repository this \
1824 conversation started in regardless - if a write elsewhere is \
1825 refused, say so plainly rather than working around it. Once you \
1826 have made an edit, say plainly what you edited. Anything bigger, \
1827 or anything still open-ended, still goes through the queue below \
1828 rather than being done here."
1829 } else {
1830 "Do not write files. Implementing a change is not this \
1831 conversation's job; a separate, blind competition of agents does \
1832 that, and a repository this conversation has already edited would \
1833 make their diffs unjudgeable."
1834 };
1835 let mut out = format!(
1836 "You are magi's standing conversation partner for its operator, who \
1837 usually has this open on a phone. Keep replies short: no preamble, \
1838 no restating what they just said.\n\n\
1839 # Repository\n\n{repo}\n\n\
1840 You may look around: read files, run shell commands, search history, \
1841 run tests - whatever answers the question. {write_policy}\n\n\
1842 A short, command-shaped message (\"list\", \"info <id>\", \"show \
1843 3cbf\") is almost always the operator asking you to look something \
1844 up, not an instruction to file - answer it yourself with `magi \
1845 list`, `magi show <id>`, `magi task list`, or the like, the same way \
1846 you would answer any other question in this conversation.\n\n\
1847 # When the operator wants something done\n\n\
1848 Run:\n\n\
1849 magi task add --solo --repo {repo} <instruction>\n\n\
1850 and tell the operator the task id it prints, so they can follow it \
1851 from the Queue. If it refuses with a duplicate warning (the \
1852 instruction names a branch, commit or pull request that an \
1853 unfinished task, run or PR already owns), do not repeat it with \
1854 --force yourself: tell the operator what it matched and let them \
1855 decide. Write <instruction> so that an implementer who has \
1856 never seen this conversation can act on it alone - it is everything \
1857 they get. Use --solo: it runs the task through one implementer \
1858 straight into review instead of the usual multi-agent competition, \
1859 which is the right shape for a change this conversation has already \
1860 settled, rather than one still worth several independent takes.\n\n\
1861 If the operator asks for something in a different repository, \
1862 --repo does not have to be a full path: --repo owner/repo (or just \
1863 repo, when that is unambiguous) is resolved against local checkouts \
1864 the same way `magi repos` lists them. If the command fails because \
1865 nothing matches or more than one checkout shares that name, ask the \
1866 operator which repository they mean (or run `magi repos` yourself \
1867 to see the candidates) rather than guessing.\n\n\
1868 The current state of the code is whatever origin/main holds, not \
1869 whatever a working tree shows: a primary checkout often lags \
1870 upstream, sits on a detached HEAD and carries uncommitted changes. \
1871 Before answering about code, run `git fetch origin` in that \
1872 repository if it is cheap, then read through \
1873 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
1874 If the working tree differs, say so; if the fetch fails, say that \
1875 too, so the operator knows the answer may be stale.\n\n\
1876 If the operator attached an image (a screenshot, say) that the task \
1877 is about, pass it with `--attach <path>`, using the absolute path \
1878 the turn's attachment note gives; repeat the flag for several. \
1879 `magi task add --solo --attach <path> <instruction>` copies the \
1880 file into the task, so the implementer receives it. Do not paste the \
1881 path into <instruction> instead: deleting this conversation deletes \
1882 its attachments, and then that path reaches no one.\n",
1883 repo = repo.display(),
1884 );
1885 out.push_str(&language_note(language));
1886 if let Some(p) = persona {
1887 out.push_str(&crate::persona::section(p));
1888 }
1889 out
1890}
1891
1892fn language_note(language: &str) -> String {
1895 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
1896 String::new()
1897 } else {
1898 format!("\nHold this conversation in {language}.\n")
1899 }
1900}
1901
1902pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
1909 let mut tasks: Vec<Task> = queue
1910 .list()
1911 .into_iter()
1912 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
1913 .collect();
1914 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
1915 tasks
1916}
1917
1918fn read_path(path: &Path) -> Result<Talk> {
1919 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1920 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
1921}
1922
1923const PUT_RETRIES: u32 = 5;
1926
1927fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1939 let mut last_err = None;
1940 for attempt in 0..PUT_RETRIES {
1941 if attempt > 0 {
1942 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
1943 }
1944 match try_write_atomic(tmp, path, body) {
1945 Ok(()) => return Ok(()),
1946 Err(e) => last_err = Some(e),
1947 }
1948 }
1949 Err(last_err.expect("the loop above always runs at least once"))
1950}
1951
1952fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1953 #[cfg(test)]
1954 if failpoint::take_forced_put_failure() {
1955 bail!("simulated write failure (test)");
1956 }
1957 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
1958 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
1959 Ok(())
1960}
1961
1962fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
1967 let dir = store.artifacts_of(id);
1968 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
1969 let path = dir.join(format!("{stem}-lost.txt"));
1970 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
1971 Ok(path)
1972}
1973
1974#[cfg(test)]
1981mod failpoint {
1982 use std::cell::Cell;
1983
1984 thread_local! {
1985 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
1986 }
1987
1988 pub(super) fn force_put_failures(count: u32) {
1991 FORCE_PUT_FAILURES.with(|c| c.set(count));
1992 }
1993
1994 pub(super) fn take_forced_put_failure() -> bool {
1997 FORCE_PUT_FAILURES.with(|c| {
1998 let n = c.get();
1999 if n == 0 {
2000 false
2001 } else {
2002 c.set(n - 1);
2003 true
2004 }
2005 })
2006 }
2007}
2008
2009fn short(id: &str) -> &str {
2010 id.split('-').next_back().unwrap_or(id)
2011}
2012
2013fn new_id() -> String {
2014 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2015 let seed = crate::rng::entropy();
2016 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2017}
2018
2019fn attachment_ext(mime: &str) -> Option<&'static str> {
2024 match mime {
2025 "image/png" => Some("png"),
2026 "image/jpeg" => Some("jpg"),
2027 "image/gif" => Some("gif"),
2028 "image/webp" => Some("webp"),
2029 _ => None,
2030 }
2031}
2032
2033pub fn valid_attachment_id(id: &str) -> bool {
2038 id.len() == 32
2039 && id
2040 .bytes()
2041 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
2042}
2043
2044fn new_attachment_id() -> String {
2048 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
2049 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
2050}
2051
2052#[cfg(test)]
2053mod tests {
2054 #[test]
2055 fn the_briefing_points_at_origin_main_not_the_working_tree() {
2056 let b = briefing(Path::new("/r"), "en", false);
2057 assert!(b.contains("origin/main"));
2058 assert!(b.contains("git show origin/main:"));
2059 }
2060 use std::collections::BTreeMap;
2061
2062 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
2063 use crate::queue::{Queue, Source, Task};
2064
2065 use super::*;
2066
2067 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
2068 AgentSpec {
2069 id: id.to_owned(),
2070 kind: AgentKind::Command,
2071 model: model.map(str::to_owned),
2072 command: Vec::new(),
2073 extra_args: Vec::new(),
2074 env: BTreeMap::new(),
2075 prompt_delivery: None,
2076 }
2077 }
2078
2079 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
2080 Talk {
2081 schema: SCHEMA,
2082 id: "20260904-014455-ab12".to_owned(),
2083 repo: PathBuf::from("."),
2084 agent: agent.to_owned(),
2085 status: TalkStatus::Open,
2086 turns,
2087 pending: String::new(),
2088 pending_attachments: Vec::new(),
2089 fallback: false,
2090 persona: String::new(),
2091 persona_dirty: false,
2092 created_at: Timestamp::now(),
2093 updated_at: Timestamp::now(),
2094 seat: SeatState::new(SEAT, agent, 1),
2095 }
2096 }
2097
2098 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
2099 Turn {
2100 who: Who::Agent,
2101 body: body.to_owned(),
2102 at: Timestamp::now(),
2103 attachments: Vec::new(),
2104 usage: usage.map(|(t, a, m)| TurnUsage {
2105 context_tokens: t,
2106 agent: a.to_owned(),
2107 model: m.map(str::to_owned),
2108 }),
2109 }
2110 }
2111
2112 fn ctx_config(windows: &[(&str, u64)]) -> Config {
2113 Config {
2114 agents: vec![
2115 ctx_agent("small", Some("small-model")),
2116 ctx_agent("big", Some("big-model")),
2117 ctx_agent("plain", None),
2118 ],
2119 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
2120 ..Config::default()
2121 }
2122 }
2123
2124 #[test]
2125 fn context_usage_computes_percent_and_warns_at_eighty() {
2126 let cfg = ctx_config(&[("small-model", 1000)]);
2127 let at = |tokens| {
2128 let t = ctx_talk(
2129 "small",
2130 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
2131 );
2132 context_usage(&t, Some(&cfg))
2133 };
2134 let u = at(799);
2135 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
2136 let u = at(800);
2137 assert_eq!((u.percent, u.warn), (Some(80), true));
2138 let u = at(1500);
2139 assert_eq!((u.percent, u.warn), (Some(150), true));
2140 assert!(!u.since_switch);
2141 }
2142
2143 #[test]
2144 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
2145 let cfg = ctx_config(&[("small-model", 1000)]);
2146 let t = ctx_talk(
2147 "small",
2148 vec![
2149 reply("old", Some((900, "small", Some("small-model")))),
2150 reply("new", None),
2151 ],
2152 );
2153 let u = context_usage(&t, Some(&cfg));
2154 assert!(u.estimated);
2156 assert_ne!(u.tokens, Some(900));
2157 assert!(u.tokens.is_some());
2158 let t = ctx_talk(
2160 "small",
2161 vec![
2162 reply("old", Some((900, "small", Some("small-model")))),
2163 reply("magi: could not run agent", None),
2164 ],
2165 );
2166 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
2167 assert_eq!(
2168 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
2169 None
2170 );
2171 }
2172
2173 #[test]
2174 fn estimate_counts_chars_both_sides_and_standing_prompt() {
2175 let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
2176 assert_eq!(estimate_context_tokens(&t, 0), Some(2)); let op = Turn {
2178 who: Who::Operator,
2179 ..reply("abcdefg", None)
2180 };
2181 t.turns.push(op);
2182 assert_eq!(estimate_context_tokens(&t, 0), Some(4));
2183 assert!(
2184 estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
2185 );
2186 let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
2188 assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
2189 let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
2191 assert_eq!(estimate_context_tokens(¬e, 1000), None);
2192 assert_eq!(
2193 estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
2194 None
2195 );
2196 }
2197
2198 #[test]
2199 fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
2200 let cfg = ctx_config(&[("small-model", 1000)]);
2201 let t = ctx_talk(
2202 "small",
2203 vec![reply(
2204 &"x".repeat(5000),
2205 Some((10, "small", Some("small-model"))),
2206 )],
2207 );
2208 let u = context_usage(&t, Some(&cfg));
2209 assert_eq!((u.tokens, u.estimated), (Some(10), false));
2210 let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
2211 let u = context_usage(&t, Some(&cfg));
2212 assert!(u.estimated && !u.since_switch);
2213 assert_eq!(u.window, Some(1000));
2214 assert!(u.warn && u.percent.unwrap() >= 80);
2215 let t = ctx_talk("small", vec![reply("hi", None)]);
2216 let u = context_usage(&t, Some(&cfg));
2217 assert!(u.estimated && u.percent.is_some());
2218 }
2219
2220 #[test]
2221 fn context_usage_without_a_window_shows_tokens_only() {
2222 let cfg = ctx_config(&[]);
2223 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
2225 let u = context_usage(&t, Some(&cfg));
2226 assert_eq!(
2227 (u.tokens, u.window, u.percent, u.warn),
2228 (Some(5000), None, None, false)
2229 );
2230 let t = ctx_talk(
2231 "small",
2232 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
2233 );
2234 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
2235 assert_eq!(context_usage(&t, None).window, None);
2237 }
2238
2239 #[test]
2240 fn context_usage_switching_model_changes_the_denominator() {
2241 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
2242 let used = reply("hi", Some((900, "small", Some("small-model"))));
2243 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
2244 assert_eq!(
2245 (before.percent, before.warn, before.since_switch),
2246 (Some(90), true, false)
2247 );
2248 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
2251 assert_eq!(after.window, Some(10_000));
2252 assert_eq!(
2253 (after.percent, after.warn, after.since_switch),
2254 (Some(9), false, true)
2255 );
2256 assert_eq!(after.model.as_deref(), Some("big-model"));
2257 }
2258
2259 #[test]
2260 fn a_turn_recorded_before_usage_existed_still_reads() {
2261 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
2262 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
2263 assert!(turn.usage.is_none());
2264 let json = serde_json::to_string(&turn).expect("serialize");
2265 assert!(
2266 !json.contains("usage"),
2267 "absent usage is not written: {json}"
2268 );
2269 }
2270
2271 fn store() -> (tempfile::TempDir, Talks) {
2273 let tmp = tempfile::tempdir().expect("tempdir");
2274 let talks = Talks::at(tmp.path().join("talks"));
2275 (tmp, talks)
2276 }
2277
2278 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
2282 let path = dir.join("mock-talk-agent.sh");
2283 std::fs::write(&path, script).expect("write mock");
2284 AgentSpec {
2285 id: "mock".to_owned(),
2286 kind: AgentKind::Command,
2287 model: None,
2288 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2289 extra_args: Vec::new(),
2290 env,
2291 prompt_delivery: None,
2292 }
2293 }
2294
2295 fn config(spec: AgentSpec) -> Config {
2296 Config {
2297 agents: vec![spec],
2298 graph: Graph {
2299 language: "en".to_owned(),
2300 ..Graph::default()
2301 },
2302 ..Config::default()
2303 }
2304 }
2305
2306 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
2308
2309 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
2311
2312 const ECHO: &str = "#!/bin/sh\ncat\n";
2315
2316 fn env(reply: &str) -> BTreeMap<String, String> {
2317 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
2318 }
2319
2320 #[test]
2321 fn the_frozen_json_field_names_round_trip_through_disk() {
2322 let (tmp, talks) = store();
2323 let mut talk = Talk {
2324 schema: SCHEMA,
2325 id: "20260904-014455-ab12".to_owned(),
2326 repo: tmp.path().to_owned(),
2327 agent: "sonnet".to_owned(),
2328 status: TalkStatus::Open,
2329 turns: Vec::new(),
2330 pending: String::new(),
2331 pending_attachments: Vec::new(),
2332 fallback: false,
2333 persona: String::new(),
2334 persona_dirty: false,
2335 created_at: Timestamp::now(),
2336 updated_at: Timestamp::now(),
2337 seat: SeatState::new(SEAT, "sonnet", 7),
2338 };
2339 talks.put(&mut talk).expect("put");
2340
2341 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
2342 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
2343 for field in [
2344 "schema",
2345 "id",
2346 "repo",
2347 "agent",
2348 "status",
2349 "turns",
2350 "created_at",
2351 "updated_at",
2352 ] {
2353 assert!(v.get(field).is_some(), "missing field `{field}`");
2354 }
2355 assert_eq!(v["schema"], 1);
2356 assert_eq!(v["status"], "open");
2357
2358 let back = talks.get(&talk.id).expect("get");
2359 assert_eq!(back.id, talk.id);
2360 assert_eq!(back.status, TalkStatus::Open);
2361 }
2362
2363 #[test]
2364 fn opening_a_talk_takes_no_agent_turn() {
2365 let (tmp, talks) = store();
2366 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2370 let cfg = config(spec);
2371
2372 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2373 assert_eq!(talk.status, TalkStatus::Open);
2374 assert!(talk.turns.is_empty(), "nothing has been said yet");
2375
2376 let on_disk = talks.get(&talk.id).expect("get");
2377 assert_eq!(on_disk.turns.len(), 0);
2378 }
2379
2380 #[test]
2388 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
2389 let (tmp, talks) = store();
2390 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2391 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2392 chatter_spec.id = "chatter-mock".to_owned();
2393
2394 let mut cfg = Config {
2395 agents: vec![first_spec.clone(), chatter_spec.clone()],
2396 graph: Graph {
2397 language: "en".to_owned(),
2398 ..Graph::default()
2399 },
2400 ..Config::default()
2401 };
2402 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
2403
2404 let talk =
2405 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
2406 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
2407
2408 cfg.roles.chatter = None;
2409 let fallback =
2410 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
2411 assert_eq!(
2412 fallback.agent, first_spec.id,
2413 "unset chatter must fall back to agent::pick's own default order"
2414 );
2415 }
2416
2417 #[test]
2420 fn a_talk_recorded_without_attachments_still_reads() {
2421 let (tmp, talks) = store();
2422 let path = talks.path_of("20260904-014455-ab12");
2423 std::fs::create_dir_all(talks.root()).expect("talks dir");
2424 std::fs::write(
2425 &path,
2426 serde_json::json!({
2427 "schema": 1,
2428 "id": "20260904-014455-ab12",
2429 "repo": tmp.path(),
2430 "agent": "sonnet",
2431 "status": "open",
2432 "turns": [
2433 { "who": "operator", "body": "still there?",
2434 "at": Timestamp::now().to_string() },
2435 ],
2436 "created_at": Timestamp::now().to_string(),
2437 "updated_at": Timestamp::now().to_string(),
2438 "seat": SeatState::new(SEAT, "sonnet", 7),
2439 })
2440 .to_string(),
2441 )
2442 .expect("write pre-attachments talk");
2443
2444 let talk = talks.get("20260904-014455-ab12").expect("must still read");
2445 assert!(talk.turns[0].attachments.is_empty());
2446 }
2447
2448 fn lease_store() -> (tempfile::TempDir, Talks, Talks) {
2449 let tmp = tempfile::TempDir::new().expect("tmp");
2450 let root = tmp.path().join("talks");
2451 (tmp, Talks::at(root.clone()), Talks::at(root))
2452 }
2453
2454 #[test]
2455 fn two_starters_on_one_talk_one_wins_and_the_other_is_refused() {
2456 let (_tmp, a, b) = lease_store();
2457 let won = a.claim_turn("t1").expect("claim").expect("first wins");
2458 assert!(
2459 b.claim_turn("t1").expect("claim").is_none(),
2460 "second is refused"
2461 );
2462 assert!(b.turn_held("t1"));
2463 assert!(
2464 b.claim_turn("t2").expect("claim").is_some(),
2465 "other talks are free"
2466 );
2467 drop(won);
2468 }
2469
2470 #[test]
2471 fn a_stale_lease_is_taken_over_and_the_old_guard_cannot_release_it() {
2472 let (_tmp, a, b) = lease_store();
2473 let old = a.claim_turn("t1").expect("claim").expect("held");
2474 let later = Timestamp::now()
2475 .checked_add(jiff::SignedDuration::from_secs(
2476 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2477 ))
2478 .expect("later");
2479 let new = b
2480 .claim_turn_at("t1", later)
2481 .expect("claim")
2482 .expect("a stale lease is taken over");
2483 drop(old);
2484 assert!(a.turn_held("t1"), "the old guard left the new lease alone");
2485 assert!(new.beat().expect("beat"), "the new owner still beats");
2486 drop(new);
2487 assert!(!a.turn_held("t1"));
2488 }
2489
2490 #[tokio::test]
2491 async fn a_turn_whose_lease_was_taken_over_is_stopped() {
2492 let (_tmp, a, b) = lease_store();
2493 let old = a.claim_turn("t1").expect("claim").expect("held");
2494 let later = Timestamp::now()
2495 .checked_add(jiff::SignedDuration::from_secs(
2496 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2497 ))
2498 .expect("later");
2499 let _new = b
2500 .claim_turn_at("t1", later)
2501 .expect("claim")
2502 .expect("taken over");
2503 let out = old
2504 .beating_every(Duration::from_millis(10), std::future::pending::<()>())
2505 .await;
2506 assert!(out.is_err(), "the displaced turn must stop, not run on");
2507 }
2508
2509 #[tokio::test]
2510 async fn a_turn_that_finishes_is_returned_and_keeps_its_lease_beating() {
2511 let (_tmp, a, _b) = lease_store();
2512 let lease = a.claim_turn("t1").expect("claim").expect("held");
2513 let out = lease
2514 .beating_every(Duration::from_millis(5), async {
2515 tokio::time::sleep(Duration::from_millis(40)).await;
2516 7
2517 })
2518 .await
2519 .expect("still ours");
2520 assert_eq!(out, 7);
2521 assert!(a.turn_held("t1"));
2522 }
2523
2524 #[test]
2525 fn an_unreadable_lease_counts_as_stale() {
2526 let (_tmp, a, b) = lease_store();
2527 std::fs::create_dir_all(a.root()).expect("dir");
2528 std::fs::write(a.turn_path("t1"), "not json").expect("write");
2529 assert!(!a.turn_held("t1"));
2530 assert!(b.claim_turn("t1").expect("claim").is_some());
2531 }
2532
2533 #[test]
2534 fn a_lease_is_released_when_the_turn_ends_or_fails() {
2535 let (_tmp, a, b) = lease_store();
2536 let lease = a.claim_turn("t1").expect("claim").expect("held");
2537 let failed: Result<()> = (|| {
2538 let _held = &lease;
2539 bail!("turn failed")
2540 })();
2541 assert!(failed.is_err());
2542 assert!(
2543 b.claim_turn("t1").expect("claim").is_none(),
2544 "held mid-turn"
2545 );
2546 drop(lease);
2547 assert!(
2548 b.claim_turn("t1").expect("claim").is_some(),
2549 "free after the turn"
2550 );
2551 }
2552
2553 #[test]
2554 fn concurrent_takeovers_of_a_stale_lease_have_one_winner() {
2555 let (_tmp, a, _b) = lease_store();
2556 drop(a.claim_turn("t1").expect("claim").expect("held"));
2557 std::fs::write(
2558 a.turn_path("t1"),
2559 serde_json::to_string(&TurnRecord {
2560 token: "gone".into(),
2561 pid: 1,
2562 beat_at: Timestamp::from_second(1).expect("ts"),
2563 })
2564 .expect("json"),
2565 )
2566 .expect("write");
2567 let wins: Vec<_> = std::thread::scope(|sc| {
2568 let hs: Vec<_> = (0..8)
2569 .map(|_| {
2570 let s = a.clone();
2571 sc.spawn(move || s.claim_turn("t1").expect("claim"))
2572 })
2573 .collect();
2574 hs.into_iter().map(|h| h.join().expect("join")).collect()
2575 });
2576 assert_eq!(wins.iter().filter(|w| w.is_some()).count(), 1);
2577 }
2578
2579 fn age_file(path: &Path) {
2580 let f = std::fs::OpenOptions::new()
2581 .write(true)
2582 .open(path)
2583 .expect("open");
2584 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
2585 .expect("age");
2586 }
2587
2588 #[test]
2589 fn a_late_taker_of_a_broken_lock_cannot_disturb_its_replacement() {
2590 let (_tmp, a, _b) = lease_store();
2591 let lease = a.turn_path("t1");
2592 let lock = lease.with_extension("turn.lock");
2593 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2594 std::fs::write(&lock, "t1-dead").expect("dead lock");
2595 age_file(&lock);
2596 let b = TurnLock::take(&lease).expect("take").expect("b wins");
2598 let fresh = std::fs::read_to_string(&lock).expect("read");
2599 assert_eq!(fresh, b.token);
2600 let ticket = std::fs::read_dir(lock.parent().expect("dir"))
2603 .expect("dir")
2604 .flatten()
2605 .map(|e| e.path())
2606 .find(|p| p.to_string_lossy().ends_with(".break.0"))
2607 .expect("ticket");
2608 assert!(!create_exclusive(&ticket, "").expect("ticket"));
2609 assert!(TurnLock::take(&lease).expect("take").is_none());
2610 assert_eq!(std::fs::read_to_string(&lock).expect("read"), fresh);
2611 age_file(&lock);
2613 age_file(&ticket);
2614 std::mem::forget(b);
2615 let c = TurnLock::take(&lease)
2616 .expect("take")
2617 .expect("next generation");
2618 assert_ne!(c.token, fresh);
2619 }
2620
2621 #[test]
2622 fn a_live_ticket_blocks_and_a_stale_one_hands_over_to_the_next_generation() {
2623 let (_tmp, a, _b) = lease_store();
2624 let lease = a.turn_path("t1");
2625 let lock = lease.with_extension("turn.lock");
2626 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2627 std::fs::write(&lock, "t1-dead").expect("dead lock");
2628 age_file(&lock);
2629 let t0 = lock.with_extension("lock.t1-dead.break.0");
2630 assert!(create_exclusive(&t0, "").expect("ticket"));
2631 assert!(TurnLock::take(&lease).expect("take").is_none());
2633 assert_eq!(std::fs::read_to_string(&lock).expect("read"), "t1-dead");
2634 age_file(&t0);
2636 let c = TurnLock::take(&lease).expect("take").expect("generation 1");
2637 assert_eq!(std::fs::read_to_string(&lock).expect("read"), c.token);
2638 assert!(lock.with_extension("lock.t1-dead.break.1").exists());
2639 assert!(TurnLock::take(&lease).expect("take").is_none());
2640 }
2641
2642 #[test]
2643 fn exhausted_ticket_generations_recover_once_the_sweep_ages_them_out() {
2644 let (_tmp, a, _b) = lease_store();
2645 let lease = a.turn_path("t1");
2646 let lock = lease.with_extension("turn.lock");
2647 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2648 std::fs::write(&lock, "t1-dead").expect("dead lock");
2649 age_file(&lock);
2650 let tickets: Vec<_> = (0..TICKET_GENERATIONS)
2651 .map(|n| lock.with_extension(format!("lock.t1-dead.break.{n}")))
2652 .collect();
2653 for t in &tickets {
2654 assert!(create_exclusive(t, "").expect("ticket"));
2655 age_file(t);
2656 }
2657 assert!(TurnLock::take(&lease).expect("take").is_none());
2659 assert!(tickets.iter().all(|t| t.exists()));
2660 for t in &tickets {
2662 let f = std::fs::OpenOptions::new()
2663 .write(true)
2664 .open(t)
2665 .expect("open");
2666 f.set_modified(std::time::SystemTime::now() - TICKET_SWEEP_AGE * 2)
2667 .expect("age");
2668 }
2669 assert!(TurnLock::take(&lease).expect("take").is_none());
2670 assert!(TurnLock::take(&lease).expect("take").is_some());
2671 }
2672
2673 #[test]
2674 fn an_aged_empty_lock_is_broken() {
2675 let (_tmp, a, _b) = lease_store();
2676 let lease = a.turn_path("t1");
2677 let lock = lease.with_extension("turn.lock");
2678 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2679 std::fs::write(&lock, "").expect("empty lock");
2680 age_file(&lock);
2681 assert!(TurnLock::take(&lease).expect("take").is_some());
2682 }
2683
2684 #[test]
2685 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
2686 let (tmp, talks) = store();
2687 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2688 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2689
2690 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
2691 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
2692 let saved = talks.get(&talk.id).expect("reload queued talk");
2693 assert_eq!(saved.pending, "first\n\nsecond");
2694 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
2695
2696 let drained = drain(&mut talk, &talks).expect("drain");
2697 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
2698 let saved = talks.get(&talk.id).expect("reload drained talk");
2699 assert!(saved.pending.is_empty());
2700 assert_eq!(saved.turns.len(), 1);
2701 assert_eq!(saved.turns[0].body, "first\n\nsecond");
2702 }
2703
2704 #[test]
2705 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
2706 let (tmp, talks) = store();
2707 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2708 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2709 let attachment = Attachment {
2710 id: "a".repeat(32),
2711 name: "shot.png".to_owned(),
2712 mime: "image/png".to_owned(),
2713 bytes: 3,
2714 };
2715
2716 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
2717 assert!(
2718 edit_pending_text(
2719 &mut talk,
2720 &talks,
2721 "corrected",
2722 "first",
2723 std::slice::from_ref(&attachment.id),
2724 )
2725 .expect("edit")
2726 );
2727 let saved = talks.get(&talk.id).expect("reload edited draft");
2728 assert_eq!(saved.pending, "corrected");
2729 assert_eq!(saved.pending_attachments, vec![attachment]);
2730
2731 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
2732 assert!(
2733 !edit_pending_text(
2734 &mut talk,
2735 &talks,
2736 "stale edit",
2737 "corrected",
2738 &["a".repeat(32)],
2739 )
2740 .expect("stale edit is a conflict")
2741 );
2742 assert_eq!(
2743 talks.get(&talk.id).expect("reload after conflict").pending,
2744 "corrected\n\nlater"
2745 );
2746 assert!(
2747 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
2748 .expect("stale clear is a conflict")
2749 );
2750 assert_eq!(
2751 talks
2752 .get(&talk.id)
2753 .expect("reload after stale clear")
2754 .pending,
2755 "corrected\n\nlater"
2756 );
2757 }
2758
2759 #[tokio::test]
2760 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
2761 let (tmp, talks) = store();
2762 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
2763 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
2764 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2765 let id = running.id.clone();
2766 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
2767
2768 let response_talks = talks.clone();
2769 let response_cfg = cfg.clone();
2770 let reply = tokio::spawn(async move {
2771 respond(&mut running, &response_talks, &response_cfg, &first).await
2772 });
2773 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
2774
2775 let mut queued = talks.get(&id).expect("queued handle");
2776 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
2777 reply.await.expect("join").expect("reply");
2778
2779 let saved = talks.get(&id).expect("reload");
2780 assert_eq!(saved.pending, "next");
2781 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
2782 }
2783
2784 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
2787 let calls = dir.join(format!("{id}.calls"));
2788 let script = format!(
2789 "#!/bin/sh\necho x >> '{}'\n{body}\n",
2790 calls.to_string_lossy()
2791 );
2792 let path = dir.join(format!("mock-{id}.sh"));
2793 std::fs::write(&path, script).expect("write mock");
2794 AgentSpec {
2795 id: id.to_owned(),
2796 kind: AgentKind::Command,
2797 model: None,
2798 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2799 extra_args: Vec::new(),
2800 env: BTreeMap::new(),
2801 prompt_delivery: None,
2802 }
2803 }
2804
2805 fn calls(dir: &Path, id: &str) -> usize {
2806 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
2807 }
2808
2809 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
2810 let mut cfg = config(specs[0].clone());
2811 cfg.agents = specs;
2812 cfg.roles.chatter = Some(AgentChoice::Chain(
2813 ids.iter().map(|s| (*s).to_owned()).collect(),
2814 ));
2815 cfg
2816 }
2817
2818 #[tokio::test]
2819 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
2820 let (tmp, talks) = store();
2821 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2822 let b = counting_agent(tmp.path(), "b", "cat");
2823 let cfg = chain_config(vec![a, b], &["a", "b"]);
2824 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2825 assert_eq!(talk.agent, "a");
2826
2827 say(&mut talk, &talks, &cfg, "hello there", Vec::new())
2828 .await
2829 .expect("turn");
2830 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
2831 assert_eq!(calls(tmp.path(), "b"), 1);
2832 assert_eq!(talk.agent, "b", "the switch persists");
2833 assert!(talks.get(&talk.id).unwrap().agent == "b");
2834 let reply = talk.turns.last().unwrap();
2835 assert!(reply.body.contains("hello there"));
2836 assert!(
2837 reply.body.contains("magi task add --solo"),
2838 "a fresh seat gets the full briefing"
2839 );
2840 assert!(
2841 talk.turns
2842 .iter()
2843 .any(|t| t.body.contains("agent changed from a to b")),
2844 "the switch is noted"
2845 );
2846 }
2847
2848 #[tokio::test]
2849 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
2850 let (tmp, talks) = store();
2851 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2852 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
2853 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
2854 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2855
2856 let err = say(&mut talk, &talks, &cfg, "hi", Vec::new())
2857 .await
2858 .expect_err("every agent failed");
2859 assert!(err.to_string().contains("`a`"), "{err:#}");
2860 assert_eq!(calls(tmp.path(), "a"), 1);
2861 assert_eq!(calls(tmp.path(), "b"), 1);
2862 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
2863 }
2864
2865 #[test]
2866 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
2867 let (tmp, talks) = store();
2868 let b = counting_agent(tmp.path(), "b", "cat");
2869 let cfg = chain_config(vec![b], &["ghost", "b"]);
2870 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2871 assert_eq!(talk.agent, "b");
2872 }
2873
2874 #[tokio::test]
2875 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
2876 let (tmp, talks) = store();
2877 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2878 let b = counting_agent(tmp.path(), "b", "cat");
2879 let cfg = chain_config(vec![a, b], &["a", "b"]);
2880 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
2881 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2882 .await
2883 .expect_err("a alone, and it fails");
2884 assert_eq!(calls(tmp.path(), "b"), 0);
2885 assert_eq!(talk.agent, "a");
2886 }
2887
2888 #[tokio::test]
2889 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
2890 let (tmp, talks) = store();
2891 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2892 let b = counting_agent(tmp.path(), "b", "cat");
2893 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
2894 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
2895 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
2896 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2897 .await
2898 .expect_err("c alone, and it fails");
2899 assert_eq!(calls(tmp.path(), "b"), 0);
2900 }
2901
2902 #[test]
2903 fn briefing_carries_a_persona_section_only_when_one_is_chosen() {
2904 let plain = briefing(Path::new("/repo"), "en", false);
2905 assert_eq!(plain, briefing_with(Path::new("/repo"), "en", false, None));
2906 assert!(!plain.contains("Persona"));
2907 let rei = crate::persona::builtin_catalog()
2908 .into_iter()
2909 .find(|p| p.id == "rei")
2910 .expect("rei");
2911 let with = briefing_with(Path::new("/repo"), "en", false, Some(&rei));
2912 assert!(with.starts_with(&plain), "the plain briefing is untouched");
2913 assert!(with.contains("# Persona (tone only)"));
2914 assert!(with.contains("TONE ONLY"));
2915 assert!(with.contains("task ids"));
2916 assert!(with.contains("write policy"));
2917 assert!(with.contains("`magi task add`"));
2918 assert!(with.contains("Rei Ayanami"));
2919 }
2920
2921 #[test]
2922 fn a_talk_written_before_personas_still_loads() {
2923 let (tmp, talks) = store();
2924 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2925 let cfg = config(spec);
2926 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2927 let path = talks.root.join(format!("{}.json", talk.id));
2928 let mut v: serde_json::Value =
2929 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
2930 v.as_object_mut().unwrap().remove("persona");
2931 v.as_object_mut().unwrap().remove("persona_dirty");
2932 std::fs::write(&path, v.to_string()).unwrap();
2933 let loaded = talks.get(&talk.id).expect("old record loads");
2934 assert_eq!(loaded.persona, "");
2935 assert!(!loaded.persona_dirty);
2936 }
2937
2938 #[tokio::test]
2939 async fn a_persona_switch_notes_marks_dirty_and_updates_the_next_turn_once() {
2940 let (tmp, talks) = store();
2941 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2942 let cfg = config(spec);
2943 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2944 say(&mut talk, &talks, &cfg, "hello", Vec::new())
2945 .await
2946 .expect("first turn");
2947 assert!(!talk.turns[1].body.contains("Persona"));
2948
2949 assert!(switch_persona(&mut talk, &talks, "misato").expect("switch"));
2950 assert!(!switch_persona(&mut talk, &talks, "misato").expect("same"));
2951 assert!(talk.persona_dirty);
2952 assert!(talk.turns.last().unwrap().body.contains("persona changed"));
2953
2954 say(&mut talk, &talks, &cfg, "next", Vec::new())
2955 .await
2956 .expect("turn");
2957 let prompt = &talk.turns.last().unwrap().body;
2958 assert!(prompt.contains("# Persona update"), "{prompt}");
2959 assert!(prompt.contains("Misato Katsuragi"));
2960 assert!(!talk.persona_dirty, "cleared after a successful turn");
2961
2962 say(&mut talk, &talks, &cfg, "again", Vec::new())
2963 .await
2964 .expect("turn");
2965 assert!(!talk.turns.last().unwrap().body.contains("# Persona update"));
2966
2967 assert!(switch_persona(&mut talk, &talks, "default").expect("back"));
2968 assert_eq!(talk.persona, "");
2969 say(&mut talk, &talks, &cfg, "plain", Vec::new())
2970 .await
2971 .expect("turn");
2972 assert!(
2973 talk.turns
2974 .last()
2975 .unwrap()
2976 .body
2977 .contains("turned the persona off")
2978 );
2979
2980 assert!(switch_persona(&mut talk, &talks, "rei").expect("rei"));
2982 let mut no_sessions = cfg.clone();
2983 no_sessions.graph.sessions = false;
2984 say(&mut talk, &talks, &no_sessions, "one", Vec::new())
2985 .await
2986 .expect("turn");
2987 say(&mut talk, &talks, &no_sessions, "two", Vec::new())
2988 .await
2989 .expect("turn");
2990 let last = &talk.turns.last().unwrap().body;
2991 assert!(!talk.persona_dirty);
2992 assert!(last.contains("# Persona (tone only)"), "{last}");
2993 assert!(last.contains("Rei Ayanami"));
2994 }
2995
2996 #[tokio::test]
2997 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
2998 let (tmp, talks) = store();
2999 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3000 let cfg = config(spec);
3001 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3002
3003 say(
3004 &mut talk,
3005 &talks,
3006 &cfg,
3007 "what does the queue module do?",
3008 Vec::new(),
3009 )
3010 .await
3011 .expect("first turn");
3012 let first_prompt = &talk.turns[1].body;
3013 assert!(first_prompt.contains("magi task add --solo"));
3014 assert!(first_prompt.contains("what does the queue module do?"));
3015
3016 say(&mut talk, &talks, &cfg, "and how is it locked?", Vec::new())
3017 .await
3018 .expect("second turn");
3019 let second_prompt = &talk.turns[3].body;
3020 assert!(
3021 !second_prompt.contains("magi task add --solo"),
3022 "the briefing is sent once, not on every turn: {second_prompt}"
3023 );
3024 assert!(second_prompt.contains("and how is it locked?"));
3025 }
3026
3027 #[tokio::test]
3028 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
3029 let (tmp, talks) = store();
3030 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3031 let mut b = a.clone();
3032 b.id = "other".to_owned();
3033 let mut cfg = config(a.clone());
3034 cfg.agents.push(b.clone());
3035 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
3036 say(&mut talk, &talks, &cfg, "remember the walrus", Vec::new())
3037 .await
3038 .expect("first turn");
3039 let old_session = talk.seat.claude_session.clone();
3040 assert_eq!(talk.seat.turns, 1);
3041
3042 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
3043 assert_eq!(talk.agent, "other");
3044 assert_eq!(talk.seat.turns, 0);
3045 assert_eq!(talk.seat.agent, "other");
3046 assert_ne!(talk.seat.claude_session, old_session);
3047 let note = talk.turns.last().expect("note");
3048 assert_eq!(note.who, Who::Agent);
3049 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3050 assert!(note.body.contains("changed from"), "{}", note.body);
3051 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
3052
3053 let before = talk.turns.len();
3054 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
3055 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
3056
3057 say(&mut talk, &talks, &cfg, "what did I say?", Vec::new())
3058 .await
3059 .expect("turn after switch");
3060 let prompt = &talk.turns.last().expect("reply").body;
3061 assert!(prompt.contains("remember the walrus"), "{prompt}");
3062 assert!(prompt.contains("## magi"), "{prompt}");
3063 assert!(prompt.contains("what did I say?"), "{prompt}");
3064 }
3065
3066 #[tokio::test]
3067 async fn say_appends_the_operator_turn_then_the_agent_turn() {
3068 let (tmp, talks) = store();
3069 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3070 let cfg = config(spec);
3071 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3072
3073 say(
3074 &mut talk,
3075 &talks,
3076 &cfg,
3077 "can I rename this function?",
3078 Vec::new(),
3079 )
3080 .await
3081 .expect("say");
3082
3083 assert_eq!(talk.turns.len(), 2);
3084 assert_eq!(talk.turns[0].who, Who::Operator);
3085 assert_eq!(talk.turns[0].body, "can I rename this function?");
3086 assert_eq!(talk.turns[1].who, Who::Agent);
3087 assert_eq!(talk.turns[1].body, "go ahead");
3088 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
3089 }
3090
3091 #[tokio::test]
3092 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
3093 let (tmp, talks) = store();
3094 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
3095 let cfg = config(spec);
3096 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3097
3098 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
3099 .await
3100 .expect_err("a turn with no answer is an error");
3101 assert!(err.to_string().contains("no answer"), "{err}");
3102
3103 let on_disk = talks.get(&talk.id).expect("get");
3104 assert_eq!(on_disk.turns.len(), 2);
3105 assert_eq!(on_disk.turns[0].body, "check the tests");
3106 let note = &on_disk.turns[1];
3107 assert_eq!(note.who, Who::Agent);
3108 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3109 assert!(note.body.contains("your message is saved"));
3110 }
3111
3112 #[tokio::test]
3118 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
3119 let (tmp, talks) = store();
3120 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3121 let cfg = config(spec);
3122 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3123
3124 let text =
3125 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
3126 failpoint::force_put_failures(PUT_RETRIES - 1);
3129 respond(&mut talk, &talks, &cfg, &text)
3130 .await
3131 .expect("respond must survive a write failure its own retries can outlast");
3132
3133 assert_eq!(talk.turns.len(), 2);
3134 assert_eq!(talk.turns[1].who, Who::Agent);
3135 assert_eq!(talk.turns[1].body, "go ahead");
3136 let on_disk = talks.get(&talk.id).expect("get");
3137 assert_eq!(
3138 on_disk.turns, talk.turns,
3139 "the reply must reach disk despite the early write failures"
3140 );
3141 }
3142
3143 #[tokio::test]
3149 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
3150 let (tmp, talks) = store();
3151 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3152 let cfg = config(spec);
3153 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3154
3155 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3156 failpoint::force_put_failures(PUT_RETRIES);
3161 let err = respond(&mut talk, &talks, &cfg, &text)
3162 .await
3163 .expect_err("a reply that cannot be saved must be reported, not swallowed");
3164 assert!(err.to_string().contains("could not be saved"), "{err}");
3165
3166 let on_disk = talks.get(&talk.id).expect("get");
3167 assert_eq!(
3168 on_disk.turns.len(),
3169 2,
3170 "the operator turn plus a visible note"
3171 );
3172 assert_eq!(on_disk.turns[0].body, "check the tests");
3173 let note = &on_disk.turns[1];
3174 assert_eq!(note.who, Who::Agent);
3175 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3176 assert!(
3177 note.body.contains("could not be saved"),
3178 "the operator must be told the reply is missing, not left staring \
3179 at a gap with no explanation: {}",
3180 note.body
3181 );
3182 assert_eq!(
3183 talk.turns, on_disk.turns,
3184 "the in-memory talk must match what actually landed on disk"
3185 );
3186
3187 let artifacts = talks.artifacts_of(&talk.id);
3190 let stash = std::fs::read_dir(&artifacts)
3191 .expect("artifacts dir")
3192 .filter_map(|e| e.ok())
3193 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
3194 .expect("a stash file for the lost reply");
3195 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
3196 assert_eq!(stashed, "go ahead");
3197
3198 assert_eq!(
3206 on_disk.seat.turns, 1,
3207 "the note's write must carry the turn the CLI actually took"
3208 );
3209 assert_eq!(
3210 on_disk.seat.claude_session, talk.seat.claude_session,
3211 "the session id handed to the CLI must survive the failed reply"
3212 );
3213 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
3214 assert!(
3215 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
3216 "the next turn must resume, not open the same session id twice"
3217 );
3218 }
3219
3220 #[tokio::test]
3225 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
3226 let (tmp, talks) = store();
3227 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3228 let cfg = config(spec);
3229 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3230
3231 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3232 failpoint::force_put_failures(PUT_RETRIES * 2);
3235 let err = respond(&mut talk, &talks, &cfg, &text)
3236 .await
3237 .expect_err("neither the reply nor the note could be saved");
3238 assert!(err.to_string().contains("could not be saved"), "{err}");
3239
3240 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
3241 let on_disk = talks.get(&talk.id).expect("get");
3242 assert_eq!(on_disk.turns.len(), 1);
3243
3244 assert_eq!(
3254 on_disk.seat.turns, 0,
3255 "an unwritable file cannot record the turn the CLI took"
3256 );
3257 assert_eq!(
3258 talk.seat.turns, 1,
3259 "the in-memory seat still reports the turn the CLI actually took"
3260 );
3261 assert_eq!(
3262 on_disk.seat.claude_session, talk.seat.claude_session,
3263 "the session id was minted at `begin` and never changes here"
3264 );
3265 }
3266
3267 #[tokio::test]
3271 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
3272 let (tmp, talks) = store();
3273 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3274 let cfg = config(spec);
3275 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3276
3277 let att = talks
3278 .put_attachment(
3279 &talk.id,
3280 "image/png",
3281 "screenshot.png",
3282 b"pretend-png-bytes",
3283 )
3284 .expect("put attachment");
3285
3286 say(&mut talk, &talks, &cfg, "", vec![att.clone()])
3287 .await
3288 .expect("an empty body with an attachment is still a turn");
3289
3290 let operator_turn = &talk.turns[0];
3291 assert_eq!(operator_turn.who, Who::Operator);
3292 assert_eq!(operator_turn.body, "");
3293 assert_eq!(operator_turn.attachments, vec![att.clone()]);
3294
3295 let prompt = &talk.turns[1].body;
3296 let expected_path = talks
3297 .attachments_dir(&talk.id)
3298 .join(format!("{}.png", att.id));
3299 assert!(
3300 prompt.contains(&expected_path.display().to_string()),
3301 "the agent must be told the attachment's absolute path: {prompt}"
3302 );
3303 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
3304 }
3305
3306 #[test]
3315 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
3316 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
3317 let att = Attachment {
3318 id: "0".repeat(32),
3319 name: "shot.png".to_owned(),
3320 mime: "image/png".to_owned(),
3321 bytes: 3,
3322 };
3323 let path = talks
3324 .attachment_path("some-talk-id", &att)
3325 .expect("a supported mime always yields a path");
3326 assert!(
3327 path.is_absolute(),
3328 "must be absolute even off a relative store root: {}",
3329 path.display()
3330 );
3331 }
3332
3333 #[tokio::test]
3334 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
3335 let (tmp, talks) = store();
3340 let slow = mock_agent(
3341 tmp.path(),
3342 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
3343 BTreeMap::new(),
3344 );
3345 let mut cfg = config(slow);
3346 cfg.graph.timeout_talk = 1;
3347 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3348
3349 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
3350 .await
3351 .expect_err("a turn that never answers is an error");
3352 assert!(
3353 err.to_string().contains("did not answer within 1s"),
3354 "{err}"
3355 );
3356
3357 let on_disk = talks.get(&talk.id).expect("get");
3358 let note = on_disk.turns.last().expect("a note turn was recorded");
3359 assert!(
3360 note.body.contains("did not answer within 1s"),
3361 "the transcript must show the configured timeout: {}",
3362 note.body
3363 );
3364 }
3365
3366 #[test]
3367 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
3368 let (tmp, talks) = store();
3369 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3370 let cfg = config(spec);
3371 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3372
3373 close(&mut talk, &talks).expect("close");
3374 assert_eq!(talk.status, TalkStatus::Closed);
3375 close(&mut talk, &talks).expect("closing twice is not an error");
3376
3377 let err =
3378 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
3379 assert!(err.to_string().contains("closed"));
3380 let _ = &cfg; }
3382
3383 #[tokio::test]
3384 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3385 let (tmp, talks) = store();
3386 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3387 let cfg = config(spec);
3388 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3391
3392 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
3396 close(&mut closed_elsewhere, &talks).expect("close");
3397 assert_eq!(
3398 talks.get(&in_flight.id).expect("reread").status,
3399 TalkStatus::Closed,
3400 "the close landed on disk before the turn finished"
3401 );
3402
3403 assert_eq!(in_flight.status, TalkStatus::Open);
3407 respond(&mut in_flight, &talks, &cfg, "one more question")
3408 .await
3409 .expect("the turn itself still completes");
3410
3411 let on_disk = talks.get(&in_flight.id).expect("reread");
3412 assert_eq!(
3413 on_disk.status,
3414 TalkStatus::Closed,
3415 "a close must stick even when a turn that started before it finishes after it"
3416 );
3417 assert!(
3420 on_disk.turns.iter().any(|t| t.body == "here you go"),
3421 "the in-flight turn's own reply is still recorded: {:?}",
3422 on_disk.turns
3423 );
3424 }
3425
3426 #[test]
3427 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
3428 let (tmp, talks) = store();
3429 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3430 let cfg = config(spec);
3431 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3434
3435 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
3438 close(&mut closed_elsewhere, &talks).expect("close");
3439 assert_eq!(
3440 talks.get(&stale.id).expect("reread").status,
3441 TalkStatus::Closed,
3442 "the close landed on disk before record was called"
3443 );
3444
3445 assert_eq!(stale.status, TalkStatus::Open);
3449 let err = record(&mut stale, &talks, "still there?", Vec::new())
3450 .expect_err("a close that landed first must be honored, not overwritten");
3451 assert!(err.to_string().contains("closed"));
3452
3453 let on_disk = talks.get(&stale.id).expect("reread");
3454 assert_eq!(
3455 on_disk.status,
3456 TalkStatus::Closed,
3457 "record must not resurrect a conversation closed while its snapshot was stale"
3458 );
3459 assert!(
3460 on_disk.turns.is_empty(),
3461 "the rejected turn must not have been appended: {:?}",
3462 on_disk.turns
3463 );
3464 let _ = &cfg; }
3466
3467 #[test]
3468 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
3469 let (tmp, talks) = store();
3470 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3471 let cfg = config(spec);
3472 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3473
3474 let held = talks.guard();
3478
3479 let talks2 = talks.clone();
3480 let id = talk.id.clone();
3481 let closing = std::thread::spawn(move || {
3482 let mut talk = talks2.get(&id).expect("get");
3483 close(&mut talk, &talks2).expect("close");
3484 });
3485
3486 std::thread::sleep(Duration::from_millis(50));
3487 assert!(
3488 !closing.is_finished(),
3489 "close must wait for the guard, not read and write while it is held - \
3490 a re-read alone narrows this window without closing it"
3491 );
3492
3493 drop(held);
3494 closing.join().expect("close thread panicked");
3495
3496 assert_eq!(
3497 talks.get(&talk.id).expect("reread").status,
3498 TalkStatus::Closed,
3499 "once the guard is free, close still lands"
3500 );
3501 let _ = &cfg; }
3503
3504 #[test]
3505 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
3506 let (tmp, talks) = store();
3507 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3508 let cfg = config(spec);
3509 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3510
3511 close(&mut talk, &talks).expect("close");
3512 assert_eq!(talk.status, TalkStatus::Closed);
3513
3514 reopen(&mut talk, &talks).expect("reopen");
3515 assert_eq!(talk.status, TalkStatus::Open);
3516 assert_eq!(
3517 talks.get(&talk.id).expect("reread").status,
3518 TalkStatus::Open
3519 );
3520
3521 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
3523 assert_eq!(talk.status, TalkStatus::Open);
3524
3525 record(&mut talk, &talks, "one more thing", Vec::new())
3526 .expect("a reopened talk takes turns again");
3527 let _ = &cfg; }
3529
3530 #[test]
3531 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
3532 let (tmp, talks) = store();
3533 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3534 let cfg = config(spec);
3535 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3536
3537 let artifacts = talks.artifacts_of(&talk.id);
3538 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
3539 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
3540
3541 talks.remove(&talk.id).expect("remove");
3542 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
3543 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
3544 assert!(
3545 talks.get(&talk.id).is_err(),
3546 "a removed talk cannot be read back"
3547 );
3548
3549 let err = talks
3550 .remove("nonexistent-id")
3551 .expect_err("unknown id refused");
3552 assert!(err.to_string().contains("no talk matches"), "{err}");
3553 let _ = &cfg; }
3555
3556 #[tokio::test]
3557 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3558 let (tmp, talks) = store();
3559 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3560 let cfg = config(spec);
3561 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3564
3565 talks.remove(&in_flight.id).expect("remove");
3566 assert!(
3567 talks.get(&in_flight.id).is_err(),
3568 "the delete landed on disk before the turn finished"
3569 );
3570
3571 respond(&mut in_flight, &talks, &cfg, "one more question")
3574 .await
3575 .expect("the turn itself still completes rather than erroring");
3576
3577 assert!(
3578 talks.get(&in_flight.id).is_err(),
3579 "a delete must stick even when a turn that started before it finishes after it"
3580 );
3581 }
3582
3583 #[test]
3584 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
3585 let (tmp, talks) = store();
3586 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3587 let cfg = config(spec);
3588 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3591
3592 talks.remove(&stale.id).expect("remove");
3593
3594 let err = record(&mut stale, &talks, "still there?", Vec::new())
3598 .expect_err("a delete that landed first must be honored, not overwritten");
3599 assert!(err.to_string().contains("deleted"), "{err}");
3600
3601 assert!(
3602 talks.get(&stale.id).is_err(),
3603 "record must not resurrect a conversation deleted while its snapshot was stale"
3604 );
3605 let _ = &cfg; }
3607
3608 #[test]
3609 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
3610 let (tmp, talks) = store();
3611 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3612 let cfg = config(spec);
3613 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3616
3617 talks.remove(&stale.id).expect("remove");
3618
3619 let err = close(&mut stale, &talks)
3623 .expect_err("a delete that landed first must be honored, not overwritten");
3624 assert!(err.to_string().contains("deleted"), "{err}");
3625
3626 assert!(
3627 talks.get(&stale.id).is_err(),
3628 "close must not resurrect a conversation deleted while its snapshot was stale"
3629 );
3630 let _ = &cfg; }
3632
3633 #[test]
3634 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
3635 let (tmp, talks) = store();
3636 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3637 let cfg = config(spec);
3638 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3641 close(&mut stale, &talks).expect("close");
3642
3643 talks.remove(&stale.id).expect("remove");
3644
3645 let err = reopen(&mut stale, &talks)
3649 .expect_err("a delete that landed first must be honored, not overwritten");
3650 assert!(err.to_string().contains("deleted"), "{err}");
3651
3652 assert!(
3653 talks.get(&stale.id).is_err(),
3654 "reopen must not resurrect a conversation deleted while its snapshot was stale"
3655 );
3656 let _ = &cfg; }
3658
3659 #[test]
3660 fn list_puts_open_talks_before_closed_ones() {
3661 let (tmp, talks) = store();
3662 let make = |id: &str, status: TalkStatus| {
3663 let mut t = Talk {
3664 schema: SCHEMA,
3665 id: id.to_owned(),
3666 repo: tmp.path().to_owned(),
3667 agent: "mock".to_owned(),
3668 status,
3669 turns: Vec::new(),
3670 pending: String::new(),
3671 pending_attachments: Vec::new(),
3672 fallback: false,
3673 persona: String::new(),
3674 persona_dirty: false,
3675 created_at: Timestamp::now(),
3676 updated_at: Timestamp::now(),
3677 seat: SeatState::new(SEAT, "mock", 7),
3678 };
3679 talks.put(&mut t).expect("put");
3680 };
3681 make("20260901-000000-0001", TalkStatus::Open);
3682 make("20260902-000000-0002", TalkStatus::Open);
3683 make("20260903-000000-0003", TalkStatus::Closed);
3684
3685 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
3686 assert_eq!(
3687 ids,
3688 [
3689 "20260902-000000-0002",
3690 "20260901-000000-0001",
3691 "20260903-000000-0003"
3692 ]
3693 );
3694 assert_eq!(talks.count_open(), 2);
3695 }
3696
3697 #[test]
3698 fn tasks_of_finds_only_this_talks_own_tasks() {
3699 let dir = tempfile::tempdir().expect("tempdir");
3700 let queue = Queue::at(dir.path().join("queue"));
3701
3702 let mut mine = Task::new(
3703 "rework the loader".to_owned(),
3704 "rework the loader".to_owned(),
3705 PathBuf::from("/repo"),
3706 Source::Agent {
3707 run: "20260904-014455-ab12".to_owned(),
3708 node: "chat".to_owned(),
3709 },
3710 );
3711 queue.put(&mut mine).expect("put mine");
3712
3713 let mut theirs = Task::new(
3714 "unrelated".to_owned(),
3715 "unrelated".to_owned(),
3716 PathBuf::from("/repo"),
3717 Source::Agent {
3718 run: "20260904-090000-zz99".to_owned(),
3719 node: "implement".to_owned(),
3720 },
3721 );
3722 queue.put(&mut theirs).expect("put theirs");
3723
3724 let mut human = Task::new(
3725 "typed by hand".to_owned(),
3726 "typed by hand".to_owned(),
3727 PathBuf::from("/repo"),
3728 Source::Human,
3729 );
3730 queue.put(&mut human).expect("put human");
3731
3732 let found = tasks_of(&queue, "20260904-014455-ab12");
3733 assert_eq!(found.len(), 1);
3734 assert_eq!(found[0].id, mine.id);
3735 }
3736
3737 #[test]
3738 fn the_briefing_names_solo_task_add() {
3739 let brief = briefing(Path::new("/repo"), "en", false);
3740 assert!(brief.contains("magi task add --solo"));
3741 assert!(brief.contains("/repo"));
3742 assert!(!brief.contains("Hold this conversation in"));
3743 }
3744
3745 #[test]
3751 fn the_briefing_explains_targeting_a_different_repository_by_name() {
3752 let brief = briefing(Path::new("/repo"), "en", false);
3753 assert!(brief.contains("--repo does not have to be a full path"));
3754 assert!(brief.contains("owner/repo"));
3755 assert!(brief.contains("magi repos"));
3756 assert!(brief.contains("ask the operator"));
3757 }
3758
3759 #[test]
3760 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
3761 let brief = briefing(Path::new("/repo"), "en", false);
3762 assert!(brief.contains("--attach <path>"), "{brief}");
3763 assert!(brief.contains("deleting this conversation"), "{brief}");
3764 }
3765
3766 #[test]
3767 fn the_briefing_names_the_language_when_it_is_not_english() {
3768 let brief = briefing(Path::new("/repo"), "Japanese", false);
3769 assert!(brief.contains("Hold this conversation in Japanese"));
3770 }
3771
3772 #[test]
3773 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
3774 let read_only = briefing(Path::new("/repo"), "en", false);
3775 assert!(read_only.contains("Do not write files"));
3776 assert!(!read_only.contains("allow_write"));
3777
3778 let writable = briefing(Path::new("/repo"), "en", true);
3779 assert!(!writable.contains("Do not write files"));
3780 assert!(writable.contains("allow_write = true"));
3781 assert!(writable.contains("magi task add --solo"));
3784 assert!(writable.contains("say plainly what you"));
3785 }
3786}