1use std::path::{Path, PathBuf};
52use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
53use std::time::Duration;
54
55use anyhow::{Context, Result, bail};
56use jiff::Timestamp;
57use serde::{Deserialize, Serialize};
58
59use crate::agent::{self, Invocation, SeatState};
60use crate::config::{AgentSpec, Config};
61use crate::queue::{Queue, Source, Task};
62
63pub const SCHEMA: u32 = 1;
65
66fn turn_timeout(cfg: &Config) -> Duration {
77 Duration::from_secs(cfg.graph.timeout_talk)
78}
79
80const SEAT: &str = "talk";
83
84const MAGI_NOTE: &str = "magi: ";
86
87#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
89#[serde(rename_all = "lowercase")]
90pub enum Who {
91 Operator,
93 Agent,
96}
97
98#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
106#[serde(deny_unknown_fields)]
107pub struct Attachment {
108 pub id: String,
110 pub name: String,
112 pub mime: String,
115 pub bytes: u64,
117}
118
119#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
121#[serde(deny_unknown_fields)]
122pub struct Turn {
123 pub who: Who,
125 pub body: String,
127 pub at: Timestamp,
129 #[serde(default)]
132 pub attachments: Vec<Attachment>,
133 #[serde(default, skip_serializing_if = "Option::is_none")]
138 pub usage: Option<TurnUsage>,
139}
140
141#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
149pub struct TurnUsage {
150 pub context_tokens: u64,
152 pub agent: String,
154 #[serde(default)]
156 pub model: Option<String>,
157}
158
159#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
164pub struct ContextUsage {
165 pub tokens: Option<u64>,
167 pub window: Option<u64>,
169 pub percent: Option<u64>,
171 pub warn: bool,
173 pub since_switch: bool,
177 pub model: Option<String>,
179 #[serde(default)]
182 pub estimated: bool,
183}
184
185const CONTEXT_WARN_PERCENT: u64 = 80;
187
188const STANDING_PROMPT_FALLBACK_CHARS: u64 = 7000;
191
192pub fn estimate_context_tokens(talk: &Talk, standing_chars: u64) -> Option<u64> {
205 let mut counted = false;
206 let mut chars = standing_chars;
207 for t in talk.turns.iter().filter(|t| !t.body.starts_with(MAGI_NOTE)) {
208 counted = true;
209 chars += t.body.chars().count() as u64;
210 }
211 counted.then(|| (chars * 2).div_ceil(7))
213}
214
215pub fn context_usage(talk: &Talk, cfg: Option<&Config>) -> ContextUsage {
228 let current = cfg.and_then(|c| c.agents.iter().find(|a| a.id == talk.agent));
229 let model = current.and_then(|a| a.model.clone());
230 let window = cfg
231 .zip(model.as_deref())
232 .and_then(|(c, m)| c.context_window(m))
233 .filter(|w| *w > 0);
234 let usage = talk
235 .turns
236 .iter()
237 .rev()
238 .find(|t| t.who == Who::Agent && !t.body.starts_with(MAGI_NOTE))
239 .and_then(|t| t.usage.as_ref());
240 let measured = usage.map(|u| u.context_tokens);
241 let tokens = measured.or_else(|| {
242 let standing = cfg.map_or(STANDING_PROMPT_FALLBACK_CHARS, |c| {
243 briefing_with(
244 &talk.repo,
245 &c.graph.language,
246 c.talk.allow_write,
247 crate::persona::active(&c.talk.personas, &talk.persona).as_ref(),
248 c.talk.operator_name(),
249 )
250 .chars()
251 .count() as u64
252 });
253 estimate_context_tokens(talk, standing)
254 });
255 let estimated = measured.is_none() && tokens.is_some();
256 let since_switch =
257 usage.is_some_and(|u| u.agent != talk.agent || (current.is_some() && u.model != model));
258 let (percent, warn) = match (tokens, window) {
259 (Some(t), Some(w)) => (
260 Some(t.saturating_mul(100) / w),
261 t.saturating_mul(100) >= w.saturating_mul(CONTEXT_WARN_PERCENT),
262 ),
263 _ => (None, false),
264 };
265 ContextUsage {
266 tokens,
267 window,
268 percent,
269 warn,
270 since_switch,
271 model,
272 estimated,
273 }
274}
275
276#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
280#[serde(rename_all = "lowercase")]
281pub enum TalkStatus {
282 Open,
285 Closed,
287}
288
289impl TalkStatus {
290 pub fn open(self) -> bool {
292 matches!(self, Self::Open)
293 }
294
295 pub fn as_str(self) -> &'static str {
297 match self {
298 Self::Open => "open",
299 Self::Closed => "closed",
300 }
301 }
302}
303
304#[derive(Debug, Clone, Serialize, Deserialize)]
306#[serde(deny_unknown_fields)]
307pub struct Talk {
308 pub schema: u32,
310 pub id: String,
312 pub repo: PathBuf,
314 pub agent: String,
316 pub status: TalkStatus,
318 pub turns: Vec<Turn>,
320 #[serde(default)]
323 pub pending: String,
324 #[serde(default)]
326 pub pending_attachments: Vec<Attachment>,
327 #[serde(default)]
332 pub fallback: bool,
333 #[serde(default)]
336 pub persona: String,
337 #[serde(default)]
341 pub persona_dirty: bool,
342 pub created_at: Timestamp,
344 pub updated_at: Timestamp,
346 seat: SeatState,
351}
352
353impl Talk {
354 pub fn short(&self) -> &str {
356 short(&self.id)
357 }
358}
359
360#[derive(Debug, Clone)]
362pub struct Talks {
363 root: PathBuf,
364 lock: Arc<Mutex<()>>,
373}
374
375impl Talks {
376 pub fn open() -> Self {
378 Self::at(crate::run::home().join("talks"))
379 }
380
381 pub fn at(root: PathBuf) -> Self {
384 Self {
385 root,
386 lock: Arc::new(Mutex::new(())),
387 }
388 }
389
390 fn guard(&self) -> Result<StoreGuard<'_>> {
405 let mutex = self.lock.lock().unwrap_or_else(PoisonError::into_inner);
406 std::fs::create_dir_all(&self.root)
407 .with_context(|| format!("create {}", self.root.display()))?;
408 let path = self.root.join(".store.turn");
409 let deadline = std::time::Instant::now() + TAKEOVER_LOCK_TTL * 2;
410 loop {
411 if let Some(file) = TurnLock::take(&path)? {
412 return Ok(StoreGuard {
413 _file: file,
414 _mutex: mutex,
415 });
416 }
417 if std::time::Instant::now() >= deadline {
418 bail!("the talk store is locked by another process");
419 }
420 std::thread::sleep(Duration::from_millis(10));
421 }
422 }
423
424 pub fn root(&self) -> &Path {
426 &self.root
427 }
428
429 pub fn path_of(&self, id: &str) -> PathBuf {
431 self.root.join(format!("{id}.json"))
432 }
433
434 pub fn artifacts_of(&self, id: &str) -> PathBuf {
437 self.root.join(format!("{id}.artifacts"))
438 }
439
440 pub fn attachments_dir(&self, id: &str) -> PathBuf {
444 self.artifacts_of(id).join("attachments")
445 }
446
447 pub fn put_attachment(
456 &self,
457 id: &str,
458 mime: &str,
459 name: &str,
460 data: &[u8],
461 ) -> Result<Attachment> {
462 let dir = self.attachments_dir(id);
463 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
464 let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
465 let att = Attachment {
466 id: new_attachment_id(),
467 name: name.to_owned(),
468 mime: mime.to_owned(),
469 bytes: data.len() as u64,
470 };
471 std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
472 .with_context(|| format!("write attachment {}", att.id))?;
473 std::fs::write(
474 dir.join(format!("{}.json", att.id)),
475 serde_json::to_string(&att).context("serialize attachment")?,
476 )
477 .with_context(|| format!("write attachment metadata {}", att.id))?;
478 Ok(att)
479 }
480
481 pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
490 if !valid_attachment_id(att_id) {
491 return Ok(None);
492 }
493 let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
494 if !meta_path.is_file() {
495 return Ok(None);
496 }
497 let att = serde_json::from_str(
498 &std::fs::read_to_string(&meta_path)
499 .with_context(|| format!("read {}", meta_path.display()))?,
500 )
501 .with_context(|| format!("parse {}", meta_path.display()))?;
502 Ok(Some(att))
503 }
504
505 pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
509 let Some(att) = self.attachment_meta(id, att_id)? else {
510 return Ok(None);
511 };
512 let ext = attachment_ext(&att.mime).with_context(|| {
513 format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
514 })?;
515 let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
516 let data =
517 std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
518 Ok(Some((att, data)))
519 }
520
521 fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
538 let ext = attachment_ext(&att.mime)?;
539 let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
540 std::path::absolute(&path).ok()
541 }
542
543 pub fn put(&self, t: &mut Talk) -> Result<()> {
552 std::fs::create_dir_all(&self.root)
553 .with_context(|| format!("create {}", self.root.display()))?;
554 t.updated_at = Timestamp::now();
555 let body = serde_json::to_string_pretty(t).context("serialize talk")?;
556 let path = self.path_of(&t.id);
557 let tmp = path.with_extension("json.tmp");
558 write_atomic(&tmp, &path, &body)
559 }
560
561 pub fn get(&self, id: &str) -> Result<Talk> {
563 let resolved = self.resolve_id(id)?;
564 read_path(&self.path_of(&resolved))
565 }
566
567 pub fn list(&self) -> Vec<Talk> {
570 self.list_counting_unreadable().0
571 }
572
573 pub fn list_counting_unreadable(&self) -> (Vec<Talk>, usize) {
576 let mut unreadable = 0;
577 let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
578 .into_iter()
579 .flatten()
580 .flatten()
581 .map(|e| e.path())
582 .filter(|p| p.extension().is_some_and(|x| x == "json"))
583 .filter_map(|p| {
584 let talk = read_path(&p).ok();
585 if talk.is_none() {
586 unreadable += 1;
587 }
588 talk
589 })
590 .collect();
591 all.sort_unstable_by(|a, b| {
592 let rank = |t: &Talk| u8::from(!t.status.open());
593 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
594 });
595 (all, unreadable)
596 }
597
598 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
600 if self.path_of(prefix).is_file() {
601 return Ok(prefix.to_owned());
602 }
603 let hits: Vec<String> = self
604 .list()
605 .into_iter()
606 .map(|t| t.id)
607 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
608 .collect();
609 match hits.len() {
610 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
611 0 => bail!("no talk matches `{prefix}`"),
612 _ => bail!(
613 "`{prefix}` matches {} talks: {}",
614 hits.len(),
615 hits.join(", ")
616 ),
617 }
618 }
619
620 pub fn revision(&self) -> u64 {
623 std::fs::read_dir(&self.root)
624 .into_iter()
625 .flatten()
626 .flatten()
627 .filter(|e| e.path().extension().is_none_or(|x| x != "turn"))
628 .filter(|e| !e.file_name().to_string_lossy().starts_with('.'))
629 .filter_map(|e| e.metadata().ok())
630 .filter_map(|m| m.modified().ok())
631 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
632 .map(|d| d.as_millis() as u64)
633 .max()
634 .unwrap_or(0)
635 }
636
637 pub fn count_open(&self) -> usize {
639 self.list().iter().filter(|t| t.status.open()).count()
640 }
641
642 pub fn turn_path(&self, id: &str) -> PathBuf {
645 self.root.join(format!("{id}.turn"))
646 }
647
648 pub fn claim_turn(&self, id: &str) -> Result<Option<TurnLease>> {
657 self.claim_turn_at(id, Timestamp::now())
658 }
659
660 fn claim_turn_at(&self, id: &str, now: Timestamp) -> Result<Option<TurnLease>> {
661 std::fs::create_dir_all(&self.root)
662 .with_context(|| format!("create {}", self.root.display()))?;
663 let path = self.turn_path(id);
664 let token = fresh_token();
665 if create_turn(&path, &token, now)? {
666 return Ok(Some(TurnLease {
667 talk: id.to_owned(),
668 path,
669 token,
670 }));
671 }
672 if lease_blocks(&path, now) {
673 return Ok(None);
674 }
675 let Some(_lock) = TurnLock::take(&path)? else {
680 return Ok(None);
681 };
682 if lease_blocks(&path, now) {
683 return Ok(None);
684 }
685 let _ = std::fs::remove_file(&path);
686 Ok(create_turn(&path, &token, now)?.then_some(TurnLease {
687 talk: id.to_owned(),
688 path,
689 token,
690 }))
691 }
692
693 pub fn turn_held(&self, id: &str) -> bool {
695 read_turn(&self.turn_path(id)).is_some_and(|r| r.fresh(Timestamp::now()))
696 }
697
698 pub fn remove(&self, id: &str) -> Result<()> {
711 let _guard = self.guard()?;
712 let resolved = self.resolve_id(id)?;
713 let path = self.path_of(&resolved);
714 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
715 let artifacts = self.artifacts_of(&resolved);
716 if artifacts.is_dir() {
717 std::fs::remove_dir_all(&artifacts)
718 .with_context(|| format!("remove {}", artifacts.display()))?;
719 }
720 let _ = std::fs::remove_file(self.turn_path(&resolved));
721 Ok(())
722 }
723}
724
725struct StoreGuard<'a> {
728 _file: TurnLock,
729 _mutex: MutexGuard<'a, ()>,
730}
731
732#[derive(Debug, Serialize, Deserialize)]
735struct TurnRecord {
736 token: String,
737 pid: u32,
738 beat_at: Timestamp,
739}
740
741impl TurnRecord {
742 fn fresh(&self, now: Timestamp) -> bool {
743 now.as_second() - self.beat_at.as_second() <= crate::ask::LEASE_TTL.as_secs() as i64
744 }
745}
746
747fn fresh_token() -> String {
751 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
752 let n = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
753 let seed = crate::rng::entropy() ^ n.wrapping_mul(0x9E37_79B9_7F4A_7C15);
754 crate::rng::SplitMix64::new(seed).uuid_v4()
755}
756
757fn read_turn(path: &Path) -> Option<TurnRecord> {
758 serde_json::from_str(&std::fs::read_to_string(path).ok()?).ok()
759}
760
761const TAKEOVER_LOCK_TTL: Duration = Duration::from_secs(10);
763
764const TICKET_TTL: Duration = Duration::from_secs(10);
767
768const TICKET_GENERATIONS: u32 = 16;
770
771const TICKET_SWEEP_AGE: Duration = Duration::from_secs(3600);
773
774fn create_exclusive(path: &Path, body: &str) -> Result<bool> {
776 use std::io::Write as _;
777 match std::fs::OpenOptions::new()
778 .write(true)
779 .create_new(true)
780 .open(path)
781 {
782 Ok(mut f) => {
783 if let Err(e) = f.write_all(body.as_bytes()) {
784 drop(f);
785 let _ = std::fs::remove_file(path);
786 return Err(e).with_context(|| format!("write {}", path.display()));
787 }
788 Ok(true)
789 }
790 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
791 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
792 }
793}
794
795fn create_turn(path: &Path, token: &str, now: Timestamp) -> Result<bool> {
796 let record = TurnRecord {
797 token: token.to_owned(),
798 pid: std::process::id(),
799 beat_at: now,
800 };
801 let body = serde_json::to_string(&record).context("serialize turn lease")?;
802 let tmp = path.with_extension(format!("turn.{token}.new"));
805 publish_exclusive(path, &tmp, &body)
806}
807
808fn lease_blocks(path: &Path, now: Timestamp) -> bool {
812 if read_turn(path).is_some_and(|r| r.fresh(now)) {
813 return true;
814 }
815 std::fs::metadata(path).is_ok_and(|m| {
816 m.len() == 0
817 && m.modified()
818 .ok()
819 .and_then(|t| t.elapsed().ok())
820 .is_some_and(|age| age <= TAKEOVER_LOCK_TTL)
821 })
822}
823
824fn link_unsupported(e: &std::io::Error) -> bool {
827 e.kind() == std::io::ErrorKind::Unsupported || (cfg!(windows) && e.raw_os_error() == Some(1))
828}
829
830fn create_in_place(path: &Path, body: &str) -> Result<bool> {
848 use std::io::Write as _;
849 let mut f = match std::fs::OpenOptions::new()
850 .write(true)
851 .create_new(true)
852 .open(path)
853 {
854 Ok(f) => f,
855 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => return Ok(false),
856 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
857 };
858 f.write_all(body.as_bytes())
859 .with_context(|| format!("write {}", path.display()))?;
860 Ok(std::fs::read_to_string(path).is_ok_and(|t| t == body))
863}
864
865fn publish_exclusive(path: &Path, tmp: &Path, body: &str) -> Result<bool> {
866 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
867 #[cfg(test)]
868 let linked = if failpoint::no_link_forced() {
869 Err(std::io::Error::from(std::io::ErrorKind::Unsupported))
870 } else {
871 std::fs::hard_link(tmp, path)
872 };
873 #[cfg(not(test))]
874 let linked = std::fs::hard_link(tmp, path);
875 let out = match linked {
876 Ok(()) => Ok(true),
877 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
878 Err(e) if link_unsupported(&e) => create_in_place(path, body),
879 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
880 };
881 let _ = std::fs::remove_file(tmp);
882 out
883}
884
885fn remove_if_carries(path: &Path, expected: &str) -> bool {
896 let gone = path.with_extension(format!("gone.{}", fresh_token()));
897 if std::fs::rename(path, &gone).is_err() {
898 return false;
899 }
900 let found = std::fs::read_to_string(&gone);
901 if found.as_ref().is_ok_and(|t| t == expected) {
902 let _ = std::fs::remove_file(&gone);
903 return true;
904 }
905 if let Ok(body) = found {
906 let back = path.with_extension(format!("back.{}", fresh_token()));
907 match publish_exclusive(path, &back, &body) {
908 Ok(true) => {}
909 Ok(false) | Err(_) => tracing::warn!(
910 "{} was replaced while a stale removal had it aside; \
911 the displaced file is dropped",
912 path.display()
913 ),
914 }
915 }
916 let _ = std::fs::remove_file(&gone);
917 false
918}
919
920struct TurnLock {
924 path: PathBuf,
925 token: String,
926}
927
928impl TurnLock {
929 fn token() -> String {
930 fresh_token()
931 }
932
933 fn publish(path: &Path, token: &str) -> Result<bool> {
937 let tmp = path.with_extension(format!("lock.{token}.new"));
938 publish_exclusive(path, &tmp, token)
939 }
940
941 fn take(lease: &Path) -> Result<Option<Self>> {
942 let path = lease.with_extension("turn.lock");
943 let token = Self::token();
944 if Self::publish(&path, &token)? {
945 return Ok(Some(Self { path, token }));
946 }
947 let Ok(seen) = std::fs::read_to_string(&path) else {
948 return Ok(None);
949 };
950 let key = if !seen.is_empty() && seen.chars().all(|c| c.is_ascii_alphanumeric() || c == '-')
953 {
954 seen.as_str()
955 } else {
956 "invalid"
957 };
958 let aged = std::fs::metadata(&path)
959 .and_then(|m| m.modified())
960 .ok()
961 .and_then(|t| t.elapsed().ok())
962 .is_some_and(|age| age > TAKEOVER_LOCK_TTL);
963 if !aged {
964 return Ok(None);
965 }
966 let mut won = false;
978 for n in 0..TICKET_GENERATIONS {
979 let ticket = path.with_extension(format!("lock.{key}.break.{n}"));
980 if create_exclusive(&ticket, "")? {
981 won = true;
982 break;
983 }
984 let stale = std::fs::metadata(&ticket)
985 .and_then(|m| m.modified())
986 .ok()
987 .and_then(|t| t.elapsed().ok())
988 .is_some_and(|age| age > TICKET_TTL);
989 if !stale {
990 return Ok(None);
991 }
992 }
993 if !won {
994 Self::sweep_tickets(&path);
999 return Ok(None);
1000 }
1001 Self::sweep_tickets(&path);
1002 if !remove_if_carries(&path, &seen) {
1008 return Ok(None);
1009 }
1010 if Self::publish(&path, &token)? {
1011 return Ok(Some(Self { path, token }));
1012 }
1013 Ok(None)
1014 }
1015
1016 fn sweep_tickets(path: &Path) {
1019 let (Some(dir), Some(name)) = (path.parent(), path.file_name().and_then(|n| n.to_str()))
1020 else {
1021 return;
1022 };
1023 let prefix = format!("{name}.");
1024 let Ok(entries) = std::fs::read_dir(dir) else {
1025 return;
1026 };
1027 for entry in entries.flatten() {
1028 let file = entry.file_name();
1029 let Some(file) = file.to_str() else { continue };
1030 if !(file.starts_with(&prefix) && file.contains(".break.")) {
1031 continue;
1032 }
1033 let old = entry
1034 .metadata()
1035 .and_then(|m| m.modified())
1036 .ok()
1037 .and_then(|t| t.elapsed().ok())
1038 .is_some_and(|age| age > TICKET_SWEEP_AGE);
1039 if old {
1040 let _ = std::fs::remove_file(entry.path());
1041 }
1042 }
1043 }
1044
1045 fn take_patiently(lease: &Path) -> Option<Self> {
1047 for _ in 0..50 {
1048 match Self::take(lease) {
1049 Ok(Some(lock)) => return Some(lock),
1050 Ok(None) => std::thread::sleep(Duration::from_millis(10)),
1051 Err(_) => return None,
1052 }
1053 }
1054 None
1055 }
1056}
1057
1058impl Drop for TurnLock {
1059 fn drop(&mut self) {
1060 remove_if_carries(&self.path, &self.token);
1062 }
1063}
1064
1065pub const TURN_BEAT: Duration = Duration::from_secs(20);
1068
1069#[derive(Debug)]
1074pub struct TurnLease {
1075 talk: String,
1076 path: PathBuf,
1077 token: String,
1078}
1079
1080impl TurnLease {
1081 pub fn holds(&self) -> bool {
1083 read_turn(&self.path).is_some_and(|r| r.token == self.token)
1084 }
1085
1086 pub fn beat(&self) -> Result<bool> {
1090 let _lock = TurnLock::take_patiently(&self.path)
1091 .with_context(|| format!("lock {} to renew it", self.path.display()))?;
1092 let Some(mut record) = read_turn(&self.path).filter(|r| r.token == self.token) else {
1093 return Ok(false);
1094 };
1095 record.beat_at = Timestamp::now();
1096 let body = serde_json::to_string(&record).context("serialize turn lease")?;
1097 let tmp = self.path.with_extension(format!("turn.{}.tmp", self.token));
1098 write_atomic(&tmp, &self.path, &body)?;
1099 Ok(true)
1100 }
1101
1102 pub async fn beating<T>(&self, fut: impl std::future::Future<Output = T>) -> Result<T> {
1107 self.beating_every(TURN_BEAT, fut).await
1108 }
1109
1110 async fn beating_every<T>(
1111 &self,
1112 period: Duration,
1113 fut: impl std::future::Future<Output = T>,
1114 ) -> Result<T> {
1115 tokio::pin!(fut);
1116 loop {
1117 match tokio::time::timeout(period, &mut fut).await {
1118 Ok(out) => return Ok(out),
1119 Err(_) => match self.beat() {
1120 Ok(true) => {}
1121 Ok(false) => bail!(
1122 "the turn lease {} was taken over; this turn is stopped",
1123 self.path.display()
1124 ),
1125 Err(e) => tracing::warn!("{e:#}"),
1126 },
1127 }
1128 }
1129 }
1130}
1131
1132impl Drop for TurnLease {
1133 fn drop(&mut self) {
1134 if let Some(_lock) = TurnLock::take_patiently(&self.path) {
1137 if read_turn(&self.path).is_some_and(|r| r.token == self.token) {
1138 let _ = std::fs::remove_file(&self.path);
1139 }
1140 }
1141 }
1142}
1143
1144pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
1154 let repo = repo.canonicalize().unwrap_or(repo);
1157 let spec = match agent {
1160 Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
1161 None => agent::pick_chain(
1162 &cfg.agents,
1163 cfg.roles.chatter.as_ref(),
1164 &agent::installed,
1165 "chatter",
1166 )?
1167 .remove(0),
1168 };
1169
1170 let now = Timestamp::now();
1171 let mut talk = Talk {
1172 schema: SCHEMA,
1173 id: new_id(),
1174 repo,
1175 agent: spec.id.clone(),
1176 status: TalkStatus::Open,
1177 turns: Vec::new(),
1178 pending: String::new(),
1179 pending_attachments: Vec::new(),
1180 fallback: agent.is_none(),
1181 persona: String::new(),
1182 persona_dirty: false,
1183 created_at: now,
1184 updated_at: now,
1185 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
1186 };
1187 store.put(&mut talk)?;
1188 Ok(talk)
1189}
1190
1191pub fn record(
1198 talk: &mut Talk,
1199 store: &Talks,
1200 text: &str,
1201 attachments: Vec<Attachment>,
1202) -> Result<String> {
1203 let _guard = store.guard()?;
1211 let Ok(fresh) = store.get(&talk.id) else {
1216 bail!("talk {} was deleted", talk.short());
1217 };
1218 talk.status = fresh.status;
1219 talk.pending = fresh.pending;
1222 talk.pending_attachments = fresh.pending_attachments;
1223 if !talk.status.open() {
1224 bail!(
1225 "talk {} is {} and takes no more turns",
1226 talk.short(),
1227 talk.status.as_str()
1228 );
1229 }
1230 let text = text.trim();
1231 if text.is_empty() && attachments.is_empty() {
1232 bail!("nothing to say");
1233 }
1234 talk.turns.push(Turn {
1235 who: Who::Operator,
1236 body: text.to_owned(),
1237 at: Timestamp::now(),
1238 attachments,
1239 usage: None,
1240 });
1241 store.put(talk)?;
1242 Ok(text.to_owned())
1243}
1244
1245pub fn queue(
1247 talk: &mut Talk,
1248 store: &Talks,
1249 text: &str,
1250 attachments: Vec<Attachment>,
1251) -> Result<()> {
1252 let text = text.trim();
1253 if text.is_empty() && attachments.is_empty() {
1254 bail!("nothing to say");
1255 }
1256 let _guard = store.guard()?;
1257 let mut fresh = store
1258 .get(&talk.id)
1259 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1260 if !fresh.status.open() {
1261 bail!(
1262 "talk {} is {} and takes no more turns",
1263 fresh.short(),
1264 fresh.status.as_str()
1265 );
1266 }
1267 if !text.is_empty() {
1268 if fresh.pending.is_empty() {
1269 fresh.pending = text.to_owned();
1270 } else {
1271 fresh.pending.push_str("\n\n");
1272 fresh.pending.push_str(text);
1273 }
1274 }
1275 fresh.pending_attachments.extend(attachments);
1276 store.put(&mut fresh)?;
1277 *talk = fresh;
1278 Ok(())
1279}
1280
1281pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
1283 let _guard = store.guard()?;
1284 let mut fresh = store
1285 .get(&talk.id)
1286 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1287 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
1288 *talk = fresh;
1289 return Ok(None);
1290 }
1291 let text = std::mem::take(&mut fresh.pending);
1292 let attachments = std::mem::take(&mut fresh.pending_attachments);
1293 fresh.turns.push(Turn {
1294 who: Who::Operator,
1295 body: text.clone(),
1296 at: Timestamp::now(),
1297 attachments,
1298 usage: None,
1299 });
1300 store.put(&mut fresh)?;
1301 *talk = fresh;
1302 Ok(Some(text))
1303}
1304
1305pub async fn say(
1308 lease: &TurnLease,
1309 talk: &mut Talk,
1310 store: &Talks,
1311 cfg: &Config,
1312 text: &str,
1313 attachments: Vec<Attachment>,
1314) -> Result<()> {
1315 check_lease(lease, talk)?;
1316 let text = record(talk, store, text, attachments)?;
1317 respond(lease, talk, store, cfg, &text).await
1318}
1319
1320pub async fn respond(
1326 lease: &TurnLease,
1327 talk: &mut Talk,
1328 store: &Talks,
1329 cfg: &Config,
1330 text: &str,
1331) -> Result<()> {
1332 check_lease(lease, talk)?;
1333 lease
1334 .beating(turn(talk, store, cfg, text))
1335 .await
1336 .and_then(|done| done)
1337}
1338
1339fn check_lease(lease: &TurnLease, talk: &Talk) -> Result<()> {
1340 if lease.talk != talk.id {
1341 bail!("the turn lease is for talk {}, not {}", lease.talk, talk.id);
1342 }
1343 if !lease.holds() {
1346 bail!("the turn lease for talk {} is no longer held", talk.short());
1347 }
1348 Ok(())
1349}
1350
1351pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
1369 let _guard = store.guard()?;
1370 let mut fresh = store
1371 .get(&talk.id)
1372 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1373 fresh.status = TalkStatus::Closed;
1374 fresh.pending.clear();
1376 fresh.pending_attachments.clear();
1377 store.put(&mut fresh)?;
1378 *talk = fresh;
1379 Ok(())
1380}
1381
1382pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
1393 let _guard = store.guard()?;
1394 let mut fresh = store
1395 .get(&talk.id)
1396 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1397 fresh.status = TalkStatus::Open;
1398 store.put(&mut fresh)?;
1399 *talk = fresh;
1400 Ok(())
1401}
1402
1403pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
1416 let _guard = store.guard()?;
1417 let mut fresh = store
1418 .get(&talk.id)
1419 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1420 if fresh.agent == spec.id {
1421 *talk = fresh;
1422 return Ok(false);
1423 }
1424 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
1425 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1426 fresh.fallback = false;
1428 fresh.turns.push(Turn {
1429 who: Who::Agent,
1430 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
1431 at: Timestamp::now(),
1432 attachments: Vec::new(),
1433 usage: None,
1434 });
1435 store.put(&mut fresh)?;
1436 *talk = fresh;
1437 Ok(true)
1438}
1439
1440pub fn switch_persona(talk: &mut Talk, store: &Talks, id: &str) -> Result<bool> {
1448 let id = if id.trim() == crate::persona::DEFAULT_ID {
1449 ""
1450 } else {
1451 id.trim()
1452 };
1453 let _guard = store.guard()?;
1454 let mut fresh = store
1455 .get(&talk.id)
1456 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1457 if fresh.persona == id {
1458 *talk = fresh;
1459 return Ok(false);
1460 }
1461 let from = std::mem::replace(&mut fresh.persona, id.to_owned());
1462 fresh.persona_dirty = true;
1463 let label = |p: &str| {
1464 if p.is_empty() {
1465 crate::persona::DEFAULT_ID.to_owned()
1466 } else {
1467 p.to_owned()
1468 }
1469 };
1470 fresh.turns.push(Turn {
1471 who: Who::Agent,
1472 body: format!(
1473 "{MAGI_NOTE}persona changed from {} to {}",
1474 label(&from),
1475 label(id)
1476 ),
1477 at: Timestamp::now(),
1478 attachments: Vec::new(),
1479 usage: None,
1480 });
1481 store.put(&mut fresh)?;
1482 *talk = fresh;
1483 Ok(true)
1484}
1485
1486pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
1488 let _guard = store.guard()?;
1489 let mut fresh = store
1490 .get(&talk.id)
1491 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1492 fresh.pending.clear();
1493 fresh.pending_attachments.clear();
1494 store.put(&mut fresh)?;
1495 *talk = fresh;
1496 Ok(())
1497}
1498
1499pub fn clear_pending_if_matches(
1501 talk: &mut Talk,
1502 store: &Talks,
1503 expected_text: &str,
1504 expected_attachments: &[String],
1505) -> Result<bool> {
1506 let _guard = store.guard()?;
1507 let mut fresh = store
1508 .get(&talk.id)
1509 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1510 if !pending_matches(&fresh, expected_text, expected_attachments) {
1511 *talk = fresh;
1512 return Ok(false);
1513 }
1514 fresh.pending.clear();
1515 fresh.pending_attachments.clear();
1516 store.put(&mut fresh)?;
1517 *talk = fresh;
1518 Ok(true)
1519}
1520
1521pub fn edit_pending_text(
1525 talk: &mut Talk,
1526 store: &Talks,
1527 text: &str,
1528 expected_text: &str,
1529 expected_attachments: &[String],
1530) -> Result<bool> {
1531 let _guard = store.guard()?;
1532 let mut fresh = store
1533 .get(&talk.id)
1534 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1535 if !pending_matches(&fresh, expected_text, expected_attachments) {
1536 *talk = fresh;
1537 return Ok(false);
1538 }
1539 fresh.pending = text.trim().to_owned();
1540 store.put(&mut fresh)?;
1541 *talk = fresh;
1542 Ok(true)
1543}
1544
1545fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
1546 talk.pending == expected_text
1547 && talk
1548 .pending_attachments
1549 .iter()
1550 .map(|attachment| &attachment.id)
1551 .eq(expected_attachments.iter())
1552}
1553
1554async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1561 let spec = cfg
1562 .agents
1563 .iter()
1564 .find(|a| a.id == talk.agent)
1565 .with_context(|| {
1566 format!(
1567 "talk {} was opened with agent `{}`, which is no longer in \
1568 the roster; restore it in magi.toml or start a new \
1569 conversation",
1570 talk.short(),
1571 talk.agent
1572 )
1573 })?;
1574
1575 let last_note = attachment_note(
1579 store,
1580 &talk.id,
1581 talk.turns
1582 .last()
1583 .map_or(&[][..], |t| t.attachments.as_slice()),
1584 );
1585
1586 let attachment_paths: Vec<PathBuf> = talk
1592 .turns
1593 .iter()
1594 .flat_map(|t| t.attachments.iter())
1595 .filter_map(|a| store.attachment_path(&talk.id, a))
1596 .collect();
1597
1598 let questions = crate::ask::Questions::open();
1606 let consulted = crate::consult::pending_consults(&questions, &talk.id);
1607 let consult_roots: Vec<PathBuf> = if consulted {
1608 vec![questions.root().to_path_buf()]
1609 } else {
1610 Vec::new()
1611 };
1612
1613 let (allow_write, unsandboxed) = turn_access(cfg.talk.allow_write, consulted);
1614
1615 let artifacts = store.artifacts_of(&talk.id);
1616 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
1619 let stem = format!("turn-{}", operator_turns.max(1));
1620 let cache_dir = cfg.cache_dir();
1623
1624 let mut chain = vec![spec.clone()];
1627 if let Some(choice) = cfg.roles.chatter.as_ref()
1628 && talk.fallback
1629 {
1630 for id in choice.ids() {
1631 if id == talk.agent || chain.iter().any(|s| s.id == id) {
1632 continue;
1633 }
1634 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1635 Ok(s) => chain.push(s),
1636 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1637 }
1638 }
1639 }
1640
1641 let persona = crate::persona::active(&cfg.talk.personas, &talk.persona);
1642 let operator_name = cfg.talk.operator_name();
1643 let persona_update = if talk.persona_dirty {
1644 format!(
1645 "{}\n\n",
1646 crate::persona::update_block_for(persona.as_ref(), operator_name)
1647 )
1648 } else {
1649 String::new()
1650 };
1651
1652 let mut outcome = None;
1653 let mut fell_back_from: Option<String> = None;
1654 let mut first_try: Option<(String, SeatState)> = None;
1657 for (n, spec) in chain.iter().enumerate() {
1658 if n > 0 {
1659 if first_try.is_none() {
1660 first_try = Some((talk.agent.clone(), talk.seat.clone()));
1661 }
1662 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1663 fell_back_from.get_or_insert_with(|| talk.agent.clone());
1666 talk.agent = spec.id.clone();
1667 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1668 }
1669 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1670 let first_ever = talk.turns.len() <= 1;
1671 let body = if talk.seat.turns == 0 && first_ever {
1672 format!(
1673 "{}\n\n# Operator\n\n{text}{last_note}",
1674 briefing_with(
1675 &talk.repo,
1676 &cfg.graph.language,
1677 cfg.talk.allow_write,
1678 persona.as_ref(),
1679 operator_name
1680 )
1681 )
1682 } else if talk.seat.turns == 0 {
1683 format!(
1686 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1687 briefing_with(
1688 &talk.repo,
1689 &cfg.graph.language,
1690 cfg.talk.allow_write,
1691 persona.as_ref(),
1692 operator_name
1693 ),
1694 transcript(talk, store)
1695 )
1696 } else if resuming {
1697 format!("{persona_update}{text}{last_note}")
1698 } else {
1699 let mut standing = match (&persona, talk.persona_dirty) {
1703 (Some(p), false) => format!("{}\n", crate::persona::section_for(p, operator_name)),
1704 _ => persona_update.clone(),
1705 };
1706 if let (None, Some(n)) = (&persona, operator_name) {
1709 standing.push_str(&format!(
1710 "# Addressing the operator\n{}\n",
1711 crate::persona::addressing(n)
1712 ));
1713 }
1714 format!("{}\n\n{standing}{text}{last_note}", transcript(talk, store))
1715 };
1716 let attempt_stem = if n == 0 {
1717 stem.clone()
1718 } else {
1719 format!("{stem}-{}", spec.id)
1720 };
1721 let inv = Invocation {
1722 cwd: &talk.repo,
1723 prompt: &body,
1724 timeout: turn_timeout(cfg),
1725 allow_write,
1730 unsandboxed,
1731 sessions: cfg.graph.sessions,
1732 artifacts: &artifacts,
1733 stem: &attempt_stem,
1734 run: &talk.id,
1737 node: crate::queue::CHAT_NODE,
1738 cache_dir: cache_dir.as_deref(),
1739 attachments: &attachment_paths,
1740 writable: &consult_roots,
1741 };
1742 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1743 let advance = agent::chain_advances(&result);
1744 if n == 0 || !advance {
1745 outcome = Some(result);
1746 } else {
1747 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1749 }
1750 if !advance {
1751 break;
1752 }
1753 }
1754 if outcome.as_ref().is_some_and(agent::chain_advances) {
1755 if let Some((id, seat)) = first_try {
1758 talk.agent = id;
1759 talk.seat = seat;
1760 fell_back_from = None;
1761 }
1762 }
1763 let outcome = outcome.expect("a chain holds at least one agent");
1764 let note = |why: String| Turn {
1765 who: Who::Agent,
1766 body: format!("{MAGI_NOTE}{why}"),
1767 at: Timestamp::now(),
1768 attachments: Vec::new(),
1769 usage: None,
1770 };
1771 let (reply, failure) = match outcome {
1772 Err(e) => (
1773 note(format!("could not run agent `{}`: {e}", talk.agent)),
1774 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1775 ),
1776 Ok(out) if out.quota_exhausted() => {
1777 let reset = out
1778 .quota
1779 .as_ref()
1780 .and_then(|q| q.reset.clone())
1781 .map_or_else(String::new, |r| format!(" (resets {r})"));
1782 let why = format!(
1783 "agent `{}` is out of quota{reset}; your message is saved, so \
1784 say it again when the window reopens",
1785 talk.agent
1786 );
1787 (note(why.clone()), Some(why))
1788 }
1789 Ok(out) if out.timed_out => {
1790 let why = format!(
1791 "agent `{}` did not answer within {}s; your message is saved",
1792 talk.agent,
1793 turn_timeout(cfg).as_secs()
1794 );
1795 (note(why.clone()), Some(why))
1796 }
1797 Ok(out) if !out.usable() => {
1798 let why = format!(
1799 "agent `{}` produced no answer (exit {}); your message is saved",
1800 talk.agent,
1801 out.exit_code
1802 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1803 );
1804 (note(why.clone()), Some(why))
1805 }
1806 Ok(out) => (
1807 Turn {
1808 who: Who::Agent,
1809 body: out.text.trim().to_owned(),
1810 at: Timestamp::now(),
1811 attachments: Vec::new(),
1812 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1815 context_tokens,
1816 agent: talk.agent.clone(),
1817 model: cfg
1818 .agents
1819 .iter()
1820 .find(|a| a.id == talk.agent)
1821 .and_then(|a| a.model.clone()),
1822 }),
1823 },
1824 None,
1825 ),
1826 };
1827
1828 let _guard = store.guard()?;
1840 let Ok(fresh) = store.get(&talk.id) else {
1846 return Ok(());
1847 };
1848 talk.status = fresh.status;
1849 talk.pending = fresh.pending;
1853 talk.pending_attachments = fresh.pending_attachments;
1854 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1855 talk.turns.push(note(format!(
1858 "agent changed from {from} to {} (fallback)",
1859 talk.agent
1860 )));
1861 }
1862 if failure.is_none() {
1865 talk.persona_dirty = false;
1866 }
1867 talk.turns.push(reply);
1868 if let Err(put_err) = store.put(talk) {
1869 let lost = talk.turns.pop().expect("just pushed above");
1878 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1879 let why = match &stash {
1880 Ok(path) => format!(
1881 "agent `{}` answered, but the reply could not be saved to \
1882 this conversation ({put_err:#}); the raw text was kept at \
1883 {} - your message is saved, ask again",
1884 talk.agent,
1885 path.display()
1886 ),
1887 Err(stash_err) => format!(
1888 "agent `{}` answered, but the reply could not be saved to \
1889 this conversation ({put_err:#}), and it could not be kept \
1890 anywhere else either ({stash_err:#}); your message is \
1891 saved, ask again",
1892 talk.agent
1893 ),
1894 };
1895 talk.turns.push(note(why.clone()));
1896 return match store.put(talk) {
1903 Ok(()) => bail!("{why}"),
1904 Err(note_err) => {
1905 talk.turns.pop();
1925 Err(note_err).context(why)
1926 }
1927 };
1928 }
1929
1930 match failure {
1931 Some(why) => bail!("{why}"),
1932 None => Ok(()),
1933 }
1934}
1935
1936fn transcript(talk: &Talk, store: &Talks) -> String {
1939 let mut out = String::from(
1940 "This conversation cannot resume on the CLI's side, so here is \
1941 everything said so far; answer only the last message.\n",
1942 );
1943 for t in &talk.turns {
1944 let who = match t.who {
1945 Who::Operator => "operator",
1946 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1947 Who::Agent => "you",
1948 };
1949 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1950 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1951 }
1952 out
1953}
1954
1955fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1960 if attachments.is_empty() {
1961 return String::new();
1962 }
1963 let mut out = String::from(
1964 "\n\nThe operator attached the image(s) below to this message. Open \
1965 and look at each one before you answer.\n",
1966 );
1967 for att in attachments {
1968 if let Some(path) = store.attachment_path(talk_id, att) {
1969 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1970 }
1971 }
1972 out.push('\n');
1973 out
1974}
1975
1976pub(crate) fn turn_access(talk_allow_write: bool, consulted: bool) -> (bool, bool) {
1987 (talk_allow_write || consulted, talk_allow_write)
1988}
1989
1990pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
2006 briefing_with(repo, language, allow_write, None, None)
2007}
2008
2009pub fn briefing_with(
2014 repo: &Path,
2015 language: &str,
2016 allow_write: bool,
2017 persona: Option<&crate::persona::Persona>,
2018 operator_name: Option<&str>,
2019) -> String {
2020 let write_policy = if allow_write {
2021 "Write access is enabled for this conversation (`allow_write = \
2022 true`), so you may write files - but only a small, \
2023 already-decided edit the operator names outright in this \
2024 conversation, not an implementation. This is a permission on the \
2025 conversation as a whole, not a property of whichever repository \
2026 it happened to start in: if the operator names a different \
2027 repository for that small edit, the policy allows it there too. \
2028 Your own tool may still confine writes to the repository this \
2029 conversation started in regardless - if a write elsewhere is \
2030 refused, say so plainly rather than working around it. Once you \
2031 have made an edit, say plainly what you edited. Anything bigger, \
2032 or anything still open-ended, still goes through the queue below \
2033 rather than being done here."
2034 } else {
2035 "Do not write files. Implementing a change is not this \
2036 conversation's job; a separate, blind competition of agents does \
2037 that, and a repository this conversation has already edited would \
2038 make their diffs unjudgeable."
2039 };
2040 let mut out = format!(
2041 "You are magi's standing conversation partner for its operator, who \
2042 usually has this open on a phone. Keep replies short: no preamble, \
2043 no restating what they just said.\n\n\
2044 # Repository\n\n{repo}\n\n\
2045 You may look around: read files, run shell commands, search history, \
2046 run tests - whatever answers the question. {write_policy}\n\n\
2047 A short, command-shaped message (\"list\", \"info <id>\", \"show \
2048 3cbf\") is almost always the operator asking you to look something \
2049 up, not an instruction to file - answer it yourself with `magi \
2050 list`, `magi show <id>`, `magi task list`, or the like, the same way \
2051 you would answer any other question in this conversation.\n\n\
2052 # When the operator wants something done\n\n\
2053 Run:\n\n\
2054 magi task add --solo --repo {repo} <instruction>\n\n\
2055 and tell the operator the task id it prints, so they can follow it \
2056 from the Queue. If it refuses with a duplicate warning (the \
2057 instruction names a branch, commit or pull request that an \
2058 unfinished task, run or PR already owns), do not repeat it with \
2059 --force yourself: tell the operator what it matched and let them \
2060 decide. Write <instruction> so that an implementer who has \
2061 never seen this conversation can act on it alone - it is everything \
2062 they get. Use --solo: it runs the task through one implementer \
2063 straight into review instead of the usual multi-agent competition, \
2064 which is the right shape for a change this conversation has already \
2065 settled, rather than one still worth several independent takes.\n\n\
2066 If the operator asks for something in a different repository, \
2067 --repo does not have to be a full path: --repo owner/repo (or just \
2068 repo, when that is unambiguous) is resolved against local checkouts \
2069 the same way `magi repos` lists them. If the command fails because \
2070 nothing matches or more than one checkout shares that name, ask the \
2071 operator which repository they mean (or run `magi repos` yourself \
2072 to see the candidates) rather than guessing.\n\n\
2073 The current state of the code is whatever origin/main holds, not \
2074 whatever a working tree shows: a primary checkout often lags \
2075 upstream, sits on a detached HEAD and carries uncommitted changes. \
2076 Before answering about code, run `git fetch origin` in that \
2077 repository if it is cheap, then read through \
2078 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
2079 If the working tree differs, say so; if the fetch fails, say that \
2080 too, so the operator knows the answer may be stale.\n\n\
2081 If the operator attached an image (a screenshot, say) that the task \
2082 is about, pass it with `--attach <path>`, using the absolute path \
2083 the turn's attachment note gives; repeat the flag for several. \
2084 `magi task add --solo --attach <path> <instruction>` copies the \
2085 file into the task, so the implementer receives it. Do not paste the \
2086 path into <instruction> instead: deleting this conversation deletes \
2087 its attachments, and then that path reaches no one.\n",
2088 repo = repo.display(),
2089 );
2090 out.push_str(&language_note(language));
2091 match (persona, operator_name) {
2092 (Some(p), n) => out.push_str(&crate::persona::section_for(p, n)),
2093 (None, Some(n)) => {
2094 out.push_str("\n# Addressing the operator\n");
2095 out.push_str(&crate::persona::addressing(n));
2096 }
2097 (None, None) => {}
2098 }
2099 out
2100}
2101
2102fn language_note(language: &str) -> String {
2105 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
2106 String::new()
2107 } else {
2108 format!("\nHold this conversation in {language}.\n")
2109 }
2110}
2111
2112pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
2119 let mut tasks: Vec<Task> = queue
2120 .list()
2121 .into_iter()
2122 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
2123 .collect();
2124 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
2125 tasks
2126}
2127
2128fn read_path(path: &Path) -> Result<Talk> {
2129 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2130 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
2131}
2132
2133const PUT_RETRIES: u32 = 5;
2136
2137fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
2149 let mut last_err = None;
2150 for attempt in 0..PUT_RETRIES {
2151 if attempt > 0 {
2152 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
2153 }
2154 match try_write_atomic(tmp, path, body) {
2155 Ok(()) => return Ok(()),
2156 Err(e) => last_err = Some(e),
2157 }
2158 }
2159 Err(last_err.expect("the loop above always runs at least once"))
2160}
2161
2162fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
2163 #[cfg(test)]
2164 if failpoint::take_forced_put_failure() {
2165 bail!("simulated write failure (test)");
2166 }
2167 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
2168 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
2169 Ok(())
2170}
2171
2172fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
2177 let dir = store.artifacts_of(id);
2178 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
2179 let path = dir.join(format!("{stem}-lost.txt"));
2180 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
2181 Ok(path)
2182}
2183
2184#[cfg(test)]
2191mod failpoint {
2192 use std::cell::Cell;
2193
2194 thread_local! {
2195 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
2196 static FORCE_NO_LINK: Cell<bool> = const { Cell::new(false) };
2197 }
2198
2199 pub(super) fn force_no_link(on: bool) {
2202 FORCE_NO_LINK.with(|c| c.set(on));
2203 }
2204
2205 pub(super) fn no_link_forced() -> bool {
2206 FORCE_NO_LINK.with(Cell::get)
2207 }
2208
2209 pub(super) fn force_put_failures(count: u32) {
2212 FORCE_PUT_FAILURES.with(|c| c.set(count));
2213 }
2214
2215 pub(super) fn take_forced_put_failure() -> bool {
2218 FORCE_PUT_FAILURES.with(|c| {
2219 let n = c.get();
2220 if n == 0 {
2221 false
2222 } else {
2223 c.set(n - 1);
2224 true
2225 }
2226 })
2227 }
2228}
2229
2230fn short(id: &str) -> &str {
2231 id.split('-').next_back().unwrap_or(id)
2232}
2233
2234fn new_id() -> String {
2235 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2236 let seed = crate::rng::entropy();
2237 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2238}
2239
2240fn attachment_ext(mime: &str) -> Option<&'static str> {
2245 match mime {
2246 "image/png" => Some("png"),
2247 "image/jpeg" => Some("jpg"),
2248 "image/gif" => Some("gif"),
2249 "image/webp" => Some("webp"),
2250 _ => None,
2251 }
2252}
2253
2254pub fn valid_attachment_id(id: &str) -> bool {
2259 id.len() == 32
2260 && id
2261 .bytes()
2262 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
2263}
2264
2265fn new_attachment_id() -> String {
2269 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
2270 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
2271}
2272
2273#[cfg(test)]
2274mod tests {
2275 #[test]
2276 fn turn_access_lifts_the_sandbox_only_for_the_talk_opt_in() {
2277 use super::turn_access;
2278 assert_eq!(turn_access(false, false), (false, false));
2279 assert_eq!(turn_access(true, false), (true, true));
2280 assert_eq!(turn_access(false, true), (true, false));
2282 assert_eq!(turn_access(true, true), (true, true));
2283 }
2284
2285 #[test]
2286 fn the_briefing_points_at_origin_main_not_the_working_tree() {
2287 let b = briefing(Path::new("/r"), "en", false);
2288 assert!(b.contains("origin/main"));
2289 assert!(b.contains("git show origin/main:"));
2290 }
2291 use std::collections::BTreeMap;
2292
2293 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
2294 use crate::queue::{Queue, Source, Task};
2295
2296 use super::*;
2297
2298 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
2299 AgentSpec {
2300 id: id.to_owned(),
2301 kind: AgentKind::Command,
2302 model: model.map(str::to_owned),
2303 command: Vec::new(),
2304 extra_args: Vec::new(),
2305 env: BTreeMap::new(),
2306 prompt_delivery: None,
2307 }
2308 }
2309
2310 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
2311 Talk {
2312 schema: SCHEMA,
2313 id: "20260904-014455-ab12".to_owned(),
2314 repo: PathBuf::from("."),
2315 agent: agent.to_owned(),
2316 status: TalkStatus::Open,
2317 turns,
2318 pending: String::new(),
2319 pending_attachments: Vec::new(),
2320 fallback: false,
2321 persona: String::new(),
2322 persona_dirty: false,
2323 created_at: Timestamp::now(),
2324 updated_at: Timestamp::now(),
2325 seat: SeatState::new(SEAT, agent, 1),
2326 }
2327 }
2328
2329 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
2330 Turn {
2331 who: Who::Agent,
2332 body: body.to_owned(),
2333 at: Timestamp::now(),
2334 attachments: Vec::new(),
2335 usage: usage.map(|(t, a, m)| TurnUsage {
2336 context_tokens: t,
2337 agent: a.to_owned(),
2338 model: m.map(str::to_owned),
2339 }),
2340 }
2341 }
2342
2343 fn ctx_config(windows: &[(&str, u64)]) -> Config {
2344 Config {
2345 agents: vec![
2346 ctx_agent("small", Some("small-model")),
2347 ctx_agent("big", Some("big-model")),
2348 ctx_agent("plain", None),
2349 ],
2350 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
2351 ..Config::default()
2352 }
2353 }
2354
2355 #[test]
2356 fn context_usage_computes_percent_and_warns_at_eighty() {
2357 let cfg = ctx_config(&[("small-model", 1000)]);
2358 let at = |tokens| {
2359 let t = ctx_talk(
2360 "small",
2361 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
2362 );
2363 context_usage(&t, Some(&cfg))
2364 };
2365 let u = at(799);
2366 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
2367 let u = at(800);
2368 assert_eq!((u.percent, u.warn), (Some(80), true));
2369 let u = at(1500);
2370 assert_eq!((u.percent, u.warn), (Some(150), true));
2371 assert!(!u.since_switch);
2372 }
2373
2374 #[test]
2375 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
2376 let cfg = ctx_config(&[("small-model", 1000)]);
2377 let t = ctx_talk(
2378 "small",
2379 vec![
2380 reply("old", Some((900, "small", Some("small-model")))),
2381 reply("new", None),
2382 ],
2383 );
2384 let u = context_usage(&t, Some(&cfg));
2385 assert!(u.estimated);
2387 assert_ne!(u.tokens, Some(900));
2388 assert!(u.tokens.is_some());
2389 let t = ctx_talk(
2391 "small",
2392 vec![
2393 reply("old", Some((900, "small", Some("small-model")))),
2394 reply("magi: could not run agent", None),
2395 ],
2396 );
2397 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
2398 assert_eq!(
2399 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
2400 None
2401 );
2402 }
2403
2404 #[test]
2405 fn estimate_counts_chars_both_sides_and_standing_prompt() {
2406 let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
2407 assert_eq!(estimate_context_tokens(&t, 0), Some(2)); let op = Turn {
2409 who: Who::Operator,
2410 ..reply("abcdefg", None)
2411 };
2412 t.turns.push(op);
2413 assert_eq!(estimate_context_tokens(&t, 0), Some(4));
2414 assert!(
2415 estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
2416 );
2417 let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
2419 assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
2420 let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
2422 assert_eq!(estimate_context_tokens(¬e, 1000), None);
2423 assert_eq!(
2424 estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
2425 None
2426 );
2427 }
2428
2429 #[test]
2430 fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
2431 let cfg = ctx_config(&[("small-model", 1000)]);
2432 let t = ctx_talk(
2433 "small",
2434 vec![reply(
2435 &"x".repeat(5000),
2436 Some((10, "small", Some("small-model"))),
2437 )],
2438 );
2439 let u = context_usage(&t, Some(&cfg));
2440 assert_eq!((u.tokens, u.estimated), (Some(10), false));
2441 let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
2442 let u = context_usage(&t, Some(&cfg));
2443 assert!(u.estimated && !u.since_switch);
2444 assert_eq!(u.window, Some(1000));
2445 assert!(u.warn && u.percent.unwrap() >= 80);
2446 let t = ctx_talk("small", vec![reply("hi", None)]);
2447 let u = context_usage(&t, Some(&cfg));
2448 assert!(u.estimated && u.percent.is_some());
2449 }
2450
2451 #[test]
2452 fn context_usage_without_a_window_shows_tokens_only() {
2453 let cfg = ctx_config(&[]);
2454 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
2456 let u = context_usage(&t, Some(&cfg));
2457 assert_eq!(
2458 (u.tokens, u.window, u.percent, u.warn),
2459 (Some(5000), None, None, false)
2460 );
2461 let t = ctx_talk(
2462 "small",
2463 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
2464 );
2465 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
2466 assert_eq!(context_usage(&t, None).window, None);
2468 }
2469
2470 #[test]
2471 fn context_usage_switching_model_changes_the_denominator() {
2472 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
2473 let used = reply("hi", Some((900, "small", Some("small-model"))));
2474 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
2475 assert_eq!(
2476 (before.percent, before.warn, before.since_switch),
2477 (Some(90), true, false)
2478 );
2479 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
2482 assert_eq!(after.window, Some(10_000));
2483 assert_eq!(
2484 (after.percent, after.warn, after.since_switch),
2485 (Some(9), false, true)
2486 );
2487 assert_eq!(after.model.as_deref(), Some("big-model"));
2488 }
2489
2490 #[test]
2491 fn a_turn_recorded_before_usage_existed_still_reads() {
2492 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
2493 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
2494 assert!(turn.usage.is_none());
2495 let json = serde_json::to_string(&turn).expect("serialize");
2496 assert!(
2497 !json.contains("usage"),
2498 "absent usage is not written: {json}"
2499 );
2500 }
2501
2502 fn store() -> (tempfile::TempDir, Talks) {
2504 let tmp = tempfile::tempdir().expect("tempdir");
2505 let talks = Talks::at(tmp.path().join("talks"));
2506 (tmp, talks)
2507 }
2508
2509 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
2513 let path = dir.join("mock-talk-agent.sh");
2514 std::fs::write(&path, script).expect("write mock");
2515 AgentSpec {
2516 id: "mock".to_owned(),
2517 kind: AgentKind::Command,
2518 model: None,
2519 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2520 extra_args: Vec::new(),
2521 env,
2522 prompt_delivery: None,
2523 }
2524 }
2525
2526 fn config(spec: AgentSpec) -> Config {
2527 Config {
2528 agents: vec![spec],
2529 graph: Graph {
2530 language: "en".to_owned(),
2531 ..Graph::default()
2532 },
2533 ..Config::default()
2534 }
2535 }
2536
2537 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
2539
2540 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
2542
2543 const ECHO: &str = "#!/bin/sh\ncat\n";
2546
2547 fn env(reply: &str) -> BTreeMap<String, String> {
2548 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
2549 }
2550
2551 #[test]
2552 fn the_frozen_json_field_names_round_trip_through_disk() {
2553 let (tmp, talks) = store();
2554 let mut talk = Talk {
2555 schema: SCHEMA,
2556 id: "20260904-014455-ab12".to_owned(),
2557 repo: tmp.path().to_owned(),
2558 agent: "sonnet".to_owned(),
2559 status: TalkStatus::Open,
2560 turns: Vec::new(),
2561 pending: String::new(),
2562 pending_attachments: Vec::new(),
2563 fallback: false,
2564 persona: String::new(),
2565 persona_dirty: false,
2566 created_at: Timestamp::now(),
2567 updated_at: Timestamp::now(),
2568 seat: SeatState::new(SEAT, "sonnet", 7),
2569 };
2570 talks.put(&mut talk).expect("put");
2571
2572 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
2573 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
2574 for field in [
2575 "schema",
2576 "id",
2577 "repo",
2578 "agent",
2579 "status",
2580 "turns",
2581 "created_at",
2582 "updated_at",
2583 ] {
2584 assert!(v.get(field).is_some(), "missing field `{field}`");
2585 }
2586 assert_eq!(v["schema"], 1);
2587 assert_eq!(v["status"], "open");
2588
2589 let back = talks.get(&talk.id).expect("get");
2590 assert_eq!(back.id, talk.id);
2591 assert_eq!(back.status, TalkStatus::Open);
2592 }
2593
2594 #[test]
2595 fn opening_a_talk_takes_no_agent_turn() {
2596 let (tmp, talks) = store();
2597 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2601 let cfg = config(spec);
2602
2603 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2604 assert_eq!(talk.status, TalkStatus::Open);
2605 assert!(talk.turns.is_empty(), "nothing has been said yet");
2606
2607 let on_disk = talks.get(&talk.id).expect("get");
2608 assert_eq!(on_disk.turns.len(), 0);
2609 }
2610
2611 #[test]
2619 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
2620 let (tmp, talks) = store();
2621 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2622 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2623 chatter_spec.id = "chatter-mock".to_owned();
2624
2625 let mut cfg = Config {
2626 agents: vec![first_spec.clone(), chatter_spec.clone()],
2627 graph: Graph {
2628 language: "en".to_owned(),
2629 ..Graph::default()
2630 },
2631 ..Config::default()
2632 };
2633 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
2634
2635 let talk =
2636 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
2637 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
2638
2639 cfg.roles.chatter = None;
2640 let fallback =
2641 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
2642 assert_eq!(
2643 fallback.agent, first_spec.id,
2644 "unset chatter must fall back to agent::pick's own default order"
2645 );
2646 }
2647
2648 #[test]
2651 fn a_talk_recorded_without_attachments_still_reads() {
2652 let (tmp, talks) = store();
2653 let path = talks.path_of("20260904-014455-ab12");
2654 std::fs::create_dir_all(talks.root()).expect("talks dir");
2655 std::fs::write(
2656 &path,
2657 serde_json::json!({
2658 "schema": 1,
2659 "id": "20260904-014455-ab12",
2660 "repo": tmp.path(),
2661 "agent": "sonnet",
2662 "status": "open",
2663 "turns": [
2664 { "who": "operator", "body": "still there?",
2665 "at": Timestamp::now().to_string() },
2666 ],
2667 "created_at": Timestamp::now().to_string(),
2668 "updated_at": Timestamp::now().to_string(),
2669 "seat": SeatState::new(SEAT, "sonnet", 7),
2670 })
2671 .to_string(),
2672 )
2673 .expect("write pre-attachments talk");
2674
2675 let talk = talks.get("20260904-014455-ab12").expect("must still read");
2676 assert!(talk.turns[0].attachments.is_empty());
2677 }
2678
2679 fn lease_store() -> (tempfile::TempDir, Talks, Talks) {
2680 let tmp = tempfile::TempDir::new().expect("tmp");
2681 let root = tmp.path().join("talks");
2682 (tmp, Talks::at(root.clone()), Talks::at(root))
2683 }
2684
2685 #[test]
2686 fn two_starters_on_one_talk_one_wins_and_the_other_is_refused() {
2687 let (_tmp, a, b) = lease_store();
2688 let won = a.claim_turn("t1").expect("claim").expect("first wins");
2689 assert!(
2690 b.claim_turn("t1").expect("claim").is_none(),
2691 "second is refused"
2692 );
2693 assert!(b.turn_held("t1"));
2694 assert!(
2695 b.claim_turn("t2").expect("claim").is_some(),
2696 "other talks are free"
2697 );
2698 drop(won);
2699 }
2700
2701 #[test]
2702 fn a_stale_lease_is_taken_over_and_the_old_guard_cannot_release_it() {
2703 let (_tmp, a, b) = lease_store();
2704 let old = a.claim_turn("t1").expect("claim").expect("held");
2705 let later = Timestamp::now()
2706 .checked_add(jiff::SignedDuration::from_secs(
2707 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2708 ))
2709 .expect("later");
2710 let new = b
2711 .claim_turn_at("t1", later)
2712 .expect("claim")
2713 .expect("a stale lease is taken over");
2714 drop(old);
2715 assert!(a.turn_held("t1"), "the old guard left the new lease alone");
2716 assert!(new.beat().expect("beat"), "the new owner still beats");
2717 drop(new);
2718 assert!(!a.turn_held("t1"));
2719 }
2720
2721 #[tokio::test]
2722 async fn a_turn_whose_lease_was_taken_over_is_stopped() {
2723 let (_tmp, a, b) = lease_store();
2724 let old = a.claim_turn("t1").expect("claim").expect("held");
2725 let later = Timestamp::now()
2726 .checked_add(jiff::SignedDuration::from_secs(
2727 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2728 ))
2729 .expect("later");
2730 let _new = b
2731 .claim_turn_at("t1", later)
2732 .expect("claim")
2733 .expect("taken over");
2734 let out = old
2735 .beating_every(Duration::from_millis(10), std::future::pending::<()>())
2736 .await;
2737 assert!(out.is_err(), "the displaced turn must stop, not run on");
2738 }
2739
2740 #[tokio::test]
2741 async fn a_turn_that_finishes_is_returned_and_keeps_its_lease_beating() {
2742 let (_tmp, a, _b) = lease_store();
2743 let lease = a.claim_turn("t1").expect("claim").expect("held");
2744 let out = lease
2745 .beating_every(Duration::from_millis(5), async {
2746 tokio::time::sleep(Duration::from_millis(40)).await;
2747 7
2748 })
2749 .await
2750 .expect("still ours");
2751 assert_eq!(out, 7);
2752 assert!(a.turn_held("t1"));
2753 }
2754
2755 #[test]
2756 fn an_unreadable_lease_counts_as_stale() {
2757 let (_tmp, a, b) = lease_store();
2758 std::fs::create_dir_all(a.root()).expect("dir");
2759 std::fs::write(a.turn_path("t1"), "not json").expect("write");
2760 assert!(!a.turn_held("t1"));
2761 assert!(b.claim_turn("t1").expect("claim").is_some());
2762 }
2763
2764 #[test]
2765 fn without_hard_links_a_claim_is_still_exclusive_and_a_young_placeholder_blocks() {
2766 let (_tmp, a, b) = lease_store();
2767 failpoint::force_no_link(true);
2768 let held = a.claim_turn("t1").expect("claim").expect("first wins");
2769 assert!(b.claim_turn("t1").expect("claim").is_none());
2770 assert!(a.turn_held("t1"));
2771 drop(held);
2772 assert!(!a.turn_held("t1"));
2773 let path = a.turn_path("t1");
2775 assert!(create_exclusive(&path, "").expect("placeholder"));
2776 assert!(b.claim_turn("t1").expect("claim").is_none(), "young: held");
2777 age_file(&path);
2778 assert!(b.claim_turn("t1").expect("claim").is_some(), "old: stale");
2779 let lease = a.turn_path("t2");
2781 let lock = TurnLock::take(&lease).expect("take").expect("lock");
2782 assert!(TurnLock::take(&lease).expect("take").is_none());
2783 drop(lock);
2784 assert!(TurnLock::take(&lease).expect("take").is_some());
2785 failpoint::force_no_link(false);
2786 }
2787
2788 #[test]
2789 fn remove_if_carries_removes_only_the_expected_content() {
2790 let (_tmp, a, _b) = lease_store();
2791 std::fs::create_dir_all(a.root()).expect("dir");
2792 let p = a.root().join("x.turn.lock");
2793 std::fs::write(&p, "mine").expect("write");
2794 assert!(!remove_if_carries(&p, "other"));
2795 assert_eq!(std::fs::read_to_string(&p).expect("kept"), "mine");
2796 assert!(remove_if_carries(&p, "mine"));
2797 assert!(!p.exists());
2798 assert!(!remove_if_carries(&p, "mine"), "absent is not a removal");
2799 }
2800
2801 #[test]
2802 fn a_dropped_lock_does_not_remove_a_lock_taken_over_since() {
2803 let (_tmp, a, _b) = lease_store();
2804 std::fs::create_dir_all(a.root()).expect("dir");
2805 let lease = a.turn_path("t1");
2806 let lock = TurnLock::take(&lease).expect("take").expect("lock");
2807 let path = lock.path.clone();
2808 std::fs::write(&path, "someone-else").expect("replace");
2809 drop(lock);
2810 assert_eq!(
2811 std::fs::read_to_string(&path).expect("kept"),
2812 "someone-else"
2813 );
2814 }
2815
2816 #[test]
2817 fn concurrent_takeovers_of_a_stale_lease_have_one_winner() {
2818 let (_tmp, a, _b) = lease_store();
2819 drop(a.claim_turn("t1").expect("claim").expect("held"));
2820 std::fs::write(
2821 a.turn_path("t1"),
2822 serde_json::to_string(&TurnRecord {
2823 token: "gone".into(),
2824 pid: 1,
2825 beat_at: Timestamp::from_second(1).expect("ts"),
2826 })
2827 .expect("json"),
2828 )
2829 .expect("write");
2830 let wins: Vec<_> = std::thread::scope(|sc| {
2831 let hs: Vec<_> = (0..8)
2832 .map(|_| {
2833 let s = a.clone();
2834 sc.spawn(move || s.claim_turn("t1").expect("claim"))
2835 })
2836 .collect();
2837 hs.into_iter().map(|h| h.join().expect("join")).collect()
2838 });
2839 assert_eq!(wins.iter().filter(|w| w.is_some()).count(), 1);
2840 }
2841
2842 fn age_file(path: &Path) {
2843 let f = std::fs::OpenOptions::new()
2844 .write(true)
2845 .open(path)
2846 .expect("open");
2847 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
2848 .expect("age");
2849 }
2850
2851 #[test]
2852 fn a_late_taker_of_a_broken_lock_cannot_disturb_its_replacement() {
2853 let (_tmp, a, _b) = lease_store();
2854 let lease = a.turn_path("t1");
2855 let lock = lease.with_extension("turn.lock");
2856 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2857 std::fs::write(&lock, "t1-dead").expect("dead lock");
2858 age_file(&lock);
2859 let b = TurnLock::take(&lease).expect("take").expect("b wins");
2861 let fresh = std::fs::read_to_string(&lock).expect("read");
2862 assert_eq!(fresh, b.token);
2863 let ticket = std::fs::read_dir(lock.parent().expect("dir"))
2866 .expect("dir")
2867 .flatten()
2868 .map(|e| e.path())
2869 .find(|p| p.to_string_lossy().ends_with(".break.0"))
2870 .expect("ticket");
2871 assert!(!create_exclusive(&ticket, "").expect("ticket"));
2872 assert!(TurnLock::take(&lease).expect("take").is_none());
2873 assert_eq!(std::fs::read_to_string(&lock).expect("read"), fresh);
2874 age_file(&lock);
2876 age_file(&ticket);
2877 std::mem::forget(b);
2878 let c = TurnLock::take(&lease)
2879 .expect("take")
2880 .expect("next generation");
2881 assert_ne!(c.token, fresh);
2882 }
2883
2884 #[test]
2885 fn a_live_ticket_blocks_and_a_stale_one_hands_over_to_the_next_generation() {
2886 let (_tmp, a, _b) = lease_store();
2887 let lease = a.turn_path("t1");
2888 let lock = lease.with_extension("turn.lock");
2889 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2890 std::fs::write(&lock, "t1-dead").expect("dead lock");
2891 age_file(&lock);
2892 let t0 = lock.with_extension("lock.t1-dead.break.0");
2893 assert!(create_exclusive(&t0, "").expect("ticket"));
2894 assert!(TurnLock::take(&lease).expect("take").is_none());
2896 assert_eq!(std::fs::read_to_string(&lock).expect("read"), "t1-dead");
2897 age_file(&t0);
2899 let c = TurnLock::take(&lease).expect("take").expect("generation 1");
2900 assert_eq!(std::fs::read_to_string(&lock).expect("read"), c.token);
2901 assert!(lock.with_extension("lock.t1-dead.break.1").exists());
2902 assert!(TurnLock::take(&lease).expect("take").is_none());
2903 }
2904
2905 #[test]
2906 fn exhausted_ticket_generations_recover_once_the_sweep_ages_them_out() {
2907 let (_tmp, a, _b) = lease_store();
2908 let lease = a.turn_path("t1");
2909 let lock = lease.with_extension("turn.lock");
2910 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2911 std::fs::write(&lock, "t1-dead").expect("dead lock");
2912 age_file(&lock);
2913 let tickets: Vec<_> = (0..TICKET_GENERATIONS)
2914 .map(|n| lock.with_extension(format!("lock.t1-dead.break.{n}")))
2915 .collect();
2916 for t in &tickets {
2917 assert!(create_exclusive(t, "").expect("ticket"));
2918 age_file(t);
2919 }
2920 assert!(TurnLock::take(&lease).expect("take").is_none());
2922 assert!(tickets.iter().all(|t| t.exists()));
2923 for t in &tickets {
2925 let f = std::fs::OpenOptions::new()
2926 .write(true)
2927 .open(t)
2928 .expect("open");
2929 f.set_modified(std::time::SystemTime::now() - TICKET_SWEEP_AGE * 2)
2930 .expect("age");
2931 }
2932 assert!(TurnLock::take(&lease).expect("take").is_none());
2933 assert!(TurnLock::take(&lease).expect("take").is_some());
2934 }
2935
2936 #[test]
2937 fn an_aged_empty_lock_is_broken() {
2938 let (_tmp, a, _b) = lease_store();
2939 let lease = a.turn_path("t1");
2940 let lock = lease.with_extension("turn.lock");
2941 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2942 std::fs::write(&lock, "").expect("empty lock");
2943 age_file(&lock);
2944 assert!(TurnLock::take(&lease).expect("take").is_some());
2945 }
2946
2947 #[test]
2948 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
2949 let (tmp, talks) = store();
2950 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2951 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2952
2953 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
2954 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
2955 let saved = talks.get(&talk.id).expect("reload queued talk");
2956 assert_eq!(saved.pending, "first\n\nsecond");
2957 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
2958
2959 let drained = drain(&mut talk, &talks).expect("drain");
2960 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
2961 let saved = talks.get(&talk.id).expect("reload drained talk");
2962 assert!(saved.pending.is_empty());
2963 assert_eq!(saved.turns.len(), 1);
2964 assert_eq!(saved.turns[0].body, "first\n\nsecond");
2965 }
2966
2967 #[test]
2968 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
2969 let (tmp, talks) = store();
2970 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2971 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2972 let attachment = Attachment {
2973 id: "a".repeat(32),
2974 name: "shot.png".to_owned(),
2975 mime: "image/png".to_owned(),
2976 bytes: 3,
2977 };
2978
2979 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
2980 assert!(
2981 edit_pending_text(
2982 &mut talk,
2983 &talks,
2984 "corrected",
2985 "first",
2986 std::slice::from_ref(&attachment.id),
2987 )
2988 .expect("edit")
2989 );
2990 let saved = talks.get(&talk.id).expect("reload edited draft");
2991 assert_eq!(saved.pending, "corrected");
2992 assert_eq!(saved.pending_attachments, vec![attachment]);
2993
2994 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
2995 assert!(
2996 !edit_pending_text(
2997 &mut talk,
2998 &talks,
2999 "stale edit",
3000 "corrected",
3001 &["a".repeat(32)],
3002 )
3003 .expect("stale edit is a conflict")
3004 );
3005 assert_eq!(
3006 talks.get(&talk.id).expect("reload after conflict").pending,
3007 "corrected\n\nlater"
3008 );
3009 assert!(
3010 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
3011 .expect("stale clear is a conflict")
3012 );
3013 assert_eq!(
3014 talks
3015 .get(&talk.id)
3016 .expect("reload after stale clear")
3017 .pending,
3018 "corrected\n\nlater"
3019 );
3020 }
3021
3022 #[tokio::test]
3023 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
3024 let (tmp, talks) = store();
3025 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
3026 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
3027 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3028 let id = running.id.clone();
3029 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
3030
3031 let response_talks = talks.clone();
3032 let response_cfg = cfg.clone();
3033 let reply = tokio::spawn(async move {
3034 respond(
3035 &response_talks.claim_turn(&running.id).unwrap().unwrap(),
3036 &mut running,
3037 &response_talks,
3038 &response_cfg,
3039 &first,
3040 )
3041 .await
3042 });
3043 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
3044
3045 let mut queued = talks.get(&id).expect("queued handle");
3046 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
3047 reply.await.expect("join").expect("reply");
3048
3049 let saved = talks.get(&id).expect("reload");
3050 assert_eq!(saved.pending, "next");
3051 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
3052 }
3053
3054 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
3057 let calls = dir.join(format!("{id}.calls"));
3058 let script = format!(
3059 "#!/bin/sh\necho x >> '{}'\n{body}\n",
3060 calls.to_string_lossy()
3061 );
3062 let path = dir.join(format!("mock-{id}.sh"));
3063 std::fs::write(&path, script).expect("write mock");
3064 AgentSpec {
3065 id: id.to_owned(),
3066 kind: AgentKind::Command,
3067 model: None,
3068 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
3069 extra_args: Vec::new(),
3070 env: BTreeMap::new(),
3071 prompt_delivery: None,
3072 }
3073 }
3074
3075 fn calls(dir: &Path, id: &str) -> usize {
3076 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
3077 }
3078
3079 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
3080 let mut cfg = config(specs[0].clone());
3081 cfg.agents = specs;
3082 cfg.roles.chatter = Some(AgentChoice::Chain(
3083 ids.iter().map(|s| (*s).to_owned()).collect(),
3084 ));
3085 cfg
3086 }
3087
3088 #[tokio::test]
3089 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
3090 let (tmp, talks) = store();
3091 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3092 let b = counting_agent(tmp.path(), "b", "cat");
3093 let cfg = chain_config(vec![a, b], &["a", "b"]);
3094 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3095 assert_eq!(talk.agent, "a");
3096
3097 say(
3098 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3099 &mut talk,
3100 &talks,
3101 &cfg,
3102 "hello there",
3103 Vec::new(),
3104 )
3105 .await
3106 .expect("turn");
3107 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
3108 assert_eq!(calls(tmp.path(), "b"), 1);
3109 assert_eq!(talk.agent, "b", "the switch persists");
3110 assert!(talks.get(&talk.id).unwrap().agent == "b");
3111 let reply = talk.turns.last().unwrap();
3112 assert!(reply.body.contains("hello there"));
3113 assert!(
3114 reply.body.contains("magi task add --solo"),
3115 "a fresh seat gets the full briefing"
3116 );
3117 assert!(
3118 talk.turns
3119 .iter()
3120 .any(|t| t.body.contains("agent changed from a to b")),
3121 "the switch is noted"
3122 );
3123 }
3124
3125 #[tokio::test]
3126 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
3127 let (tmp, talks) = store();
3128 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3129 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
3130 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
3131 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3132
3133 let err = say(
3134 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3135 &mut talk,
3136 &talks,
3137 &cfg,
3138 "hi",
3139 Vec::new(),
3140 )
3141 .await
3142 .expect_err("every agent failed");
3143 assert!(err.to_string().contains("`a`"), "{err:#}");
3144 assert_eq!(calls(tmp.path(), "a"), 1);
3145 assert_eq!(calls(tmp.path(), "b"), 1);
3146 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
3147 }
3148
3149 #[test]
3150 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
3151 let (tmp, talks) = store();
3152 let b = counting_agent(tmp.path(), "b", "cat");
3153 let cfg = chain_config(vec![b], &["ghost", "b"]);
3154 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3155 assert_eq!(talk.agent, "b");
3156 }
3157
3158 #[tokio::test]
3159 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
3160 let (tmp, talks) = store();
3161 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3162 let b = counting_agent(tmp.path(), "b", "cat");
3163 let cfg = chain_config(vec![a, b], &["a", "b"]);
3164 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
3165 say(
3166 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3167 &mut talk,
3168 &talks,
3169 &cfg,
3170 "hi",
3171 Vec::new(),
3172 )
3173 .await
3174 .expect_err("a alone, and it fails");
3175 assert_eq!(calls(tmp.path(), "b"), 0);
3176 assert_eq!(talk.agent, "a");
3177 }
3178
3179 #[tokio::test]
3180 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
3181 let (tmp, talks) = store();
3182 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3183 let b = counting_agent(tmp.path(), "b", "cat");
3184 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
3185 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
3186 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
3187 say(
3188 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3189 &mut talk,
3190 &talks,
3191 &cfg,
3192 "hi",
3193 Vec::new(),
3194 )
3195 .await
3196 .expect_err("c alone, and it fails");
3197 assert_eq!(calls(tmp.path(), "b"), 0);
3198 }
3199
3200 #[test]
3201 fn briefing_names_the_operator_with_or_without_a_persona() {
3202 let plain = briefing(Path::new("/repo"), "en", false);
3203 let named = briefing_with(Path::new("/repo"), "en", false, None, Some("Commander"));
3204 assert!(named.starts_with(&plain));
3205 assert!(named.contains("# Addressing the operator"));
3206 assert!(named.contains("\"Commander\""));
3207 let rei = crate::persona::builtin_catalog()
3208 .into_iter()
3209 .find(|p| p.id == "rei")
3210 .expect("rei");
3211 let with = briefing_with(
3212 Path::new("/repo"),
3213 "en",
3214 false,
3215 Some(&rei),
3216 Some("Commander"),
3217 );
3218 assert!(with.contains("\"Commander\""));
3219 assert!(!with.contains("# Addressing the operator"));
3220 }
3221
3222 #[test]
3223 fn briefing_carries_a_persona_section_only_when_one_is_chosen() {
3224 let plain = briefing(Path::new("/repo"), "en", false);
3225 assert_eq!(
3226 plain,
3227 briefing_with(Path::new("/repo"), "en", false, None, None)
3228 );
3229 assert!(!plain.contains("Persona"));
3230 let rei = crate::persona::builtin_catalog()
3231 .into_iter()
3232 .find(|p| p.id == "rei")
3233 .expect("rei");
3234 let with = briefing_with(Path::new("/repo"), "en", false, Some(&rei), None);
3235 assert!(with.starts_with(&plain), "the plain briefing is untouched");
3236 assert!(with.contains("# Persona (tone only)"));
3237 assert!(with.contains("TONE ONLY"));
3238 assert!(with.contains("task ids"));
3239 assert!(with.contains("write policy"));
3240 assert!(with.contains("`magi task add`"));
3241 assert!(with.contains("Rei Ayanami"));
3242 }
3243
3244 #[test]
3245 fn a_talk_written_before_personas_still_loads() {
3246 let (tmp, talks) = store();
3247 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3248 let cfg = config(spec);
3249 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3250 let path = talks.root.join(format!("{}.json", talk.id));
3251 let mut v: serde_json::Value =
3252 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3253 v.as_object_mut().unwrap().remove("persona");
3254 v.as_object_mut().unwrap().remove("persona_dirty");
3255 std::fs::write(&path, v.to_string()).unwrap();
3256 let loaded = talks.get(&talk.id).expect("old record loads");
3257 assert_eq!(loaded.persona, "");
3258 assert!(!loaded.persona_dirty);
3259 }
3260
3261 #[tokio::test]
3262 async fn a_persona_switch_notes_marks_dirty_and_updates_the_next_turn_once() {
3263 let (tmp, talks) = store();
3264 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3265 let cfg = config(spec);
3266 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3267 let lease = talks.claim_turn(&talk.id).unwrap().unwrap();
3268 say(&lease, &mut talk, &talks, &cfg, "hello", Vec::new())
3269 .await
3270 .expect("first turn");
3271 assert!(!talk.turns[1].body.contains("Persona"));
3272
3273 assert!(switch_persona(&mut talk, &talks, "misato").expect("switch"));
3274 assert!(!switch_persona(&mut talk, &talks, "misato").expect("same"));
3275 assert!(talk.persona_dirty);
3276 assert!(talk.turns.last().unwrap().body.contains("persona changed"));
3277
3278 say(&lease, &mut talk, &talks, &cfg, "next", Vec::new())
3279 .await
3280 .expect("turn");
3281 let prompt = &talk.turns.last().unwrap().body;
3282 assert!(prompt.contains("# Persona update"), "{prompt}");
3283 assert!(prompt.contains("Misato Katsuragi"));
3284 assert!(!talk.persona_dirty, "cleared after a successful turn");
3285
3286 say(&lease, &mut talk, &talks, &cfg, "again", Vec::new())
3287 .await
3288 .expect("turn");
3289 assert!(!talk.turns.last().unwrap().body.contains("# Persona update"));
3290
3291 assert!(switch_persona(&mut talk, &talks, "default").expect("back"));
3292 assert_eq!(talk.persona, "");
3293 say(&lease, &mut talk, &talks, &cfg, "plain", Vec::new())
3294 .await
3295 .expect("turn");
3296 assert!(
3297 talk.turns
3298 .last()
3299 .unwrap()
3300 .body
3301 .contains("turned the persona off")
3302 );
3303
3304 assert!(switch_persona(&mut talk, &talks, "rei").expect("rei"));
3306 let mut no_sessions = cfg.clone();
3307 no_sessions.graph.sessions = false;
3308 say(&lease, &mut talk, &talks, &no_sessions, "one", Vec::new())
3309 .await
3310 .expect("turn");
3311 say(&lease, &mut talk, &talks, &no_sessions, "two", Vec::new())
3312 .await
3313 .expect("turn");
3314 let last = &talk.turns.last().unwrap().body;
3315 assert!(!talk.persona_dirty);
3316 assert!(last.contains("# Persona (tone only)"), "{last}");
3317 assert!(last.contains("Rei Ayanami"));
3318 }
3319
3320 #[tokio::test]
3321 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
3322 let (tmp, talks) = store();
3323 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3324 let cfg = config(spec);
3325 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3326
3327 say(
3328 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3329 &mut talk,
3330 &talks,
3331 &cfg,
3332 "what does the queue module do?",
3333 Vec::new(),
3334 )
3335 .await
3336 .expect("first turn");
3337 let first_prompt = &talk.turns[1].body;
3338 assert!(first_prompt.contains("magi task add --solo"));
3339 assert!(first_prompt.contains("what does the queue module do?"));
3340
3341 say(
3342 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3343 &mut talk,
3344 &talks,
3345 &cfg,
3346 "and how is it locked?",
3347 Vec::new(),
3348 )
3349 .await
3350 .expect("second turn");
3351 let second_prompt = &talk.turns[3].body;
3352 assert!(
3353 !second_prompt.contains("magi task add --solo"),
3354 "the briefing is sent once, not on every turn: {second_prompt}"
3355 );
3356 assert!(second_prompt.contains("and how is it locked?"));
3357 }
3358
3359 #[tokio::test]
3360 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
3361 let (tmp, talks) = store();
3362 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3363 let mut b = a.clone();
3364 b.id = "other".to_owned();
3365 let mut cfg = config(a.clone());
3366 cfg.agents.push(b.clone());
3367 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
3368 say(
3369 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3370 &mut talk,
3371 &talks,
3372 &cfg,
3373 "remember the walrus",
3374 Vec::new(),
3375 )
3376 .await
3377 .expect("first turn");
3378 let old_session = talk.seat.claude_session.clone();
3379 assert_eq!(talk.seat.turns, 1);
3380
3381 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
3382 assert_eq!(talk.agent, "other");
3383 assert_eq!(talk.seat.turns, 0);
3384 assert_eq!(talk.seat.agent, "other");
3385 assert_ne!(talk.seat.claude_session, old_session);
3386 let note = talk.turns.last().expect("note");
3387 assert_eq!(note.who, Who::Agent);
3388 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3389 assert!(note.body.contains("changed from"), "{}", note.body);
3390 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
3391
3392 let before = talk.turns.len();
3393 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
3394 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
3395
3396 say(
3397 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3398 &mut talk,
3399 &talks,
3400 &cfg,
3401 "what did I say?",
3402 Vec::new(),
3403 )
3404 .await
3405 .expect("turn after switch");
3406 let prompt = &talk.turns.last().expect("reply").body;
3407 assert!(prompt.contains("remember the walrus"), "{prompt}");
3408 assert!(prompt.contains("## magi"), "{prompt}");
3409 assert!(prompt.contains("what did I say?"), "{prompt}");
3410 }
3411
3412 #[tokio::test]
3413 async fn say_appends_the_operator_turn_then_the_agent_turn() {
3414 let (tmp, talks) = store();
3415 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3416 let cfg = config(spec);
3417 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3418
3419 say(
3420 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3421 &mut talk,
3422 &talks,
3423 &cfg,
3424 "can I rename this function?",
3425 Vec::new(),
3426 )
3427 .await
3428 .expect("say");
3429
3430 assert_eq!(talk.turns.len(), 2);
3431 assert_eq!(talk.turns[0].who, Who::Operator);
3432 assert_eq!(talk.turns[0].body, "can I rename this function?");
3433 assert_eq!(talk.turns[1].who, Who::Agent);
3434 assert_eq!(talk.turns[1].body, "go ahead");
3435 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
3436 }
3437
3438 #[tokio::test]
3439 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
3440 let (tmp, talks) = store();
3441 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
3442 let cfg = config(spec);
3443 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3444
3445 let err = say(
3446 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3447 &mut talk,
3448 &talks,
3449 &cfg,
3450 "check the tests",
3451 Vec::new(),
3452 )
3453 .await
3454 .expect_err("a turn with no answer is an error");
3455 assert!(err.to_string().contains("no answer"), "{err}");
3456
3457 let on_disk = talks.get(&talk.id).expect("get");
3458 assert_eq!(on_disk.turns.len(), 2);
3459 assert_eq!(on_disk.turns[0].body, "check the tests");
3460 let note = &on_disk.turns[1];
3461 assert_eq!(note.who, Who::Agent);
3462 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3463 assert!(note.body.contains("your message is saved"));
3464 }
3465
3466 #[tokio::test]
3472 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
3473 let (tmp, talks) = store();
3474 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3475 let cfg = config(spec);
3476 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3477
3478 let text =
3479 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
3480 failpoint::force_put_failures(PUT_RETRIES - 1);
3483 respond(
3484 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3485 &mut talk,
3486 &talks,
3487 &cfg,
3488 &text,
3489 )
3490 .await
3491 .expect("respond must survive a write failure its own retries can outlast");
3492
3493 assert_eq!(talk.turns.len(), 2);
3494 assert_eq!(talk.turns[1].who, Who::Agent);
3495 assert_eq!(talk.turns[1].body, "go ahead");
3496 let on_disk = talks.get(&talk.id).expect("get");
3497 assert_eq!(
3498 on_disk.turns, talk.turns,
3499 "the reply must reach disk despite the early write failures"
3500 );
3501 }
3502
3503 #[tokio::test]
3509 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
3510 let (tmp, talks) = store();
3511 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3512 let cfg = config(spec);
3513 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3514
3515 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3516 failpoint::force_put_failures(PUT_RETRIES);
3521 let err = respond(
3522 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3523 &mut talk,
3524 &talks,
3525 &cfg,
3526 &text,
3527 )
3528 .await
3529 .expect_err("a reply that cannot be saved must be reported, not swallowed");
3530 assert!(err.to_string().contains("could not be saved"), "{err}");
3531
3532 let on_disk = talks.get(&talk.id).expect("get");
3533 assert_eq!(
3534 on_disk.turns.len(),
3535 2,
3536 "the operator turn plus a visible note"
3537 );
3538 assert_eq!(on_disk.turns[0].body, "check the tests");
3539 let note = &on_disk.turns[1];
3540 assert_eq!(note.who, Who::Agent);
3541 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3542 assert!(
3543 note.body.contains("could not be saved"),
3544 "the operator must be told the reply is missing, not left staring \
3545 at a gap with no explanation: {}",
3546 note.body
3547 );
3548 assert_eq!(
3549 talk.turns, on_disk.turns,
3550 "the in-memory talk must match what actually landed on disk"
3551 );
3552
3553 let artifacts = talks.artifacts_of(&talk.id);
3556 let stash = std::fs::read_dir(&artifacts)
3557 .expect("artifacts dir")
3558 .filter_map(|e| e.ok())
3559 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
3560 .expect("a stash file for the lost reply");
3561 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
3562 assert_eq!(stashed, "go ahead");
3563
3564 assert_eq!(
3572 on_disk.seat.turns, 1,
3573 "the note's write must carry the turn the CLI actually took"
3574 );
3575 assert_eq!(
3576 on_disk.seat.claude_session, talk.seat.claude_session,
3577 "the session id handed to the CLI must survive the failed reply"
3578 );
3579 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
3580 assert!(
3581 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
3582 "the next turn must resume, not open the same session id twice"
3583 );
3584 }
3585
3586 #[tokio::test]
3591 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
3592 let (tmp, talks) = store();
3593 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3594 let cfg = config(spec);
3595 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3596
3597 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3598 failpoint::force_put_failures(PUT_RETRIES * 2);
3601 let err = respond(
3602 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3603 &mut talk,
3604 &talks,
3605 &cfg,
3606 &text,
3607 )
3608 .await
3609 .expect_err("neither the reply nor the note could be saved");
3610 assert!(err.to_string().contains("could not be saved"), "{err}");
3611
3612 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
3613 let on_disk = talks.get(&talk.id).expect("get");
3614 assert_eq!(on_disk.turns.len(), 1);
3615
3616 assert_eq!(
3626 on_disk.seat.turns, 0,
3627 "an unwritable file cannot record the turn the CLI took"
3628 );
3629 assert_eq!(
3630 talk.seat.turns, 1,
3631 "the in-memory seat still reports the turn the CLI actually took"
3632 );
3633 assert_eq!(
3634 on_disk.seat.claude_session, talk.seat.claude_session,
3635 "the session id was minted at `begin` and never changes here"
3636 );
3637 }
3638
3639 #[tokio::test]
3643 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
3644 let (tmp, talks) = store();
3645 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3646 let cfg = config(spec);
3647 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3648
3649 let att = talks
3650 .put_attachment(
3651 &talk.id,
3652 "image/png",
3653 "screenshot.png",
3654 b"pretend-png-bytes",
3655 )
3656 .expect("put attachment");
3657
3658 say(
3659 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3660 &mut talk,
3661 &talks,
3662 &cfg,
3663 "",
3664 vec![att.clone()],
3665 )
3666 .await
3667 .expect("an empty body with an attachment is still a turn");
3668
3669 let operator_turn = &talk.turns[0];
3670 assert_eq!(operator_turn.who, Who::Operator);
3671 assert_eq!(operator_turn.body, "");
3672 assert_eq!(operator_turn.attachments, vec![att.clone()]);
3673
3674 let prompt = &talk.turns[1].body;
3675 let expected_path = talks
3676 .attachments_dir(&talk.id)
3677 .join(format!("{}.png", att.id));
3678 assert!(
3679 prompt.contains(&expected_path.display().to_string()),
3680 "the agent must be told the attachment's absolute path: {prompt}"
3681 );
3682 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
3683 }
3684
3685 #[test]
3694 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
3695 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
3696 let att = Attachment {
3697 id: "0".repeat(32),
3698 name: "shot.png".to_owned(),
3699 mime: "image/png".to_owned(),
3700 bytes: 3,
3701 };
3702 let path = talks
3703 .attachment_path("some-talk-id", &att)
3704 .expect("a supported mime always yields a path");
3705 assert!(
3706 path.is_absolute(),
3707 "must be absolute even off a relative store root: {}",
3708 path.display()
3709 );
3710 }
3711
3712 #[tokio::test]
3713 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
3714 let (tmp, talks) = store();
3719 let slow = mock_agent(
3720 tmp.path(),
3721 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
3722 BTreeMap::new(),
3723 );
3724 let mut cfg = config(slow);
3725 cfg.graph.timeout_talk = 1;
3726 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3727
3728 let err = say(
3729 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3730 &mut talk,
3731 &talks,
3732 &cfg,
3733 "check the tests",
3734 Vec::new(),
3735 )
3736 .await
3737 .expect_err("a turn that never answers is an error");
3738 assert!(
3739 err.to_string().contains("did not answer within 1s"),
3740 "{err}"
3741 );
3742
3743 let on_disk = talks.get(&talk.id).expect("get");
3744 let note = on_disk.turns.last().expect("a note turn was recorded");
3745 assert!(
3746 note.body.contains("did not answer within 1s"),
3747 "the transcript must show the configured timeout: {}",
3748 note.body
3749 );
3750 }
3751
3752 #[test]
3753 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
3754 let (tmp, talks) = store();
3755 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3756 let cfg = config(spec);
3757 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3758
3759 close(&mut talk, &talks).expect("close");
3760 assert_eq!(talk.status, TalkStatus::Closed);
3761 close(&mut talk, &talks).expect("closing twice is not an error");
3762
3763 let err =
3764 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
3765 assert!(err.to_string().contains("closed"));
3766 let _ = &cfg; }
3768
3769 #[tokio::test]
3770 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3771 let (tmp, talks) = store();
3772 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3773 let cfg = config(spec);
3774 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3777
3778 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
3782 close(&mut closed_elsewhere, &talks).expect("close");
3783 assert_eq!(
3784 talks.get(&in_flight.id).expect("reread").status,
3785 TalkStatus::Closed,
3786 "the close landed on disk before the turn finished"
3787 );
3788
3789 assert_eq!(in_flight.status, TalkStatus::Open);
3793 respond(
3794 &talks.claim_turn(&in_flight.id).unwrap().unwrap(),
3795 &mut in_flight,
3796 &talks,
3797 &cfg,
3798 "one more question",
3799 )
3800 .await
3801 .expect("the turn itself still completes");
3802
3803 let on_disk = talks.get(&in_flight.id).expect("reread");
3804 assert_eq!(
3805 on_disk.status,
3806 TalkStatus::Closed,
3807 "a close must stick even when a turn that started before it finishes after it"
3808 );
3809 assert!(
3812 on_disk.turns.iter().any(|t| t.body == "here you go"),
3813 "the in-flight turn's own reply is still recorded: {:?}",
3814 on_disk.turns
3815 );
3816 }
3817
3818 #[test]
3819 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
3820 let (tmp, talks) = store();
3821 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3822 let cfg = config(spec);
3823 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3826
3827 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
3830 close(&mut closed_elsewhere, &talks).expect("close");
3831 assert_eq!(
3832 talks.get(&stale.id).expect("reread").status,
3833 TalkStatus::Closed,
3834 "the close landed on disk before record was called"
3835 );
3836
3837 assert_eq!(stale.status, TalkStatus::Open);
3841 let err = record(&mut stale, &talks, "still there?", Vec::new())
3842 .expect_err("a close that landed first must be honored, not overwritten");
3843 assert!(err.to_string().contains("closed"));
3844
3845 let on_disk = talks.get(&stale.id).expect("reread");
3846 assert_eq!(
3847 on_disk.status,
3848 TalkStatus::Closed,
3849 "record must not resurrect a conversation closed while its snapshot was stale"
3850 );
3851 assert!(
3852 on_disk.turns.is_empty(),
3853 "the rejected turn must not have been appended: {:?}",
3854 on_disk.turns
3855 );
3856 let _ = &cfg; }
3858
3859 #[test]
3860 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
3861 let (tmp, talks) = store();
3862 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3863 let cfg = config(spec);
3864 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3865
3866 let held = talks.guard().unwrap();
3870
3871 let talks2 = talks.clone();
3872 let id = talk.id.clone();
3873 let closing = std::thread::spawn(move || {
3874 let mut talk = talks2.get(&id).expect("get");
3875 close(&mut talk, &talks2).expect("close");
3876 });
3877
3878 std::thread::sleep(Duration::from_millis(50));
3879 assert!(
3880 !closing.is_finished(),
3881 "close must wait for the guard, not read and write while it is held - \
3882 a re-read alone narrows this window without closing it"
3883 );
3884
3885 drop(held);
3886 closing.join().expect("close thread panicked");
3887
3888 assert_eq!(
3889 talks.get(&talk.id).expect("reread").status,
3890 TalkStatus::Closed,
3891 "once the guard is free, close still lands"
3892 );
3893 let _ = &cfg; }
3895
3896 #[test]
3897 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
3898 let (tmp, talks) = store();
3899 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3900 let cfg = config(spec);
3901 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3902
3903 close(&mut talk, &talks).expect("close");
3904 assert_eq!(talk.status, TalkStatus::Closed);
3905
3906 reopen(&mut talk, &talks).expect("reopen");
3907 assert_eq!(talk.status, TalkStatus::Open);
3908 assert_eq!(
3909 talks.get(&talk.id).expect("reread").status,
3910 TalkStatus::Open
3911 );
3912
3913 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
3915 assert_eq!(talk.status, TalkStatus::Open);
3916
3917 record(&mut talk, &talks, "one more thing", Vec::new())
3918 .expect("a reopened talk takes turns again");
3919 let _ = &cfg; }
3921
3922 #[test]
3923 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
3924 let (tmp, talks) = store();
3925 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3926 let cfg = config(spec);
3927 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3928
3929 let artifacts = talks.artifacts_of(&talk.id);
3930 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
3931 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
3932
3933 talks.remove(&talk.id).expect("remove");
3934 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
3935 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
3936 assert!(
3937 talks.get(&talk.id).is_err(),
3938 "a removed talk cannot be read back"
3939 );
3940
3941 let err = talks
3942 .remove("nonexistent-id")
3943 .expect_err("unknown id refused");
3944 assert!(err.to_string().contains("no talk matches"), "{err}");
3945 let _ = &cfg; }
3947
3948 #[tokio::test]
3949 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3950 let (tmp, talks) = store();
3951 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3952 let cfg = config(spec);
3953 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3956
3957 talks.remove(&in_flight.id).expect("remove");
3958 assert!(
3959 talks.get(&in_flight.id).is_err(),
3960 "the delete landed on disk before the turn finished"
3961 );
3962
3963 respond(
3966 &talks.claim_turn(&in_flight.id).unwrap().unwrap(),
3967 &mut in_flight,
3968 &talks,
3969 &cfg,
3970 "one more question",
3971 )
3972 .await
3973 .expect("the turn itself still completes rather than erroring");
3974
3975 assert!(
3976 talks.get(&in_flight.id).is_err(),
3977 "a delete must stick even when a turn that started before it finishes after it"
3978 );
3979 }
3980
3981 #[test]
3982 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
3983 let (tmp, talks) = store();
3984 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3985 let cfg = config(spec);
3986 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3989
3990 talks.remove(&stale.id).expect("remove");
3991
3992 let err = record(&mut stale, &talks, "still there?", Vec::new())
3996 .expect_err("a delete that landed first must be honored, not overwritten");
3997 assert!(err.to_string().contains("deleted"), "{err}");
3998
3999 assert!(
4000 talks.get(&stale.id).is_err(),
4001 "record must not resurrect a conversation deleted while its snapshot was stale"
4002 );
4003 let _ = &cfg; }
4005
4006 #[test]
4007 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
4008 let (tmp, talks) = store();
4009 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4010 let cfg = config(spec);
4011 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4014
4015 talks.remove(&stale.id).expect("remove");
4016
4017 let err = close(&mut stale, &talks)
4021 .expect_err("a delete that landed first must be honored, not overwritten");
4022 assert!(err.to_string().contains("deleted"), "{err}");
4023
4024 assert!(
4025 talks.get(&stale.id).is_err(),
4026 "close must not resurrect a conversation deleted while its snapshot was stale"
4027 );
4028 let _ = &cfg; }
4030
4031 #[test]
4032 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
4033 let (tmp, talks) = store();
4034 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4035 let cfg = config(spec);
4036 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4039 close(&mut stale, &talks).expect("close");
4040
4041 talks.remove(&stale.id).expect("remove");
4042
4043 let err = reopen(&mut stale, &talks)
4047 .expect_err("a delete that landed first must be honored, not overwritten");
4048 assert!(err.to_string().contains("deleted"), "{err}");
4049
4050 assert!(
4051 talks.get(&stale.id).is_err(),
4052 "reopen must not resurrect a conversation deleted while its snapshot was stale"
4053 );
4054 let _ = &cfg; }
4056
4057 #[test]
4058 fn list_puts_open_talks_before_closed_ones() {
4059 let (tmp, talks) = store();
4060 let make = |id: &str, status: TalkStatus| {
4061 let mut t = Talk {
4062 schema: SCHEMA,
4063 id: id.to_owned(),
4064 repo: tmp.path().to_owned(),
4065 agent: "mock".to_owned(),
4066 status,
4067 turns: Vec::new(),
4068 pending: String::new(),
4069 pending_attachments: Vec::new(),
4070 fallback: false,
4071 persona: String::new(),
4072 persona_dirty: false,
4073 created_at: Timestamp::now(),
4074 updated_at: Timestamp::now(),
4075 seat: SeatState::new(SEAT, "mock", 7),
4076 };
4077 talks.put(&mut t).expect("put");
4078 };
4079 make("20260901-000000-0001", TalkStatus::Open);
4080 make("20260902-000000-0002", TalkStatus::Open);
4081 make("20260903-000000-0003", TalkStatus::Closed);
4082
4083 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
4084 assert_eq!(
4085 ids,
4086 [
4087 "20260902-000000-0002",
4088 "20260901-000000-0001",
4089 "20260903-000000-0003"
4090 ]
4091 );
4092 assert_eq!(talks.count_open(), 2);
4093 }
4094
4095 #[test]
4096 fn tasks_of_finds_only_this_talks_own_tasks() {
4097 let dir = tempfile::tempdir().expect("tempdir");
4098 let queue = Queue::at(dir.path().join("queue"));
4099
4100 let mut mine = Task::new(
4101 "rework the loader".to_owned(),
4102 "rework the loader".to_owned(),
4103 PathBuf::from("/repo"),
4104 Source::Agent {
4105 run: "20260904-014455-ab12".to_owned(),
4106 node: "chat".to_owned(),
4107 },
4108 );
4109 queue.put(&mut mine).expect("put mine");
4110
4111 let mut theirs = Task::new(
4112 "unrelated".to_owned(),
4113 "unrelated".to_owned(),
4114 PathBuf::from("/repo"),
4115 Source::Agent {
4116 run: "20260904-090000-zz99".to_owned(),
4117 node: "implement".to_owned(),
4118 },
4119 );
4120 queue.put(&mut theirs).expect("put theirs");
4121
4122 let mut human = Task::new(
4123 "typed by hand".to_owned(),
4124 "typed by hand".to_owned(),
4125 PathBuf::from("/repo"),
4126 Source::Human,
4127 );
4128 queue.put(&mut human).expect("put human");
4129
4130 let found = tasks_of(&queue, "20260904-014455-ab12");
4131 assert_eq!(found.len(), 1);
4132 assert_eq!(found[0].id, mine.id);
4133 }
4134
4135 #[test]
4136 fn the_briefing_names_solo_task_add() {
4137 let brief = briefing(Path::new("/repo"), "en", false);
4138 assert!(brief.contains("magi task add --solo"));
4139 assert!(brief.contains("/repo"));
4140 assert!(!brief.contains("Hold this conversation in"));
4141 }
4142
4143 #[test]
4149 fn the_briefing_explains_targeting_a_different_repository_by_name() {
4150 let brief = briefing(Path::new("/repo"), "en", false);
4151 assert!(brief.contains("--repo does not have to be a full path"));
4152 assert!(brief.contains("owner/repo"));
4153 assert!(brief.contains("magi repos"));
4154 assert!(brief.contains("ask the operator"));
4155 }
4156
4157 #[test]
4158 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
4159 let brief = briefing(Path::new("/repo"), "en", false);
4160 assert!(brief.contains("--attach <path>"), "{brief}");
4161 assert!(brief.contains("deleting this conversation"), "{brief}");
4162 }
4163
4164 #[test]
4165 fn the_briefing_names_the_language_when_it_is_not_english() {
4166 let brief = briefing(Path::new("/repo"), "Japanese", false);
4167 assert!(brief.contains("Hold this conversation in Japanese"));
4168 }
4169
4170 #[test]
4171 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
4172 let read_only = briefing(Path::new("/repo"), "en", false);
4173 assert!(read_only.contains("Do not write files"));
4174 assert!(!read_only.contains("allow_write"));
4175
4176 let writable = briefing(Path::new("/repo"), "en", true);
4177 assert!(!writable.contains("Do not write files"));
4178 assert!(writable.contains("allow_write = true"));
4179 assert!(writable.contains("magi task add --solo"));
4182 assert!(writable.contains("say plainly what you"));
4183 }
4184}