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 =
1608 crate::run::try_home().map(|home| crate::ask::Questions::at(home.join("questions")));
1609 let consulted = questions
1610 .as_ref()
1611 .is_some_and(|q| crate::consult::pending_consults(q, &talk.id));
1612 let consult_roots: Vec<PathBuf> = match &questions {
1613 Some(q) if consulted => vec![q.root().to_path_buf()],
1614 _ => Vec::new(),
1615 };
1616
1617 let (allow_write, unsandboxed) = turn_access(cfg.talk.allow_write, consulted);
1618
1619 let artifacts = store.artifacts_of(&talk.id);
1620 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
1623 let stem = format!("turn-{}", operator_turns.max(1));
1624 let cache_dir = cfg.cache_dir();
1627
1628 let mut chain = vec![spec.clone()];
1631 if let Some(choice) = cfg.roles.chatter.as_ref()
1632 && talk.fallback
1633 {
1634 for id in choice.ids() {
1635 if id == talk.agent || chain.iter().any(|s| s.id == id) {
1636 continue;
1637 }
1638 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1639 Ok(s) => chain.push(s),
1640 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1641 }
1642 }
1643 }
1644
1645 let persona = crate::persona::active(&cfg.talk.personas, &talk.persona);
1646 let operator_name = cfg.talk.operator_name();
1647 let persona_update = if talk.persona_dirty {
1648 format!(
1649 "{}\n\n",
1650 crate::persona::update_block_for(persona.as_ref(), operator_name)
1651 )
1652 } else {
1653 String::new()
1654 };
1655
1656 let mut outcome = None;
1657 let mut fell_back_from: Option<String> = None;
1658 let mut first_try: Option<(String, SeatState)> = None;
1661 for (n, spec) in chain.iter().enumerate() {
1662 if n > 0 {
1663 if first_try.is_none() {
1664 first_try = Some((talk.agent.clone(), talk.seat.clone()));
1665 }
1666 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1667 fell_back_from.get_or_insert_with(|| talk.agent.clone());
1670 talk.agent = spec.id.clone();
1671 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1672 }
1673 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1674 let first_ever = talk.turns.len() <= 1;
1675 let body = if talk.seat.turns == 0 && first_ever {
1676 format!(
1677 "{}\n\n# Operator\n\n{text}{last_note}",
1678 briefing_with(
1679 &talk.repo,
1680 &cfg.graph.language,
1681 cfg.talk.allow_write,
1682 persona.as_ref(),
1683 operator_name
1684 )
1685 )
1686 } else if talk.seat.turns == 0 {
1687 format!(
1690 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1691 briefing_with(
1692 &talk.repo,
1693 &cfg.graph.language,
1694 cfg.talk.allow_write,
1695 persona.as_ref(),
1696 operator_name
1697 ),
1698 transcript(talk, store)
1699 )
1700 } else if resuming {
1701 format!("{persona_update}{text}{last_note}")
1702 } else {
1703 let mut standing = match (&persona, talk.persona_dirty) {
1707 (Some(p), false) => format!("{}\n", crate::persona::section_for(p, operator_name)),
1708 _ => persona_update.clone(),
1709 };
1710 if let (None, Some(n)) = (&persona, operator_name) {
1713 standing.push_str(&format!(
1714 "# Addressing the operator\n{}\n",
1715 crate::persona::addressing(n)
1716 ));
1717 }
1718 format!("{}\n\n{standing}{text}{last_note}", transcript(talk, store))
1719 };
1720 let attempt_stem = if n == 0 {
1721 stem.clone()
1722 } else {
1723 format!("{stem}-{}", spec.id)
1724 };
1725 let inv = Invocation {
1726 cwd: &talk.repo,
1727 prompt: &body,
1728 timeout: turn_timeout(cfg),
1729 allow_write,
1734 unsandboxed,
1735 sessions: cfg.graph.sessions,
1736 artifacts: &artifacts,
1737 stem: &attempt_stem,
1738 run: &talk.id,
1741 node: crate::queue::CHAT_NODE,
1742 cache_dir: cache_dir.as_deref(),
1743 attachments: &attachment_paths,
1744 writable: &consult_roots,
1745 };
1746 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1747 let advance = agent::chain_advances(&result);
1748 if n == 0 || !advance {
1749 outcome = Some(result);
1750 } else {
1751 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1753 }
1754 if !advance {
1755 break;
1756 }
1757 }
1758 if outcome.as_ref().is_some_and(agent::chain_advances) {
1759 if let Some((id, seat)) = first_try {
1762 talk.agent = id;
1763 talk.seat = seat;
1764 fell_back_from = None;
1765 }
1766 }
1767 let outcome = outcome.expect("a chain holds at least one agent");
1768 let note = |why: String| Turn {
1769 who: Who::Agent,
1770 body: format!("{MAGI_NOTE}{why}"),
1771 at: Timestamp::now(),
1772 attachments: Vec::new(),
1773 usage: None,
1774 };
1775 let (reply, failure) = match outcome {
1776 Err(e) => (
1777 note(format!("could not run agent `{}`: {e}", talk.agent)),
1778 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1779 ),
1780 Ok(out) if out.quota_exhausted() => {
1781 let reset = out
1782 .quota
1783 .as_ref()
1784 .and_then(|q| q.reset.clone())
1785 .map_or_else(String::new, |r| format!(" (resets {r})"));
1786 let why = format!(
1787 "agent `{}` is out of quota{reset}; your message is saved, so \
1788 say it again when the window reopens",
1789 talk.agent
1790 );
1791 (note(why.clone()), Some(why))
1792 }
1793 Ok(out) if out.timed_out => {
1794 let why = format!(
1795 "agent `{}` did not answer within {}s; your message is saved",
1796 talk.agent,
1797 turn_timeout(cfg).as_secs()
1798 );
1799 (note(why.clone()), Some(why))
1800 }
1801 Ok(out) if !out.usable() => {
1802 let why = format!(
1803 "agent `{}` produced no answer (exit {}); your message is saved",
1804 talk.agent,
1805 out.exit_code
1806 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1807 );
1808 (note(why.clone()), Some(why))
1809 }
1810 Ok(out) => (
1811 Turn {
1812 who: Who::Agent,
1813 body: out.text.trim().to_owned(),
1814 at: Timestamp::now(),
1815 attachments: Vec::new(),
1816 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1819 context_tokens,
1820 agent: talk.agent.clone(),
1821 model: cfg
1822 .agents
1823 .iter()
1824 .find(|a| a.id == talk.agent)
1825 .and_then(|a| a.model.clone()),
1826 }),
1827 },
1828 None,
1829 ),
1830 };
1831
1832 let _guard = store.guard()?;
1844 let Ok(fresh) = store.get(&talk.id) else {
1850 return Ok(());
1851 };
1852 talk.status = fresh.status;
1853 talk.pending = fresh.pending;
1857 talk.pending_attachments = fresh.pending_attachments;
1858 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1859 talk.turns.push(note(format!(
1862 "agent changed from {from} to {} (fallback)",
1863 talk.agent
1864 )));
1865 }
1866 if failure.is_none() {
1869 talk.persona_dirty = false;
1870 }
1871 talk.turns.push(reply);
1872 if let Err(put_err) = store.put(talk) {
1873 let lost = talk.turns.pop().expect("just pushed above");
1882 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1883 let why = match &stash {
1884 Ok(path) => format!(
1885 "agent `{}` answered, but the reply could not be saved to \
1886 this conversation ({put_err:#}); the raw text was kept at \
1887 {} - your message is saved, ask again",
1888 talk.agent,
1889 path.display()
1890 ),
1891 Err(stash_err) => format!(
1892 "agent `{}` answered, but the reply could not be saved to \
1893 this conversation ({put_err:#}), and it could not be kept \
1894 anywhere else either ({stash_err:#}); your message is \
1895 saved, ask again",
1896 talk.agent
1897 ),
1898 };
1899 talk.turns.push(note(why.clone()));
1900 return match store.put(talk) {
1907 Ok(()) => bail!("{why}"),
1908 Err(note_err) => {
1909 talk.turns.pop();
1929 Err(note_err).context(why)
1930 }
1931 };
1932 }
1933
1934 match failure {
1935 Some(why) => bail!("{why}"),
1936 None => Ok(()),
1937 }
1938}
1939
1940fn transcript(talk: &Talk, store: &Talks) -> String {
1943 let mut out = String::from(
1944 "This conversation cannot resume on the CLI's side, so here is \
1945 everything said so far; answer only the last message.\n",
1946 );
1947 for t in &talk.turns {
1948 let who = match t.who {
1949 Who::Operator => "operator",
1950 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1951 Who::Agent => "you",
1952 };
1953 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1954 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1955 }
1956 out
1957}
1958
1959fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1964 if attachments.is_empty() {
1965 return String::new();
1966 }
1967 let mut out = String::from(
1968 "\n\nThe operator attached the image(s) below to this message. Open \
1969 and look at each one before you answer.\n",
1970 );
1971 for att in attachments {
1972 if let Some(path) = store.attachment_path(talk_id, att) {
1973 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1974 }
1975 }
1976 out.push('\n');
1977 out
1978}
1979
1980pub(crate) fn turn_access(talk_allow_write: bool, consulted: bool) -> (bool, bool) {
1991 (talk_allow_write || consulted, talk_allow_write)
1992}
1993
1994pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
2010 briefing_with(repo, language, allow_write, None, None)
2011}
2012
2013pub fn briefing_with(
2018 repo: &Path,
2019 language: &str,
2020 allow_write: bool,
2021 persona: Option<&crate::persona::Persona>,
2022 operator_name: Option<&str>,
2023) -> String {
2024 let write_policy = if allow_write {
2025 "Write access is enabled for this conversation (`allow_write = \
2026 true`), so you may write files - but only a small, \
2027 already-decided edit the operator names outright in this \
2028 conversation, not an implementation. This is a permission on the \
2029 conversation as a whole, not a property of whichever repository \
2030 it happened to start in: if the operator names a different \
2031 repository for that small edit, the policy allows it there too. \
2032 Your own tool may still confine writes to the repository this \
2033 conversation started in regardless - if a write elsewhere is \
2034 refused, say so plainly rather than working around it. Once you \
2035 have made an edit, say plainly what you edited. Anything bigger, \
2036 or anything still open-ended, still goes through the queue below \
2037 rather than being done here."
2038 } else {
2039 "Do not write files. Implementing a change is not this \
2040 conversation's job; a separate, blind competition of agents does \
2041 that, and a repository this conversation has already edited would \
2042 make their diffs unjudgeable."
2043 };
2044 let mut out = format!(
2045 "You are magi's standing conversation partner for its operator, who \
2046 usually has this open on a phone. Keep replies short: no preamble, \
2047 no restating what they just said.\n\n\
2048 # Repository\n\n{repo}\n\n\
2049 You may look around: read files, run shell commands, search history, \
2050 run tests - whatever answers the question. {write_policy}\n\n\
2051 A short, command-shaped message (\"list\", \"info <id>\", \"show \
2052 3cbf\") is almost always the operator asking you to look something \
2053 up, not an instruction to file - answer it yourself with `magi \
2054 list`, `magi show <id>`, `magi task list`, or the like, the same way \
2055 you would answer any other question in this conversation.\n\n\
2056 # When the operator wants something done\n\n\
2057 Run:\n\n\
2058 magi task add --solo --repo {repo} <instruction>\n\n\
2059 and tell the operator the task id it prints, so they can follow it \
2060 from the Queue. If it refuses with a duplicate warning (the \
2061 instruction names a branch, commit or pull request that an \
2062 unfinished task, run or PR already owns), do not repeat it with \
2063 --force yourself: tell the operator what it matched and let them \
2064 decide. Write <instruction> so that an implementer who has \
2065 never seen this conversation can act on it alone - it is everything \
2066 they get. Use --solo: it runs the task through one implementer \
2067 straight into review instead of the usual multi-agent competition, \
2068 which is the right shape for a change this conversation has already \
2069 settled, rather than one still worth several independent takes.\n\n\
2070 If the operator asks for something in a different repository, \
2071 --repo does not have to be a full path: --repo owner/repo (or just \
2072 repo, when that is unambiguous) is resolved against local checkouts \
2073 the same way `magi repos` lists them. If the command fails because \
2074 nothing matches or more than one checkout shares that name, ask the \
2075 operator which repository they mean (or run `magi repos` yourself \
2076 to see the candidates) rather than guessing.\n\n\
2077 The current state of the code is whatever origin/main holds, not \
2078 whatever a working tree shows: a primary checkout often lags \
2079 upstream, sits on a detached HEAD and carries uncommitted changes. \
2080 Before answering about code, run `git fetch origin` in that \
2081 repository if it is cheap, then read through \
2082 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
2083 If the working tree differs, say so; if the fetch fails, say that \
2084 too, so the operator knows the answer may be stale.\n\n\
2085 If the operator attached an image (a screenshot, say) that the task \
2086 is about, pass it with `--attach <path>`, using the absolute path \
2087 the turn's attachment note gives; repeat the flag for several. \
2088 `magi task add --solo --attach <path> <instruction>` copies the \
2089 file into the task, so the implementer receives it. Do not paste the \
2090 path into <instruction> instead: deleting this conversation deletes \
2091 its attachments, and then that path reaches no one.\n",
2092 repo = repo.display(),
2093 );
2094 out.push_str(&language_note(language));
2095 match (persona, operator_name) {
2096 (Some(p), n) => out.push_str(&crate::persona::section_for(p, n)),
2097 (None, Some(n)) => {
2098 out.push_str("\n# Addressing the operator\n");
2099 out.push_str(&crate::persona::addressing(n));
2100 }
2101 (None, None) => {}
2102 }
2103 out
2104}
2105
2106fn language_note(language: &str) -> String {
2109 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
2110 String::new()
2111 } else {
2112 format!("\nHold this conversation in {language}.\n")
2113 }
2114}
2115
2116pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
2123 let mut tasks: Vec<Task> = queue
2124 .list()
2125 .into_iter()
2126 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
2127 .collect();
2128 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
2129 tasks
2130}
2131
2132fn read_path(path: &Path) -> Result<Talk> {
2133 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2134 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
2135}
2136
2137const PUT_RETRIES: u32 = 5;
2140
2141fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
2153 let mut last_err = None;
2154 for attempt in 0..PUT_RETRIES {
2155 if attempt > 0 {
2156 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
2157 }
2158 match try_write_atomic(tmp, path, body) {
2159 Ok(()) => return Ok(()),
2160 Err(e) => last_err = Some(e),
2161 }
2162 }
2163 Err(last_err.expect("the loop above always runs at least once"))
2164}
2165
2166fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
2167 #[cfg(test)]
2168 if failpoint::take_forced_put_failure() {
2169 bail!("simulated write failure (test)");
2170 }
2171 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
2172 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
2173 Ok(())
2174}
2175
2176fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
2181 let dir = store.artifacts_of(id);
2182 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
2183 let path = dir.join(format!("{stem}-lost.txt"));
2184 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
2185 Ok(path)
2186}
2187
2188#[cfg(test)]
2195mod failpoint {
2196 use std::cell::Cell;
2197
2198 thread_local! {
2199 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
2200 static FORCE_NO_LINK: Cell<bool> = const { Cell::new(false) };
2201 }
2202
2203 pub(super) fn force_no_link(on: bool) {
2206 FORCE_NO_LINK.with(|c| c.set(on));
2207 }
2208
2209 pub(super) fn no_link_forced() -> bool {
2210 FORCE_NO_LINK.with(Cell::get)
2211 }
2212
2213 pub(super) fn force_put_failures(count: u32) {
2216 FORCE_PUT_FAILURES.with(|c| c.set(count));
2217 }
2218
2219 pub(super) fn take_forced_put_failure() -> bool {
2222 FORCE_PUT_FAILURES.with(|c| {
2223 let n = c.get();
2224 if n == 0 {
2225 false
2226 } else {
2227 c.set(n - 1);
2228 true
2229 }
2230 })
2231 }
2232}
2233
2234fn short(id: &str) -> &str {
2235 id.split('-').next_back().unwrap_or(id)
2236}
2237
2238fn new_id() -> String {
2239 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2240 let seed = crate::rng::entropy();
2241 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2242}
2243
2244fn attachment_ext(mime: &str) -> Option<&'static str> {
2249 match mime {
2250 "image/png" => Some("png"),
2251 "image/jpeg" => Some("jpg"),
2252 "image/gif" => Some("gif"),
2253 "image/webp" => Some("webp"),
2254 _ => None,
2255 }
2256}
2257
2258pub fn valid_attachment_id(id: &str) -> bool {
2263 id.len() == 32
2264 && id
2265 .bytes()
2266 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
2267}
2268
2269fn new_attachment_id() -> String {
2273 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
2274 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
2275}
2276
2277#[cfg(test)]
2278mod tests {
2279 #[test]
2280 fn turn_access_lifts_the_sandbox_only_for_the_talk_opt_in() {
2281 use super::turn_access;
2282 assert_eq!(turn_access(false, false), (false, false));
2283 assert_eq!(turn_access(true, false), (true, true));
2284 assert_eq!(turn_access(false, true), (true, false));
2286 assert_eq!(turn_access(true, true), (true, true));
2287 }
2288
2289 #[test]
2290 fn the_briefing_points_at_origin_main_not_the_working_tree() {
2291 let b = briefing(Path::new("/r"), "en", false);
2292 assert!(b.contains("origin/main"));
2293 assert!(b.contains("git show origin/main:"));
2294 }
2295 use std::collections::BTreeMap;
2296
2297 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
2298 use crate::queue::{Queue, Source, Task};
2299
2300 use super::*;
2301
2302 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
2303 AgentSpec {
2304 id: id.to_owned(),
2305 kind: AgentKind::Command,
2306 model: model.map(str::to_owned),
2307 command: Vec::new(),
2308 extra_args: Vec::new(),
2309 env: BTreeMap::new(),
2310 prompt_delivery: None,
2311 }
2312 }
2313
2314 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
2315 Talk {
2316 schema: SCHEMA,
2317 id: "20260904-014455-ab12".to_owned(),
2318 repo: PathBuf::from("."),
2319 agent: agent.to_owned(),
2320 status: TalkStatus::Open,
2321 turns,
2322 pending: String::new(),
2323 pending_attachments: Vec::new(),
2324 fallback: false,
2325 persona: String::new(),
2326 persona_dirty: false,
2327 created_at: Timestamp::now(),
2328 updated_at: Timestamp::now(),
2329 seat: SeatState::new(SEAT, agent, 1),
2330 }
2331 }
2332
2333 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
2334 Turn {
2335 who: Who::Agent,
2336 body: body.to_owned(),
2337 at: Timestamp::now(),
2338 attachments: Vec::new(),
2339 usage: usage.map(|(t, a, m)| TurnUsage {
2340 context_tokens: t,
2341 agent: a.to_owned(),
2342 model: m.map(str::to_owned),
2343 }),
2344 }
2345 }
2346
2347 fn ctx_config(windows: &[(&str, u64)]) -> Config {
2348 Config {
2349 agents: vec![
2350 ctx_agent("small", Some("small-model")),
2351 ctx_agent("big", Some("big-model")),
2352 ctx_agent("plain", None),
2353 ],
2354 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
2355 ..Config::default()
2356 }
2357 }
2358
2359 #[test]
2360 fn context_usage_computes_percent_and_warns_at_eighty() {
2361 let cfg = ctx_config(&[("small-model", 1000)]);
2362 let at = |tokens| {
2363 let t = ctx_talk(
2364 "small",
2365 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
2366 );
2367 context_usage(&t, Some(&cfg))
2368 };
2369 let u = at(799);
2370 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
2371 let u = at(800);
2372 assert_eq!((u.percent, u.warn), (Some(80), true));
2373 let u = at(1500);
2374 assert_eq!((u.percent, u.warn), (Some(150), true));
2375 assert!(!u.since_switch);
2376 }
2377
2378 #[test]
2379 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
2380 let cfg = ctx_config(&[("small-model", 1000)]);
2381 let t = ctx_talk(
2382 "small",
2383 vec![
2384 reply("old", Some((900, "small", Some("small-model")))),
2385 reply("new", None),
2386 ],
2387 );
2388 let u = context_usage(&t, Some(&cfg));
2389 assert!(u.estimated);
2391 assert_ne!(u.tokens, Some(900));
2392 assert!(u.tokens.is_some());
2393 let t = ctx_talk(
2395 "small",
2396 vec![
2397 reply("old", Some((900, "small", Some("small-model")))),
2398 reply("magi: could not run agent", None),
2399 ],
2400 );
2401 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
2402 assert_eq!(
2403 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
2404 None
2405 );
2406 }
2407
2408 #[test]
2409 fn estimate_counts_chars_both_sides_and_standing_prompt() {
2410 let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
2411 assert_eq!(estimate_context_tokens(&t, 0), Some(2)); let op = Turn {
2413 who: Who::Operator,
2414 ..reply("abcdefg", None)
2415 };
2416 t.turns.push(op);
2417 assert_eq!(estimate_context_tokens(&t, 0), Some(4));
2418 assert!(
2419 estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
2420 );
2421 let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
2423 assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
2424 let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
2426 assert_eq!(estimate_context_tokens(¬e, 1000), None);
2427 assert_eq!(
2428 estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
2429 None
2430 );
2431 }
2432
2433 #[test]
2434 fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
2435 let cfg = ctx_config(&[("small-model", 1000)]);
2436 let t = ctx_talk(
2437 "small",
2438 vec![reply(
2439 &"x".repeat(5000),
2440 Some((10, "small", Some("small-model"))),
2441 )],
2442 );
2443 let u = context_usage(&t, Some(&cfg));
2444 assert_eq!((u.tokens, u.estimated), (Some(10), false));
2445 let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
2446 let u = context_usage(&t, Some(&cfg));
2447 assert!(u.estimated && !u.since_switch);
2448 assert_eq!(u.window, Some(1000));
2449 assert!(u.warn && u.percent.unwrap() >= 80);
2450 let t = ctx_talk("small", vec![reply("hi", None)]);
2451 let u = context_usage(&t, Some(&cfg));
2452 assert!(u.estimated && u.percent.is_some());
2453 }
2454
2455 #[test]
2456 fn context_usage_without_a_window_shows_tokens_only() {
2457 let cfg = ctx_config(&[]);
2458 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
2460 let u = context_usage(&t, Some(&cfg));
2461 assert_eq!(
2462 (u.tokens, u.window, u.percent, u.warn),
2463 (Some(5000), None, None, false)
2464 );
2465 let t = ctx_talk(
2466 "small",
2467 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
2468 );
2469 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
2470 assert_eq!(context_usage(&t, None).window, None);
2472 }
2473
2474 #[test]
2475 fn context_usage_switching_model_changes_the_denominator() {
2476 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
2477 let used = reply("hi", Some((900, "small", Some("small-model"))));
2478 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
2479 assert_eq!(
2480 (before.percent, before.warn, before.since_switch),
2481 (Some(90), true, false)
2482 );
2483 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
2486 assert_eq!(after.window, Some(10_000));
2487 assert_eq!(
2488 (after.percent, after.warn, after.since_switch),
2489 (Some(9), false, true)
2490 );
2491 assert_eq!(after.model.as_deref(), Some("big-model"));
2492 }
2493
2494 #[test]
2495 fn a_turn_recorded_before_usage_existed_still_reads() {
2496 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
2497 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
2498 assert!(turn.usage.is_none());
2499 let json = serde_json::to_string(&turn).expect("serialize");
2500 assert!(
2501 !json.contains("usage"),
2502 "absent usage is not written: {json}"
2503 );
2504 }
2505
2506 fn store() -> (tempfile::TempDir, Talks) {
2508 let tmp = tempfile::tempdir().expect("tempdir");
2509 let talks = Talks::at(tmp.path().join("talks"));
2510 (tmp, talks)
2511 }
2512
2513 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
2517 let path = dir.join("mock-talk-agent.sh");
2518 std::fs::write(&path, script).expect("write mock");
2519 AgentSpec {
2520 id: "mock".to_owned(),
2521 kind: AgentKind::Command,
2522 model: None,
2523 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2524 extra_args: Vec::new(),
2525 env,
2526 prompt_delivery: None,
2527 }
2528 }
2529
2530 fn config(spec: AgentSpec) -> Config {
2531 Config {
2532 agents: vec![spec],
2533 graph: Graph {
2534 language: "en".to_owned(),
2535 ..Graph::default()
2536 },
2537 ..Config::default()
2538 }
2539 }
2540
2541 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
2543
2544 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
2546
2547 const ECHO: &str = "#!/bin/sh\ncat\n";
2550
2551 fn env(reply: &str) -> BTreeMap<String, String> {
2552 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
2553 }
2554
2555 #[test]
2556 fn the_frozen_json_field_names_round_trip_through_disk() {
2557 let (tmp, talks) = store();
2558 let mut talk = Talk {
2559 schema: SCHEMA,
2560 id: "20260904-014455-ab12".to_owned(),
2561 repo: tmp.path().to_owned(),
2562 agent: "sonnet".to_owned(),
2563 status: TalkStatus::Open,
2564 turns: Vec::new(),
2565 pending: String::new(),
2566 pending_attachments: Vec::new(),
2567 fallback: false,
2568 persona: String::new(),
2569 persona_dirty: false,
2570 created_at: Timestamp::now(),
2571 updated_at: Timestamp::now(),
2572 seat: SeatState::new(SEAT, "sonnet", 7),
2573 };
2574 talks.put(&mut talk).expect("put");
2575
2576 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
2577 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
2578 for field in [
2579 "schema",
2580 "id",
2581 "repo",
2582 "agent",
2583 "status",
2584 "turns",
2585 "created_at",
2586 "updated_at",
2587 ] {
2588 assert!(v.get(field).is_some(), "missing field `{field}`");
2589 }
2590 assert_eq!(v["schema"], 1);
2591 assert_eq!(v["status"], "open");
2592
2593 let back = talks.get(&talk.id).expect("get");
2594 assert_eq!(back.id, talk.id);
2595 assert_eq!(back.status, TalkStatus::Open);
2596 }
2597
2598 #[test]
2599 fn opening_a_talk_takes_no_agent_turn() {
2600 let (tmp, talks) = store();
2601 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2605 let cfg = config(spec);
2606
2607 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2608 assert_eq!(talk.status, TalkStatus::Open);
2609 assert!(talk.turns.is_empty(), "nothing has been said yet");
2610
2611 let on_disk = talks.get(&talk.id).expect("get");
2612 assert_eq!(on_disk.turns.len(), 0);
2613 }
2614
2615 #[test]
2623 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
2624 let (tmp, talks) = store();
2625 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2626 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2627 chatter_spec.id = "chatter-mock".to_owned();
2628
2629 let mut cfg = Config {
2630 agents: vec![first_spec.clone(), chatter_spec.clone()],
2631 graph: Graph {
2632 language: "en".to_owned(),
2633 ..Graph::default()
2634 },
2635 ..Config::default()
2636 };
2637 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
2638
2639 let talk =
2640 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
2641 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
2642
2643 cfg.roles.chatter = None;
2644 let fallback =
2645 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
2646 assert_eq!(
2647 fallback.agent, first_spec.id,
2648 "unset chatter must fall back to agent::pick's own default order"
2649 );
2650 }
2651
2652 #[test]
2655 fn a_talk_recorded_without_attachments_still_reads() {
2656 let (tmp, talks) = store();
2657 let path = talks.path_of("20260904-014455-ab12");
2658 std::fs::create_dir_all(talks.root()).expect("talks dir");
2659 std::fs::write(
2660 &path,
2661 serde_json::json!({
2662 "schema": 1,
2663 "id": "20260904-014455-ab12",
2664 "repo": tmp.path(),
2665 "agent": "sonnet",
2666 "status": "open",
2667 "turns": [
2668 { "who": "operator", "body": "still there?",
2669 "at": Timestamp::now().to_string() },
2670 ],
2671 "created_at": Timestamp::now().to_string(),
2672 "updated_at": Timestamp::now().to_string(),
2673 "seat": SeatState::new(SEAT, "sonnet", 7),
2674 })
2675 .to_string(),
2676 )
2677 .expect("write pre-attachments talk");
2678
2679 let talk = talks.get("20260904-014455-ab12").expect("must still read");
2680 assert!(talk.turns[0].attachments.is_empty());
2681 }
2682
2683 fn lease_store() -> (tempfile::TempDir, Talks, Talks) {
2684 let tmp = tempfile::TempDir::new().expect("tmp");
2685 let root = tmp.path().join("talks");
2686 (tmp, Talks::at(root.clone()), Talks::at(root))
2687 }
2688
2689 #[test]
2690 fn two_starters_on_one_talk_one_wins_and_the_other_is_refused() {
2691 let (_tmp, a, b) = lease_store();
2692 let won = a.claim_turn("t1").expect("claim").expect("first wins");
2693 assert!(
2694 b.claim_turn("t1").expect("claim").is_none(),
2695 "second is refused"
2696 );
2697 assert!(b.turn_held("t1"));
2698 assert!(
2699 b.claim_turn("t2").expect("claim").is_some(),
2700 "other talks are free"
2701 );
2702 drop(won);
2703 }
2704
2705 #[test]
2706 fn a_stale_lease_is_taken_over_and_the_old_guard_cannot_release_it() {
2707 let (_tmp, a, b) = lease_store();
2708 let old = a.claim_turn("t1").expect("claim").expect("held");
2709 let later = Timestamp::now()
2710 .checked_add(jiff::SignedDuration::from_secs(
2711 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2712 ))
2713 .expect("later");
2714 let new = b
2715 .claim_turn_at("t1", later)
2716 .expect("claim")
2717 .expect("a stale lease is taken over");
2718 drop(old);
2719 assert!(a.turn_held("t1"), "the old guard left the new lease alone");
2720 assert!(new.beat().expect("beat"), "the new owner still beats");
2721 drop(new);
2722 assert!(!a.turn_held("t1"));
2723 }
2724
2725 #[tokio::test]
2726 async fn a_turn_whose_lease_was_taken_over_is_stopped() {
2727 let (_tmp, a, b) = lease_store();
2728 let old = a.claim_turn("t1").expect("claim").expect("held");
2729 let later = Timestamp::now()
2730 .checked_add(jiff::SignedDuration::from_secs(
2731 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2732 ))
2733 .expect("later");
2734 let _new = b
2735 .claim_turn_at("t1", later)
2736 .expect("claim")
2737 .expect("taken over");
2738 let out = old
2739 .beating_every(Duration::from_millis(10), std::future::pending::<()>())
2740 .await;
2741 assert!(out.is_err(), "the displaced turn must stop, not run on");
2742 }
2743
2744 #[tokio::test]
2745 async fn a_turn_that_finishes_is_returned_and_keeps_its_lease_beating() {
2746 let (_tmp, a, _b) = lease_store();
2747 let lease = a.claim_turn("t1").expect("claim").expect("held");
2748 let out = lease
2749 .beating_every(Duration::from_millis(5), async {
2750 tokio::time::sleep(Duration::from_millis(40)).await;
2751 7
2752 })
2753 .await
2754 .expect("still ours");
2755 assert_eq!(out, 7);
2756 assert!(a.turn_held("t1"));
2757 }
2758
2759 #[test]
2760 fn an_unreadable_lease_counts_as_stale() {
2761 let (_tmp, a, b) = lease_store();
2762 std::fs::create_dir_all(a.root()).expect("dir");
2763 std::fs::write(a.turn_path("t1"), "not json").expect("write");
2764 assert!(!a.turn_held("t1"));
2765 assert!(b.claim_turn("t1").expect("claim").is_some());
2766 }
2767
2768 #[test]
2769 fn without_hard_links_a_claim_is_still_exclusive_and_a_young_placeholder_blocks() {
2770 let (_tmp, a, b) = lease_store();
2771 failpoint::force_no_link(true);
2772 let held = a.claim_turn("t1").expect("claim").expect("first wins");
2773 assert!(b.claim_turn("t1").expect("claim").is_none());
2774 assert!(a.turn_held("t1"));
2775 drop(held);
2776 assert!(!a.turn_held("t1"));
2777 let path = a.turn_path("t1");
2779 assert!(create_exclusive(&path, "").expect("placeholder"));
2780 assert!(b.claim_turn("t1").expect("claim").is_none(), "young: held");
2781 age_file(&path);
2782 assert!(b.claim_turn("t1").expect("claim").is_some(), "old: stale");
2783 let lease = a.turn_path("t2");
2785 let lock = TurnLock::take(&lease).expect("take").expect("lock");
2786 assert!(TurnLock::take(&lease).expect("take").is_none());
2787 drop(lock);
2788 assert!(TurnLock::take(&lease).expect("take").is_some());
2789 failpoint::force_no_link(false);
2790 }
2791
2792 #[test]
2793 fn remove_if_carries_removes_only_the_expected_content() {
2794 let (_tmp, a, _b) = lease_store();
2795 std::fs::create_dir_all(a.root()).expect("dir");
2796 let p = a.root().join("x.turn.lock");
2797 std::fs::write(&p, "mine").expect("write");
2798 assert!(!remove_if_carries(&p, "other"));
2799 assert_eq!(std::fs::read_to_string(&p).expect("kept"), "mine");
2800 assert!(remove_if_carries(&p, "mine"));
2801 assert!(!p.exists());
2802 assert!(!remove_if_carries(&p, "mine"), "absent is not a removal");
2803 }
2804
2805 #[test]
2806 fn a_dropped_lock_does_not_remove_a_lock_taken_over_since() {
2807 let (_tmp, a, _b) = lease_store();
2808 std::fs::create_dir_all(a.root()).expect("dir");
2809 let lease = a.turn_path("t1");
2810 let lock = TurnLock::take(&lease).expect("take").expect("lock");
2811 let path = lock.path.clone();
2812 std::fs::write(&path, "someone-else").expect("replace");
2813 drop(lock);
2814 assert_eq!(
2815 std::fs::read_to_string(&path).expect("kept"),
2816 "someone-else"
2817 );
2818 }
2819
2820 #[test]
2821 fn concurrent_takeovers_of_a_stale_lease_have_one_winner() {
2822 let (_tmp, a, _b) = lease_store();
2823 drop(a.claim_turn("t1").expect("claim").expect("held"));
2824 std::fs::write(
2825 a.turn_path("t1"),
2826 serde_json::to_string(&TurnRecord {
2827 token: "gone".into(),
2828 pid: 1,
2829 beat_at: Timestamp::from_second(1).expect("ts"),
2830 })
2831 .expect("json"),
2832 )
2833 .expect("write");
2834 let wins: Vec<_> = std::thread::scope(|sc| {
2835 let hs: Vec<_> = (0..8)
2836 .map(|_| {
2837 let s = a.clone();
2838 sc.spawn(move || s.claim_turn("t1").expect("claim"))
2839 })
2840 .collect();
2841 hs.into_iter().map(|h| h.join().expect("join")).collect()
2842 });
2843 assert_eq!(wins.iter().filter(|w| w.is_some()).count(), 1);
2844 }
2845
2846 fn age_file(path: &Path) {
2847 let f = std::fs::OpenOptions::new()
2848 .write(true)
2849 .open(path)
2850 .expect("open");
2851 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
2852 .expect("age");
2853 }
2854
2855 #[test]
2856 fn a_late_taker_of_a_broken_lock_cannot_disturb_its_replacement() {
2857 let (_tmp, a, _b) = lease_store();
2858 let lease = a.turn_path("t1");
2859 let lock = lease.with_extension("turn.lock");
2860 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2861 std::fs::write(&lock, "t1-dead").expect("dead lock");
2862 age_file(&lock);
2863 let b = TurnLock::take(&lease).expect("take").expect("b wins");
2865 let fresh = std::fs::read_to_string(&lock).expect("read");
2866 assert_eq!(fresh, b.token);
2867 let ticket = std::fs::read_dir(lock.parent().expect("dir"))
2870 .expect("dir")
2871 .flatten()
2872 .map(|e| e.path())
2873 .find(|p| p.to_string_lossy().ends_with(".break.0"))
2874 .expect("ticket");
2875 assert!(!create_exclusive(&ticket, "").expect("ticket"));
2876 assert!(TurnLock::take(&lease).expect("take").is_none());
2877 assert_eq!(std::fs::read_to_string(&lock).expect("read"), fresh);
2878 age_file(&lock);
2880 age_file(&ticket);
2881 std::mem::forget(b);
2882 let c = TurnLock::take(&lease)
2883 .expect("take")
2884 .expect("next generation");
2885 assert_ne!(c.token, fresh);
2886 }
2887
2888 #[test]
2889 fn a_live_ticket_blocks_and_a_stale_one_hands_over_to_the_next_generation() {
2890 let (_tmp, a, _b) = lease_store();
2891 let lease = a.turn_path("t1");
2892 let lock = lease.with_extension("turn.lock");
2893 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2894 std::fs::write(&lock, "t1-dead").expect("dead lock");
2895 age_file(&lock);
2896 let t0 = lock.with_extension("lock.t1-dead.break.0");
2897 assert!(create_exclusive(&t0, "").expect("ticket"));
2898 assert!(TurnLock::take(&lease).expect("take").is_none());
2900 assert_eq!(std::fs::read_to_string(&lock).expect("read"), "t1-dead");
2901 age_file(&t0);
2903 let c = TurnLock::take(&lease).expect("take").expect("generation 1");
2904 assert_eq!(std::fs::read_to_string(&lock).expect("read"), c.token);
2905 assert!(lock.with_extension("lock.t1-dead.break.1").exists());
2906 assert!(TurnLock::take(&lease).expect("take").is_none());
2907 }
2908
2909 #[test]
2910 fn exhausted_ticket_generations_recover_once_the_sweep_ages_them_out() {
2911 let (_tmp, a, _b) = lease_store();
2912 let lease = a.turn_path("t1");
2913 let lock = lease.with_extension("turn.lock");
2914 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2915 std::fs::write(&lock, "t1-dead").expect("dead lock");
2916 age_file(&lock);
2917 let tickets: Vec<_> = (0..TICKET_GENERATIONS)
2918 .map(|n| lock.with_extension(format!("lock.t1-dead.break.{n}")))
2919 .collect();
2920 for t in &tickets {
2921 assert!(create_exclusive(t, "").expect("ticket"));
2922 age_file(t);
2923 }
2924 assert!(TurnLock::take(&lease).expect("take").is_none());
2926 assert!(tickets.iter().all(|t| t.exists()));
2927 for t in &tickets {
2929 let f = std::fs::OpenOptions::new()
2930 .write(true)
2931 .open(t)
2932 .expect("open");
2933 f.set_modified(std::time::SystemTime::now() - TICKET_SWEEP_AGE * 2)
2934 .expect("age");
2935 }
2936 assert!(TurnLock::take(&lease).expect("take").is_none());
2937 assert!(TurnLock::take(&lease).expect("take").is_some());
2938 }
2939
2940 #[test]
2941 fn an_aged_empty_lock_is_broken() {
2942 let (_tmp, a, _b) = lease_store();
2943 let lease = a.turn_path("t1");
2944 let lock = lease.with_extension("turn.lock");
2945 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2946 std::fs::write(&lock, "").expect("empty lock");
2947 age_file(&lock);
2948 assert!(TurnLock::take(&lease).expect("take").is_some());
2949 }
2950
2951 #[test]
2952 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
2953 let (tmp, talks) = store();
2954 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2955 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2956
2957 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
2958 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
2959 let saved = talks.get(&talk.id).expect("reload queued talk");
2960 assert_eq!(saved.pending, "first\n\nsecond");
2961 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
2962
2963 let drained = drain(&mut talk, &talks).expect("drain");
2964 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
2965 let saved = talks.get(&talk.id).expect("reload drained talk");
2966 assert!(saved.pending.is_empty());
2967 assert_eq!(saved.turns.len(), 1);
2968 assert_eq!(saved.turns[0].body, "first\n\nsecond");
2969 }
2970
2971 #[test]
2972 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
2973 let (tmp, talks) = store();
2974 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2975 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2976 let attachment = Attachment {
2977 id: "a".repeat(32),
2978 name: "shot.png".to_owned(),
2979 mime: "image/png".to_owned(),
2980 bytes: 3,
2981 };
2982
2983 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
2984 assert!(
2985 edit_pending_text(
2986 &mut talk,
2987 &talks,
2988 "corrected",
2989 "first",
2990 std::slice::from_ref(&attachment.id),
2991 )
2992 .expect("edit")
2993 );
2994 let saved = talks.get(&talk.id).expect("reload edited draft");
2995 assert_eq!(saved.pending, "corrected");
2996 assert_eq!(saved.pending_attachments, vec![attachment]);
2997
2998 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
2999 assert!(
3000 !edit_pending_text(
3001 &mut talk,
3002 &talks,
3003 "stale edit",
3004 "corrected",
3005 &["a".repeat(32)],
3006 )
3007 .expect("stale edit is a conflict")
3008 );
3009 assert_eq!(
3010 talks.get(&talk.id).expect("reload after conflict").pending,
3011 "corrected\n\nlater"
3012 );
3013 assert!(
3014 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
3015 .expect("stale clear is a conflict")
3016 );
3017 assert_eq!(
3018 talks
3019 .get(&talk.id)
3020 .expect("reload after stale clear")
3021 .pending,
3022 "corrected\n\nlater"
3023 );
3024 }
3025
3026 #[tokio::test]
3027 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
3028 let (tmp, talks) = store();
3029 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
3030 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
3031 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3032 let id = running.id.clone();
3033 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
3034
3035 let response_talks = talks.clone();
3036 let response_cfg = cfg.clone();
3037 let reply = tokio::spawn(async move {
3038 respond(
3039 &response_talks.claim_turn(&running.id).unwrap().unwrap(),
3040 &mut running,
3041 &response_talks,
3042 &response_cfg,
3043 &first,
3044 )
3045 .await
3046 });
3047 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
3048
3049 let mut queued = talks.get(&id).expect("queued handle");
3050 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
3051 reply.await.expect("join").expect("reply");
3052
3053 let saved = talks.get(&id).expect("reload");
3054 assert_eq!(saved.pending, "next");
3055 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
3056 }
3057
3058 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
3061 let calls = dir.join(format!("{id}.calls"));
3062 let script = format!(
3063 "#!/bin/sh\necho x >> '{}'\n{body}\n",
3064 calls.to_string_lossy()
3065 );
3066 let path = dir.join(format!("mock-{id}.sh"));
3067 std::fs::write(&path, script).expect("write mock");
3068 AgentSpec {
3069 id: id.to_owned(),
3070 kind: AgentKind::Command,
3071 model: None,
3072 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
3073 extra_args: Vec::new(),
3074 env: BTreeMap::new(),
3075 prompt_delivery: None,
3076 }
3077 }
3078
3079 fn calls(dir: &Path, id: &str) -> usize {
3080 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
3081 }
3082
3083 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
3084 let mut cfg = config(specs[0].clone());
3085 cfg.agents = specs;
3086 cfg.roles.chatter = Some(AgentChoice::Chain(
3087 ids.iter().map(|s| (*s).to_owned()).collect(),
3088 ));
3089 cfg
3090 }
3091
3092 #[tokio::test]
3093 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
3094 let (tmp, talks) = store();
3095 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3096 let b = counting_agent(tmp.path(), "b", "cat");
3097 let cfg = chain_config(vec![a, b], &["a", "b"]);
3098 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3099 assert_eq!(talk.agent, "a");
3100
3101 say(
3102 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3103 &mut talk,
3104 &talks,
3105 &cfg,
3106 "hello there",
3107 Vec::new(),
3108 )
3109 .await
3110 .expect("turn");
3111 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
3112 assert_eq!(calls(tmp.path(), "b"), 1);
3113 assert_eq!(talk.agent, "b", "the switch persists");
3114 assert!(talks.get(&talk.id).unwrap().agent == "b");
3115 let reply = talk.turns.last().unwrap();
3116 assert!(reply.body.contains("hello there"));
3117 assert!(
3118 reply.body.contains("magi task add --solo"),
3119 "a fresh seat gets the full briefing"
3120 );
3121 assert!(
3122 talk.turns
3123 .iter()
3124 .any(|t| t.body.contains("agent changed from a to b")),
3125 "the switch is noted"
3126 );
3127 }
3128
3129 #[tokio::test]
3130 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
3131 let (tmp, talks) = store();
3132 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3133 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
3134 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
3135 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3136
3137 let err = say(
3138 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3139 &mut talk,
3140 &talks,
3141 &cfg,
3142 "hi",
3143 Vec::new(),
3144 )
3145 .await
3146 .expect_err("every agent failed");
3147 assert!(err.to_string().contains("`a`"), "{err:#}");
3148 assert_eq!(calls(tmp.path(), "a"), 1);
3149 assert_eq!(calls(tmp.path(), "b"), 1);
3150 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
3151 }
3152
3153 #[test]
3154 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
3155 let (tmp, talks) = store();
3156 let b = counting_agent(tmp.path(), "b", "cat");
3157 let cfg = chain_config(vec![b], &["ghost", "b"]);
3158 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3159 assert_eq!(talk.agent, "b");
3160 }
3161
3162 #[tokio::test]
3163 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
3164 let (tmp, talks) = store();
3165 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3166 let b = counting_agent(tmp.path(), "b", "cat");
3167 let cfg = chain_config(vec![a, b], &["a", "b"]);
3168 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
3169 say(
3170 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3171 &mut talk,
3172 &talks,
3173 &cfg,
3174 "hi",
3175 Vec::new(),
3176 )
3177 .await
3178 .expect_err("a alone, and it fails");
3179 assert_eq!(calls(tmp.path(), "b"), 0);
3180 assert_eq!(talk.agent, "a");
3181 }
3182
3183 #[tokio::test]
3184 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
3185 let (tmp, talks) = store();
3186 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3187 let b = counting_agent(tmp.path(), "b", "cat");
3188 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
3189 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
3190 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
3191 say(
3192 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3193 &mut talk,
3194 &talks,
3195 &cfg,
3196 "hi",
3197 Vec::new(),
3198 )
3199 .await
3200 .expect_err("c alone, and it fails");
3201 assert_eq!(calls(tmp.path(), "b"), 0);
3202 }
3203
3204 #[test]
3205 fn briefing_names_the_operator_with_or_without_a_persona() {
3206 let plain = briefing(Path::new("/repo"), "en", false);
3207 let named = briefing_with(Path::new("/repo"), "en", false, None, Some("Commander"));
3208 assert!(named.starts_with(&plain));
3209 assert!(named.contains("# Addressing the operator"));
3210 assert!(named.contains("\"Commander\""));
3211 let rei = crate::persona::builtin_catalog()
3212 .into_iter()
3213 .find(|p| p.id == "rei")
3214 .expect("rei");
3215 let with = briefing_with(
3216 Path::new("/repo"),
3217 "en",
3218 false,
3219 Some(&rei),
3220 Some("Commander"),
3221 );
3222 assert!(with.contains("\"Commander\""));
3223 assert!(!with.contains("# Addressing the operator"));
3224 }
3225
3226 #[test]
3227 fn briefing_carries_a_persona_section_only_when_one_is_chosen() {
3228 let plain = briefing(Path::new("/repo"), "en", false);
3229 assert_eq!(
3230 plain,
3231 briefing_with(Path::new("/repo"), "en", false, None, None)
3232 );
3233 assert!(!plain.contains("Persona"));
3234 let rei = crate::persona::builtin_catalog()
3235 .into_iter()
3236 .find(|p| p.id == "rei")
3237 .expect("rei");
3238 let with = briefing_with(Path::new("/repo"), "en", false, Some(&rei), None);
3239 assert!(with.starts_with(&plain), "the plain briefing is untouched");
3240 assert!(with.contains("# Persona (tone only)"));
3241 assert!(with.contains("TONE ONLY"));
3242 assert!(with.contains("task ids"));
3243 assert!(with.contains("write policy"));
3244 assert!(with.contains("`magi task add`"));
3245 assert!(with.contains("Rei Ayanami"));
3246 }
3247
3248 #[test]
3249 fn a_talk_written_before_personas_still_loads() {
3250 let (tmp, talks) = store();
3251 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3252 let cfg = config(spec);
3253 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3254 let path = talks.root.join(format!("{}.json", talk.id));
3255 let mut v: serde_json::Value =
3256 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3257 v.as_object_mut().unwrap().remove("persona");
3258 v.as_object_mut().unwrap().remove("persona_dirty");
3259 std::fs::write(&path, v.to_string()).unwrap();
3260 let loaded = talks.get(&talk.id).expect("old record loads");
3261 assert_eq!(loaded.persona, "");
3262 assert!(!loaded.persona_dirty);
3263 }
3264
3265 #[tokio::test]
3266 async fn a_persona_switch_notes_marks_dirty_and_updates_the_next_turn_once() {
3267 let (tmp, talks) = store();
3268 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3269 let cfg = config(spec);
3270 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3271 let lease = talks.claim_turn(&talk.id).unwrap().unwrap();
3272 say(&lease, &mut talk, &talks, &cfg, "hello", Vec::new())
3273 .await
3274 .expect("first turn");
3275 assert!(!talk.turns[1].body.contains("Persona"));
3276
3277 assert!(switch_persona(&mut talk, &talks, "misato").expect("switch"));
3278 assert!(!switch_persona(&mut talk, &talks, "misato").expect("same"));
3279 assert!(talk.persona_dirty);
3280 assert!(talk.turns.last().unwrap().body.contains("persona changed"));
3281
3282 say(&lease, &mut talk, &talks, &cfg, "next", Vec::new())
3283 .await
3284 .expect("turn");
3285 let prompt = &talk.turns.last().unwrap().body;
3286 assert!(prompt.contains("# Persona update"), "{prompt}");
3287 assert!(prompt.contains("Misato Katsuragi"));
3288 assert!(!talk.persona_dirty, "cleared after a successful turn");
3289
3290 say(&lease, &mut talk, &talks, &cfg, "again", Vec::new())
3291 .await
3292 .expect("turn");
3293 assert!(!talk.turns.last().unwrap().body.contains("# Persona update"));
3294
3295 assert!(switch_persona(&mut talk, &talks, "default").expect("back"));
3296 assert_eq!(talk.persona, "");
3297 say(&lease, &mut talk, &talks, &cfg, "plain", Vec::new())
3298 .await
3299 .expect("turn");
3300 assert!(
3301 talk.turns
3302 .last()
3303 .unwrap()
3304 .body
3305 .contains("turned the persona off")
3306 );
3307
3308 assert!(switch_persona(&mut talk, &talks, "rei").expect("rei"));
3310 let mut no_sessions = cfg.clone();
3311 no_sessions.graph.sessions = false;
3312 say(&lease, &mut talk, &talks, &no_sessions, "one", Vec::new())
3313 .await
3314 .expect("turn");
3315 say(&lease, &mut talk, &talks, &no_sessions, "two", Vec::new())
3316 .await
3317 .expect("turn");
3318 let last = &talk.turns.last().unwrap().body;
3319 assert!(!talk.persona_dirty);
3320 assert!(last.contains("# Persona (tone only)"), "{last}");
3321 assert!(last.contains("Rei Ayanami"));
3322 }
3323
3324 #[tokio::test]
3325 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
3326 let (tmp, talks) = store();
3327 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3328 let cfg = config(spec);
3329 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3330
3331 say(
3332 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3333 &mut talk,
3334 &talks,
3335 &cfg,
3336 "what does the queue module do?",
3337 Vec::new(),
3338 )
3339 .await
3340 .expect("first turn");
3341 let first_prompt = &talk.turns[1].body;
3342 assert!(first_prompt.contains("magi task add --solo"));
3343 assert!(first_prompt.contains("what does the queue module do?"));
3344
3345 say(
3346 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3347 &mut talk,
3348 &talks,
3349 &cfg,
3350 "and how is it locked?",
3351 Vec::new(),
3352 )
3353 .await
3354 .expect("second turn");
3355 let second_prompt = &talk.turns[3].body;
3356 assert!(
3357 !second_prompt.contains("magi task add --solo"),
3358 "the briefing is sent once, not on every turn: {second_prompt}"
3359 );
3360 assert!(second_prompt.contains("and how is it locked?"));
3361 }
3362
3363 #[tokio::test]
3364 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
3365 let (tmp, talks) = store();
3366 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3367 let mut b = a.clone();
3368 b.id = "other".to_owned();
3369 let mut cfg = config(a.clone());
3370 cfg.agents.push(b.clone());
3371 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
3372 say(
3373 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3374 &mut talk,
3375 &talks,
3376 &cfg,
3377 "remember the walrus",
3378 Vec::new(),
3379 )
3380 .await
3381 .expect("first turn");
3382 let old_session = talk.seat.claude_session.clone();
3383 assert_eq!(talk.seat.turns, 1);
3384
3385 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
3386 assert_eq!(talk.agent, "other");
3387 assert_eq!(talk.seat.turns, 0);
3388 assert_eq!(talk.seat.agent, "other");
3389 assert_ne!(talk.seat.claude_session, old_session);
3390 let note = talk.turns.last().expect("note");
3391 assert_eq!(note.who, Who::Agent);
3392 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3393 assert!(note.body.contains("changed from"), "{}", note.body);
3394 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
3395
3396 let before = talk.turns.len();
3397 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
3398 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
3399
3400 say(
3401 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3402 &mut talk,
3403 &talks,
3404 &cfg,
3405 "what did I say?",
3406 Vec::new(),
3407 )
3408 .await
3409 .expect("turn after switch");
3410 let prompt = &talk.turns.last().expect("reply").body;
3411 assert!(prompt.contains("remember the walrus"), "{prompt}");
3412 assert!(prompt.contains("## magi"), "{prompt}");
3413 assert!(prompt.contains("what did I say?"), "{prompt}");
3414 }
3415
3416 #[tokio::test]
3417 async fn say_appends_the_operator_turn_then_the_agent_turn() {
3418 let (tmp, talks) = store();
3419 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3420 let cfg = config(spec);
3421 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3422
3423 say(
3424 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3425 &mut talk,
3426 &talks,
3427 &cfg,
3428 "can I rename this function?",
3429 Vec::new(),
3430 )
3431 .await
3432 .expect("say");
3433
3434 assert_eq!(talk.turns.len(), 2);
3435 assert_eq!(talk.turns[0].who, Who::Operator);
3436 assert_eq!(talk.turns[0].body, "can I rename this function?");
3437 assert_eq!(talk.turns[1].who, Who::Agent);
3438 assert_eq!(talk.turns[1].body, "go ahead");
3439 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
3440 }
3441
3442 #[tokio::test]
3443 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
3444 let (tmp, talks) = store();
3445 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
3446 let cfg = config(spec);
3447 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3448
3449 let err = say(
3450 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3451 &mut talk,
3452 &talks,
3453 &cfg,
3454 "check the tests",
3455 Vec::new(),
3456 )
3457 .await
3458 .expect_err("a turn with no answer is an error");
3459 assert!(err.to_string().contains("no answer"), "{err}");
3460
3461 let on_disk = talks.get(&talk.id).expect("get");
3462 assert_eq!(on_disk.turns.len(), 2);
3463 assert_eq!(on_disk.turns[0].body, "check the tests");
3464 let note = &on_disk.turns[1];
3465 assert_eq!(note.who, Who::Agent);
3466 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3467 assert!(note.body.contains("your message is saved"));
3468 }
3469
3470 #[tokio::test]
3476 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
3477 let (tmp, talks) = store();
3478 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3479 let cfg = config(spec);
3480 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3481
3482 let text =
3483 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
3484 failpoint::force_put_failures(PUT_RETRIES - 1);
3487 respond(
3488 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3489 &mut talk,
3490 &talks,
3491 &cfg,
3492 &text,
3493 )
3494 .await
3495 .expect("respond must survive a write failure its own retries can outlast");
3496
3497 assert_eq!(talk.turns.len(), 2);
3498 assert_eq!(talk.turns[1].who, Who::Agent);
3499 assert_eq!(talk.turns[1].body, "go ahead");
3500 let on_disk = talks.get(&talk.id).expect("get");
3501 assert_eq!(
3502 on_disk.turns, talk.turns,
3503 "the reply must reach disk despite the early write failures"
3504 );
3505 }
3506
3507 #[tokio::test]
3513 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
3514 let (tmp, talks) = store();
3515 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3516 let cfg = config(spec);
3517 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3518
3519 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3520 failpoint::force_put_failures(PUT_RETRIES);
3525 let err = respond(
3526 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3527 &mut talk,
3528 &talks,
3529 &cfg,
3530 &text,
3531 )
3532 .await
3533 .expect_err("a reply that cannot be saved must be reported, not swallowed");
3534 assert!(err.to_string().contains("could not be saved"), "{err}");
3535
3536 let on_disk = talks.get(&talk.id).expect("get");
3537 assert_eq!(
3538 on_disk.turns.len(),
3539 2,
3540 "the operator turn plus a visible note"
3541 );
3542 assert_eq!(on_disk.turns[0].body, "check the tests");
3543 let note = &on_disk.turns[1];
3544 assert_eq!(note.who, Who::Agent);
3545 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3546 assert!(
3547 note.body.contains("could not be saved"),
3548 "the operator must be told the reply is missing, not left staring \
3549 at a gap with no explanation: {}",
3550 note.body
3551 );
3552 assert_eq!(
3553 talk.turns, on_disk.turns,
3554 "the in-memory talk must match what actually landed on disk"
3555 );
3556
3557 let artifacts = talks.artifacts_of(&talk.id);
3560 let stash = std::fs::read_dir(&artifacts)
3561 .expect("artifacts dir")
3562 .filter_map(|e| e.ok())
3563 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
3564 .expect("a stash file for the lost reply");
3565 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
3566 assert_eq!(stashed, "go ahead");
3567
3568 assert_eq!(
3576 on_disk.seat.turns, 1,
3577 "the note's write must carry the turn the CLI actually took"
3578 );
3579 assert_eq!(
3580 on_disk.seat.claude_session, talk.seat.claude_session,
3581 "the session id handed to the CLI must survive the failed reply"
3582 );
3583 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
3584 assert!(
3585 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
3586 "the next turn must resume, not open the same session id twice"
3587 );
3588 }
3589
3590 #[tokio::test]
3595 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
3596 let (tmp, talks) = store();
3597 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3598 let cfg = config(spec);
3599 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3600
3601 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3602 failpoint::force_put_failures(PUT_RETRIES * 2);
3605 let err = respond(
3606 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3607 &mut talk,
3608 &talks,
3609 &cfg,
3610 &text,
3611 )
3612 .await
3613 .expect_err("neither the reply nor the note could be saved");
3614 assert!(err.to_string().contains("could not be saved"), "{err}");
3615
3616 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
3617 let on_disk = talks.get(&talk.id).expect("get");
3618 assert_eq!(on_disk.turns.len(), 1);
3619
3620 assert_eq!(
3630 on_disk.seat.turns, 0,
3631 "an unwritable file cannot record the turn the CLI took"
3632 );
3633 assert_eq!(
3634 talk.seat.turns, 1,
3635 "the in-memory seat still reports the turn the CLI actually took"
3636 );
3637 assert_eq!(
3638 on_disk.seat.claude_session, talk.seat.claude_session,
3639 "the session id was minted at `begin` and never changes here"
3640 );
3641 }
3642
3643 #[tokio::test]
3647 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
3648 let (tmp, talks) = store();
3649 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3650 let cfg = config(spec);
3651 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3652
3653 let att = talks
3654 .put_attachment(
3655 &talk.id,
3656 "image/png",
3657 "screenshot.png",
3658 b"pretend-png-bytes",
3659 )
3660 .expect("put attachment");
3661
3662 say(
3663 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3664 &mut talk,
3665 &talks,
3666 &cfg,
3667 "",
3668 vec![att.clone()],
3669 )
3670 .await
3671 .expect("an empty body with an attachment is still a turn");
3672
3673 let operator_turn = &talk.turns[0];
3674 assert_eq!(operator_turn.who, Who::Operator);
3675 assert_eq!(operator_turn.body, "");
3676 assert_eq!(operator_turn.attachments, vec![att.clone()]);
3677
3678 let prompt = &talk.turns[1].body;
3679 let expected_path = talks
3680 .attachments_dir(&talk.id)
3681 .join(format!("{}.png", att.id));
3682 assert!(
3683 prompt.contains(&expected_path.display().to_string()),
3684 "the agent must be told the attachment's absolute path: {prompt}"
3685 );
3686 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
3687 }
3688
3689 #[test]
3698 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
3699 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
3700 let att = Attachment {
3701 id: "0".repeat(32),
3702 name: "shot.png".to_owned(),
3703 mime: "image/png".to_owned(),
3704 bytes: 3,
3705 };
3706 let path = talks
3707 .attachment_path("some-talk-id", &att)
3708 .expect("a supported mime always yields a path");
3709 assert!(
3710 path.is_absolute(),
3711 "must be absolute even off a relative store root: {}",
3712 path.display()
3713 );
3714 }
3715
3716 #[tokio::test]
3717 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
3718 let (tmp, talks) = store();
3723 let slow = mock_agent(
3724 tmp.path(),
3725 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
3726 BTreeMap::new(),
3727 );
3728 let mut cfg = config(slow);
3729 cfg.graph.timeout_talk = 1;
3730 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3731
3732 let err = say(
3733 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3734 &mut talk,
3735 &talks,
3736 &cfg,
3737 "check the tests",
3738 Vec::new(),
3739 )
3740 .await
3741 .expect_err("a turn that never answers is an error");
3742 assert!(
3743 err.to_string().contains("did not answer within 1s"),
3744 "{err}"
3745 );
3746
3747 let on_disk = talks.get(&talk.id).expect("get");
3748 let note = on_disk.turns.last().expect("a note turn was recorded");
3749 assert!(
3750 note.body.contains("did not answer within 1s"),
3751 "the transcript must show the configured timeout: {}",
3752 note.body
3753 );
3754 }
3755
3756 #[test]
3757 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
3758 let (tmp, talks) = store();
3759 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3760 let cfg = config(spec);
3761 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3762
3763 close(&mut talk, &talks).expect("close");
3764 assert_eq!(talk.status, TalkStatus::Closed);
3765 close(&mut talk, &talks).expect("closing twice is not an error");
3766
3767 let err =
3768 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
3769 assert!(err.to_string().contains("closed"));
3770 let _ = &cfg; }
3772
3773 #[tokio::test]
3774 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3775 let (tmp, talks) = store();
3776 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3777 let cfg = config(spec);
3778 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3781
3782 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
3786 close(&mut closed_elsewhere, &talks).expect("close");
3787 assert_eq!(
3788 talks.get(&in_flight.id).expect("reread").status,
3789 TalkStatus::Closed,
3790 "the close landed on disk before the turn finished"
3791 );
3792
3793 assert_eq!(in_flight.status, TalkStatus::Open);
3797 respond(
3798 &talks.claim_turn(&in_flight.id).unwrap().unwrap(),
3799 &mut in_flight,
3800 &talks,
3801 &cfg,
3802 "one more question",
3803 )
3804 .await
3805 .expect("the turn itself still completes");
3806
3807 let on_disk = talks.get(&in_flight.id).expect("reread");
3808 assert_eq!(
3809 on_disk.status,
3810 TalkStatus::Closed,
3811 "a close must stick even when a turn that started before it finishes after it"
3812 );
3813 assert!(
3816 on_disk.turns.iter().any(|t| t.body == "here you go"),
3817 "the in-flight turn's own reply is still recorded: {:?}",
3818 on_disk.turns
3819 );
3820 }
3821
3822 #[test]
3823 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
3824 let (tmp, talks) = store();
3825 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3826 let cfg = config(spec);
3827 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3830
3831 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
3834 close(&mut closed_elsewhere, &talks).expect("close");
3835 assert_eq!(
3836 talks.get(&stale.id).expect("reread").status,
3837 TalkStatus::Closed,
3838 "the close landed on disk before record was called"
3839 );
3840
3841 assert_eq!(stale.status, TalkStatus::Open);
3845 let err = record(&mut stale, &talks, "still there?", Vec::new())
3846 .expect_err("a close that landed first must be honored, not overwritten");
3847 assert!(err.to_string().contains("closed"));
3848
3849 let on_disk = talks.get(&stale.id).expect("reread");
3850 assert_eq!(
3851 on_disk.status,
3852 TalkStatus::Closed,
3853 "record must not resurrect a conversation closed while its snapshot was stale"
3854 );
3855 assert!(
3856 on_disk.turns.is_empty(),
3857 "the rejected turn must not have been appended: {:?}",
3858 on_disk.turns
3859 );
3860 let _ = &cfg; }
3862
3863 #[test]
3864 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
3865 let (tmp, talks) = store();
3866 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3867 let cfg = config(spec);
3868 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3869
3870 let held = talks.guard().unwrap();
3874
3875 let talks2 = talks.clone();
3876 let id = talk.id.clone();
3877 let closing = std::thread::spawn(move || {
3878 let mut talk = talks2.get(&id).expect("get");
3879 close(&mut talk, &talks2).expect("close");
3880 });
3881
3882 std::thread::sleep(Duration::from_millis(50));
3883 assert!(
3884 !closing.is_finished(),
3885 "close must wait for the guard, not read and write while it is held - \
3886 a re-read alone narrows this window without closing it"
3887 );
3888
3889 drop(held);
3890 closing.join().expect("close thread panicked");
3891
3892 assert_eq!(
3893 talks.get(&talk.id).expect("reread").status,
3894 TalkStatus::Closed,
3895 "once the guard is free, close still lands"
3896 );
3897 let _ = &cfg; }
3899
3900 #[test]
3901 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
3902 let (tmp, talks) = store();
3903 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3904 let cfg = config(spec);
3905 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3906
3907 close(&mut talk, &talks).expect("close");
3908 assert_eq!(talk.status, TalkStatus::Closed);
3909
3910 reopen(&mut talk, &talks).expect("reopen");
3911 assert_eq!(talk.status, TalkStatus::Open);
3912 assert_eq!(
3913 talks.get(&talk.id).expect("reread").status,
3914 TalkStatus::Open
3915 );
3916
3917 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
3919 assert_eq!(talk.status, TalkStatus::Open);
3920
3921 record(&mut talk, &talks, "one more thing", Vec::new())
3922 .expect("a reopened talk takes turns again");
3923 let _ = &cfg; }
3925
3926 #[test]
3927 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
3928 let (tmp, talks) = store();
3929 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3930 let cfg = config(spec);
3931 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3932
3933 let artifacts = talks.artifacts_of(&talk.id);
3934 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
3935 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
3936
3937 talks.remove(&talk.id).expect("remove");
3938 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
3939 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
3940 assert!(
3941 talks.get(&talk.id).is_err(),
3942 "a removed talk cannot be read back"
3943 );
3944
3945 let err = talks
3946 .remove("nonexistent-id")
3947 .expect_err("unknown id refused");
3948 assert!(err.to_string().contains("no talk matches"), "{err}");
3949 let _ = &cfg; }
3951
3952 #[tokio::test]
3953 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3954 let (tmp, talks) = store();
3955 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3956 let cfg = config(spec);
3957 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3960
3961 talks.remove(&in_flight.id).expect("remove");
3962 assert!(
3963 talks.get(&in_flight.id).is_err(),
3964 "the delete landed on disk before the turn finished"
3965 );
3966
3967 respond(
3970 &talks.claim_turn(&in_flight.id).unwrap().unwrap(),
3971 &mut in_flight,
3972 &talks,
3973 &cfg,
3974 "one more question",
3975 )
3976 .await
3977 .expect("the turn itself still completes rather than erroring");
3978
3979 assert!(
3980 talks.get(&in_flight.id).is_err(),
3981 "a delete must stick even when a turn that started before it finishes after it"
3982 );
3983 }
3984
3985 #[test]
3986 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
3987 let (tmp, talks) = store();
3988 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3989 let cfg = config(spec);
3990 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3993
3994 talks.remove(&stale.id).expect("remove");
3995
3996 let err = record(&mut stale, &talks, "still there?", Vec::new())
4000 .expect_err("a delete that landed first must be honored, not overwritten");
4001 assert!(err.to_string().contains("deleted"), "{err}");
4002
4003 assert!(
4004 talks.get(&stale.id).is_err(),
4005 "record must not resurrect a conversation deleted while its snapshot was stale"
4006 );
4007 let _ = &cfg; }
4009
4010 #[test]
4011 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
4012 let (tmp, talks) = store();
4013 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4014 let cfg = config(spec);
4015 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4018
4019 talks.remove(&stale.id).expect("remove");
4020
4021 let err = close(&mut stale, &talks)
4025 .expect_err("a delete that landed first must be honored, not overwritten");
4026 assert!(err.to_string().contains("deleted"), "{err}");
4027
4028 assert!(
4029 talks.get(&stale.id).is_err(),
4030 "close must not resurrect a conversation deleted while its snapshot was stale"
4031 );
4032 let _ = &cfg; }
4034
4035 #[test]
4036 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
4037 let (tmp, talks) = store();
4038 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4039 let cfg = config(spec);
4040 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4043 close(&mut stale, &talks).expect("close");
4044
4045 talks.remove(&stale.id).expect("remove");
4046
4047 let err = reopen(&mut stale, &talks)
4051 .expect_err("a delete that landed first must be honored, not overwritten");
4052 assert!(err.to_string().contains("deleted"), "{err}");
4053
4054 assert!(
4055 talks.get(&stale.id).is_err(),
4056 "reopen must not resurrect a conversation deleted while its snapshot was stale"
4057 );
4058 let _ = &cfg; }
4060
4061 #[test]
4062 fn list_puts_open_talks_before_closed_ones() {
4063 let (tmp, talks) = store();
4064 let make = |id: &str, status: TalkStatus| {
4065 let mut t = Talk {
4066 schema: SCHEMA,
4067 id: id.to_owned(),
4068 repo: tmp.path().to_owned(),
4069 agent: "mock".to_owned(),
4070 status,
4071 turns: Vec::new(),
4072 pending: String::new(),
4073 pending_attachments: Vec::new(),
4074 fallback: false,
4075 persona: String::new(),
4076 persona_dirty: false,
4077 created_at: Timestamp::now(),
4078 updated_at: Timestamp::now(),
4079 seat: SeatState::new(SEAT, "mock", 7),
4080 };
4081 talks.put(&mut t).expect("put");
4082 };
4083 make("20260901-000000-0001", TalkStatus::Open);
4084 make("20260902-000000-0002", TalkStatus::Open);
4085 make("20260903-000000-0003", TalkStatus::Closed);
4086
4087 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
4088 assert_eq!(
4089 ids,
4090 [
4091 "20260902-000000-0002",
4092 "20260901-000000-0001",
4093 "20260903-000000-0003"
4094 ]
4095 );
4096 assert_eq!(talks.count_open(), 2);
4097 }
4098
4099 #[test]
4100 fn tasks_of_finds_only_this_talks_own_tasks() {
4101 let dir = tempfile::tempdir().expect("tempdir");
4102 let queue = Queue::at(dir.path().join("queue"));
4103
4104 let mut mine = Task::new(
4105 "rework the loader".to_owned(),
4106 "rework the loader".to_owned(),
4107 PathBuf::from("/repo"),
4108 Source::Agent {
4109 run: "20260904-014455-ab12".to_owned(),
4110 node: "chat".to_owned(),
4111 },
4112 );
4113 queue.put(&mut mine).expect("put mine");
4114
4115 let mut theirs = Task::new(
4116 "unrelated".to_owned(),
4117 "unrelated".to_owned(),
4118 PathBuf::from("/repo"),
4119 Source::Agent {
4120 run: "20260904-090000-zz99".to_owned(),
4121 node: "implement".to_owned(),
4122 },
4123 );
4124 queue.put(&mut theirs).expect("put theirs");
4125
4126 let mut human = Task::new(
4127 "typed by hand".to_owned(),
4128 "typed by hand".to_owned(),
4129 PathBuf::from("/repo"),
4130 Source::Human,
4131 );
4132 queue.put(&mut human).expect("put human");
4133
4134 let found = tasks_of(&queue, "20260904-014455-ab12");
4135 assert_eq!(found.len(), 1);
4136 assert_eq!(found[0].id, mine.id);
4137 }
4138
4139 #[test]
4140 fn the_briefing_names_solo_task_add() {
4141 let brief = briefing(Path::new("/repo"), "en", false);
4142 assert!(brief.contains("magi task add --solo"));
4143 assert!(brief.contains("/repo"));
4144 assert!(!brief.contains("Hold this conversation in"));
4145 }
4146
4147 #[test]
4153 fn the_briefing_explains_targeting_a_different_repository_by_name() {
4154 let brief = briefing(Path::new("/repo"), "en", false);
4155 assert!(brief.contains("--repo does not have to be a full path"));
4156 assert!(brief.contains("owner/repo"));
4157 assert!(brief.contains("magi repos"));
4158 assert!(brief.contains("ask the operator"));
4159 }
4160
4161 #[test]
4162 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
4163 let brief = briefing(Path::new("/repo"), "en", false);
4164 assert!(brief.contains("--attach <path>"), "{brief}");
4165 assert!(brief.contains("deleting this conversation"), "{brief}");
4166 }
4167
4168 #[test]
4169 fn the_briefing_names_the_language_when_it_is_not_english() {
4170 let brief = briefing(Path::new("/repo"), "Japanese", false);
4171 assert!(brief.contains("Hold this conversation in Japanese"));
4172 }
4173
4174 #[test]
4175 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
4176 let read_only = briefing(Path::new("/repo"), "en", false);
4177 assert!(read_only.contains("Do not write files"));
4178 assert!(!read_only.contains("allow_write"));
4179
4180 let writable = briefing(Path::new("/repo"), "en", true);
4181 assert!(!writable.contains("Do not write files"));
4182 assert!(writable.contains("allow_write = true"));
4183 assert!(writable.contains("magi task add --solo"));
4186 assert!(writable.contains("say plainly what you"));
4187 }
4188}