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 #[serde(default, skip_serializing_if = "Option::is_none")]
146 pub breaks: Option<Vec<usize>>,
147}
148
149#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
157pub struct TurnUsage {
158 pub context_tokens: u64,
160 pub agent: String,
162 #[serde(default)]
164 pub model: Option<String>,
165}
166
167#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
172pub struct ContextUsage {
173 pub tokens: Option<u64>,
175 pub window: Option<u64>,
177 pub percent: Option<u64>,
179 pub warn: bool,
181 pub since_switch: bool,
185 pub model: Option<String>,
187 #[serde(default)]
190 pub estimated: bool,
191}
192
193const CONTEXT_WARN_PERCENT: u64 = 80;
195
196const STANDING_PROMPT_FALLBACK_CHARS: u64 = 7000;
199
200pub fn estimate_context_tokens(talk: &Talk, standing_chars: u64) -> Option<u64> {
213 let mut counted = false;
214 let mut chars = standing_chars;
215 for t in talk.turns.iter().filter(|t| !t.body.starts_with(MAGI_NOTE)) {
216 counted = true;
217 chars += t.body.chars().count() as u64;
218 }
219 counted.then(|| (chars * 2).div_ceil(7))
221}
222
223pub fn context_usage(talk: &Talk, cfg: Option<&Config>) -> ContextUsage {
236 let current = cfg.and_then(|c| c.agents.iter().find(|a| a.id == talk.agent));
237 let model = current.and_then(|a| a.model.clone());
238 let window = cfg
239 .zip(model.as_deref())
240 .and_then(|(c, m)| c.context_window(m))
241 .filter(|w| *w > 0);
242 let usage = talk
243 .turns
244 .iter()
245 .rev()
246 .find(|t| t.who == Who::Agent && !t.body.starts_with(MAGI_NOTE))
247 .and_then(|t| t.usage.as_ref());
248 let measured = usage.map(|u| u.context_tokens);
249 let tokens = measured.or_else(|| {
250 let standing = cfg.map_or(STANDING_PROMPT_FALLBACK_CHARS, |c| {
251 briefing_with(
252 &talk.repo,
253 &c.graph.language,
254 c.talk.allow_write,
255 crate::persona::active(&c.talk.personas, &talk.persona).as_ref(),
256 c.talk.operator_name(),
257 )
258 .chars()
259 .count() as u64
260 });
261 estimate_context_tokens(talk, standing)
262 });
263 let estimated = measured.is_none() && tokens.is_some();
264 let since_switch =
265 usage.is_some_and(|u| u.agent != talk.agent || (current.is_some() && u.model != model));
266 let (percent, warn) = match (tokens, window) {
267 (Some(t), Some(w)) => (
268 Some(t.saturating_mul(100) / w),
269 t.saturating_mul(100) >= w.saturating_mul(CONTEXT_WARN_PERCENT),
270 ),
271 _ => (None, false),
272 };
273 ContextUsage {
274 tokens,
275 window,
276 percent,
277 warn,
278 since_switch,
279 model,
280 estimated,
281 }
282}
283
284#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
288#[serde(rename_all = "lowercase")]
289pub enum TalkStatus {
290 Open,
293 Closed,
295}
296
297impl TalkStatus {
298 pub fn open(self) -> bool {
300 matches!(self, Self::Open)
301 }
302
303 pub fn as_str(self) -> &'static str {
305 match self {
306 Self::Open => "open",
307 Self::Closed => "closed",
308 }
309 }
310}
311
312#[derive(Debug, Clone, Serialize, Deserialize)]
314#[serde(deny_unknown_fields)]
315pub struct Talk {
316 pub schema: u32,
318 pub id: String,
320 pub repo: PathBuf,
322 pub agent: String,
324 pub status: TalkStatus,
326 pub turns: Vec<Turn>,
328 #[serde(default)]
331 pub pending: String,
332 #[serde(default)]
334 pub pending_attachments: Vec<Attachment>,
335 #[serde(default, skip_serializing_if = "Option::is_none")]
337 pub pending_breaks: Option<Vec<usize>>,
338 #[serde(default)]
343 pub fallback: bool,
344 #[serde(default)]
347 pub persona: String,
348 #[serde(default)]
352 pub persona_dirty: bool,
353 pub created_at: Timestamp,
355 pub updated_at: Timestamp,
357 seat: SeatState,
362}
363
364impl Talk {
365 pub fn short(&self) -> &str {
367 short(&self.id)
368 }
369}
370
371#[derive(Debug, Clone)]
373pub struct Talks {
374 root: PathBuf,
375 lock: Arc<Mutex<()>>,
384}
385
386impl Talks {
387 pub fn open() -> Self {
389 Self::at(crate::run::home().join("talks"))
390 }
391
392 pub fn at(root: PathBuf) -> Self {
395 Self {
396 root,
397 lock: Arc::new(Mutex::new(())),
398 }
399 }
400
401 fn guard(&self) -> Result<StoreGuard<'_>> {
416 let mutex = self.lock.lock().unwrap_or_else(PoisonError::into_inner);
417 std::fs::create_dir_all(&self.root)
418 .with_context(|| format!("create {}", self.root.display()))?;
419 let path = self.root.join(".store.turn");
420 let deadline = std::time::Instant::now() + TAKEOVER_LOCK_TTL * 2;
421 loop {
422 if let Some(file) = TurnLock::take(&path)? {
423 return Ok(StoreGuard {
424 _file: file,
425 _mutex: mutex,
426 });
427 }
428 if std::time::Instant::now() >= deadline {
429 bail!("the talk store is locked by another process");
430 }
431 std::thread::sleep(Duration::from_millis(10));
432 }
433 }
434
435 pub fn root(&self) -> &Path {
437 &self.root
438 }
439
440 pub fn path_of(&self, id: &str) -> PathBuf {
442 self.root.join(format!("{id}.json"))
443 }
444
445 pub fn artifacts_of(&self, id: &str) -> PathBuf {
448 self.root.join(format!("{id}.artifacts"))
449 }
450
451 pub fn attachments_dir(&self, id: &str) -> PathBuf {
455 self.artifacts_of(id).join("attachments")
456 }
457
458 pub fn put_attachment(
467 &self,
468 id: &str,
469 mime: &str,
470 name: &str,
471 data: &[u8],
472 ) -> Result<Attachment> {
473 let dir = self.attachments_dir(id);
474 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
475 let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
476 let att = Attachment {
477 id: new_attachment_id(),
478 name: name.to_owned(),
479 mime: mime.to_owned(),
480 bytes: data.len() as u64,
481 };
482 std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
483 .with_context(|| format!("write attachment {}", att.id))?;
484 std::fs::write(
485 dir.join(format!("{}.json", att.id)),
486 serde_json::to_string(&att).context("serialize attachment")?,
487 )
488 .with_context(|| format!("write attachment metadata {}", att.id))?;
489 Ok(att)
490 }
491
492 pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
501 if !valid_attachment_id(att_id) {
502 return Ok(None);
503 }
504 let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
505 if !meta_path.is_file() {
506 return Ok(None);
507 }
508 let att = serde_json::from_str(
509 &std::fs::read_to_string(&meta_path)
510 .with_context(|| format!("read {}", meta_path.display()))?,
511 )
512 .with_context(|| format!("parse {}", meta_path.display()))?;
513 Ok(Some(att))
514 }
515
516 pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
520 let Some(att) = self.attachment_meta(id, att_id)? else {
521 return Ok(None);
522 };
523 let ext = attachment_ext(&att.mime).with_context(|| {
524 format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
525 })?;
526 let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
527 let data =
528 std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
529 Ok(Some((att, data)))
530 }
531
532 fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
549 let ext = attachment_ext(&att.mime)?;
550 let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
551 std::path::absolute(&path).ok()
552 }
553
554 pub fn put(&self, t: &mut Talk) -> Result<()> {
563 std::fs::create_dir_all(&self.root)
564 .with_context(|| format!("create {}", self.root.display()))?;
565 t.updated_at = Timestamp::now();
566 let body = serde_json::to_string_pretty(t).context("serialize talk")?;
567 let path = self.path_of(&t.id);
568 let tmp = path.with_extension("json.tmp");
569 write_atomic(&tmp, &path, &body)
570 }
571
572 pub fn get(&self, id: &str) -> Result<Talk> {
574 let resolved = self.resolve_id(id)?;
575 read_path(&self.path_of(&resolved))
576 }
577
578 pub fn list(&self) -> Vec<Talk> {
581 self.list_counting_unreadable().0
582 }
583
584 pub fn list_counting_unreadable(&self) -> (Vec<Talk>, usize) {
587 let mut unreadable = 0;
588 let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
589 .into_iter()
590 .flatten()
591 .flatten()
592 .map(|e| e.path())
593 .filter(|p| p.extension().is_some_and(|x| x == "json"))
594 .filter_map(|p| {
595 let talk = read_path(&p).ok();
596 if talk.is_none() {
597 unreadable += 1;
598 }
599 talk
600 })
601 .collect();
602 all.sort_unstable_by(|a, b| {
603 let rank = |t: &Talk| u8::from(!t.status.open());
604 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
605 });
606 (all, unreadable)
607 }
608
609 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
611 if self.path_of(prefix).is_file() {
612 return Ok(prefix.to_owned());
613 }
614 let hits: Vec<String> = self
615 .list()
616 .into_iter()
617 .map(|t| t.id)
618 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
619 .collect();
620 match hits.len() {
621 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
622 0 => bail!("no talk matches `{prefix}`"),
623 _ => bail!(
624 "`{prefix}` matches {} talks: {}",
625 hits.len(),
626 hits.join(", ")
627 ),
628 }
629 }
630
631 pub fn revision(&self) -> u64 {
634 std::fs::read_dir(&self.root)
635 .into_iter()
636 .flatten()
637 .flatten()
638 .filter(|e| e.path().extension().is_none_or(|x| x != "turn"))
639 .filter(|e| !e.file_name().to_string_lossy().starts_with('.'))
640 .filter_map(|e| e.metadata().ok())
641 .filter_map(|m| m.modified().ok())
642 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
643 .map(|d| d.as_millis() as u64)
644 .max()
645 .unwrap_or(0)
646 }
647
648 pub fn count_open(&self) -> usize {
650 self.list().iter().filter(|t| t.status.open()).count()
651 }
652
653 pub fn turn_path(&self, id: &str) -> PathBuf {
656 self.root.join(format!("{id}.turn"))
657 }
658
659 pub fn claim_turn(&self, id: &str) -> Result<Option<TurnLease>> {
668 self.claim_turn_at(id, Timestamp::now())
669 }
670
671 fn claim_turn_at(&self, id: &str, now: Timestamp) -> Result<Option<TurnLease>> {
672 std::fs::create_dir_all(&self.root)
673 .with_context(|| format!("create {}", self.root.display()))?;
674 let path = self.turn_path(id);
675 let token = fresh_token();
676 if create_turn(&path, &token, now)? {
677 return Ok(Some(TurnLease {
678 talk: id.to_owned(),
679 path,
680 token,
681 }));
682 }
683 if lease_blocks(&path, now) {
684 return Ok(None);
685 }
686 let Some(_lock) = TurnLock::take(&path)? else {
691 return Ok(None);
692 };
693 if lease_blocks(&path, now) {
694 return Ok(None);
695 }
696 let _ = std::fs::remove_file(&path);
697 Ok(create_turn(&path, &token, now)?.then_some(TurnLease {
698 talk: id.to_owned(),
699 path,
700 token,
701 }))
702 }
703
704 pub fn turn_held(&self, id: &str) -> bool {
706 read_turn(&self.turn_path(id)).is_some_and(|r| r.fresh(Timestamp::now()))
707 }
708
709 pub fn remove(&self, id: &str) -> Result<()> {
722 let _guard = self.guard()?;
723 let resolved = self.resolve_id(id)?;
724 let path = self.path_of(&resolved);
725 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
726 let artifacts = self.artifacts_of(&resolved);
727 if artifacts.is_dir() {
728 std::fs::remove_dir_all(&artifacts)
729 .with_context(|| format!("remove {}", artifacts.display()))?;
730 }
731 let _ = std::fs::remove_file(self.turn_path(&resolved));
732 Ok(())
733 }
734}
735
736struct StoreGuard<'a> {
739 _file: TurnLock,
740 _mutex: MutexGuard<'a, ()>,
741}
742
743#[derive(Debug, Serialize, Deserialize)]
746struct TurnRecord {
747 token: String,
748 pid: u32,
749 beat_at: Timestamp,
750}
751
752impl TurnRecord {
753 fn fresh(&self, now: Timestamp) -> bool {
754 now.as_second() - self.beat_at.as_second() <= crate::ask::LEASE_TTL.as_secs() as i64
755 }
756}
757
758fn fresh_token() -> String {
762 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
763 let n = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
764 let seed = crate::rng::entropy() ^ n.wrapping_mul(0x9E37_79B9_7F4A_7C15);
765 crate::rng::SplitMix64::new(seed).uuid_v4()
766}
767
768fn read_turn(path: &Path) -> Option<TurnRecord> {
769 serde_json::from_str(&std::fs::read_to_string(path).ok()?).ok()
770}
771
772const TAKEOVER_LOCK_TTL: Duration = Duration::from_secs(10);
774
775const TICKET_TTL: Duration = Duration::from_secs(10);
778
779const TICKET_GENERATIONS: u32 = 16;
781
782const TICKET_SWEEP_AGE: Duration = Duration::from_secs(3600);
784
785fn create_exclusive(path: &Path, body: &str) -> Result<bool> {
787 use std::io::Write as _;
788 match std::fs::OpenOptions::new()
789 .write(true)
790 .create_new(true)
791 .open(path)
792 {
793 Ok(mut f) => {
794 if let Err(e) = f.write_all(body.as_bytes()) {
795 drop(f);
796 let _ = std::fs::remove_file(path);
797 return Err(e).with_context(|| format!("write {}", path.display()));
798 }
799 Ok(true)
800 }
801 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
802 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
803 }
804}
805
806fn create_turn(path: &Path, token: &str, now: Timestamp) -> Result<bool> {
807 let record = TurnRecord {
808 token: token.to_owned(),
809 pid: std::process::id(),
810 beat_at: now,
811 };
812 let body = serde_json::to_string(&record).context("serialize turn lease")?;
813 let tmp = path.with_extension(format!("turn.{token}.new"));
816 publish_exclusive(path, &tmp, &body)
817}
818
819fn lease_blocks(path: &Path, now: Timestamp) -> bool {
823 if read_turn(path).is_some_and(|r| r.fresh(now)) {
824 return true;
825 }
826 std::fs::metadata(path).is_ok_and(|m| {
827 m.len() == 0
828 && m.modified()
829 .ok()
830 .and_then(|t| t.elapsed().ok())
831 .is_some_and(|age| age <= TAKEOVER_LOCK_TTL)
832 })
833}
834
835fn link_unsupported(e: &std::io::Error) -> bool {
838 e.kind() == std::io::ErrorKind::Unsupported || (cfg!(windows) && e.raw_os_error() == Some(1))
839}
840
841fn create_in_place(path: &Path, body: &str) -> Result<bool> {
859 use std::io::Write as _;
860 let mut f = match std::fs::OpenOptions::new()
861 .write(true)
862 .create_new(true)
863 .open(path)
864 {
865 Ok(f) => f,
866 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => return Ok(false),
867 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
868 };
869 f.write_all(body.as_bytes())
870 .with_context(|| format!("write {}", path.display()))?;
871 Ok(std::fs::read_to_string(path).is_ok_and(|t| t == body))
874}
875
876fn publish_exclusive(path: &Path, tmp: &Path, body: &str) -> Result<bool> {
877 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
878 #[cfg(test)]
879 let linked = if failpoint::no_link_forced() {
880 Err(std::io::Error::from(std::io::ErrorKind::Unsupported))
881 } else {
882 std::fs::hard_link(tmp, path)
883 };
884 #[cfg(not(test))]
885 let linked = std::fs::hard_link(tmp, path);
886 let out = match linked {
887 Ok(()) => Ok(true),
888 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
889 Err(e) if link_unsupported(&e) => create_in_place(path, body),
890 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
891 };
892 let _ = std::fs::remove_file(tmp);
893 out
894}
895
896fn remove_if_carries(path: &Path, expected: &str) -> bool {
907 let gone = path.with_extension(format!("gone.{}", fresh_token()));
908 if std::fs::rename(path, &gone).is_err() {
909 return false;
910 }
911 let found = std::fs::read_to_string(&gone);
912 if found.as_ref().is_ok_and(|t| t == expected) {
913 let _ = std::fs::remove_file(&gone);
914 return true;
915 }
916 if let Ok(body) = found {
917 let back = path.with_extension(format!("back.{}", fresh_token()));
918 match publish_exclusive(path, &back, &body) {
919 Ok(true) => {}
920 Ok(false) | Err(_) => tracing::warn!(
921 "{} was replaced while a stale removal had it aside; \
922 the displaced file is dropped",
923 path.display()
924 ),
925 }
926 }
927 let _ = std::fs::remove_file(&gone);
928 false
929}
930
931struct TurnLock {
935 path: PathBuf,
936 token: String,
937}
938
939impl TurnLock {
940 fn token() -> String {
941 fresh_token()
942 }
943
944 fn publish(path: &Path, token: &str) -> Result<bool> {
948 let tmp = path.with_extension(format!("lock.{token}.new"));
949 publish_exclusive(path, &tmp, token)
950 }
951
952 fn take(lease: &Path) -> Result<Option<Self>> {
953 let path = lease.with_extension("turn.lock");
954 let token = Self::token();
955 if Self::publish(&path, &token)? {
956 return Ok(Some(Self { path, token }));
957 }
958 let Ok(seen) = std::fs::read_to_string(&path) else {
959 return Ok(None);
960 };
961 let key = if !seen.is_empty() && seen.chars().all(|c| c.is_ascii_alphanumeric() || c == '-')
964 {
965 seen.as_str()
966 } else {
967 "invalid"
968 };
969 let aged = std::fs::metadata(&path)
970 .and_then(|m| m.modified())
971 .ok()
972 .and_then(|t| t.elapsed().ok())
973 .is_some_and(|age| age > TAKEOVER_LOCK_TTL);
974 if !aged {
975 return Ok(None);
976 }
977 let mut won = false;
989 for n in 0..TICKET_GENERATIONS {
990 let ticket = path.with_extension(format!("lock.{key}.break.{n}"));
991 if create_exclusive(&ticket, "")? {
992 won = true;
993 break;
994 }
995 let stale = std::fs::metadata(&ticket)
996 .and_then(|m| m.modified())
997 .ok()
998 .and_then(|t| t.elapsed().ok())
999 .is_some_and(|age| age > TICKET_TTL);
1000 if !stale {
1001 return Ok(None);
1002 }
1003 }
1004 if !won {
1005 Self::sweep_tickets(&path);
1010 return Ok(None);
1011 }
1012 Self::sweep_tickets(&path);
1013 if !remove_if_carries(&path, &seen) {
1019 return Ok(None);
1020 }
1021 if Self::publish(&path, &token)? {
1022 return Ok(Some(Self { path, token }));
1023 }
1024 Ok(None)
1025 }
1026
1027 fn sweep_tickets(path: &Path) {
1030 let (Some(dir), Some(name)) = (path.parent(), path.file_name().and_then(|n| n.to_str()))
1031 else {
1032 return;
1033 };
1034 let prefix = format!("{name}.");
1035 let Ok(entries) = std::fs::read_dir(dir) else {
1036 return;
1037 };
1038 for entry in entries.flatten() {
1039 let file = entry.file_name();
1040 let Some(file) = file.to_str() else { continue };
1041 if !(file.starts_with(&prefix) && file.contains(".break.")) {
1042 continue;
1043 }
1044 let old = entry
1045 .metadata()
1046 .and_then(|m| m.modified())
1047 .ok()
1048 .and_then(|t| t.elapsed().ok())
1049 .is_some_and(|age| age > TICKET_SWEEP_AGE);
1050 if old {
1051 let _ = std::fs::remove_file(entry.path());
1052 }
1053 }
1054 }
1055
1056 fn take_patiently(lease: &Path) -> Option<Self> {
1058 for _ in 0..50 {
1059 match Self::take(lease) {
1060 Ok(Some(lock)) => return Some(lock),
1061 Ok(None) => std::thread::sleep(Duration::from_millis(10)),
1062 Err(_) => return None,
1063 }
1064 }
1065 None
1066 }
1067}
1068
1069impl Drop for TurnLock {
1070 fn drop(&mut self) {
1071 remove_if_carries(&self.path, &self.token);
1073 }
1074}
1075
1076pub const TURN_BEAT: Duration = Duration::from_secs(20);
1079
1080#[derive(Debug)]
1085pub struct TurnLease {
1086 talk: String,
1087 path: PathBuf,
1088 token: String,
1089}
1090
1091impl TurnLease {
1092 pub fn holds(&self) -> bool {
1094 read_turn(&self.path).is_some_and(|r| r.token == self.token)
1095 }
1096
1097 pub fn beat(&self) -> Result<bool> {
1101 let _lock = TurnLock::take_patiently(&self.path)
1102 .with_context(|| format!("lock {} to renew it", self.path.display()))?;
1103 let Some(mut record) = read_turn(&self.path).filter(|r| r.token == self.token) else {
1104 return Ok(false);
1105 };
1106 record.beat_at = Timestamp::now();
1107 let body = serde_json::to_string(&record).context("serialize turn lease")?;
1108 let tmp = self.path.with_extension(format!("turn.{}.tmp", self.token));
1109 write_atomic(&tmp, &self.path, &body)?;
1110 Ok(true)
1111 }
1112
1113 pub async fn beating<T>(&self, fut: impl std::future::Future<Output = T>) -> Result<T> {
1118 self.beating_every(TURN_BEAT, fut).await
1119 }
1120
1121 async fn beating_every<T>(
1122 &self,
1123 period: Duration,
1124 fut: impl std::future::Future<Output = T>,
1125 ) -> Result<T> {
1126 tokio::pin!(fut);
1127 loop {
1128 match tokio::time::timeout(period, &mut fut).await {
1129 Ok(out) => return Ok(out),
1130 Err(_) => match self.beat() {
1131 Ok(true) => {}
1132 Ok(false) => bail!(
1133 "the turn lease {} was taken over; this turn is stopped",
1134 self.path.display()
1135 ),
1136 Err(e) => tracing::warn!("{e:#}"),
1137 },
1138 }
1139 }
1140 }
1141}
1142
1143impl Drop for TurnLease {
1144 fn drop(&mut self) {
1145 if let Some(_lock) = TurnLock::take_patiently(&self.path) {
1148 if read_turn(&self.path).is_some_and(|r| r.token == self.token) {
1149 let _ = std::fs::remove_file(&self.path);
1150 }
1151 }
1152 }
1153}
1154
1155pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
1165 let repo = repo.canonicalize().unwrap_or(repo);
1168 let spec = match agent {
1171 Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
1172 None => agent::pick_chain(
1173 &cfg.agents,
1174 cfg.roles.chatter.as_ref(),
1175 &agent::installed,
1176 "chatter",
1177 )?
1178 .remove(0),
1179 };
1180
1181 let now = Timestamp::now();
1182 let mut talk = Talk {
1183 pending_breaks: None,
1184 schema: SCHEMA,
1185 id: new_id(),
1186 repo,
1187 agent: spec.id.clone(),
1188 status: TalkStatus::Open,
1189 turns: Vec::new(),
1190 pending: String::new(),
1191 pending_attachments: Vec::new(),
1192 fallback: agent.is_none(),
1193 persona: String::new(),
1194 persona_dirty: false,
1195 created_at: now,
1196 updated_at: now,
1197 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
1198 };
1199 store.put(&mut talk)?;
1200 Ok(talk)
1201}
1202
1203pub fn record(
1210 talk: &mut Talk,
1211 store: &Talks,
1212 text: &str,
1213 attachments: Vec<Attachment>,
1214) -> Result<String> {
1215 let _guard = store.guard()?;
1223 let Ok(fresh) = store.get(&talk.id) else {
1228 bail!("talk {} was deleted", talk.short());
1229 };
1230 talk.status = fresh.status;
1231 talk.pending = fresh.pending;
1234 talk.pending_breaks = fresh.pending_breaks;
1235 talk.pending_attachments = fresh.pending_attachments;
1236 if !talk.status.open() {
1237 bail!(
1238 "talk {} is {} and takes no more turns",
1239 talk.short(),
1240 talk.status.as_str()
1241 );
1242 }
1243 let text = text.trim();
1244 if text.is_empty() && attachments.is_empty() {
1245 bail!("nothing to say");
1246 }
1247 talk.turns.push(Turn {
1248 who: Who::Operator,
1249 body: text.to_owned(),
1250 at: Timestamp::now(),
1251 attachments,
1252 usage: None,
1253 breaks: Some(Vec::new()),
1254 });
1255 store.put(talk)?;
1256 Ok(text.to_owned())
1257}
1258
1259pub fn queue(
1261 talk: &mut Talk,
1262 store: &Talks,
1263 text: &str,
1264 attachments: Vec<Attachment>,
1265) -> Result<()> {
1266 let text = text.trim();
1267 if text.is_empty() && attachments.is_empty() {
1268 bail!("nothing to say");
1269 }
1270 let _guard = store.guard()?;
1271 let mut fresh = store
1272 .get(&talk.id)
1273 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1274 if !fresh.status.open() {
1275 bail!(
1276 "talk {} is {} and takes no more turns",
1277 fresh.short(),
1278 fresh.status.as_str()
1279 );
1280 }
1281 if !text.is_empty() {
1282 if fresh.pending.is_empty() {
1283 fresh.pending = text.to_owned();
1284 fresh.pending_breaks = Some(Vec::new());
1285 } else {
1286 fresh.pending.push_str("\n\n");
1287 if let Some(b) = fresh.pending_breaks.as_mut() {
1288 b.push(fresh.pending.len());
1289 }
1290 fresh.pending.push_str(text);
1291 }
1292 }
1293 fresh.pending_attachments.extend(attachments);
1294 store.put(&mut fresh)?;
1295 *talk = fresh;
1296 Ok(())
1297}
1298
1299pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
1301 let _guard = store.guard()?;
1302 let mut fresh = store
1303 .get(&talk.id)
1304 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1305 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
1306 *talk = fresh;
1307 return Ok(None);
1308 }
1309 let text = std::mem::take(&mut fresh.pending);
1310 let attachments = std::mem::take(&mut fresh.pending_attachments);
1311 let breaks = fresh.pending_breaks.take();
1312 fresh.turns.push(Turn {
1313 who: Who::Operator,
1314 body: text.clone(),
1315 at: Timestamp::now(),
1316 attachments,
1317 usage: None,
1318 breaks,
1319 });
1320 store.put(&mut fresh)?;
1321 *talk = fresh;
1322 Ok(Some(text))
1323}
1324
1325pub async fn say(
1328 lease: &TurnLease,
1329 talk: &mut Talk,
1330 store: &Talks,
1331 cfg: &Config,
1332 text: &str,
1333 attachments: Vec<Attachment>,
1334) -> Result<()> {
1335 check_lease(lease, talk)?;
1336 let text = record(talk, store, text, attachments)?;
1337 respond(lease, talk, store, cfg, &text).await
1338}
1339
1340pub async fn respond(
1346 lease: &TurnLease,
1347 talk: &mut Talk,
1348 store: &Talks,
1349 cfg: &Config,
1350 text: &str,
1351) -> Result<()> {
1352 check_lease(lease, talk)?;
1353 lease
1354 .beating(turn(talk, store, cfg, text))
1355 .await
1356 .and_then(|done| done)
1357}
1358
1359fn check_lease(lease: &TurnLease, talk: &Talk) -> Result<()> {
1360 if lease.talk != talk.id {
1361 bail!("the turn lease is for talk {}, not {}", lease.talk, talk.id);
1362 }
1363 if !lease.holds() {
1366 bail!("the turn lease for talk {} is no longer held", talk.short());
1367 }
1368 Ok(())
1369}
1370
1371pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
1389 let _guard = store.guard()?;
1390 let mut fresh = store
1391 .get(&talk.id)
1392 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1393 fresh.status = TalkStatus::Closed;
1394 fresh.pending.clear();
1396 fresh.pending_breaks = None;
1397 fresh.pending_attachments.clear();
1398 store.put(&mut fresh)?;
1399 *talk = fresh;
1400 Ok(())
1401}
1402
1403pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
1414 let _guard = store.guard()?;
1415 let mut fresh = store
1416 .get(&talk.id)
1417 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1418 fresh.status = TalkStatus::Open;
1419 store.put(&mut fresh)?;
1420 *talk = fresh;
1421 Ok(())
1422}
1423
1424pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
1437 let _guard = store.guard()?;
1438 let mut fresh = store
1439 .get(&talk.id)
1440 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1441 if fresh.agent == spec.id {
1442 *talk = fresh;
1443 return Ok(false);
1444 }
1445 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
1446 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1447 fresh.fallback = false;
1449 fresh.turns.push(Turn {
1450 breaks: None,
1451 who: Who::Agent,
1452 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
1453 at: Timestamp::now(),
1454 attachments: Vec::new(),
1455 usage: None,
1456 });
1457 store.put(&mut fresh)?;
1458 *talk = fresh;
1459 Ok(true)
1460}
1461
1462pub fn switch_persona(talk: &mut Talk, store: &Talks, id: &str) -> Result<bool> {
1470 let id = if id.trim() == crate::persona::DEFAULT_ID {
1471 ""
1472 } else {
1473 id.trim()
1474 };
1475 let _guard = store.guard()?;
1476 let mut fresh = store
1477 .get(&talk.id)
1478 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1479 if fresh.persona == id {
1480 *talk = fresh;
1481 return Ok(false);
1482 }
1483 let from = std::mem::replace(&mut fresh.persona, id.to_owned());
1484 fresh.persona_dirty = true;
1485 let label = |p: &str| {
1486 if p.is_empty() {
1487 crate::persona::DEFAULT_ID.to_owned()
1488 } else {
1489 p.to_owned()
1490 }
1491 };
1492 fresh.turns.push(Turn {
1493 breaks: None,
1494 who: Who::Agent,
1495 body: format!(
1496 "{MAGI_NOTE}persona changed from {} to {}",
1497 label(&from),
1498 label(id)
1499 ),
1500 at: Timestamp::now(),
1501 attachments: Vec::new(),
1502 usage: None,
1503 });
1504 store.put(&mut fresh)?;
1505 *talk = fresh;
1506 Ok(true)
1507}
1508
1509pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
1511 let _guard = store.guard()?;
1512 let mut fresh = store
1513 .get(&talk.id)
1514 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1515 fresh.pending.clear();
1516 fresh.pending_breaks = None;
1517 fresh.pending_attachments.clear();
1518 store.put(&mut fresh)?;
1519 *talk = fresh;
1520 Ok(())
1521}
1522
1523pub fn clear_pending_if_matches(
1525 talk: &mut Talk,
1526 store: &Talks,
1527 expected_text: &str,
1528 expected_attachments: &[String],
1529) -> Result<bool> {
1530 let _guard = store.guard()?;
1531 let mut fresh = store
1532 .get(&talk.id)
1533 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1534 if !pending_matches(&fresh, expected_text, expected_attachments) {
1535 *talk = fresh;
1536 return Ok(false);
1537 }
1538 fresh.pending.clear();
1539 fresh.pending_breaks = None;
1540 fresh.pending_attachments.clear();
1541 store.put(&mut fresh)?;
1542 *talk = fresh;
1543 Ok(true)
1544}
1545
1546pub fn edit_pending_text(
1550 talk: &mut Talk,
1551 store: &Talks,
1552 text: &str,
1553 expected_text: &str,
1554 expected_attachments: &[String],
1555) -> Result<bool> {
1556 let _guard = store.guard()?;
1557 let mut fresh = store
1558 .get(&talk.id)
1559 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1560 if !pending_matches(&fresh, expected_text, expected_attachments) {
1561 *talk = fresh;
1562 return Ok(false);
1563 }
1564 fresh.pending = text.trim().to_owned();
1565 fresh.pending_breaks = Some(Vec::new());
1566 store.put(&mut fresh)?;
1567 *talk = fresh;
1568 Ok(true)
1569}
1570
1571fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
1572 talk.pending == expected_text
1573 && talk
1574 .pending_attachments
1575 .iter()
1576 .map(|attachment| &attachment.id)
1577 .eq(expected_attachments.iter())
1578}
1579
1580async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1587 let spec = cfg
1588 .agents
1589 .iter()
1590 .find(|a| a.id == talk.agent)
1591 .with_context(|| {
1592 format!(
1593 "talk {} was opened with agent `{}`, which is no longer in \
1594 the roster; restore it in magi.toml or start a new \
1595 conversation",
1596 talk.short(),
1597 talk.agent
1598 )
1599 })?;
1600
1601 let last_note = attachment_note(
1605 store,
1606 &talk.id,
1607 talk.turns
1608 .last()
1609 .map_or(&[][..], |t| t.attachments.as_slice()),
1610 );
1611
1612 let attachment_paths: Vec<PathBuf> = talk
1618 .turns
1619 .iter()
1620 .flat_map(|t| t.attachments.iter())
1621 .filter_map(|a| store.attachment_path(&talk.id, a))
1622 .collect();
1623
1624 let questions =
1634 crate::run::try_home().map(|home| crate::ask::Questions::at(home.join("questions")));
1635 let consulted = questions
1636 .as_ref()
1637 .is_some_and(|q| crate::consult::pending_consults(q, &talk.id));
1638 let consult_roots: Vec<PathBuf> = match &questions {
1639 Some(q) if consulted => vec![q.root().to_path_buf()],
1640 _ => Vec::new(),
1641 };
1642
1643 let (allow_write, unsandboxed) = turn_access(cfg.talk.allow_write, consulted);
1644
1645 let artifacts = store.artifacts_of(&talk.id);
1646 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
1649 let stem = format!("turn-{}", operator_turns.max(1));
1650 let cache_dir = cfg.cache_dir();
1653
1654 let mut chain = vec![spec.clone()];
1657 if let Some(choice) = cfg.roles.chatter.as_ref()
1658 && talk.fallback
1659 {
1660 for id in choice.ids() {
1661 if id == talk.agent || chain.iter().any(|s| s.id == id) {
1662 continue;
1663 }
1664 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1665 Ok(s) => chain.push(s),
1666 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1667 }
1668 }
1669 }
1670
1671 let persona = crate::persona::active(&cfg.talk.personas, &talk.persona);
1672 let operator_name = cfg.talk.operator_name();
1673 let persona_update = if talk.persona_dirty {
1674 format!(
1675 "{}\n\n",
1676 crate::persona::update_block_for(persona.as_ref(), operator_name)
1677 )
1678 } else {
1679 String::new()
1680 };
1681
1682 let mut outcome = None;
1683 let mut fell_back_from: Option<String> = None;
1684 let mut first_try: Option<(String, SeatState)> = None;
1687 for (n, spec) in chain.iter().enumerate() {
1688 if n > 0 {
1689 if first_try.is_none() {
1690 first_try = Some((talk.agent.clone(), talk.seat.clone()));
1691 }
1692 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1693 fell_back_from.get_or_insert_with(|| talk.agent.clone());
1696 talk.agent = spec.id.clone();
1697 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1698 }
1699 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1700 let first_ever = talk.turns.len() <= 1;
1701 let body = if talk.seat.turns == 0 && first_ever {
1702 format!(
1703 "{}\n\n# Operator\n\n{text}{last_note}",
1704 briefing_with(
1705 &talk.repo,
1706 &cfg.graph.language,
1707 cfg.talk.allow_write,
1708 persona.as_ref(),
1709 operator_name
1710 )
1711 )
1712 } else if talk.seat.turns == 0 {
1713 format!(
1716 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1717 briefing_with(
1718 &talk.repo,
1719 &cfg.graph.language,
1720 cfg.talk.allow_write,
1721 persona.as_ref(),
1722 operator_name
1723 ),
1724 transcript(talk, store)
1725 )
1726 } else if resuming {
1727 format!("{persona_update}{text}{last_note}")
1728 } else {
1729 let mut standing = match (&persona, talk.persona_dirty) {
1733 (Some(p), false) => format!("{}\n", crate::persona::section_for(p, operator_name)),
1734 _ => persona_update.clone(),
1735 };
1736 if let (None, Some(n)) = (&persona, operator_name) {
1739 standing.push_str(&format!(
1740 "# Addressing the operator\n{}\n",
1741 crate::persona::addressing(n)
1742 ));
1743 }
1744 format!("{}\n\n{standing}{text}{last_note}", transcript(talk, store))
1745 };
1746 let attempt_stem = if n == 0 {
1747 stem.clone()
1748 } else {
1749 format!("{stem}-{}", spec.id)
1750 };
1751 let inv = Invocation {
1752 cwd: &talk.repo,
1753 prompt: &body,
1754 timeout: turn_timeout(cfg),
1755 allow_write,
1760 unsandboxed,
1761 sessions: cfg.graph.sessions,
1762 artifacts: &artifacts,
1763 stem: &attempt_stem,
1764 run: &talk.id,
1767 node: crate::queue::CHAT_NODE,
1768 cache_dir: cache_dir.as_deref(),
1769 attachments: &attachment_paths,
1770 writable: &consult_roots,
1771 };
1772 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1773 let advance = agent::chain_advances(&result);
1774 if n == 0 || !advance {
1775 outcome = Some(result);
1776 } else {
1777 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1779 }
1780 if !advance {
1781 break;
1782 }
1783 }
1784 if outcome.as_ref().is_some_and(agent::chain_advances) {
1785 if let Some((id, seat)) = first_try {
1788 talk.agent = id;
1789 talk.seat = seat;
1790 fell_back_from = None;
1791 }
1792 }
1793 let outcome = outcome.expect("a chain holds at least one agent");
1794 let note = |why: String| Turn {
1795 breaks: None,
1796 who: Who::Agent,
1797 body: format!("{MAGI_NOTE}{why}"),
1798 at: Timestamp::now(),
1799 attachments: Vec::new(),
1800 usage: None,
1801 };
1802 let (reply, failure) = match outcome {
1803 Err(e) => (
1804 note(format!("could not run agent `{}`: {e}", talk.agent)),
1805 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1806 ),
1807 Ok(out) if out.quota_exhausted() => {
1808 let reset = out
1809 .quota
1810 .as_ref()
1811 .and_then(|q| q.reset.clone())
1812 .map_or_else(String::new, |r| format!(" (resets {r})"));
1813 let why = format!(
1814 "agent `{}` is out of quota{reset}; your message is saved, so \
1815 say it again when the window reopens",
1816 talk.agent
1817 );
1818 (note(why.clone()), Some(why))
1819 }
1820 Ok(out) if out.timed_out => {
1821 let why = format!(
1822 "agent `{}` did not answer within {}s; your message is saved",
1823 talk.agent,
1824 turn_timeout(cfg).as_secs()
1825 );
1826 (note(why.clone()), Some(why))
1827 }
1828 Ok(out) if !out.usable() => {
1829 let why = format!(
1830 "agent `{}` produced no answer (exit {}); your message is saved",
1831 talk.agent,
1832 out.exit_code
1833 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1834 );
1835 (note(why.clone()), Some(why))
1836 }
1837 Ok(out) => (
1838 Turn {
1839 breaks: None,
1840 who: Who::Agent,
1841 body: out.text.trim().to_owned(),
1842 at: Timestamp::now(),
1843 attachments: Vec::new(),
1844 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1847 context_tokens,
1848 agent: talk.agent.clone(),
1849 model: cfg
1850 .agents
1851 .iter()
1852 .find(|a| a.id == talk.agent)
1853 .and_then(|a| a.model.clone()),
1854 }),
1855 },
1856 None,
1857 ),
1858 };
1859
1860 let _guard = store.guard()?;
1872 let Ok(fresh) = store.get(&talk.id) else {
1878 return Ok(());
1879 };
1880 talk.status = fresh.status;
1881 talk.pending = fresh.pending;
1885 talk.pending_breaks = fresh.pending_breaks;
1886 talk.pending_attachments = fresh.pending_attachments;
1887 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1888 talk.turns.push(note(format!(
1891 "agent changed from {from} to {} (fallback)",
1892 talk.agent
1893 )));
1894 }
1895 if failure.is_none() {
1898 talk.persona_dirty = false;
1899 }
1900 talk.turns.push(reply);
1901 if let Err(put_err) = store.put(talk) {
1902 let lost = talk.turns.pop().expect("just pushed above");
1911 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1912 let why = match &stash {
1913 Ok(path) => format!(
1914 "agent `{}` answered, but the reply could not be saved to \
1915 this conversation ({put_err:#}); the raw text was kept at \
1916 {} - your message is saved, ask again",
1917 talk.agent,
1918 path.display()
1919 ),
1920 Err(stash_err) => format!(
1921 "agent `{}` answered, but the reply could not be saved to \
1922 this conversation ({put_err:#}), and it could not be kept \
1923 anywhere else either ({stash_err:#}); your message is \
1924 saved, ask again",
1925 talk.agent
1926 ),
1927 };
1928 talk.turns.push(note(why.clone()));
1929 return match store.put(talk) {
1936 Ok(()) => bail!("{why}"),
1937 Err(note_err) => {
1938 talk.turns.pop();
1958 Err(note_err).context(why)
1959 }
1960 };
1961 }
1962
1963 match failure {
1964 Some(why) => bail!("{why}"),
1965 None => Ok(()),
1966 }
1967}
1968
1969fn transcript(talk: &Talk, store: &Talks) -> String {
1972 let mut out = String::from(
1973 "This conversation cannot resume on the CLI's side, so here is \
1974 everything said so far; answer only the last message.\n",
1975 );
1976 for t in &talk.turns {
1977 let who = match t.who {
1978 Who::Operator => "operator",
1979 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1980 Who::Agent => "you",
1981 };
1982 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1983 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1984 }
1985 out
1986}
1987
1988fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1993 if attachments.is_empty() {
1994 return String::new();
1995 }
1996 let mut out = String::from(
1997 "\n\nThe operator attached the image(s) below to this message. Open \
1998 and look at each one before you answer.\n",
1999 );
2000 for att in attachments {
2001 if let Some(path) = store.attachment_path(talk_id, att) {
2002 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
2003 }
2004 }
2005 out.push('\n');
2006 out
2007}
2008
2009pub(crate) fn turn_access(talk_allow_write: bool, consulted: bool) -> (bool, bool) {
2020 (talk_allow_write || consulted, talk_allow_write)
2021}
2022
2023pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
2039 briefing_with(repo, language, allow_write, None, None)
2040}
2041
2042pub fn briefing_with(
2047 repo: &Path,
2048 language: &str,
2049 allow_write: bool,
2050 persona: Option<&crate::persona::Persona>,
2051 operator_name: Option<&str>,
2052) -> String {
2053 let write_policy = if allow_write {
2054 "Write access is enabled for this conversation (`allow_write = \
2055 true`), so you may write files - but only a small, \
2056 already-decided edit the operator names outright in this \
2057 conversation, not an implementation. This is a permission on the \
2058 conversation as a whole, not a property of whichever repository \
2059 it happened to start in: if the operator names a different \
2060 repository for that small edit, the policy allows it there too. \
2061 Your own tool may still confine writes to the repository this \
2062 conversation started in regardless - if a write elsewhere is \
2063 refused, say so plainly rather than working around it. Once you \
2064 have made an edit, say plainly what you edited. Anything bigger, \
2065 or anything still open-ended, still goes through the queue below \
2066 rather than being done here."
2067 } else {
2068 "Do not write files. Implementing a change is not this \
2069 conversation's job; a separate, blind competition of agents does \
2070 that, and a repository this conversation has already edited would \
2071 make their diffs unjudgeable."
2072 };
2073 let mut out = format!(
2074 "You are magi's standing conversation partner for its operator, who \
2075 usually has this open on a phone. Keep replies short: no preamble, \
2076 no restating what they just said.\n\n\
2077 # Repository\n\n{repo}\n\n\
2078 You may look around: read files, run shell commands, search history, \
2079 run tests - whatever answers the question. {write_policy}\n\n\
2080 A short, command-shaped message (\"list\", \"info <id>\", \"show \
2081 3cbf\") is almost always the operator asking you to look something \
2082 up, not an instruction to file - answer it yourself with `magi \
2083 list`, `magi show <id>`, `magi task list`, or the like, the same way \
2084 you would answer any other question in this conversation.\n\n\
2085 # When the operator wants something done\n\n\
2086 Run:\n\n\
2087 magi task add --solo --repo {repo} <instruction>\n\n\
2088 and tell the operator the task id it prints, so they can follow it \
2089 from the Queue. If it refuses with a duplicate warning (the \
2090 instruction names a branch, commit or pull request that an \
2091 unfinished task, run or PR already owns), do not repeat it with \
2092 --force yourself: tell the operator what it matched and let them \
2093 decide. Write <instruction> so that an implementer who has \
2094 never seen this conversation can act on it alone - it is everything \
2095 they get. Use --solo: it runs the task through one implementer \
2096 straight into review instead of the usual multi-agent competition, \
2097 which is the right shape for a change this conversation has already \
2098 settled, rather than one still worth several independent takes.\n\n\
2099 If the operator asks for something in a different repository, \
2100 --repo does not have to be a full path: --repo owner/repo (or just \
2101 repo, when that is unambiguous) is resolved against local checkouts \
2102 the same way `magi repos` lists them. If the command fails because \
2103 nothing matches or more than one checkout shares that name, ask the \
2104 operator which repository they mean (or run `magi repos` yourself \
2105 to see the candidates) rather than guessing.\n\n\
2106 The current state of the code is whatever origin/main holds, not \
2107 whatever a working tree shows: a primary checkout often lags \
2108 upstream, sits on a detached HEAD and carries uncommitted changes. \
2109 Before answering about code, run `git fetch origin` in that \
2110 repository if it is cheap, then read through \
2111 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
2112 If the working tree differs, say so; if the fetch fails, say that \
2113 too, so the operator knows the answer may be stale.\n\n\
2114 If the operator attached an image (a screenshot, say) that the task \
2115 is about, pass it with `--attach <path>`, using the absolute path \
2116 the turn's attachment note gives; repeat the flag for several. \
2117 `magi task add --solo --attach <path> <instruction>` copies the \
2118 file into the task, so the implementer receives it. Do not paste the \
2119 path into <instruction> instead: deleting this conversation deletes \
2120 its attachments, and then that path reaches no one.\n",
2121 repo = repo.display(),
2122 );
2123 out.push_str(&language_note(language));
2124 match (persona, operator_name) {
2125 (Some(p), n) => out.push_str(&crate::persona::section_for(p, n)),
2126 (None, Some(n)) => {
2127 out.push_str("\n# Addressing the operator\n");
2128 out.push_str(&crate::persona::addressing(n));
2129 }
2130 (None, None) => {}
2131 }
2132 out
2133}
2134
2135fn language_note(language: &str) -> String {
2138 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
2139 String::new()
2140 } else {
2141 format!("\nHold this conversation in {language}.\n")
2142 }
2143}
2144
2145pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
2152 let mut tasks: Vec<Task> = queue
2153 .list()
2154 .into_iter()
2155 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
2156 .collect();
2157 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
2158 tasks
2159}
2160
2161fn read_path(path: &Path) -> Result<Talk> {
2162 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2163 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
2164}
2165
2166const PUT_RETRIES: u32 = 5;
2169
2170fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
2182 let mut last_err = None;
2183 for attempt in 0..PUT_RETRIES {
2184 if attempt > 0 {
2185 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
2186 }
2187 match try_write_atomic(tmp, path, body) {
2188 Ok(()) => return Ok(()),
2189 Err(e) => last_err = Some(e),
2190 }
2191 }
2192 Err(last_err.expect("the loop above always runs at least once"))
2193}
2194
2195fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
2196 #[cfg(test)]
2197 if failpoint::take_forced_put_failure() {
2198 bail!("simulated write failure (test)");
2199 }
2200 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
2201 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
2202 Ok(())
2203}
2204
2205fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
2210 let dir = store.artifacts_of(id);
2211 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
2212 let path = dir.join(format!("{stem}-lost.txt"));
2213 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
2214 Ok(path)
2215}
2216
2217#[cfg(test)]
2224mod failpoint {
2225 use std::cell::Cell;
2226
2227 thread_local! {
2228 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
2229 static FORCE_NO_LINK: Cell<bool> = const { Cell::new(false) };
2230 }
2231
2232 pub(super) fn force_no_link(on: bool) {
2235 FORCE_NO_LINK.with(|c| c.set(on));
2236 }
2237
2238 pub(super) fn no_link_forced() -> bool {
2239 FORCE_NO_LINK.with(Cell::get)
2240 }
2241
2242 pub(super) fn force_put_failures(count: u32) {
2245 FORCE_PUT_FAILURES.with(|c| c.set(count));
2246 }
2247
2248 pub(super) fn take_forced_put_failure() -> bool {
2251 FORCE_PUT_FAILURES.with(|c| {
2252 let n = c.get();
2253 if n == 0 {
2254 false
2255 } else {
2256 c.set(n - 1);
2257 true
2258 }
2259 })
2260 }
2261}
2262
2263fn short(id: &str) -> &str {
2264 id.split('-').next_back().unwrap_or(id)
2265}
2266
2267fn new_id() -> String {
2268 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2269 let seed = crate::rng::entropy();
2270 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2271}
2272
2273fn attachment_ext(mime: &str) -> Option<&'static str> {
2278 match mime {
2279 "image/png" => Some("png"),
2280 "image/jpeg" => Some("jpg"),
2281 "image/gif" => Some("gif"),
2282 "image/webp" => Some("webp"),
2283 _ => None,
2284 }
2285}
2286
2287pub fn valid_attachment_id(id: &str) -> bool {
2292 id.len() == 32
2293 && id
2294 .bytes()
2295 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
2296}
2297
2298fn new_attachment_id() -> String {
2302 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
2303 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
2304}
2305
2306#[cfg(test)]
2307mod tests {
2308 #[test]
2309 fn turn_access_lifts_the_sandbox_only_for_the_talk_opt_in() {
2310 use super::turn_access;
2311 assert_eq!(turn_access(false, false), (false, false));
2312 assert_eq!(turn_access(true, false), (true, true));
2313 assert_eq!(turn_access(false, true), (true, false));
2315 assert_eq!(turn_access(true, true), (true, true));
2316 }
2317
2318 #[test]
2319 fn the_briefing_points_at_origin_main_not_the_working_tree() {
2320 let b = briefing(Path::new("/r"), "en", false);
2321 assert!(b.contains("origin/main"));
2322 assert!(b.contains("git show origin/main:"));
2323 }
2324 use std::collections::BTreeMap;
2325
2326 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
2327 use crate::queue::{Queue, Source, Task};
2328
2329 use super::*;
2330
2331 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
2332 AgentSpec {
2333 id: id.to_owned(),
2334 kind: AgentKind::Command,
2335 model: model.map(str::to_owned),
2336 command: Vec::new(),
2337 extra_args: Vec::new(),
2338 env: BTreeMap::new(),
2339 prompt_delivery: None,
2340 }
2341 }
2342
2343 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
2344 Talk {
2345 pending_breaks: None,
2346 schema: SCHEMA,
2347 id: "20260904-014455-ab12".to_owned(),
2348 repo: PathBuf::from("."),
2349 agent: agent.to_owned(),
2350 status: TalkStatus::Open,
2351 turns,
2352 pending: String::new(),
2353 pending_attachments: Vec::new(),
2354 fallback: false,
2355 persona: String::new(),
2356 persona_dirty: false,
2357 created_at: Timestamp::now(),
2358 updated_at: Timestamp::now(),
2359 seat: SeatState::new(SEAT, agent, 1),
2360 }
2361 }
2362
2363 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
2364 Turn {
2365 breaks: None,
2366 who: Who::Agent,
2367 body: body.to_owned(),
2368 at: Timestamp::now(),
2369 attachments: Vec::new(),
2370 usage: usage.map(|(t, a, m)| TurnUsage {
2371 context_tokens: t,
2372 agent: a.to_owned(),
2373 model: m.map(str::to_owned),
2374 }),
2375 }
2376 }
2377
2378 fn ctx_config(windows: &[(&str, u64)]) -> Config {
2379 Config {
2380 agents: vec![
2381 ctx_agent("small", Some("small-model")),
2382 ctx_agent("big", Some("big-model")),
2383 ctx_agent("plain", None),
2384 ],
2385 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
2386 ..Config::default()
2387 }
2388 }
2389
2390 #[test]
2391 fn context_usage_computes_percent_and_warns_at_eighty() {
2392 let cfg = ctx_config(&[("small-model", 1000)]);
2393 let at = |tokens| {
2394 let t = ctx_talk(
2395 "small",
2396 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
2397 );
2398 context_usage(&t, Some(&cfg))
2399 };
2400 let u = at(799);
2401 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
2402 let u = at(800);
2403 assert_eq!((u.percent, u.warn), (Some(80), true));
2404 let u = at(1500);
2405 assert_eq!((u.percent, u.warn), (Some(150), true));
2406 assert!(!u.since_switch);
2407 }
2408
2409 #[test]
2410 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
2411 let cfg = ctx_config(&[("small-model", 1000)]);
2412 let t = ctx_talk(
2413 "small",
2414 vec![
2415 reply("old", Some((900, "small", Some("small-model")))),
2416 reply("new", None),
2417 ],
2418 );
2419 let u = context_usage(&t, Some(&cfg));
2420 assert!(u.estimated);
2422 assert_ne!(u.tokens, Some(900));
2423 assert!(u.tokens.is_some());
2424 let t = ctx_talk(
2426 "small",
2427 vec![
2428 reply("old", Some((900, "small", Some("small-model")))),
2429 reply("magi: could not run agent", None),
2430 ],
2431 );
2432 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
2433 assert_eq!(
2434 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
2435 None
2436 );
2437 }
2438
2439 #[test]
2440 fn estimate_counts_chars_both_sides_and_standing_prompt() {
2441 let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
2442 assert_eq!(estimate_context_tokens(&t, 0), Some(2)); let op = Turn {
2444 who: Who::Operator,
2445 ..reply("abcdefg", None)
2446 };
2447 t.turns.push(op);
2448 assert_eq!(estimate_context_tokens(&t, 0), Some(4));
2449 assert!(
2450 estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
2451 );
2452 let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
2454 assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
2455 let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
2457 assert_eq!(estimate_context_tokens(¬e, 1000), None);
2458 assert_eq!(
2459 estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
2460 None
2461 );
2462 }
2463
2464 #[test]
2465 fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
2466 let cfg = ctx_config(&[("small-model", 1000)]);
2467 let t = ctx_talk(
2468 "small",
2469 vec![reply(
2470 &"x".repeat(5000),
2471 Some((10, "small", Some("small-model"))),
2472 )],
2473 );
2474 let u = context_usage(&t, Some(&cfg));
2475 assert_eq!((u.tokens, u.estimated), (Some(10), false));
2476 let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
2477 let u = context_usage(&t, Some(&cfg));
2478 assert!(u.estimated && !u.since_switch);
2479 assert_eq!(u.window, Some(1000));
2480 assert!(u.warn && u.percent.unwrap() >= 80);
2481 let t = ctx_talk("small", vec![reply("hi", None)]);
2482 let u = context_usage(&t, Some(&cfg));
2483 assert!(u.estimated && u.percent.is_some());
2484 }
2485
2486 #[test]
2487 fn context_usage_without_a_window_shows_tokens_only() {
2488 let cfg = ctx_config(&[]);
2489 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
2491 let u = context_usage(&t, Some(&cfg));
2492 assert_eq!(
2493 (u.tokens, u.window, u.percent, u.warn),
2494 (Some(5000), None, None, false)
2495 );
2496 let t = ctx_talk(
2497 "small",
2498 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
2499 );
2500 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
2501 assert_eq!(context_usage(&t, None).window, None);
2503 }
2504
2505 #[test]
2506 fn context_usage_switching_model_changes_the_denominator() {
2507 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
2508 let used = reply("hi", Some((900, "small", Some("small-model"))));
2509 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
2510 assert_eq!(
2511 (before.percent, before.warn, before.since_switch),
2512 (Some(90), true, false)
2513 );
2514 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
2517 assert_eq!(after.window, Some(10_000));
2518 assert_eq!(
2519 (after.percent, after.warn, after.since_switch),
2520 (Some(9), false, true)
2521 );
2522 assert_eq!(after.model.as_deref(), Some("big-model"));
2523 }
2524
2525 #[test]
2526 fn a_turn_recorded_before_usage_existed_still_reads() {
2527 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
2528 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
2529 assert!(turn.usage.is_none());
2530 let json = serde_json::to_string(&turn).expect("serialize");
2531 assert!(
2532 !json.contains("usage"),
2533 "absent usage is not written: {json}"
2534 );
2535 }
2536
2537 fn store() -> (tempfile::TempDir, Talks) {
2539 let tmp = tempfile::tempdir().expect("tempdir");
2540 let talks = Talks::at(tmp.path().join("talks"));
2541 (tmp, talks)
2542 }
2543
2544 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
2548 let path = dir.join("mock-talk-agent.sh");
2549 std::fs::write(&path, script).expect("write mock");
2550 AgentSpec {
2551 id: "mock".to_owned(),
2552 kind: AgentKind::Command,
2553 model: None,
2554 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2555 extra_args: Vec::new(),
2556 env,
2557 prompt_delivery: None,
2558 }
2559 }
2560
2561 fn config(spec: AgentSpec) -> Config {
2562 Config {
2563 agents: vec![spec],
2564 graph: Graph {
2565 language: "en".to_owned(),
2566 ..Graph::default()
2567 },
2568 ..Config::default()
2569 }
2570 }
2571
2572 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
2574
2575 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
2577
2578 const ECHO: &str = "#!/bin/sh\ncat\n";
2581
2582 fn env(reply: &str) -> BTreeMap<String, String> {
2583 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
2584 }
2585
2586 #[test]
2587 fn the_frozen_json_field_names_round_trip_through_disk() {
2588 let (tmp, talks) = store();
2589 let mut talk = Talk {
2590 pending_breaks: None,
2591 schema: SCHEMA,
2592 id: "20260904-014455-ab12".to_owned(),
2593 repo: tmp.path().to_owned(),
2594 agent: "sonnet".to_owned(),
2595 status: TalkStatus::Open,
2596 turns: Vec::new(),
2597 pending: String::new(),
2598 pending_attachments: Vec::new(),
2599 fallback: false,
2600 persona: String::new(),
2601 persona_dirty: false,
2602 created_at: Timestamp::now(),
2603 updated_at: Timestamp::now(),
2604 seat: SeatState::new(SEAT, "sonnet", 7),
2605 };
2606 talks.put(&mut talk).expect("put");
2607
2608 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
2609 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
2610 for field in [
2611 "schema",
2612 "id",
2613 "repo",
2614 "agent",
2615 "status",
2616 "turns",
2617 "created_at",
2618 "updated_at",
2619 ] {
2620 assert!(v.get(field).is_some(), "missing field `{field}`");
2621 }
2622 assert_eq!(v["schema"], 1);
2623 assert_eq!(v["status"], "open");
2624
2625 let back = talks.get(&talk.id).expect("get");
2626 assert_eq!(back.id, talk.id);
2627 assert_eq!(back.status, TalkStatus::Open);
2628 }
2629
2630 #[test]
2631 fn opening_a_talk_takes_no_agent_turn() {
2632 let (tmp, talks) = store();
2633 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2637 let cfg = config(spec);
2638
2639 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2640 assert_eq!(talk.status, TalkStatus::Open);
2641 assert!(talk.turns.is_empty(), "nothing has been said yet");
2642
2643 let on_disk = talks.get(&talk.id).expect("get");
2644 assert_eq!(on_disk.turns.len(), 0);
2645 }
2646
2647 #[test]
2655 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
2656 let (tmp, talks) = store();
2657 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2658 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2659 chatter_spec.id = "chatter-mock".to_owned();
2660
2661 let mut cfg = Config {
2662 agents: vec![first_spec.clone(), chatter_spec.clone()],
2663 graph: Graph {
2664 language: "en".to_owned(),
2665 ..Graph::default()
2666 },
2667 ..Config::default()
2668 };
2669 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
2670
2671 let talk =
2672 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
2673 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
2674
2675 cfg.roles.chatter = None;
2676 let fallback =
2677 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
2678 assert_eq!(
2679 fallback.agent, first_spec.id,
2680 "unset chatter must fall back to agent::pick's own default order"
2681 );
2682 }
2683
2684 #[test]
2687 fn a_talk_recorded_without_attachments_still_reads() {
2688 let (tmp, talks) = store();
2689 let path = talks.path_of("20260904-014455-ab12");
2690 std::fs::create_dir_all(talks.root()).expect("talks dir");
2691 std::fs::write(
2692 &path,
2693 serde_json::json!({
2694 "schema": 1,
2695 "id": "20260904-014455-ab12",
2696 "repo": tmp.path(),
2697 "agent": "sonnet",
2698 "status": "open",
2699 "turns": [
2700 { "who": "operator", "body": "still there?",
2701 "at": Timestamp::now().to_string() },
2702 ],
2703 "created_at": Timestamp::now().to_string(),
2704 "updated_at": Timestamp::now().to_string(),
2705 "seat": SeatState::new(SEAT, "sonnet", 7),
2706 })
2707 .to_string(),
2708 )
2709 .expect("write pre-attachments talk");
2710
2711 let talk = talks.get("20260904-014455-ab12").expect("must still read");
2712 assert!(talk.turns[0].attachments.is_empty());
2713 }
2714
2715 fn lease_store() -> (tempfile::TempDir, Talks, Talks) {
2716 let tmp = tempfile::TempDir::new().expect("tmp");
2717 let root = tmp.path().join("talks");
2718 (tmp, Talks::at(root.clone()), Talks::at(root))
2719 }
2720
2721 #[test]
2722 fn two_starters_on_one_talk_one_wins_and_the_other_is_refused() {
2723 let (_tmp, a, b) = lease_store();
2724 let won = a.claim_turn("t1").expect("claim").expect("first wins");
2725 assert!(
2726 b.claim_turn("t1").expect("claim").is_none(),
2727 "second is refused"
2728 );
2729 assert!(b.turn_held("t1"));
2730 assert!(
2731 b.claim_turn("t2").expect("claim").is_some(),
2732 "other talks are free"
2733 );
2734 drop(won);
2735 }
2736
2737 #[test]
2738 fn a_stale_lease_is_taken_over_and_the_old_guard_cannot_release_it() {
2739 let (_tmp, a, b) = lease_store();
2740 let old = a.claim_turn("t1").expect("claim").expect("held");
2741 let later = Timestamp::now()
2742 .checked_add(jiff::SignedDuration::from_secs(
2743 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2744 ))
2745 .expect("later");
2746 let new = b
2747 .claim_turn_at("t1", later)
2748 .expect("claim")
2749 .expect("a stale lease is taken over");
2750 drop(old);
2751 assert!(a.turn_held("t1"), "the old guard left the new lease alone");
2752 assert!(new.beat().expect("beat"), "the new owner still beats");
2753 drop(new);
2754 assert!(!a.turn_held("t1"));
2755 }
2756
2757 #[tokio::test]
2758 async fn a_turn_whose_lease_was_taken_over_is_stopped() {
2759 let (_tmp, a, b) = lease_store();
2760 let old = a.claim_turn("t1").expect("claim").expect("held");
2761 let later = Timestamp::now()
2762 .checked_add(jiff::SignedDuration::from_secs(
2763 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2764 ))
2765 .expect("later");
2766 let _new = b
2767 .claim_turn_at("t1", later)
2768 .expect("claim")
2769 .expect("taken over");
2770 let out = old
2771 .beating_every(Duration::from_millis(10), std::future::pending::<()>())
2772 .await;
2773 assert!(out.is_err(), "the displaced turn must stop, not run on");
2774 }
2775
2776 #[tokio::test]
2777 async fn a_turn_that_finishes_is_returned_and_keeps_its_lease_beating() {
2778 let (_tmp, a, _b) = lease_store();
2779 let lease = a.claim_turn("t1").expect("claim").expect("held");
2780 let out = lease
2781 .beating_every(Duration::from_millis(5), async {
2782 tokio::time::sleep(Duration::from_millis(40)).await;
2783 7
2784 })
2785 .await
2786 .expect("still ours");
2787 assert_eq!(out, 7);
2788 assert!(a.turn_held("t1"));
2789 }
2790
2791 #[test]
2792 fn an_unreadable_lease_counts_as_stale() {
2793 let (_tmp, a, b) = lease_store();
2794 std::fs::create_dir_all(a.root()).expect("dir");
2795 std::fs::write(a.turn_path("t1"), "not json").expect("write");
2796 assert!(!a.turn_held("t1"));
2797 assert!(b.claim_turn("t1").expect("claim").is_some());
2798 }
2799
2800 #[test]
2801 fn without_hard_links_a_claim_is_still_exclusive_and_a_young_placeholder_blocks() {
2802 let (_tmp, a, b) = lease_store();
2803 failpoint::force_no_link(true);
2804 let held = a.claim_turn("t1").expect("claim").expect("first wins");
2805 assert!(b.claim_turn("t1").expect("claim").is_none());
2806 assert!(a.turn_held("t1"));
2807 drop(held);
2808 assert!(!a.turn_held("t1"));
2809 let path = a.turn_path("t1");
2811 assert!(create_exclusive(&path, "").expect("placeholder"));
2812 assert!(b.claim_turn("t1").expect("claim").is_none(), "young: held");
2813 age_file(&path);
2814 assert!(b.claim_turn("t1").expect("claim").is_some(), "old: stale");
2815 let lease = a.turn_path("t2");
2817 let lock = TurnLock::take(&lease).expect("take").expect("lock");
2818 assert!(TurnLock::take(&lease).expect("take").is_none());
2819 drop(lock);
2820 assert!(TurnLock::take(&lease).expect("take").is_some());
2821 failpoint::force_no_link(false);
2822 }
2823
2824 #[test]
2825 fn remove_if_carries_removes_only_the_expected_content() {
2826 let (_tmp, a, _b) = lease_store();
2827 std::fs::create_dir_all(a.root()).expect("dir");
2828 let p = a.root().join("x.turn.lock");
2829 std::fs::write(&p, "mine").expect("write");
2830 assert!(!remove_if_carries(&p, "other"));
2831 assert_eq!(std::fs::read_to_string(&p).expect("kept"), "mine");
2832 assert!(remove_if_carries(&p, "mine"));
2833 assert!(!p.exists());
2834 assert!(!remove_if_carries(&p, "mine"), "absent is not a removal");
2835 }
2836
2837 #[test]
2838 fn a_dropped_lock_does_not_remove_a_lock_taken_over_since() {
2839 let (_tmp, a, _b) = lease_store();
2840 std::fs::create_dir_all(a.root()).expect("dir");
2841 let lease = a.turn_path("t1");
2842 let lock = TurnLock::take(&lease).expect("take").expect("lock");
2843 let path = lock.path.clone();
2844 std::fs::write(&path, "someone-else").expect("replace");
2845 drop(lock);
2846 assert_eq!(
2847 std::fs::read_to_string(&path).expect("kept"),
2848 "someone-else"
2849 );
2850 }
2851
2852 #[test]
2853 fn concurrent_takeovers_of_a_stale_lease_have_one_winner() {
2854 let (_tmp, a, _b) = lease_store();
2855 drop(a.claim_turn("t1").expect("claim").expect("held"));
2856 std::fs::write(
2857 a.turn_path("t1"),
2858 serde_json::to_string(&TurnRecord {
2859 token: "gone".into(),
2860 pid: 1,
2861 beat_at: Timestamp::from_second(1).expect("ts"),
2862 })
2863 .expect("json"),
2864 )
2865 .expect("write");
2866 let wins: Vec<_> = std::thread::scope(|sc| {
2867 let hs: Vec<_> = (0..8)
2868 .map(|_| {
2869 let s = a.clone();
2870 sc.spawn(move || s.claim_turn("t1").expect("claim"))
2871 })
2872 .collect();
2873 hs.into_iter().map(|h| h.join().expect("join")).collect()
2874 });
2875 assert_eq!(wins.iter().filter(|w| w.is_some()).count(), 1);
2876 }
2877
2878 fn age_file(path: &Path) {
2879 let f = std::fs::OpenOptions::new()
2880 .write(true)
2881 .open(path)
2882 .expect("open");
2883 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
2884 .expect("age");
2885 }
2886
2887 #[test]
2888 fn a_late_taker_of_a_broken_lock_cannot_disturb_its_replacement() {
2889 let (_tmp, a, _b) = lease_store();
2890 let lease = a.turn_path("t1");
2891 let lock = lease.with_extension("turn.lock");
2892 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2893 std::fs::write(&lock, "t1-dead").expect("dead lock");
2894 age_file(&lock);
2895 let b = TurnLock::take(&lease).expect("take").expect("b wins");
2897 let fresh = std::fs::read_to_string(&lock).expect("read");
2898 assert_eq!(fresh, b.token);
2899 let ticket = std::fs::read_dir(lock.parent().expect("dir"))
2902 .expect("dir")
2903 .flatten()
2904 .map(|e| e.path())
2905 .find(|p| p.to_string_lossy().ends_with(".break.0"))
2906 .expect("ticket");
2907 assert!(!create_exclusive(&ticket, "").expect("ticket"));
2908 assert!(TurnLock::take(&lease).expect("take").is_none());
2909 assert_eq!(std::fs::read_to_string(&lock).expect("read"), fresh);
2910 age_file(&lock);
2912 age_file(&ticket);
2913 std::mem::forget(b);
2914 let c = TurnLock::take(&lease)
2915 .expect("take")
2916 .expect("next generation");
2917 assert_ne!(c.token, fresh);
2918 }
2919
2920 #[test]
2921 fn a_live_ticket_blocks_and_a_stale_one_hands_over_to_the_next_generation() {
2922 let (_tmp, a, _b) = lease_store();
2923 let lease = a.turn_path("t1");
2924 let lock = lease.with_extension("turn.lock");
2925 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2926 std::fs::write(&lock, "t1-dead").expect("dead lock");
2927 age_file(&lock);
2928 let t0 = lock.with_extension("lock.t1-dead.break.0");
2929 assert!(create_exclusive(&t0, "").expect("ticket"));
2930 assert!(TurnLock::take(&lease).expect("take").is_none());
2932 assert_eq!(std::fs::read_to_string(&lock).expect("read"), "t1-dead");
2933 age_file(&t0);
2935 let c = TurnLock::take(&lease).expect("take").expect("generation 1");
2936 assert_eq!(std::fs::read_to_string(&lock).expect("read"), c.token);
2937 assert!(lock.with_extension("lock.t1-dead.break.1").exists());
2938 assert!(TurnLock::take(&lease).expect("take").is_none());
2939 }
2940
2941 #[test]
2942 fn exhausted_ticket_generations_recover_once_the_sweep_ages_them_out() {
2943 let (_tmp, a, _b) = lease_store();
2944 let lease = a.turn_path("t1");
2945 let lock = lease.with_extension("turn.lock");
2946 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2947 std::fs::write(&lock, "t1-dead").expect("dead lock");
2948 age_file(&lock);
2949 let tickets: Vec<_> = (0..TICKET_GENERATIONS)
2950 .map(|n| lock.with_extension(format!("lock.t1-dead.break.{n}")))
2951 .collect();
2952 for t in &tickets {
2953 assert!(create_exclusive(t, "").expect("ticket"));
2954 age_file(t);
2955 }
2956 assert!(TurnLock::take(&lease).expect("take").is_none());
2958 assert!(tickets.iter().all(|t| t.exists()));
2959 for t in &tickets {
2961 let f = std::fs::OpenOptions::new()
2962 .write(true)
2963 .open(t)
2964 .expect("open");
2965 f.set_modified(std::time::SystemTime::now() - TICKET_SWEEP_AGE * 2)
2966 .expect("age");
2967 }
2968 assert!(TurnLock::take(&lease).expect("take").is_none());
2969 assert!(TurnLock::take(&lease).expect("take").is_some());
2970 }
2971
2972 #[test]
2973 fn an_aged_empty_lock_is_broken() {
2974 let (_tmp, a, _b) = lease_store();
2975 let lease = a.turn_path("t1");
2976 let lock = lease.with_extension("turn.lock");
2977 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2978 std::fs::write(&lock, "").expect("empty lock");
2979 age_file(&lock);
2980 assert!(TurnLock::take(&lease).expect("take").is_some());
2981 }
2982
2983 #[test]
2984 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
2985 let (tmp, talks) = store();
2986 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2987 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2988
2989 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
2990 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
2991 let saved = talks.get(&talk.id).expect("reload queued talk");
2992 assert_eq!(saved.pending, "first\n\nsecond");
2993 assert_eq!(saved.pending_breaks, Some(vec!["first\n\n".len()]));
2994 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
2995
2996 let drained = drain(&mut talk, &talks).expect("drain");
2997 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
2998 let saved = talks.get(&talk.id).expect("reload drained talk");
2999 assert!(saved.pending.is_empty());
3000 assert_eq!(saved.turns.len(), 1);
3001 assert_eq!(saved.turns[0].body, "first\n\nsecond");
3002 assert_eq!(saved.turns[0].breaks, Some(vec!["first\n\n".len()]));
3003 assert_eq!(saved.pending_breaks, None);
3004 }
3005
3006 #[test]
3007 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
3008 let (tmp, talks) = store();
3009 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
3010 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3011 let attachment = Attachment {
3012 id: "a".repeat(32),
3013 name: "shot.png".to_owned(),
3014 mime: "image/png".to_owned(),
3015 bytes: 3,
3016 };
3017
3018 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
3019 assert!(
3020 edit_pending_text(
3021 &mut talk,
3022 &talks,
3023 "corrected",
3024 "first",
3025 std::slice::from_ref(&attachment.id),
3026 )
3027 .expect("edit")
3028 );
3029 let saved = talks.get(&talk.id).expect("reload edited draft");
3030 assert_eq!(saved.pending, "corrected");
3031 assert_eq!(saved.pending_attachments, vec![attachment]);
3032
3033 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
3034 assert!(
3035 !edit_pending_text(
3036 &mut talk,
3037 &talks,
3038 "stale edit",
3039 "corrected",
3040 &["a".repeat(32)],
3041 )
3042 .expect("stale edit is a conflict")
3043 );
3044 assert_eq!(
3045 talks.get(&talk.id).expect("reload after conflict").pending,
3046 "corrected\n\nlater"
3047 );
3048 assert!(
3049 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
3050 .expect("stale clear is a conflict")
3051 );
3052 assert_eq!(
3053 talks
3054 .get(&talk.id)
3055 .expect("reload after stale clear")
3056 .pending,
3057 "corrected\n\nlater"
3058 );
3059 }
3060
3061 #[tokio::test]
3062 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
3063 let (tmp, talks) = store();
3064 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
3065 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
3066 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3067 let id = running.id.clone();
3068 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
3069
3070 let response_talks = talks.clone();
3071 let response_cfg = cfg.clone();
3072 let reply = tokio::spawn(async move {
3073 respond(
3074 &response_talks.claim_turn(&running.id).unwrap().unwrap(),
3075 &mut running,
3076 &response_talks,
3077 &response_cfg,
3078 &first,
3079 )
3080 .await
3081 });
3082 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
3083
3084 let mut queued = talks.get(&id).expect("queued handle");
3085 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
3086 reply.await.expect("join").expect("reply");
3087
3088 let saved = talks.get(&id).expect("reload");
3089 assert_eq!(saved.pending, "next");
3090 assert_eq!(saved.pending_breaks, Some(Vec::new()));
3091 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
3092 }
3093
3094 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
3097 let calls = dir.join(format!("{id}.calls"));
3098 let script = format!(
3099 "#!/bin/sh\necho x >> '{}'\n{body}\n",
3100 calls.to_string_lossy()
3101 );
3102 let path = dir.join(format!("mock-{id}.sh"));
3103 std::fs::write(&path, script).expect("write mock");
3104 AgentSpec {
3105 id: id.to_owned(),
3106 kind: AgentKind::Command,
3107 model: None,
3108 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
3109 extra_args: Vec::new(),
3110 env: BTreeMap::new(),
3111 prompt_delivery: None,
3112 }
3113 }
3114
3115 fn calls(dir: &Path, id: &str) -> usize {
3116 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
3117 }
3118
3119 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
3120 let mut cfg = config(specs[0].clone());
3121 cfg.agents = specs;
3122 cfg.roles.chatter = Some(AgentChoice::Chain(
3123 ids.iter().map(|s| (*s).to_owned()).collect(),
3124 ));
3125 cfg
3126 }
3127
3128 #[tokio::test]
3129 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
3130 let (tmp, talks) = store();
3131 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3132 let b = counting_agent(tmp.path(), "b", "cat");
3133 let cfg = chain_config(vec![a, b], &["a", "b"]);
3134 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3135 assert_eq!(talk.agent, "a");
3136
3137 say(
3138 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3139 &mut talk,
3140 &talks,
3141 &cfg,
3142 "hello there",
3143 Vec::new(),
3144 )
3145 .await
3146 .expect("turn");
3147 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
3148 assert_eq!(calls(tmp.path(), "b"), 1);
3149 assert_eq!(talk.agent, "b", "the switch persists");
3150 assert!(talks.get(&talk.id).unwrap().agent == "b");
3151 let reply = talk.turns.last().unwrap();
3152 assert!(reply.body.contains("hello there"));
3153 assert!(
3154 reply.body.contains("magi task add --solo"),
3155 "a fresh seat gets the full briefing"
3156 );
3157 assert!(
3158 talk.turns
3159 .iter()
3160 .any(|t| t.body.contains("agent changed from a to b")),
3161 "the switch is noted"
3162 );
3163 }
3164
3165 #[tokio::test]
3166 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
3167 let (tmp, talks) = store();
3168 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3169 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
3170 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
3171 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3172
3173 let err = say(
3174 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3175 &mut talk,
3176 &talks,
3177 &cfg,
3178 "hi",
3179 Vec::new(),
3180 )
3181 .await
3182 .expect_err("every agent failed");
3183 assert!(err.to_string().contains("`a`"), "{err:#}");
3184 assert_eq!(calls(tmp.path(), "a"), 1);
3185 assert_eq!(calls(tmp.path(), "b"), 1);
3186 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
3187 }
3188
3189 #[test]
3190 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
3191 let (tmp, talks) = store();
3192 let b = counting_agent(tmp.path(), "b", "cat");
3193 let cfg = chain_config(vec![b], &["ghost", "b"]);
3194 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3195 assert_eq!(talk.agent, "b");
3196 }
3197
3198 #[tokio::test]
3199 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
3200 let (tmp, talks) = store();
3201 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3202 let b = counting_agent(tmp.path(), "b", "cat");
3203 let cfg = chain_config(vec![a, b], &["a", "b"]);
3204 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
3205 say(
3206 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3207 &mut talk,
3208 &talks,
3209 &cfg,
3210 "hi",
3211 Vec::new(),
3212 )
3213 .await
3214 .expect_err("a alone, and it fails");
3215 assert_eq!(calls(tmp.path(), "b"), 0);
3216 assert_eq!(talk.agent, "a");
3217 }
3218
3219 #[tokio::test]
3220 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
3221 let (tmp, talks) = store();
3222 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3223 let b = counting_agent(tmp.path(), "b", "cat");
3224 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
3225 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
3226 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
3227 say(
3228 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3229 &mut talk,
3230 &talks,
3231 &cfg,
3232 "hi",
3233 Vec::new(),
3234 )
3235 .await
3236 .expect_err("c alone, and it fails");
3237 assert_eq!(calls(tmp.path(), "b"), 0);
3238 }
3239
3240 #[test]
3241 fn briefing_names_the_operator_with_or_without_a_persona() {
3242 let plain = briefing(Path::new("/repo"), "en", false);
3243 let named = briefing_with(Path::new("/repo"), "en", false, None, Some("Commander"));
3244 assert!(named.starts_with(&plain));
3245 assert!(named.contains("# Addressing the operator"));
3246 assert!(named.contains("\"Commander\""));
3247 let rei = crate::persona::builtin_catalog()
3248 .into_iter()
3249 .find(|p| p.id == "rei")
3250 .expect("rei");
3251 let with = briefing_with(
3252 Path::new("/repo"),
3253 "en",
3254 false,
3255 Some(&rei),
3256 Some("Commander"),
3257 );
3258 assert!(with.contains("\"Commander\""));
3259 assert!(!with.contains("# Addressing the operator"));
3260 }
3261
3262 #[test]
3263 fn briefing_carries_a_persona_section_only_when_one_is_chosen() {
3264 let plain = briefing(Path::new("/repo"), "en", false);
3265 assert_eq!(
3266 plain,
3267 briefing_with(Path::new("/repo"), "en", false, None, None)
3268 );
3269 assert!(!plain.contains("Persona"));
3270 let rei = crate::persona::builtin_catalog()
3271 .into_iter()
3272 .find(|p| p.id == "rei")
3273 .expect("rei");
3274 let with = briefing_with(Path::new("/repo"), "en", false, Some(&rei), None);
3275 assert!(with.starts_with(&plain), "the plain briefing is untouched");
3276 assert!(with.contains("# Persona (tone only)"));
3277 assert!(with.contains("TONE ONLY"));
3278 assert!(with.contains("task ids"));
3279 assert!(with.contains("write policy"));
3280 assert!(with.contains("`magi task add`"));
3281 assert!(with.contains("Rei Ayanami"));
3282 }
3283
3284 #[test]
3285 fn a_talk_written_before_personas_still_loads() {
3286 let (tmp, talks) = store();
3287 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3288 let cfg = config(spec);
3289 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3290 let path = talks.root.join(format!("{}.json", talk.id));
3291 let mut v: serde_json::Value =
3292 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3293 v.as_object_mut().unwrap().remove("persona");
3294 v.as_object_mut().unwrap().remove("persona_dirty");
3295 std::fs::write(&path, v.to_string()).unwrap();
3296 let loaded = talks.get(&talk.id).expect("old record loads");
3297 assert_eq!(loaded.persona, "");
3298 assert!(!loaded.persona_dirty);
3299 }
3300
3301 #[tokio::test]
3302 async fn a_persona_switch_notes_marks_dirty_and_updates_the_next_turn_once() {
3303 let (tmp, talks) = store();
3304 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3305 let cfg = config(spec);
3306 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3307 let lease = talks.claim_turn(&talk.id).unwrap().unwrap();
3308 say(&lease, &mut talk, &talks, &cfg, "hello", Vec::new())
3309 .await
3310 .expect("first turn");
3311 assert!(!talk.turns[1].body.contains("Persona"));
3312
3313 assert!(switch_persona(&mut talk, &talks, "misato").expect("switch"));
3314 assert!(!switch_persona(&mut talk, &talks, "misato").expect("same"));
3315 assert!(talk.persona_dirty);
3316 assert!(talk.turns.last().unwrap().body.contains("persona changed"));
3317
3318 say(&lease, &mut talk, &talks, &cfg, "next", Vec::new())
3319 .await
3320 .expect("turn");
3321 let prompt = &talk.turns.last().unwrap().body;
3322 assert!(prompt.contains("# Persona update"), "{prompt}");
3323 assert!(prompt.contains("Misato Katsuragi"));
3324 assert!(!talk.persona_dirty, "cleared after a successful turn");
3325
3326 say(&lease, &mut talk, &talks, &cfg, "again", Vec::new())
3327 .await
3328 .expect("turn");
3329 assert!(!talk.turns.last().unwrap().body.contains("# Persona update"));
3330
3331 assert!(switch_persona(&mut talk, &talks, "default").expect("back"));
3332 assert_eq!(talk.persona, "");
3333 say(&lease, &mut talk, &talks, &cfg, "plain", Vec::new())
3334 .await
3335 .expect("turn");
3336 assert!(
3337 talk.turns
3338 .last()
3339 .unwrap()
3340 .body
3341 .contains("turned the persona off")
3342 );
3343
3344 assert!(switch_persona(&mut talk, &talks, "rei").expect("rei"));
3346 let mut no_sessions = cfg.clone();
3347 no_sessions.graph.sessions = false;
3348 say(&lease, &mut talk, &talks, &no_sessions, "one", Vec::new())
3349 .await
3350 .expect("turn");
3351 say(&lease, &mut talk, &talks, &no_sessions, "two", Vec::new())
3352 .await
3353 .expect("turn");
3354 let last = &talk.turns.last().unwrap().body;
3355 assert!(!talk.persona_dirty);
3356 assert!(last.contains("# Persona (tone only)"), "{last}");
3357 assert!(last.contains("Rei Ayanami"));
3358 }
3359
3360 #[tokio::test]
3361 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
3362 let (tmp, talks) = store();
3363 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3364 let cfg = config(spec);
3365 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3366
3367 say(
3368 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3369 &mut talk,
3370 &talks,
3371 &cfg,
3372 "what does the queue module do?",
3373 Vec::new(),
3374 )
3375 .await
3376 .expect("first turn");
3377 let first_prompt = &talk.turns[1].body;
3378 assert!(first_prompt.contains("magi task add --solo"));
3379 assert!(first_prompt.contains("what does the queue module do?"));
3380
3381 say(
3382 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3383 &mut talk,
3384 &talks,
3385 &cfg,
3386 "and how is it locked?",
3387 Vec::new(),
3388 )
3389 .await
3390 .expect("second turn");
3391 let second_prompt = &talk.turns[3].body;
3392 assert!(
3393 !second_prompt.contains("magi task add --solo"),
3394 "the briefing is sent once, not on every turn: {second_prompt}"
3395 );
3396 assert!(second_prompt.contains("and how is it locked?"));
3397 }
3398
3399 #[tokio::test]
3400 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
3401 let (tmp, talks) = store();
3402 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3403 let mut b = a.clone();
3404 b.id = "other".to_owned();
3405 let mut cfg = config(a.clone());
3406 cfg.agents.push(b.clone());
3407 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
3408 say(
3409 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3410 &mut talk,
3411 &talks,
3412 &cfg,
3413 "remember the walrus",
3414 Vec::new(),
3415 )
3416 .await
3417 .expect("first turn");
3418 let old_session = talk.seat.claude_session.clone();
3419 assert_eq!(talk.seat.turns, 1);
3420
3421 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
3422 assert_eq!(talk.agent, "other");
3423 assert_eq!(talk.seat.turns, 0);
3424 assert_eq!(talk.seat.agent, "other");
3425 assert_ne!(talk.seat.claude_session, old_session);
3426 let note = talk.turns.last().expect("note");
3427 assert_eq!(note.who, Who::Agent);
3428 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3429 assert!(note.body.contains("changed from"), "{}", note.body);
3430 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
3431
3432 let before = talk.turns.len();
3433 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
3434 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
3435
3436 say(
3437 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3438 &mut talk,
3439 &talks,
3440 &cfg,
3441 "what did I say?",
3442 Vec::new(),
3443 )
3444 .await
3445 .expect("turn after switch");
3446 let prompt = &talk.turns.last().expect("reply").body;
3447 assert!(prompt.contains("remember the walrus"), "{prompt}");
3448 assert!(prompt.contains("## magi"), "{prompt}");
3449 assert!(prompt.contains("what did I say?"), "{prompt}");
3450 }
3451
3452 #[tokio::test]
3453 async fn say_appends_the_operator_turn_then_the_agent_turn() {
3454 let (tmp, talks) = store();
3455 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3456 let cfg = config(spec);
3457 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3458
3459 say(
3460 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3461 &mut talk,
3462 &talks,
3463 &cfg,
3464 "can I rename this function?",
3465 Vec::new(),
3466 )
3467 .await
3468 .expect("say");
3469
3470 assert_eq!(talk.turns.len(), 2);
3471 assert_eq!(talk.turns[0].who, Who::Operator);
3472 assert_eq!(talk.turns[0].body, "can I rename this function?");
3473 assert_eq!(talk.turns[1].who, Who::Agent);
3474 assert_eq!(talk.turns[1].body, "go ahead");
3475 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
3476 }
3477
3478 #[tokio::test]
3479 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
3480 let (tmp, talks) = store();
3481 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
3482 let cfg = config(spec);
3483 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3484
3485 let err = say(
3486 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3487 &mut talk,
3488 &talks,
3489 &cfg,
3490 "check the tests",
3491 Vec::new(),
3492 )
3493 .await
3494 .expect_err("a turn with no answer is an error");
3495 assert!(err.to_string().contains("no answer"), "{err}");
3496
3497 let on_disk = talks.get(&talk.id).expect("get");
3498 assert_eq!(on_disk.turns.len(), 2);
3499 assert_eq!(on_disk.turns[0].body, "check the tests");
3500 let note = &on_disk.turns[1];
3501 assert_eq!(note.who, Who::Agent);
3502 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3503 assert!(note.body.contains("your message is saved"));
3504 }
3505
3506 #[tokio::test]
3512 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
3513 let (tmp, talks) = store();
3514 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3515 let cfg = config(spec);
3516 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3517
3518 let text =
3519 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
3520 failpoint::force_put_failures(PUT_RETRIES - 1);
3523 respond(
3524 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3525 &mut talk,
3526 &talks,
3527 &cfg,
3528 &text,
3529 )
3530 .await
3531 .expect("respond must survive a write failure its own retries can outlast");
3532
3533 assert_eq!(talk.turns.len(), 2);
3534 assert_eq!(talk.turns[1].who, Who::Agent);
3535 assert_eq!(talk.turns[1].body, "go ahead");
3536 let on_disk = talks.get(&talk.id).expect("get");
3537 assert_eq!(
3538 on_disk.turns, talk.turns,
3539 "the reply must reach disk despite the early write failures"
3540 );
3541 }
3542
3543 #[tokio::test]
3549 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
3550 let (tmp, talks) = store();
3551 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3552 let cfg = config(spec);
3553 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3554
3555 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3556 failpoint::force_put_failures(PUT_RETRIES);
3561 let err = respond(
3562 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3563 &mut talk,
3564 &talks,
3565 &cfg,
3566 &text,
3567 )
3568 .await
3569 .expect_err("a reply that cannot be saved must be reported, not swallowed");
3570 assert!(err.to_string().contains("could not be saved"), "{err}");
3571
3572 let on_disk = talks.get(&talk.id).expect("get");
3573 assert_eq!(
3574 on_disk.turns.len(),
3575 2,
3576 "the operator turn plus a visible note"
3577 );
3578 assert_eq!(on_disk.turns[0].body, "check the tests");
3579 let note = &on_disk.turns[1];
3580 assert_eq!(note.who, Who::Agent);
3581 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3582 assert!(
3583 note.body.contains("could not be saved"),
3584 "the operator must be told the reply is missing, not left staring \
3585 at a gap with no explanation: {}",
3586 note.body
3587 );
3588 assert_eq!(
3589 talk.turns, on_disk.turns,
3590 "the in-memory talk must match what actually landed on disk"
3591 );
3592
3593 let artifacts = talks.artifacts_of(&talk.id);
3596 let stash = std::fs::read_dir(&artifacts)
3597 .expect("artifacts dir")
3598 .filter_map(|e| e.ok())
3599 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
3600 .expect("a stash file for the lost reply");
3601 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
3602 assert_eq!(stashed, "go ahead");
3603
3604 assert_eq!(
3612 on_disk.seat.turns, 1,
3613 "the note's write must carry the turn the CLI actually took"
3614 );
3615 assert_eq!(
3616 on_disk.seat.claude_session, talk.seat.claude_session,
3617 "the session id handed to the CLI must survive the failed reply"
3618 );
3619 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
3620 assert!(
3621 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
3622 "the next turn must resume, not open the same session id twice"
3623 );
3624 }
3625
3626 #[tokio::test]
3631 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
3632 let (tmp, talks) = store();
3633 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3634 let cfg = config(spec);
3635 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3636
3637 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3638 failpoint::force_put_failures(PUT_RETRIES * 2);
3641 let err = respond(
3642 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3643 &mut talk,
3644 &talks,
3645 &cfg,
3646 &text,
3647 )
3648 .await
3649 .expect_err("neither the reply nor the note could be saved");
3650 assert!(err.to_string().contains("could not be saved"), "{err}");
3651
3652 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
3653 let on_disk = talks.get(&talk.id).expect("get");
3654 assert_eq!(on_disk.turns.len(), 1);
3655
3656 assert_eq!(
3666 on_disk.seat.turns, 0,
3667 "an unwritable file cannot record the turn the CLI took"
3668 );
3669 assert_eq!(
3670 talk.seat.turns, 1,
3671 "the in-memory seat still reports the turn the CLI actually took"
3672 );
3673 assert_eq!(
3674 on_disk.seat.claude_session, talk.seat.claude_session,
3675 "the session id was minted at `begin` and never changes here"
3676 );
3677 }
3678
3679 #[tokio::test]
3683 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
3684 let (tmp, talks) = store();
3685 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3686 let cfg = config(spec);
3687 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3688
3689 let att = talks
3690 .put_attachment(
3691 &talk.id,
3692 "image/png",
3693 "screenshot.png",
3694 b"pretend-png-bytes",
3695 )
3696 .expect("put attachment");
3697
3698 say(
3699 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3700 &mut talk,
3701 &talks,
3702 &cfg,
3703 "",
3704 vec![att.clone()],
3705 )
3706 .await
3707 .expect("an empty body with an attachment is still a turn");
3708
3709 let operator_turn = &talk.turns[0];
3710 assert_eq!(operator_turn.who, Who::Operator);
3711 assert_eq!(operator_turn.body, "");
3712 assert_eq!(operator_turn.attachments, vec![att.clone()]);
3713
3714 let prompt = &talk.turns[1].body;
3715 let expected_path = talks
3716 .attachments_dir(&talk.id)
3717 .join(format!("{}.png", att.id));
3718 assert!(
3719 prompt.contains(&expected_path.display().to_string()),
3720 "the agent must be told the attachment's absolute path: {prompt}"
3721 );
3722 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
3723 }
3724
3725 #[test]
3734 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
3735 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
3736 let att = Attachment {
3737 id: "0".repeat(32),
3738 name: "shot.png".to_owned(),
3739 mime: "image/png".to_owned(),
3740 bytes: 3,
3741 };
3742 let path = talks
3743 .attachment_path("some-talk-id", &att)
3744 .expect("a supported mime always yields a path");
3745 assert!(
3746 path.is_absolute(),
3747 "must be absolute even off a relative store root: {}",
3748 path.display()
3749 );
3750 }
3751
3752 #[tokio::test]
3753 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
3754 let (tmp, talks) = store();
3759 let slow = mock_agent(
3760 tmp.path(),
3761 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
3762 BTreeMap::new(),
3763 );
3764 let mut cfg = config(slow);
3765 cfg.graph.timeout_talk = 1;
3766 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3767
3768 let err = say(
3769 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3770 &mut talk,
3771 &talks,
3772 &cfg,
3773 "check the tests",
3774 Vec::new(),
3775 )
3776 .await
3777 .expect_err("a turn that never answers is an error");
3778 assert!(
3779 err.to_string().contains("did not answer within 1s"),
3780 "{err}"
3781 );
3782
3783 let on_disk = talks.get(&talk.id).expect("get");
3784 let note = on_disk.turns.last().expect("a note turn was recorded");
3785 assert!(
3786 note.body.contains("did not answer within 1s"),
3787 "the transcript must show the configured timeout: {}",
3788 note.body
3789 );
3790 }
3791
3792 #[test]
3793 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
3794 let (tmp, talks) = store();
3795 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3796 let cfg = config(spec);
3797 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3798
3799 close(&mut talk, &talks).expect("close");
3800 assert_eq!(talk.status, TalkStatus::Closed);
3801 close(&mut talk, &talks).expect("closing twice is not an error");
3802
3803 let err =
3804 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
3805 assert!(err.to_string().contains("closed"));
3806 let _ = &cfg; }
3808
3809 #[tokio::test]
3810 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3811 let (tmp, talks) = store();
3812 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3813 let cfg = config(spec);
3814 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3817
3818 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
3822 close(&mut closed_elsewhere, &talks).expect("close");
3823 assert_eq!(
3824 talks.get(&in_flight.id).expect("reread").status,
3825 TalkStatus::Closed,
3826 "the close landed on disk before the turn finished"
3827 );
3828
3829 assert_eq!(in_flight.status, TalkStatus::Open);
3833 respond(
3834 &talks.claim_turn(&in_flight.id).unwrap().unwrap(),
3835 &mut in_flight,
3836 &talks,
3837 &cfg,
3838 "one more question",
3839 )
3840 .await
3841 .expect("the turn itself still completes");
3842
3843 let on_disk = talks.get(&in_flight.id).expect("reread");
3844 assert_eq!(
3845 on_disk.status,
3846 TalkStatus::Closed,
3847 "a close must stick even when a turn that started before it finishes after it"
3848 );
3849 assert!(
3852 on_disk.turns.iter().any(|t| t.body == "here you go"),
3853 "the in-flight turn's own reply is still recorded: {:?}",
3854 on_disk.turns
3855 );
3856 }
3857
3858 #[test]
3859 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
3860 let (tmp, talks) = store();
3861 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3862 let cfg = config(spec);
3863 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3866
3867 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
3870 close(&mut closed_elsewhere, &talks).expect("close");
3871 assert_eq!(
3872 talks.get(&stale.id).expect("reread").status,
3873 TalkStatus::Closed,
3874 "the close landed on disk before record was called"
3875 );
3876
3877 assert_eq!(stale.status, TalkStatus::Open);
3881 let err = record(&mut stale, &talks, "still there?", Vec::new())
3882 .expect_err("a close that landed first must be honored, not overwritten");
3883 assert!(err.to_string().contains("closed"));
3884
3885 let on_disk = talks.get(&stale.id).expect("reread");
3886 assert_eq!(
3887 on_disk.status,
3888 TalkStatus::Closed,
3889 "record must not resurrect a conversation closed while its snapshot was stale"
3890 );
3891 assert!(
3892 on_disk.turns.is_empty(),
3893 "the rejected turn must not have been appended: {:?}",
3894 on_disk.turns
3895 );
3896 let _ = &cfg; }
3898
3899 #[test]
3900 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
3901 let (tmp, talks) = store();
3902 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3903 let cfg = config(spec);
3904 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3905
3906 let held = talks.guard().unwrap();
3910
3911 let talks2 = talks.clone();
3912 let id = talk.id.clone();
3913 let closing = std::thread::spawn(move || {
3914 let mut talk = talks2.get(&id).expect("get");
3915 close(&mut talk, &talks2).expect("close");
3916 });
3917
3918 std::thread::sleep(Duration::from_millis(50));
3919 assert!(
3920 !closing.is_finished(),
3921 "close must wait for the guard, not read and write while it is held - \
3922 a re-read alone narrows this window without closing it"
3923 );
3924
3925 drop(held);
3926 closing.join().expect("close thread panicked");
3927
3928 assert_eq!(
3929 talks.get(&talk.id).expect("reread").status,
3930 TalkStatus::Closed,
3931 "once the guard is free, close still lands"
3932 );
3933 let _ = &cfg; }
3935
3936 #[test]
3937 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
3938 let (tmp, talks) = store();
3939 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3940 let cfg = config(spec);
3941 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3942
3943 close(&mut talk, &talks).expect("close");
3944 assert_eq!(talk.status, TalkStatus::Closed);
3945
3946 reopen(&mut talk, &talks).expect("reopen");
3947 assert_eq!(talk.status, TalkStatus::Open);
3948 assert_eq!(
3949 talks.get(&talk.id).expect("reread").status,
3950 TalkStatus::Open
3951 );
3952
3953 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
3955 assert_eq!(talk.status, TalkStatus::Open);
3956
3957 record(&mut talk, &talks, "one more thing", Vec::new())
3958 .expect("a reopened talk takes turns again");
3959 let _ = &cfg; }
3961
3962 #[test]
3963 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
3964 let (tmp, talks) = store();
3965 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3966 let cfg = config(spec);
3967 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3968
3969 let artifacts = talks.artifacts_of(&talk.id);
3970 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
3971 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
3972
3973 talks.remove(&talk.id).expect("remove");
3974 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
3975 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
3976 assert!(
3977 talks.get(&talk.id).is_err(),
3978 "a removed talk cannot be read back"
3979 );
3980
3981 let err = talks
3982 .remove("nonexistent-id")
3983 .expect_err("unknown id refused");
3984 assert!(err.to_string().contains("no talk matches"), "{err}");
3985 let _ = &cfg; }
3987
3988 #[tokio::test]
3989 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3990 let (tmp, talks) = store();
3991 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3992 let cfg = config(spec);
3993 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3996
3997 talks.remove(&in_flight.id).expect("remove");
3998 assert!(
3999 talks.get(&in_flight.id).is_err(),
4000 "the delete landed on disk before the turn finished"
4001 );
4002
4003 respond(
4006 &talks.claim_turn(&in_flight.id).unwrap().unwrap(),
4007 &mut in_flight,
4008 &talks,
4009 &cfg,
4010 "one more question",
4011 )
4012 .await
4013 .expect("the turn itself still completes rather than erroring");
4014
4015 assert!(
4016 talks.get(&in_flight.id).is_err(),
4017 "a delete must stick even when a turn that started before it finishes after it"
4018 );
4019 }
4020
4021 #[test]
4022 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
4023 let (tmp, talks) = store();
4024 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4025 let cfg = config(spec);
4026 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4029
4030 talks.remove(&stale.id).expect("remove");
4031
4032 let err = record(&mut stale, &talks, "still there?", Vec::new())
4036 .expect_err("a delete that landed first must be honored, not overwritten");
4037 assert!(err.to_string().contains("deleted"), "{err}");
4038
4039 assert!(
4040 talks.get(&stale.id).is_err(),
4041 "record must not resurrect a conversation deleted while its snapshot was stale"
4042 );
4043 let _ = &cfg; }
4045
4046 #[test]
4047 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
4048 let (tmp, talks) = store();
4049 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4050 let cfg = config(spec);
4051 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4054
4055 talks.remove(&stale.id).expect("remove");
4056
4057 let err = close(&mut stale, &talks)
4061 .expect_err("a delete that landed first must be honored, not overwritten");
4062 assert!(err.to_string().contains("deleted"), "{err}");
4063
4064 assert!(
4065 talks.get(&stale.id).is_err(),
4066 "close must not resurrect a conversation deleted while its snapshot was stale"
4067 );
4068 let _ = &cfg; }
4070
4071 #[test]
4072 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
4073 let (tmp, talks) = store();
4074 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4075 let cfg = config(spec);
4076 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4079 close(&mut stale, &talks).expect("close");
4080
4081 talks.remove(&stale.id).expect("remove");
4082
4083 let err = reopen(&mut stale, &talks)
4087 .expect_err("a delete that landed first must be honored, not overwritten");
4088 assert!(err.to_string().contains("deleted"), "{err}");
4089
4090 assert!(
4091 talks.get(&stale.id).is_err(),
4092 "reopen must not resurrect a conversation deleted while its snapshot was stale"
4093 );
4094 let _ = &cfg; }
4096
4097 #[test]
4098 fn list_puts_open_talks_before_closed_ones() {
4099 let (tmp, talks) = store();
4100 let make = |id: &str, status: TalkStatus| {
4101 let mut t = Talk {
4102 pending_breaks: None,
4103 schema: SCHEMA,
4104 id: id.to_owned(),
4105 repo: tmp.path().to_owned(),
4106 agent: "mock".to_owned(),
4107 status,
4108 turns: Vec::new(),
4109 pending: String::new(),
4110 pending_attachments: Vec::new(),
4111 fallback: false,
4112 persona: String::new(),
4113 persona_dirty: false,
4114 created_at: Timestamp::now(),
4115 updated_at: Timestamp::now(),
4116 seat: SeatState::new(SEAT, "mock", 7),
4117 };
4118 talks.put(&mut t).expect("put");
4119 };
4120 make("20260901-000000-0001", TalkStatus::Open);
4121 make("20260902-000000-0002", TalkStatus::Open);
4122 make("20260903-000000-0003", TalkStatus::Closed);
4123
4124 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
4125 assert_eq!(
4126 ids,
4127 [
4128 "20260902-000000-0002",
4129 "20260901-000000-0001",
4130 "20260903-000000-0003"
4131 ]
4132 );
4133 assert_eq!(talks.count_open(), 2);
4134 }
4135
4136 #[test]
4137 fn tasks_of_finds_only_this_talks_own_tasks() {
4138 let dir = tempfile::tempdir().expect("tempdir");
4139 let queue = Queue::at(dir.path().join("queue"));
4140
4141 let mut mine = Task::new(
4142 "rework the loader".to_owned(),
4143 "rework the loader".to_owned(),
4144 PathBuf::from("/repo"),
4145 Source::Agent {
4146 run: "20260904-014455-ab12".to_owned(),
4147 node: "chat".to_owned(),
4148 },
4149 );
4150 queue.put(&mut mine).expect("put mine");
4151
4152 let mut theirs = Task::new(
4153 "unrelated".to_owned(),
4154 "unrelated".to_owned(),
4155 PathBuf::from("/repo"),
4156 Source::Agent {
4157 run: "20260904-090000-zz99".to_owned(),
4158 node: "implement".to_owned(),
4159 },
4160 );
4161 queue.put(&mut theirs).expect("put theirs");
4162
4163 let mut human = Task::new(
4164 "typed by hand".to_owned(),
4165 "typed by hand".to_owned(),
4166 PathBuf::from("/repo"),
4167 Source::Human,
4168 );
4169 queue.put(&mut human).expect("put human");
4170
4171 let found = tasks_of(&queue, "20260904-014455-ab12");
4172 assert_eq!(found.len(), 1);
4173 assert_eq!(found[0].id, mine.id);
4174 }
4175
4176 #[test]
4177 fn the_briefing_names_solo_task_add() {
4178 let brief = briefing(Path::new("/repo"), "en", false);
4179 assert!(brief.contains("magi task add --solo"));
4180 assert!(brief.contains("/repo"));
4181 assert!(!brief.contains("Hold this conversation in"));
4182 }
4183
4184 #[test]
4190 fn the_briefing_explains_targeting_a_different_repository_by_name() {
4191 let brief = briefing(Path::new("/repo"), "en", false);
4192 assert!(brief.contains("--repo does not have to be a full path"));
4193 assert!(brief.contains("owner/repo"));
4194 assert!(brief.contains("magi repos"));
4195 assert!(brief.contains("ask the operator"));
4196 }
4197
4198 #[test]
4199 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
4200 let brief = briefing(Path::new("/repo"), "en", false);
4201 assert!(brief.contains("--attach <path>"), "{brief}");
4202 assert!(brief.contains("deleting this conversation"), "{brief}");
4203 }
4204
4205 #[test]
4206 fn the_briefing_names_the_language_when_it_is_not_english() {
4207 let brief = briefing(Path::new("/repo"), "Japanese", false);
4208 assert!(brief.contains("Hold this conversation in Japanese"));
4209 }
4210
4211 #[test]
4212 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
4213 let read_only = briefing(Path::new("/repo"), "en", false);
4214 assert!(read_only.contains("Do not write files"));
4215 assert!(!read_only.contains("allow_write"));
4216
4217 let writable = briefing(Path::new("/repo"), "en", true);
4218 assert!(!writable.contains("Do not write files"));
4219 assert!(writable.contains("allow_write = true"));
4220 assert!(writable.contains("magi task add --solo"));
4223 assert!(writable.contains("say plainly what you"));
4224 }
4225}