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_for(
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 talk.implementers,
258 )
259 .chars()
260 .count() as u64
261 });
262 estimate_context_tokens(talk, standing)
263 });
264 let estimated = measured.is_none() && tokens.is_some();
265 let since_switch =
266 usage.is_some_and(|u| u.agent != talk.agent || (current.is_some() && u.model != model));
267 let (percent, warn) = match (tokens, window) {
268 (Some(t), Some(w)) => (
269 Some(t.saturating_mul(100) / w),
270 t.saturating_mul(100) >= w.saturating_mul(CONTEXT_WARN_PERCENT),
271 ),
272 _ => (None, false),
273 };
274 ContextUsage {
275 tokens,
276 window,
277 percent,
278 warn,
279 since_switch,
280 model,
281 estimated,
282 }
283}
284
285#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
289#[serde(rename_all = "lowercase")]
290pub enum TalkStatus {
291 Open,
294 Closed,
296}
297
298impl TalkStatus {
299 pub fn open(self) -> bool {
301 matches!(self, Self::Open)
302 }
303
304 pub fn as_str(self) -> &'static str {
306 match self {
307 Self::Open => "open",
308 Self::Closed => "closed",
309 }
310 }
311}
312
313#[derive(Debug, Clone, Serialize, Deserialize)]
315#[serde(deny_unknown_fields)]
316pub struct Talk {
317 pub schema: u32,
319 pub id: String,
321 pub repo: PathBuf,
323 pub agent: String,
325 pub status: TalkStatus,
327 pub turns: Vec<Turn>,
329 #[serde(default)]
332 pub pending: String,
333 #[serde(default)]
335 pub pending_attachments: Vec<Attachment>,
336 #[serde(default, skip_serializing_if = "Option::is_none")]
338 pub pending_breaks: Option<Vec<usize>>,
339 #[serde(default)]
344 pub fallback: bool,
345 #[serde(default)]
348 pub persona: String,
349 #[serde(default)]
353 pub persona_dirty: bool,
354 #[serde(default = "solo_implementers")]
357 pub implementers: u8,
358 #[serde(default)]
361 pub implementers_dirty: bool,
362 pub created_at: Timestamp,
364 pub updated_at: Timestamp,
366 seat: SeatState,
371}
372
373impl Talk {
374 pub fn short(&self) -> &str {
376 short(&self.id)
377 }
378}
379
380#[derive(Debug, Clone)]
382pub struct Talks {
383 root: PathBuf,
384 lock: Arc<Mutex<()>>,
393}
394
395impl Talks {
396 pub fn open() -> Self {
398 Self::at(crate::run::home().join("talks"))
399 }
400
401 pub fn at(root: PathBuf) -> Self {
404 Self {
405 root,
406 lock: Arc::new(Mutex::new(())),
407 }
408 }
409
410 fn guard(&self) -> Result<StoreGuard<'_>> {
425 let mutex = self.lock.lock().unwrap_or_else(PoisonError::into_inner);
426 std::fs::create_dir_all(&self.root)
427 .with_context(|| format!("create {}", self.root.display()))?;
428 let path = self.root.join(".store.turn");
429 let deadline = std::time::Instant::now() + TAKEOVER_LOCK_TTL * 2;
430 loop {
431 if let Some(file) = TurnLock::take(&path)? {
432 return Ok(StoreGuard {
433 _file: file,
434 _mutex: mutex,
435 });
436 }
437 if std::time::Instant::now() >= deadline {
438 bail!("the talk store is locked by another process");
439 }
440 std::thread::sleep(Duration::from_millis(10));
441 }
442 }
443
444 pub fn root(&self) -> &Path {
446 &self.root
447 }
448
449 pub fn path_of(&self, id: &str) -> PathBuf {
451 self.root.join(format!("{id}.json"))
452 }
453
454 pub fn artifacts_of(&self, id: &str) -> PathBuf {
457 self.root.join(format!("{id}.artifacts"))
458 }
459
460 pub fn attachments_dir(&self, id: &str) -> PathBuf {
464 self.artifacts_of(id).join("attachments")
465 }
466
467 pub fn put_attachment(
476 &self,
477 id: &str,
478 mime: &str,
479 name: &str,
480 data: &[u8],
481 ) -> Result<Attachment> {
482 let dir = self.attachments_dir(id);
483 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
484 let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
485 let att = Attachment {
486 id: new_attachment_id(),
487 name: name.to_owned(),
488 mime: mime.to_owned(),
489 bytes: data.len() as u64,
490 };
491 std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
492 .with_context(|| format!("write attachment {}", att.id))?;
493 std::fs::write(
494 dir.join(format!("{}.json", att.id)),
495 serde_json::to_string(&att).context("serialize attachment")?,
496 )
497 .with_context(|| format!("write attachment metadata {}", att.id))?;
498 Ok(att)
499 }
500
501 pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
510 if !valid_attachment_id(att_id) {
511 return Ok(None);
512 }
513 let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
514 if !meta_path.is_file() {
515 return Ok(None);
516 }
517 let att = serde_json::from_str(
518 &std::fs::read_to_string(&meta_path)
519 .with_context(|| format!("read {}", meta_path.display()))?,
520 )
521 .with_context(|| format!("parse {}", meta_path.display()))?;
522 Ok(Some(att))
523 }
524
525 pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
529 let Some(att) = self.attachment_meta(id, att_id)? else {
530 return Ok(None);
531 };
532 let ext = attachment_ext(&att.mime).with_context(|| {
533 format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
534 })?;
535 let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
536 let data =
537 std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
538 Ok(Some((att, data)))
539 }
540
541 fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
558 let ext = attachment_ext(&att.mime)?;
559 let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
560 std::path::absolute(&path).ok()
561 }
562
563 pub fn put(&self, t: &mut Talk) -> Result<()> {
572 std::fs::create_dir_all(&self.root)
573 .with_context(|| format!("create {}", self.root.display()))?;
574 t.updated_at = Timestamp::now();
575 let body = serde_json::to_string_pretty(t).context("serialize talk")?;
576 let path = self.path_of(&t.id);
577 let tmp = path.with_extension("json.tmp");
578 write_atomic(&tmp, &path, &body)
579 }
580
581 pub fn get(&self, id: &str) -> Result<Talk> {
583 let resolved = self.resolve_id(id)?;
584 read_path(&self.path_of(&resolved))
585 }
586
587 pub fn list(&self) -> Vec<Talk> {
590 self.list_counting_unreadable().0
591 }
592
593 pub fn list_counting_unreadable(&self) -> (Vec<Talk>, usize) {
596 let mut unreadable = 0;
597 let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
598 .into_iter()
599 .flatten()
600 .flatten()
601 .map(|e| e.path())
602 .filter(|p| p.extension().is_some_and(|x| x == "json"))
603 .filter_map(|p| {
604 let talk = read_path(&p).ok();
605 if talk.is_none() {
606 unreadable += 1;
607 }
608 talk
609 })
610 .collect();
611 all.sort_unstable_by(|a, b| {
612 let rank = |t: &Talk| u8::from(!t.status.open());
613 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
614 });
615 (all, unreadable)
616 }
617
618 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
620 if self.path_of(prefix).is_file() {
621 return Ok(prefix.to_owned());
622 }
623 let hits: Vec<String> = self
624 .list()
625 .into_iter()
626 .map(|t| t.id)
627 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
628 .collect();
629 match hits.len() {
630 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
631 0 => bail!("no talk matches `{prefix}`"),
632 _ => bail!(
633 "`{prefix}` matches {} talks: {}",
634 hits.len(),
635 hits.join(", ")
636 ),
637 }
638 }
639
640 pub fn revision(&self) -> u64 {
643 std::fs::read_dir(&self.root)
644 .into_iter()
645 .flatten()
646 .flatten()
647 .filter(|e| e.path().extension().is_none_or(|x| x != "turn"))
648 .filter(|e| !e.file_name().to_string_lossy().starts_with('.'))
649 .filter_map(|e| e.metadata().ok())
650 .filter_map(|m| m.modified().ok())
651 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
652 .map(|d| d.as_millis() as u64)
653 .max()
654 .unwrap_or(0)
655 }
656
657 pub fn count_open(&self) -> usize {
659 self.list().iter().filter(|t| t.status.open()).count()
660 }
661
662 pub fn turn_path(&self, id: &str) -> PathBuf {
665 self.root.join(format!("{id}.turn"))
666 }
667
668 pub fn claim_turn(&self, id: &str) -> Result<Option<TurnLease>> {
677 self.claim_turn_at(id, Timestamp::now())
678 }
679
680 fn claim_turn_at(&self, id: &str, now: Timestamp) -> Result<Option<TurnLease>> {
681 std::fs::create_dir_all(&self.root)
682 .with_context(|| format!("create {}", self.root.display()))?;
683 let path = self.turn_path(id);
684 let token = fresh_token();
685 if create_turn(&path, &token, now)? {
686 return Ok(Some(TurnLease {
687 talk: id.to_owned(),
688 path,
689 token,
690 }));
691 }
692 if lease_blocks(&path, now) {
693 return Ok(None);
694 }
695 let Some(_lock) = TurnLock::take(&path)? else {
700 return Ok(None);
701 };
702 if lease_blocks(&path, now) {
703 return Ok(None);
704 }
705 let _ = std::fs::remove_file(&path);
706 Ok(create_turn(&path, &token, now)?.then_some(TurnLease {
707 talk: id.to_owned(),
708 path,
709 token,
710 }))
711 }
712
713 pub fn turn_held(&self, id: &str) -> bool {
715 read_turn(&self.turn_path(id)).is_some_and(|r| r.fresh(Timestamp::now()))
716 }
717
718 pub fn remove(&self, id: &str) -> Result<()> {
731 let _guard = self.guard()?;
732 let resolved = self.resolve_id(id)?;
733 let path = self.path_of(&resolved);
734 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
735 let artifacts = self.artifacts_of(&resolved);
736 if artifacts.is_dir() {
737 std::fs::remove_dir_all(&artifacts)
738 .with_context(|| format!("remove {}", artifacts.display()))?;
739 }
740 let _ = std::fs::remove_file(self.turn_path(&resolved));
741 Ok(())
742 }
743}
744
745struct StoreGuard<'a> {
748 _file: TurnLock,
749 _mutex: MutexGuard<'a, ()>,
750}
751
752#[derive(Debug, Serialize, Deserialize)]
755struct TurnRecord {
756 token: String,
757 pid: u32,
758 beat_at: Timestamp,
759}
760
761impl TurnRecord {
762 fn fresh(&self, now: Timestamp) -> bool {
763 now.as_second() - self.beat_at.as_second() <= crate::ask::LEASE_TTL.as_secs() as i64
764 }
765}
766
767fn fresh_token() -> String {
771 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
772 let n = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
773 let seed = crate::rng::entropy() ^ n.wrapping_mul(0x9E37_79B9_7F4A_7C15);
774 crate::rng::SplitMix64::new(seed).uuid_v4()
775}
776
777fn read_turn(path: &Path) -> Option<TurnRecord> {
778 serde_json::from_str(&std::fs::read_to_string(path).ok()?).ok()
779}
780
781const TAKEOVER_LOCK_TTL: Duration = Duration::from_secs(10);
783
784const TICKET_TTL: Duration = Duration::from_secs(10);
787
788const TICKET_GENERATIONS: u32 = 16;
790
791const TICKET_SWEEP_AGE: Duration = Duration::from_secs(3600);
793
794fn create_exclusive(path: &Path, body: &str) -> Result<bool> {
796 use std::io::Write as _;
797 match std::fs::OpenOptions::new()
798 .write(true)
799 .create_new(true)
800 .open(path)
801 {
802 Ok(mut f) => {
803 if let Err(e) = f.write_all(body.as_bytes()) {
804 drop(f);
805 let _ = std::fs::remove_file(path);
806 return Err(e).with_context(|| format!("write {}", path.display()));
807 }
808 Ok(true)
809 }
810 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
811 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
812 }
813}
814
815fn create_turn(path: &Path, token: &str, now: Timestamp) -> Result<bool> {
816 let record = TurnRecord {
817 token: token.to_owned(),
818 pid: std::process::id(),
819 beat_at: now,
820 };
821 let body = serde_json::to_string(&record).context("serialize turn lease")?;
822 let tmp = path.with_extension(format!("turn.{token}.new"));
825 publish_exclusive(path, &tmp, &body)
826}
827
828fn lease_blocks(path: &Path, now: Timestamp) -> bool {
832 if read_turn(path).is_some_and(|r| r.fresh(now)) {
833 return true;
834 }
835 std::fs::metadata(path).is_ok_and(|m| {
836 m.len() == 0
837 && m.modified()
838 .ok()
839 .and_then(|t| t.elapsed().ok())
840 .is_some_and(|age| age <= TAKEOVER_LOCK_TTL)
841 })
842}
843
844fn link_unsupported(e: &std::io::Error) -> bool {
847 e.kind() == std::io::ErrorKind::Unsupported || (cfg!(windows) && e.raw_os_error() == Some(1))
848}
849
850fn create_in_place(path: &Path, body: &str) -> Result<bool> {
868 use std::io::Write as _;
869 let mut f = match std::fs::OpenOptions::new()
870 .write(true)
871 .create_new(true)
872 .open(path)
873 {
874 Ok(f) => f,
875 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => return Ok(false),
876 Err(e) => return Err(e).with_context(|| format!("create {}", path.display())),
877 };
878 f.write_all(body.as_bytes())
879 .with_context(|| format!("write {}", path.display()))?;
880 Ok(std::fs::read_to_string(path).is_ok_and(|t| t == body))
883}
884
885fn publish_exclusive(path: &Path, tmp: &Path, body: &str) -> Result<bool> {
886 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
887 #[cfg(test)]
888 let linked = if failpoint::no_link_forced() {
889 Err(std::io::Error::from(std::io::ErrorKind::Unsupported))
890 } else {
891 std::fs::hard_link(tmp, path)
892 };
893 #[cfg(not(test))]
894 let linked = std::fs::hard_link(tmp, path);
895 let out = match linked {
896 Ok(()) => Ok(true),
897 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
898 Err(e) if link_unsupported(&e) => create_in_place(path, body),
899 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
900 };
901 let _ = std::fs::remove_file(tmp);
902 out
903}
904
905fn remove_if_carries(path: &Path, expected: &str) -> bool {
916 let gone = path.with_extension(format!("gone.{}", fresh_token()));
917 if std::fs::rename(path, &gone).is_err() {
918 return false;
919 }
920 let found = std::fs::read_to_string(&gone);
921 if found.as_ref().is_ok_and(|t| t == expected) {
922 let _ = std::fs::remove_file(&gone);
923 return true;
924 }
925 if let Ok(body) = found {
926 let back = path.with_extension(format!("back.{}", fresh_token()));
927 match publish_exclusive(path, &back, &body) {
928 Ok(true) => {}
929 Ok(false) | Err(_) => tracing::warn!(
930 "{} was replaced while a stale removal had it aside; \
931 the displaced file is dropped",
932 path.display()
933 ),
934 }
935 }
936 let _ = std::fs::remove_file(&gone);
937 false
938}
939
940struct TurnLock {
944 path: PathBuf,
945 token: String,
946}
947
948impl TurnLock {
949 fn token() -> String {
950 fresh_token()
951 }
952
953 fn publish(path: &Path, token: &str) -> Result<bool> {
957 let tmp = path.with_extension(format!("lock.{token}.new"));
958 publish_exclusive(path, &tmp, token)
959 }
960
961 fn take(lease: &Path) -> Result<Option<Self>> {
962 let path = lease.with_extension("turn.lock");
963 let token = Self::token();
964 if Self::publish(&path, &token)? {
965 return Ok(Some(Self { path, token }));
966 }
967 let Ok(seen) = std::fs::read_to_string(&path) else {
968 return Ok(None);
969 };
970 let key = if !seen.is_empty() && seen.chars().all(|c| c.is_ascii_alphanumeric() || c == '-')
973 {
974 seen.as_str()
975 } else {
976 "invalid"
977 };
978 let aged = std::fs::metadata(&path)
979 .and_then(|m| m.modified())
980 .ok()
981 .and_then(|t| t.elapsed().ok())
982 .is_some_and(|age| age > TAKEOVER_LOCK_TTL);
983 if !aged {
984 return Ok(None);
985 }
986 let mut won = false;
998 for n in 0..TICKET_GENERATIONS {
999 let ticket = path.with_extension(format!("lock.{key}.break.{n}"));
1000 if create_exclusive(&ticket, "")? {
1001 won = true;
1002 break;
1003 }
1004 let stale = std::fs::metadata(&ticket)
1005 .and_then(|m| m.modified())
1006 .ok()
1007 .and_then(|t| t.elapsed().ok())
1008 .is_some_and(|age| age > TICKET_TTL);
1009 if !stale {
1010 return Ok(None);
1011 }
1012 }
1013 if !won {
1014 Self::sweep_tickets(&path);
1019 return Ok(None);
1020 }
1021 Self::sweep_tickets(&path);
1022 if !remove_if_carries(&path, &seen) {
1028 return Ok(None);
1029 }
1030 if Self::publish(&path, &token)? {
1031 return Ok(Some(Self { path, token }));
1032 }
1033 Ok(None)
1034 }
1035
1036 fn sweep_tickets(path: &Path) {
1039 let (Some(dir), Some(name)) = (path.parent(), path.file_name().and_then(|n| n.to_str()))
1040 else {
1041 return;
1042 };
1043 let prefix = format!("{name}.");
1044 let Ok(entries) = std::fs::read_dir(dir) else {
1045 return;
1046 };
1047 for entry in entries.flatten() {
1048 let file = entry.file_name();
1049 let Some(file) = file.to_str() else { continue };
1050 if !(file.starts_with(&prefix) && file.contains(".break.")) {
1051 continue;
1052 }
1053 let old = entry
1054 .metadata()
1055 .and_then(|m| m.modified())
1056 .ok()
1057 .and_then(|t| t.elapsed().ok())
1058 .is_some_and(|age| age > TICKET_SWEEP_AGE);
1059 if old {
1060 let _ = std::fs::remove_file(entry.path());
1061 }
1062 }
1063 }
1064
1065 fn take_patiently(lease: &Path) -> Option<Self> {
1067 for _ in 0..50 {
1068 match Self::take(lease) {
1069 Ok(Some(lock)) => return Some(lock),
1070 Ok(None) => std::thread::sleep(Duration::from_millis(10)),
1071 Err(_) => return None,
1072 }
1073 }
1074 None
1075 }
1076}
1077
1078impl Drop for TurnLock {
1079 fn drop(&mut self) {
1080 remove_if_carries(&self.path, &self.token);
1082 }
1083}
1084
1085pub const TURN_BEAT: Duration = Duration::from_secs(20);
1088
1089#[derive(Debug)]
1094pub struct TurnLease {
1095 talk: String,
1096 path: PathBuf,
1097 token: String,
1098}
1099
1100impl TurnLease {
1101 pub fn holds(&self) -> bool {
1103 read_turn(&self.path).is_some_and(|r| r.token == self.token)
1104 }
1105
1106 pub fn beat(&self) -> Result<bool> {
1110 let _lock = TurnLock::take_patiently(&self.path)
1111 .with_context(|| format!("lock {} to renew it", self.path.display()))?;
1112 let Some(mut record) = read_turn(&self.path).filter(|r| r.token == self.token) else {
1113 return Ok(false);
1114 };
1115 record.beat_at = Timestamp::now();
1116 let body = serde_json::to_string(&record).context("serialize turn lease")?;
1117 let tmp = self.path.with_extension(format!("turn.{}.tmp", self.token));
1118 write_atomic(&tmp, &self.path, &body)?;
1119 Ok(true)
1120 }
1121
1122 pub async fn beating<T>(&self, fut: impl std::future::Future<Output = T>) -> Result<T> {
1127 self.beating_every(TURN_BEAT, fut).await
1128 }
1129
1130 async fn beating_every<T>(
1131 &self,
1132 period: Duration,
1133 fut: impl std::future::Future<Output = T>,
1134 ) -> Result<T> {
1135 tokio::pin!(fut);
1136 loop {
1137 match tokio::time::timeout(period, &mut fut).await {
1138 Ok(out) => return Ok(out),
1139 Err(_) => match self.beat() {
1140 Ok(true) => {}
1141 Ok(false) => bail!(
1142 "the turn lease {} was taken over; this turn is stopped",
1143 self.path.display()
1144 ),
1145 Err(e) => tracing::warn!("{e:#}"),
1146 },
1147 }
1148 }
1149 }
1150}
1151
1152impl Drop for TurnLease {
1153 fn drop(&mut self) {
1154 if let Some(_lock) = TurnLock::take_patiently(&self.path) {
1157 if read_turn(&self.path).is_some_and(|r| r.token == self.token) {
1158 let _ = std::fs::remove_file(&self.path);
1159 }
1160 }
1161 }
1162}
1163
1164pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
1174 let repo = repo.canonicalize().unwrap_or(repo);
1177 let spec = match agent {
1180 Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
1181 None => agent::pick_chain(
1182 &cfg.agents,
1183 cfg.roles.chatter.as_ref(),
1184 &agent::installed,
1185 "chatter",
1186 )?
1187 .remove(0),
1188 };
1189
1190 let now = Timestamp::now();
1191 let mut talk = Talk {
1192 pending_breaks: None,
1193 schema: SCHEMA,
1194 id: new_id(),
1195 repo,
1196 agent: spec.id.clone(),
1197 status: TalkStatus::Open,
1198 turns: Vec::new(),
1199 pending: String::new(),
1200 pending_attachments: Vec::new(),
1201 fallback: agent.is_none(),
1202 persona: String::new(),
1203 persona_dirty: false,
1204 implementers: 1,
1205 implementers_dirty: false,
1206 created_at: now,
1207 updated_at: now,
1208 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
1209 };
1210 store.put(&mut talk)?;
1211 Ok(talk)
1212}
1213
1214pub fn record(
1221 talk: &mut Talk,
1222 store: &Talks,
1223 text: &str,
1224 attachments: Vec<Attachment>,
1225) -> Result<String> {
1226 let _guard = store.guard()?;
1234 let Ok(fresh) = store.get(&talk.id) else {
1239 bail!("talk {} was deleted", talk.short());
1240 };
1241 talk.status = fresh.status;
1242 talk.pending = fresh.pending;
1245 talk.pending_breaks = fresh.pending_breaks;
1246 talk.pending_attachments = fresh.pending_attachments;
1247 if !talk.status.open() {
1248 bail!(
1249 "talk {} is {} and takes no more turns",
1250 talk.short(),
1251 talk.status.as_str()
1252 );
1253 }
1254 let text = text.trim();
1255 if text.is_empty() && attachments.is_empty() {
1256 bail!("nothing to say");
1257 }
1258 talk.turns.push(Turn {
1259 who: Who::Operator,
1260 body: text.to_owned(),
1261 at: Timestamp::now(),
1262 attachments,
1263 usage: None,
1264 breaks: Some(Vec::new()),
1265 });
1266 store.put(talk)?;
1267 Ok(text.to_owned())
1268}
1269
1270pub fn queue(
1272 talk: &mut Talk,
1273 store: &Talks,
1274 text: &str,
1275 attachments: Vec<Attachment>,
1276) -> Result<()> {
1277 let text = text.trim();
1278 if text.is_empty() && attachments.is_empty() {
1279 bail!("nothing to say");
1280 }
1281 let _guard = store.guard()?;
1282 let mut fresh = store
1283 .get(&talk.id)
1284 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1285 if !fresh.status.open() {
1286 bail!(
1287 "talk {} is {} and takes no more turns",
1288 fresh.short(),
1289 fresh.status.as_str()
1290 );
1291 }
1292 if !text.is_empty() {
1293 if fresh.pending.is_empty() {
1294 fresh.pending = text.to_owned();
1295 fresh.pending_breaks = Some(Vec::new());
1296 } else {
1297 fresh.pending.push_str("\n\n");
1298 if let Some(b) = fresh.pending_breaks.as_mut() {
1299 b.push(fresh.pending.len());
1300 }
1301 fresh.pending.push_str(text);
1302 }
1303 }
1304 fresh.pending_attachments.extend(attachments);
1305 store.put(&mut fresh)?;
1306 *talk = fresh;
1307 Ok(())
1308}
1309
1310pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
1312 let _guard = store.guard()?;
1313 let mut fresh = store
1314 .get(&talk.id)
1315 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1316 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
1317 *talk = fresh;
1318 return Ok(None);
1319 }
1320 let text = std::mem::take(&mut fresh.pending);
1321 let attachments = std::mem::take(&mut fresh.pending_attachments);
1322 let breaks = fresh.pending_breaks.take();
1323 fresh.turns.push(Turn {
1324 who: Who::Operator,
1325 body: text.clone(),
1326 at: Timestamp::now(),
1327 attachments,
1328 usage: None,
1329 breaks,
1330 });
1331 store.put(&mut fresh)?;
1332 *talk = fresh;
1333 Ok(Some(text))
1334}
1335
1336pub async fn say(
1339 lease: &TurnLease,
1340 talk: &mut Talk,
1341 store: &Talks,
1342 cfg: &Config,
1343 text: &str,
1344 attachments: Vec<Attachment>,
1345) -> Result<()> {
1346 check_lease(lease, talk)?;
1347 let text = record(talk, store, text, attachments)?;
1348 respond(lease, talk, store, cfg, &text).await
1349}
1350
1351pub async fn respond(
1357 lease: &TurnLease,
1358 talk: &mut Talk,
1359 store: &Talks,
1360 cfg: &Config,
1361 text: &str,
1362) -> Result<()> {
1363 check_lease(lease, talk)?;
1364 lease
1365 .beating(turn(talk, store, cfg, text))
1366 .await
1367 .and_then(|done| done)
1368}
1369
1370fn check_lease(lease: &TurnLease, talk: &Talk) -> Result<()> {
1371 if lease.talk != talk.id {
1372 bail!("the turn lease is for talk {}, not {}", lease.talk, talk.id);
1373 }
1374 if !lease.holds() {
1377 bail!("the turn lease for talk {} is no longer held", talk.short());
1378 }
1379 Ok(())
1380}
1381
1382pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
1400 let _guard = store.guard()?;
1401 let mut fresh = store
1402 .get(&talk.id)
1403 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1404 fresh.status = TalkStatus::Closed;
1405 fresh.pending.clear();
1407 fresh.pending_breaks = None;
1408 fresh.pending_attachments.clear();
1409 store.put(&mut fresh)?;
1410 *talk = fresh;
1411 Ok(())
1412}
1413
1414pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
1425 let _guard = store.guard()?;
1426 let mut fresh = store
1427 .get(&talk.id)
1428 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1429 fresh.status = TalkStatus::Open;
1430 store.put(&mut fresh)?;
1431 *talk = fresh;
1432 Ok(())
1433}
1434
1435pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
1448 let _guard = store.guard()?;
1449 let mut fresh = store
1450 .get(&talk.id)
1451 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1452 if fresh.agent == spec.id {
1453 *talk = fresh;
1454 return Ok(false);
1455 }
1456 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
1457 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1458 fresh.fallback = false;
1460 fresh.turns.push(Turn {
1461 breaks: None,
1462 who: Who::Agent,
1463 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
1464 at: Timestamp::now(),
1465 attachments: Vec::new(),
1466 usage: None,
1467 });
1468 store.put(&mut fresh)?;
1469 *talk = fresh;
1470 Ok(true)
1471}
1472
1473pub fn switch_persona(talk: &mut Talk, store: &Talks, id: &str) -> Result<bool> {
1481 let id = if id.trim() == crate::persona::DEFAULT_ID {
1482 ""
1483 } else {
1484 id.trim()
1485 };
1486 let _guard = store.guard()?;
1487 let mut fresh = store
1488 .get(&talk.id)
1489 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1490 if fresh.persona == id {
1491 *talk = fresh;
1492 return Ok(false);
1493 }
1494 let from = std::mem::replace(&mut fresh.persona, id.to_owned());
1495 fresh.persona_dirty = true;
1496 let label = |p: &str| {
1497 if p.is_empty() {
1498 crate::persona::DEFAULT_ID.to_owned()
1499 } else {
1500 p.to_owned()
1501 }
1502 };
1503 fresh.turns.push(Turn {
1504 breaks: None,
1505 who: Who::Agent,
1506 body: format!(
1507 "{MAGI_NOTE}persona changed from {} to {}",
1508 label(&from),
1509 label(id)
1510 ),
1511 at: Timestamp::now(),
1512 attachments: Vec::new(),
1513 usage: None,
1514 });
1515 store.put(&mut fresh)?;
1516 *talk = fresh;
1517 Ok(true)
1518}
1519
1520pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
1522 let _guard = store.guard()?;
1523 let mut fresh = store
1524 .get(&talk.id)
1525 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1526 fresh.pending.clear();
1527 fresh.pending_breaks = None;
1528 fresh.pending_attachments.clear();
1529 store.put(&mut fresh)?;
1530 *talk = fresh;
1531 Ok(())
1532}
1533
1534pub fn clear_pending_if_matches(
1536 talk: &mut Talk,
1537 store: &Talks,
1538 expected_text: &str,
1539 expected_attachments: &[String],
1540) -> Result<bool> {
1541 let _guard = store.guard()?;
1542 let mut fresh = store
1543 .get(&talk.id)
1544 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1545 if !pending_matches(&fresh, expected_text, expected_attachments) {
1546 *talk = fresh;
1547 return Ok(false);
1548 }
1549 fresh.pending.clear();
1550 fresh.pending_breaks = None;
1551 fresh.pending_attachments.clear();
1552 store.put(&mut fresh)?;
1553 *talk = fresh;
1554 Ok(true)
1555}
1556
1557pub fn edit_pending_text(
1561 talk: &mut Talk,
1562 store: &Talks,
1563 text: &str,
1564 expected_text: &str,
1565 expected_attachments: &[String],
1566) -> Result<bool> {
1567 let _guard = store.guard()?;
1568 let mut fresh = store
1569 .get(&talk.id)
1570 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1571 if !pending_matches(&fresh, expected_text, expected_attachments) {
1572 *talk = fresh;
1573 return Ok(false);
1574 }
1575 fresh.pending = text.trim().to_owned();
1576 fresh.pending_breaks = Some(Vec::new());
1577 store.put(&mut fresh)?;
1578 *talk = fresh;
1579 Ok(true)
1580}
1581
1582fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
1583 talk.pending == expected_text
1584 && talk
1585 .pending_attachments
1586 .iter()
1587 .map(|attachment| &attachment.id)
1588 .eq(expected_attachments.iter())
1589}
1590
1591async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1598 let spec = cfg
1599 .agents
1600 .iter()
1601 .find(|a| a.id == talk.agent)
1602 .with_context(|| {
1603 format!(
1604 "talk {} was opened with agent `{}`, which is no longer in \
1605 the roster; restore it in magi.toml or start a new \
1606 conversation",
1607 talk.short(),
1608 talk.agent
1609 )
1610 })?;
1611
1612 let last_note = attachment_note(
1616 store,
1617 &talk.id,
1618 talk.turns
1619 .last()
1620 .map_or(&[][..], |t| t.attachments.as_slice()),
1621 );
1622
1623 let attachment_paths: Vec<PathBuf> = talk
1629 .turns
1630 .iter()
1631 .flat_map(|t| t.attachments.iter())
1632 .filter_map(|a| store.attachment_path(&talk.id, a))
1633 .collect();
1634
1635 let questions =
1645 crate::run::try_home().map(|home| crate::ask::Questions::at(home.join("questions")));
1646 let consulted = questions
1647 .as_ref()
1648 .is_some_and(|q| crate::consult::pending_consults(q, &talk.id));
1649 let consult_roots: Vec<PathBuf> = match &questions {
1650 Some(q) if consulted => vec![q.root().to_path_buf()],
1651 _ => Vec::new(),
1652 };
1653
1654 let (allow_write, unsandboxed) = turn_access(cfg.talk.allow_write, consulted);
1655
1656 let artifacts = store.artifacts_of(&talk.id);
1657 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
1660 let stem = format!("turn-{}", operator_turns.max(1));
1661 let cache_dir = cfg.cache_dir();
1664
1665 let mut chain = vec![spec.clone()];
1668 if let Some(choice) = cfg.roles.chatter.as_ref()
1669 && talk.fallback
1670 {
1671 for id in choice.ids() {
1672 if id == talk.agent || chain.iter().any(|s| s.id == id) {
1673 continue;
1674 }
1675 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1676 Ok(s) => chain.push(s),
1677 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1678 }
1679 }
1680 }
1681
1682 let persona = crate::persona::active(&cfg.talk.personas, &talk.persona);
1683 let operator_name = cfg.talk.operator_name();
1684 let persona_update = if talk.persona_dirty {
1685 format!(
1686 "{}\n\n",
1687 crate::persona::update_block_for(persona.as_ref(), operator_name)
1688 )
1689 } else {
1690 String::new()
1691 };
1692
1693 let filing_update = if talk.implementers_dirty {
1696 format!("{}\n\n", filing_update_block(talk.implementers))
1697 } else {
1698 String::new()
1699 };
1700
1701 let mut outcome = None;
1702 let mut fell_back_from: Option<String> = None;
1703 let mut first_try: Option<(String, SeatState)> = None;
1706 for (n, spec) in chain.iter().enumerate() {
1707 if n > 0 {
1708 if first_try.is_none() {
1709 first_try = Some((talk.agent.clone(), talk.seat.clone()));
1710 }
1711 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1712 fell_back_from.get_or_insert_with(|| talk.agent.clone());
1715 talk.agent = spec.id.clone();
1716 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1717 }
1718 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1719 let first_ever = talk.turns.len() <= 1;
1720 let body = if talk.seat.turns == 0 && first_ever {
1721 format!(
1722 "{}\n\n# Operator\n\n{text}{last_note}",
1723 briefing_for(
1724 &talk.repo,
1725 &cfg.graph.language,
1726 cfg.talk.allow_write,
1727 persona.as_ref(),
1728 operator_name,
1729 talk.implementers
1730 )
1731 )
1732 } else if talk.seat.turns == 0 {
1733 format!(
1736 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1737 briefing_for(
1738 &talk.repo,
1739 &cfg.graph.language,
1740 cfg.talk.allow_write,
1741 persona.as_ref(),
1742 operator_name,
1743 talk.implementers
1744 ),
1745 transcript(talk, store)
1746 )
1747 } else if resuming {
1748 format!("{persona_update}{filing_update}{text}{last_note}")
1749 } else {
1750 let mut standing = match (&persona, talk.persona_dirty) {
1754 (Some(p), false) => format!("{}\n", crate::persona::section_for(p, operator_name)),
1755 _ => persona_update.clone(),
1756 };
1757 if let (None, Some(n)) = (&persona, operator_name) {
1760 standing.push_str(&format!(
1761 "# Addressing the operator\n{}\n",
1762 crate::persona::addressing(n)
1763 ));
1764 }
1765 if talk.implementers_dirty || talk.implementers != 1 {
1766 standing.push_str(&format!("{}\n\n", filing_update_block(talk.implementers)));
1767 }
1768 format!("{}\n\n{standing}{text}{last_note}", transcript(talk, store))
1769 };
1770 let attempt_stem = if n == 0 {
1771 stem.clone()
1772 } else {
1773 format!("{stem}-{}", spec.id)
1774 };
1775 let inv = Invocation {
1776 cwd: &talk.repo,
1777 prompt: &body,
1778 timeout: turn_timeout(cfg),
1779 allow_write,
1784 unsandboxed,
1785 sessions: cfg.graph.sessions,
1786 artifacts: &artifacts,
1787 stem: &attempt_stem,
1788 run: &talk.id,
1791 node: crate::queue::CHAT_NODE,
1792 cache_dir: cache_dir.as_deref(),
1793 attachments: &attachment_paths,
1794 writable: &consult_roots,
1795 };
1796 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1797 let advance = agent::chain_advances(&result);
1798 if n == 0 || !advance {
1799 outcome = Some(result);
1800 } else {
1801 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1803 }
1804 if !advance {
1805 break;
1806 }
1807 }
1808 if outcome.as_ref().is_some_and(agent::chain_advances) {
1809 if let Some((id, seat)) = first_try {
1812 talk.agent = id;
1813 talk.seat = seat;
1814 fell_back_from = None;
1815 }
1816 }
1817 let outcome = outcome.expect("a chain holds at least one agent");
1818 let note = |why: String| Turn {
1819 breaks: None,
1820 who: Who::Agent,
1821 body: format!("{MAGI_NOTE}{why}"),
1822 at: Timestamp::now(),
1823 attachments: Vec::new(),
1824 usage: None,
1825 };
1826 let (reply, failure) = match outcome {
1827 Err(e) => (
1828 note(format!("could not run agent `{}`: {e}", talk.agent)),
1829 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1830 ),
1831 Ok(out) if out.quota_exhausted() => {
1832 let reset = out
1833 .quota
1834 .as_ref()
1835 .and_then(|q| q.reset.clone())
1836 .map_or_else(String::new, |r| format!(" (resets {r})"));
1837 let why = format!(
1838 "agent `{}` is out of quota{reset}; your message is saved, so \
1839 say it again when the window reopens",
1840 talk.agent
1841 );
1842 (note(why.clone()), Some(why))
1843 }
1844 Ok(out) if out.timed_out => {
1845 let why = format!(
1846 "agent `{}` did not answer within {}s; your message is saved",
1847 talk.agent,
1848 turn_timeout(cfg).as_secs()
1849 );
1850 (note(why.clone()), Some(why))
1851 }
1852 Ok(out) if !out.usable() => {
1853 let why = format!(
1854 "agent `{}` produced no answer (exit {}); your message is saved",
1855 talk.agent,
1856 out.exit_code
1857 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1858 );
1859 (note(why.clone()), Some(why))
1860 }
1861 Ok(out) => (
1862 Turn {
1863 breaks: None,
1864 who: Who::Agent,
1865 body: out.text.trim().to_owned(),
1866 at: Timestamp::now(),
1867 attachments: Vec::new(),
1868 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1871 context_tokens,
1872 agent: talk.agent.clone(),
1873 model: cfg
1874 .agents
1875 .iter()
1876 .find(|a| a.id == talk.agent)
1877 .and_then(|a| a.model.clone()),
1878 }),
1879 },
1880 None,
1881 ),
1882 };
1883
1884 let _guard = store.guard()?;
1896 let Ok(fresh) = store.get(&talk.id) else {
1902 return Ok(());
1903 };
1904 talk.status = fresh.status;
1905 talk.pending = fresh.pending;
1909 talk.pending_breaks = fresh.pending_breaks;
1910 talk.pending_attachments = fresh.pending_attachments;
1911 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1912 talk.turns.push(note(format!(
1915 "agent changed from {from} to {} (fallback)",
1916 talk.agent
1917 )));
1918 }
1919 if failure.is_none() {
1922 talk.persona_dirty = false;
1923 talk.implementers_dirty = false;
1924 }
1925 talk.turns.push(reply);
1926 if let Err(put_err) = store.put(talk) {
1927 let lost = talk.turns.pop().expect("just pushed above");
1936 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1937 let why = match &stash {
1938 Ok(path) => format!(
1939 "agent `{}` answered, but the reply could not be saved to \
1940 this conversation ({put_err:#}); the raw text was kept at \
1941 {} - your message is saved, ask again",
1942 talk.agent,
1943 path.display()
1944 ),
1945 Err(stash_err) => format!(
1946 "agent `{}` answered, but the reply could not be saved to \
1947 this conversation ({put_err:#}), and it could not be kept \
1948 anywhere else either ({stash_err:#}); your message is \
1949 saved, ask again",
1950 talk.agent
1951 ),
1952 };
1953 talk.turns.push(note(why.clone()));
1954 return match store.put(talk) {
1961 Ok(()) => bail!("{why}"),
1962 Err(note_err) => {
1963 talk.turns.pop();
1983 Err(note_err).context(why)
1984 }
1985 };
1986 }
1987
1988 match failure {
1989 Some(why) => bail!("{why}"),
1990 None => Ok(()),
1991 }
1992}
1993
1994fn transcript(talk: &Talk, store: &Talks) -> String {
1997 let mut out = String::from(
1998 "This conversation cannot resume on the CLI's side, so here is \
1999 everything said so far; answer only the last message.\n",
2000 );
2001 for t in &talk.turns {
2002 let who = match t.who {
2003 Who::Operator => "operator",
2004 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
2005 Who::Agent => "you",
2006 };
2007 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
2008 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
2009 }
2010 out
2011}
2012
2013fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
2018 if attachments.is_empty() {
2019 return String::new();
2020 }
2021 let mut out = String::from(
2022 "\n\nThe operator attached the image(s) below to this message. Open \
2023 and look at each one before you answer.\n",
2024 );
2025 for att in attachments {
2026 if let Some(path) = store.attachment_path(talk_id, att) {
2027 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
2028 }
2029 }
2030 out.push('\n');
2031 out
2032}
2033
2034pub(crate) fn turn_access(talk_allow_write: bool, consulted: bool) -> (bool, bool) {
2045 (talk_allow_write || consulted, talk_allow_write)
2046}
2047
2048pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
2064 briefing_with(repo, language, allow_write, None, None)
2065}
2066
2067pub fn briefing_with(
2072 repo: &Path,
2073 language: &str,
2074 allow_write: bool,
2075 persona: Option<&crate::persona::Persona>,
2076 operator_name: Option<&str>,
2077) -> String {
2078 briefing_for(repo, language, allow_write, persona, operator_name, 1)
2079}
2080
2081pub const MAX_IMPLEMENTERS: u8 = 3;
2083
2084fn solo_implementers() -> u8 {
2085 1
2086}
2087
2088fn filing_flag(n: u8) -> String {
2091 if n <= 1 {
2092 "--solo".to_owned()
2093 } else {
2094 format!("--implementers {n}")
2095 }
2096}
2097
2098pub fn filing_update_block(n: u8) -> String {
2101 let flag = filing_flag(n);
2102 if n <= 1 {
2103 format!(
2104 "# Task filing update\n\nThe operator changed how many implementers tasks use. \
2105 From now on file with `magi task add {flag} --repo <repo> <instruction>` \
2106 (one implementer, straight into review).\n"
2107 )
2108 } else {
2109 format!(
2110 "# Task filing update\n\nThe operator changed how many implementers tasks use. \
2111 From now on file with `magi task add {flag} --repo <repo> <instruction>` \
2112 ({n} independent implementations compete). Do not pass `--solo` \
2113 together with `--implementers`; drop any earlier `--solo`.\n"
2114 )
2115 }
2116}
2117
2118pub fn check_implementers(n: u8, cfg: &crate::config::Config) -> std::result::Result<u8, String> {
2121 if !(1..=MAX_IMPLEMENTERS).contains(&n) {
2122 return Err(format!(
2123 "implementers must be 1..={MAX_IMPLEMENTERS}, got {n}"
2124 ));
2125 }
2126 if n > 1 {
2127 let mut cfg = cfg.clone();
2128 crate::queue::RunOverrides {
2129 candidates: Some(usize::from(n)),
2130 ..Default::default()
2131 }
2132 .apply(&mut cfg);
2133 cfg.resolve_roles()
2134 .map_err(|e| format!("{n} implementers is not allowed here: {e:#}"))?;
2135 }
2136 Ok(n)
2137}
2138
2139pub fn switch_implementers(talk: &mut Talk, store: &Talks, n: u8) -> Result<bool> {
2142 let _guard = store.guard()?;
2143 let mut fresh = store
2144 .get(&talk.id)
2145 .with_context(|| format!("talk {} was deleted", talk.short()))?;
2146 if fresh.implementers == n {
2147 *talk = fresh;
2148 return Ok(false);
2149 }
2150 let from = std::mem::replace(&mut fresh.implementers, n);
2151 fresh.implementers_dirty = true;
2152 fresh.turns.push(Turn {
2153 breaks: None,
2154 who: Who::Agent,
2155 body: format!("{MAGI_NOTE}implementers changed from {from} to {n}"),
2156 at: Timestamp::now(),
2157 attachments: Vec::new(),
2158 usage: None,
2159 });
2160 store.put(&mut fresh)?;
2161 *talk = fresh;
2162 Ok(true)
2163}
2164
2165pub fn briefing_for(
2168 repo: &Path,
2169 language: &str,
2170 allow_write: bool,
2171 persona: Option<&crate::persona::Persona>,
2172 operator_name: Option<&str>,
2173 implementers: u8,
2174) -> String {
2175 let flag = filing_flag(implementers);
2176 let shape = if implementers <= 1 {
2177 "Use --solo: it runs the task through one implementer \
2178 straight into review instead of the usual multi-agent competition, \
2179 which is the right shape for a change this conversation has already \
2180 settled, rather than one still worth several independent takes."
2181 .to_owned()
2182 } else {
2183 format!(
2184 "Use --implementers {implementers}: it has {implementers} \
2185 implementers work on the task independently and compete, which \
2186 is the right shape when several independent takes are worth \
2187 having. Never add --solo to it; the two do not go together."
2188 )
2189 };
2190 let write_policy = if allow_write {
2191 "Write access is enabled for this conversation (`allow_write = \
2192 true`), so you may write files - but only a small, \
2193 already-decided edit the operator names outright in this \
2194 conversation, not an implementation. This is a permission on the \
2195 conversation as a whole, not a property of whichever repository \
2196 it happened to start in: if the operator names a different \
2197 repository for that small edit, the policy allows it there too. \
2198 Your own tool may still confine writes to the repository this \
2199 conversation started in regardless - if a write elsewhere is \
2200 refused, say so plainly rather than working around it. Once you \
2201 have made an edit, say plainly what you edited. Anything bigger, \
2202 or anything still open-ended, still goes through the queue below \
2203 rather than being done here."
2204 } else {
2205 "Do not write files. Implementing a change is not this \
2206 conversation's job; a separate, blind competition of agents does \
2207 that, and a repository this conversation has already edited would \
2208 make their diffs unjudgeable."
2209 };
2210 let mut out = format!(
2211 "You are magi's standing conversation partner for its operator, who \
2212 usually has this open on a phone. Keep replies short: no preamble, \
2213 no restating what they just said.\n\n\
2214 # Repository\n\n{repo}\n\n\
2215 You may look around: read files, run shell commands, search history, \
2216 run tests - whatever answers the question. {write_policy}\n\n\
2217 A short, command-shaped message (\"list\", \"info <id>\", \"show \
2218 3cbf\") is almost always the operator asking you to look something \
2219 up, not an instruction to file - answer it yourself with `magi \
2220 list`, `magi show <id>`, `magi task list`, or the like, the same way \
2221 you would answer any other question in this conversation.\n\n\
2222 # When the operator wants something done\n\n\
2223 Run:\n\n\
2224 magi task add {flag} --repo {repo} <instruction>\n\n\
2225 and tell the operator the task id it prints, so they can follow it \
2226 from the Queue. If it refuses with a duplicate warning (the \
2227 instruction names a branch, commit or pull request that an \
2228 unfinished task, run or PR already owns), do not repeat it with \
2229 --force yourself: tell the operator what it matched and let them \
2230 decide. Write <instruction> so that an implementer who has \
2231 never seen this conversation can act on it alone - it is everything \
2232 they get. {shape}\n\n\
2233 If the operator asks for something in a different repository, \
2234 --repo does not have to be a full path: --repo owner/repo (or just \
2235 repo, when that is unambiguous) is resolved against local checkouts \
2236 the same way `magi repos` lists them. If the command fails because \
2237 nothing matches or more than one checkout shares that name, ask the \
2238 operator which repository they mean (or run `magi repos` yourself \
2239 to see the candidates) rather than guessing.\n\n\
2240 The current state of the code is whatever origin/main holds, not \
2241 whatever a working tree shows: a primary checkout often lags \
2242 upstream, sits on a detached HEAD and carries uncommitted changes. \
2243 Before answering about code, run `git fetch origin` in that \
2244 repository if it is cheap, then read through \
2245 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
2246 If the working tree differs, say so; if the fetch fails, say that \
2247 too, so the operator knows the answer may be stale.\n\n\
2248 If the operator attached an image (a screenshot, say) that the task \
2249 is about, pass it with `--attach <path>`, using the absolute path \
2250 the turn's attachment note gives; repeat the flag for several. \
2251 `magi task add {flag} --attach <path> <instruction>` copies the \
2252 file into the task, so the implementer receives it. Do not paste the \
2253 path into <instruction> instead: deleting this conversation deletes \
2254 its attachments, and then that path reaches no one.\n",
2255 repo = repo.display(),
2256 );
2257 out.push_str(&language_note(language));
2258 match (persona, operator_name) {
2259 (Some(p), n) => out.push_str(&crate::persona::section_for(p, n)),
2260 (None, Some(n)) => {
2261 out.push_str("\n# Addressing the operator\n");
2262 out.push_str(&crate::persona::addressing(n));
2263 }
2264 (None, None) => {}
2265 }
2266 out
2267}
2268
2269fn language_note(language: &str) -> String {
2272 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
2273 String::new()
2274 } else {
2275 format!("\nHold this conversation in {language}.\n")
2276 }
2277}
2278
2279pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
2286 let mut tasks: Vec<Task> = queue
2287 .list()
2288 .into_iter()
2289 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
2290 .collect();
2291 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
2292 tasks
2293}
2294
2295fn read_path(path: &Path) -> Result<Talk> {
2296 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
2297 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
2298}
2299
2300const PUT_RETRIES: u32 = 5;
2303
2304fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
2316 let mut last_err = None;
2317 for attempt in 0..PUT_RETRIES {
2318 if attempt > 0 {
2319 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
2320 }
2321 match try_write_atomic(tmp, path, body) {
2322 Ok(()) => return Ok(()),
2323 Err(e) => last_err = Some(e),
2324 }
2325 }
2326 Err(last_err.expect("the loop above always runs at least once"))
2327}
2328
2329fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
2330 #[cfg(test)]
2331 if failpoint::take_forced_put_failure() {
2332 bail!("simulated write failure (test)");
2333 }
2334 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
2335 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
2336 Ok(())
2337}
2338
2339fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
2344 let dir = store.artifacts_of(id);
2345 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
2346 let path = dir.join(format!("{stem}-lost.txt"));
2347 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
2348 Ok(path)
2349}
2350
2351#[cfg(test)]
2358mod failpoint {
2359 use std::cell::Cell;
2360
2361 thread_local! {
2362 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
2363 static FORCE_NO_LINK: Cell<bool> = const { Cell::new(false) };
2364 }
2365
2366 pub(super) fn force_no_link(on: bool) {
2369 FORCE_NO_LINK.with(|c| c.set(on));
2370 }
2371
2372 pub(super) fn no_link_forced() -> bool {
2373 FORCE_NO_LINK.with(Cell::get)
2374 }
2375
2376 pub(super) fn force_put_failures(count: u32) {
2379 FORCE_PUT_FAILURES.with(|c| c.set(count));
2380 }
2381
2382 pub(super) fn take_forced_put_failure() -> bool {
2385 FORCE_PUT_FAILURES.with(|c| {
2386 let n = c.get();
2387 if n == 0 {
2388 false
2389 } else {
2390 c.set(n - 1);
2391 true
2392 }
2393 })
2394 }
2395}
2396
2397fn short(id: &str) -> &str {
2398 id.split('-').next_back().unwrap_or(id)
2399}
2400
2401fn new_id() -> String {
2402 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2403 let seed = crate::rng::entropy();
2404 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2405}
2406
2407fn attachment_ext(mime: &str) -> Option<&'static str> {
2412 match mime {
2413 "image/png" => Some("png"),
2414 "image/jpeg" => Some("jpg"),
2415 "image/gif" => Some("gif"),
2416 "image/webp" => Some("webp"),
2417 _ => None,
2418 }
2419}
2420
2421pub fn valid_attachment_id(id: &str) -> bool {
2426 id.len() == 32
2427 && id
2428 .bytes()
2429 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
2430}
2431
2432fn new_attachment_id() -> String {
2436 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
2437 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
2438}
2439
2440#[cfg(test)]
2441mod tests {
2442 #[test]
2443 fn turn_access_lifts_the_sandbox_only_for_the_talk_opt_in() {
2444 use super::turn_access;
2445 assert_eq!(turn_access(false, false), (false, false));
2446 assert_eq!(turn_access(true, false), (true, true));
2447 assert_eq!(turn_access(false, true), (true, false));
2449 assert_eq!(turn_access(true, true), (true, true));
2450 }
2451
2452 #[test]
2453 fn the_briefing_points_at_origin_main_not_the_working_tree() {
2454 let b = briefing(Path::new("/r"), "en", false);
2455 assert!(b.contains("origin/main"));
2456 assert!(b.contains("git show origin/main:"));
2457 }
2458 use std::collections::BTreeMap;
2459
2460 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
2461 use crate::queue::{Queue, Source, Task};
2462
2463 use super::*;
2464
2465 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
2466 AgentSpec {
2467 id: id.to_owned(),
2468 kind: AgentKind::Command,
2469 model: model.map(str::to_owned),
2470 command: Vec::new(),
2471 extra_args: Vec::new(),
2472 env: BTreeMap::new(),
2473 prompt_delivery: None,
2474 }
2475 }
2476
2477 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
2478 Talk {
2479 pending_breaks: None,
2480 schema: SCHEMA,
2481 id: "20260904-014455-ab12".to_owned(),
2482 repo: PathBuf::from("."),
2483 agent: agent.to_owned(),
2484 status: TalkStatus::Open,
2485 turns,
2486 pending: String::new(),
2487 pending_attachments: Vec::new(),
2488 fallback: false,
2489 persona: String::new(),
2490 persona_dirty: false,
2491 implementers: 1,
2492 implementers_dirty: false,
2493 created_at: Timestamp::now(),
2494 updated_at: Timestamp::now(),
2495 seat: SeatState::new(SEAT, agent, 1),
2496 }
2497 }
2498
2499 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
2500 Turn {
2501 breaks: None,
2502 who: Who::Agent,
2503 body: body.to_owned(),
2504 at: Timestamp::now(),
2505 attachments: Vec::new(),
2506 usage: usage.map(|(t, a, m)| TurnUsage {
2507 context_tokens: t,
2508 agent: a.to_owned(),
2509 model: m.map(str::to_owned),
2510 }),
2511 }
2512 }
2513
2514 fn ctx_config(windows: &[(&str, u64)]) -> Config {
2515 Config {
2516 agents: vec![
2517 ctx_agent("small", Some("small-model")),
2518 ctx_agent("big", Some("big-model")),
2519 ctx_agent("plain", None),
2520 ],
2521 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
2522 ..Config::default()
2523 }
2524 }
2525
2526 #[test]
2527 fn context_usage_computes_percent_and_warns_at_eighty() {
2528 let cfg = ctx_config(&[("small-model", 1000)]);
2529 let at = |tokens| {
2530 let t = ctx_talk(
2531 "small",
2532 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
2533 );
2534 context_usage(&t, Some(&cfg))
2535 };
2536 let u = at(799);
2537 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
2538 let u = at(800);
2539 assert_eq!((u.percent, u.warn), (Some(80), true));
2540 let u = at(1500);
2541 assert_eq!((u.percent, u.warn), (Some(150), true));
2542 assert!(!u.since_switch);
2543 }
2544
2545 #[test]
2546 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
2547 let cfg = ctx_config(&[("small-model", 1000)]);
2548 let t = ctx_talk(
2549 "small",
2550 vec![
2551 reply("old", Some((900, "small", Some("small-model")))),
2552 reply("new", None),
2553 ],
2554 );
2555 let u = context_usage(&t, Some(&cfg));
2556 assert!(u.estimated);
2558 assert_ne!(u.tokens, Some(900));
2559 assert!(u.tokens.is_some());
2560 let t = ctx_talk(
2562 "small",
2563 vec![
2564 reply("old", Some((900, "small", Some("small-model")))),
2565 reply("magi: could not run agent", None),
2566 ],
2567 );
2568 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
2569 assert_eq!(
2570 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
2571 None
2572 );
2573 }
2574
2575 #[test]
2576 fn estimate_counts_chars_both_sides_and_standing_prompt() {
2577 let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
2578 assert_eq!(estimate_context_tokens(&t, 0), Some(2)); let op = Turn {
2580 who: Who::Operator,
2581 ..reply("abcdefg", None)
2582 };
2583 t.turns.push(op);
2584 assert_eq!(estimate_context_tokens(&t, 0), Some(4));
2585 assert!(
2586 estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
2587 );
2588 let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
2590 assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
2591 let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
2593 assert_eq!(estimate_context_tokens(¬e, 1000), None);
2594 assert_eq!(
2595 estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
2596 None
2597 );
2598 }
2599
2600 #[test]
2601 fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
2602 let cfg = ctx_config(&[("small-model", 1000)]);
2603 let t = ctx_talk(
2604 "small",
2605 vec![reply(
2606 &"x".repeat(5000),
2607 Some((10, "small", Some("small-model"))),
2608 )],
2609 );
2610 let u = context_usage(&t, Some(&cfg));
2611 assert_eq!((u.tokens, u.estimated), (Some(10), false));
2612 let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
2613 let u = context_usage(&t, Some(&cfg));
2614 assert!(u.estimated && !u.since_switch);
2615 assert_eq!(u.window, Some(1000));
2616 assert!(u.warn && u.percent.unwrap() >= 80);
2617 let t = ctx_talk("small", vec![reply("hi", None)]);
2618 let u = context_usage(&t, Some(&cfg));
2619 assert!(u.estimated && u.percent.is_some());
2620 }
2621
2622 #[test]
2623 fn context_usage_without_a_window_shows_tokens_only() {
2624 let cfg = ctx_config(&[]);
2625 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
2627 let u = context_usage(&t, Some(&cfg));
2628 assert_eq!(
2629 (u.tokens, u.window, u.percent, u.warn),
2630 (Some(5000), None, None, false)
2631 );
2632 let t = ctx_talk(
2633 "small",
2634 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
2635 );
2636 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
2637 assert_eq!(context_usage(&t, None).window, None);
2639 }
2640
2641 #[test]
2642 fn context_usage_switching_model_changes_the_denominator() {
2643 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
2644 let used = reply("hi", Some((900, "small", Some("small-model"))));
2645 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
2646 assert_eq!(
2647 (before.percent, before.warn, before.since_switch),
2648 (Some(90), true, false)
2649 );
2650 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
2653 assert_eq!(after.window, Some(10_000));
2654 assert_eq!(
2655 (after.percent, after.warn, after.since_switch),
2656 (Some(9), false, true)
2657 );
2658 assert_eq!(after.model.as_deref(), Some("big-model"));
2659 }
2660
2661 #[test]
2662 fn a_turn_recorded_before_usage_existed_still_reads() {
2663 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
2664 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
2665 assert!(turn.usage.is_none());
2666 let json = serde_json::to_string(&turn).expect("serialize");
2667 assert!(
2668 !json.contains("usage"),
2669 "absent usage is not written: {json}"
2670 );
2671 }
2672
2673 fn store() -> (tempfile::TempDir, Talks) {
2675 let tmp = tempfile::tempdir().expect("tempdir");
2676 let talks = Talks::at(tmp.path().join("talks"));
2677 (tmp, talks)
2678 }
2679
2680 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
2684 let path = dir.join("mock-talk-agent.sh");
2685 std::fs::write(&path, script).expect("write mock");
2686 AgentSpec {
2687 id: "mock".to_owned(),
2688 kind: AgentKind::Command,
2689 model: None,
2690 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2691 extra_args: Vec::new(),
2692 env,
2693 prompt_delivery: None,
2694 }
2695 }
2696
2697 fn config(spec: AgentSpec) -> Config {
2698 Config {
2699 agents: vec![spec],
2700 graph: Graph {
2701 language: "en".to_owned(),
2702 ..Graph::default()
2703 },
2704 ..Config::default()
2705 }
2706 }
2707
2708 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
2710
2711 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
2713
2714 const ECHO: &str = "#!/bin/sh\ncat\n";
2717
2718 fn env(reply: &str) -> BTreeMap<String, String> {
2719 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
2720 }
2721
2722 #[test]
2723 fn the_frozen_json_field_names_round_trip_through_disk() {
2724 let (tmp, talks) = store();
2725 let mut talk = Talk {
2726 pending_breaks: None,
2727 schema: SCHEMA,
2728 id: "20260904-014455-ab12".to_owned(),
2729 repo: tmp.path().to_owned(),
2730 agent: "sonnet".to_owned(),
2731 status: TalkStatus::Open,
2732 turns: Vec::new(),
2733 pending: String::new(),
2734 pending_attachments: Vec::new(),
2735 fallback: false,
2736 persona: String::new(),
2737 persona_dirty: false,
2738 implementers: 1,
2739 implementers_dirty: false,
2740 created_at: Timestamp::now(),
2741 updated_at: Timestamp::now(),
2742 seat: SeatState::new(SEAT, "sonnet", 7),
2743 };
2744 talks.put(&mut talk).expect("put");
2745
2746 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
2747 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
2748 for field in [
2749 "schema",
2750 "id",
2751 "repo",
2752 "agent",
2753 "status",
2754 "turns",
2755 "created_at",
2756 "updated_at",
2757 ] {
2758 assert!(v.get(field).is_some(), "missing field `{field}`");
2759 }
2760 assert_eq!(v["schema"], 1);
2761 assert_eq!(v["status"], "open");
2762
2763 let back = talks.get(&talk.id).expect("get");
2764 assert_eq!(back.id, talk.id);
2765 assert_eq!(back.status, TalkStatus::Open);
2766 }
2767
2768 #[test]
2769 fn opening_a_talk_takes_no_agent_turn() {
2770 let (tmp, talks) = store();
2771 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2775 let cfg = config(spec);
2776
2777 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2778 assert_eq!(talk.status, TalkStatus::Open);
2779 assert!(talk.turns.is_empty(), "nothing has been said yet");
2780
2781 let on_disk = talks.get(&talk.id).expect("get");
2782 assert_eq!(on_disk.turns.len(), 0);
2783 }
2784
2785 #[test]
2793 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
2794 let (tmp, talks) = store();
2795 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2796 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2797 chatter_spec.id = "chatter-mock".to_owned();
2798
2799 let mut cfg = Config {
2800 agents: vec![first_spec.clone(), chatter_spec.clone()],
2801 graph: Graph {
2802 language: "en".to_owned(),
2803 ..Graph::default()
2804 },
2805 ..Config::default()
2806 };
2807 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
2808
2809 let talk =
2810 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
2811 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
2812
2813 cfg.roles.chatter = None;
2814 let fallback =
2815 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
2816 assert_eq!(
2817 fallback.agent, first_spec.id,
2818 "unset chatter must fall back to agent::pick's own default order"
2819 );
2820 }
2821
2822 #[test]
2825 fn a_talk_recorded_without_attachments_still_reads() {
2826 let (tmp, talks) = store();
2827 let path = talks.path_of("20260904-014455-ab12");
2828 std::fs::create_dir_all(talks.root()).expect("talks dir");
2829 std::fs::write(
2830 &path,
2831 serde_json::json!({
2832 "schema": 1,
2833 "id": "20260904-014455-ab12",
2834 "repo": tmp.path(),
2835 "agent": "sonnet",
2836 "status": "open",
2837 "turns": [
2838 { "who": "operator", "body": "still there?",
2839 "at": Timestamp::now().to_string() },
2840 ],
2841 "created_at": Timestamp::now().to_string(),
2842 "updated_at": Timestamp::now().to_string(),
2843 "seat": SeatState::new(SEAT, "sonnet", 7),
2844 })
2845 .to_string(),
2846 )
2847 .expect("write pre-attachments talk");
2848
2849 let talk = talks.get("20260904-014455-ab12").expect("must still read");
2850 assert!(talk.turns[0].attachments.is_empty());
2851 }
2852
2853 fn lease_store() -> (tempfile::TempDir, Talks, Talks) {
2854 let tmp = tempfile::TempDir::new().expect("tmp");
2855 let root = tmp.path().join("talks");
2856 (tmp, Talks::at(root.clone()), Talks::at(root))
2857 }
2858
2859 #[test]
2860 fn two_starters_on_one_talk_one_wins_and_the_other_is_refused() {
2861 let (_tmp, a, b) = lease_store();
2862 let won = a.claim_turn("t1").expect("claim").expect("first wins");
2863 assert!(
2864 b.claim_turn("t1").expect("claim").is_none(),
2865 "second is refused"
2866 );
2867 assert!(b.turn_held("t1"));
2868 assert!(
2869 b.claim_turn("t2").expect("claim").is_some(),
2870 "other talks are free"
2871 );
2872 drop(won);
2873 }
2874
2875 #[test]
2876 fn a_stale_lease_is_taken_over_and_the_old_guard_cannot_release_it() {
2877 let (_tmp, a, b) = lease_store();
2878 let old = a.claim_turn("t1").expect("claim").expect("held");
2879 let later = Timestamp::now()
2880 .checked_add(jiff::SignedDuration::from_secs(
2881 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2882 ))
2883 .expect("later");
2884 let new = b
2885 .claim_turn_at("t1", later)
2886 .expect("claim")
2887 .expect("a stale lease is taken over");
2888 drop(old);
2889 assert!(a.turn_held("t1"), "the old guard left the new lease alone");
2890 assert!(new.beat().expect("beat"), "the new owner still beats");
2891 drop(new);
2892 assert!(!a.turn_held("t1"));
2893 }
2894
2895 #[tokio::test]
2896 async fn a_turn_whose_lease_was_taken_over_is_stopped() {
2897 let (_tmp, a, b) = lease_store();
2898 let old = a.claim_turn("t1").expect("claim").expect("held");
2899 let later = Timestamp::now()
2900 .checked_add(jiff::SignedDuration::from_secs(
2901 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2902 ))
2903 .expect("later");
2904 let _new = b
2905 .claim_turn_at("t1", later)
2906 .expect("claim")
2907 .expect("taken over");
2908 let out = old
2909 .beating_every(Duration::from_millis(10), std::future::pending::<()>())
2910 .await;
2911 assert!(out.is_err(), "the displaced turn must stop, not run on");
2912 }
2913
2914 #[tokio::test]
2915 async fn a_turn_that_finishes_is_returned_and_keeps_its_lease_beating() {
2916 let (_tmp, a, _b) = lease_store();
2917 let lease = a.claim_turn("t1").expect("claim").expect("held");
2918 let out = lease
2919 .beating_every(Duration::from_millis(5), async {
2920 tokio::time::sleep(Duration::from_millis(40)).await;
2921 7
2922 })
2923 .await
2924 .expect("still ours");
2925 assert_eq!(out, 7);
2926 assert!(a.turn_held("t1"));
2927 }
2928
2929 #[test]
2930 fn an_unreadable_lease_counts_as_stale() {
2931 let (_tmp, a, b) = lease_store();
2932 std::fs::create_dir_all(a.root()).expect("dir");
2933 std::fs::write(a.turn_path("t1"), "not json").expect("write");
2934 assert!(!a.turn_held("t1"));
2935 assert!(b.claim_turn("t1").expect("claim").is_some());
2936 }
2937
2938 #[test]
2939 fn without_hard_links_a_claim_is_still_exclusive_and_a_young_placeholder_blocks() {
2940 let (_tmp, a, b) = lease_store();
2941 failpoint::force_no_link(true);
2942 let held = a.claim_turn("t1").expect("claim").expect("first wins");
2943 assert!(b.claim_turn("t1").expect("claim").is_none());
2944 assert!(a.turn_held("t1"));
2945 drop(held);
2946 assert!(!a.turn_held("t1"));
2947 let path = a.turn_path("t1");
2949 assert!(create_exclusive(&path, "").expect("placeholder"));
2950 assert!(b.claim_turn("t1").expect("claim").is_none(), "young: held");
2951 age_file(&path);
2952 assert!(b.claim_turn("t1").expect("claim").is_some(), "old: stale");
2953 let lease = a.turn_path("t2");
2955 let lock = TurnLock::take(&lease).expect("take").expect("lock");
2956 assert!(TurnLock::take(&lease).expect("take").is_none());
2957 drop(lock);
2958 assert!(TurnLock::take(&lease).expect("take").is_some());
2959 failpoint::force_no_link(false);
2960 }
2961
2962 #[test]
2963 fn remove_if_carries_removes_only_the_expected_content() {
2964 let (_tmp, a, _b) = lease_store();
2965 std::fs::create_dir_all(a.root()).expect("dir");
2966 let p = a.root().join("x.turn.lock");
2967 std::fs::write(&p, "mine").expect("write");
2968 assert!(!remove_if_carries(&p, "other"));
2969 assert_eq!(std::fs::read_to_string(&p).expect("kept"), "mine");
2970 assert!(remove_if_carries(&p, "mine"));
2971 assert!(!p.exists());
2972 assert!(!remove_if_carries(&p, "mine"), "absent is not a removal");
2973 }
2974
2975 #[test]
2976 fn a_dropped_lock_does_not_remove_a_lock_taken_over_since() {
2977 let (_tmp, a, _b) = lease_store();
2978 std::fs::create_dir_all(a.root()).expect("dir");
2979 let lease = a.turn_path("t1");
2980 let lock = TurnLock::take(&lease).expect("take").expect("lock");
2981 let path = lock.path.clone();
2982 std::fs::write(&path, "someone-else").expect("replace");
2983 drop(lock);
2984 assert_eq!(
2985 std::fs::read_to_string(&path).expect("kept"),
2986 "someone-else"
2987 );
2988 }
2989
2990 #[test]
2991 fn concurrent_takeovers_of_a_stale_lease_have_one_winner() {
2992 let (_tmp, a, _b) = lease_store();
2993 drop(a.claim_turn("t1").expect("claim").expect("held"));
2994 std::fs::write(
2995 a.turn_path("t1"),
2996 serde_json::to_string(&TurnRecord {
2997 token: "gone".into(),
2998 pid: 1,
2999 beat_at: Timestamp::from_second(1).expect("ts"),
3000 })
3001 .expect("json"),
3002 )
3003 .expect("write");
3004 let wins: Vec<_> = std::thread::scope(|sc| {
3005 let hs: Vec<_> = (0..8)
3006 .map(|_| {
3007 let s = a.clone();
3008 sc.spawn(move || s.claim_turn("t1").expect("claim"))
3009 })
3010 .collect();
3011 hs.into_iter().map(|h| h.join().expect("join")).collect()
3012 });
3013 assert_eq!(wins.iter().filter(|w| w.is_some()).count(), 1);
3014 }
3015
3016 fn age_file(path: &Path) {
3017 let f = std::fs::OpenOptions::new()
3018 .write(true)
3019 .open(path)
3020 .expect("open");
3021 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
3022 .expect("age");
3023 }
3024
3025 #[test]
3026 fn a_late_taker_of_a_broken_lock_cannot_disturb_its_replacement() {
3027 let (_tmp, a, _b) = lease_store();
3028 let lease = a.turn_path("t1");
3029 let lock = lease.with_extension("turn.lock");
3030 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
3031 std::fs::write(&lock, "t1-dead").expect("dead lock");
3032 age_file(&lock);
3033 let b = TurnLock::take(&lease).expect("take").expect("b wins");
3035 let fresh = std::fs::read_to_string(&lock).expect("read");
3036 assert_eq!(fresh, b.token);
3037 let ticket = std::fs::read_dir(lock.parent().expect("dir"))
3040 .expect("dir")
3041 .flatten()
3042 .map(|e| e.path())
3043 .find(|p| p.to_string_lossy().ends_with(".break.0"))
3044 .expect("ticket");
3045 assert!(!create_exclusive(&ticket, "").expect("ticket"));
3046 assert!(TurnLock::take(&lease).expect("take").is_none());
3047 assert_eq!(std::fs::read_to_string(&lock).expect("read"), fresh);
3048 age_file(&lock);
3050 age_file(&ticket);
3051 std::mem::forget(b);
3052 let c = TurnLock::take(&lease)
3053 .expect("take")
3054 .expect("next generation");
3055 assert_ne!(c.token, fresh);
3056 }
3057
3058 #[test]
3059 fn a_live_ticket_blocks_and_a_stale_one_hands_over_to_the_next_generation() {
3060 let (_tmp, a, _b) = lease_store();
3061 let lease = a.turn_path("t1");
3062 let lock = lease.with_extension("turn.lock");
3063 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
3064 std::fs::write(&lock, "t1-dead").expect("dead lock");
3065 age_file(&lock);
3066 let t0 = lock.with_extension("lock.t1-dead.break.0");
3067 assert!(create_exclusive(&t0, "").expect("ticket"));
3068 assert!(TurnLock::take(&lease).expect("take").is_none());
3070 assert_eq!(std::fs::read_to_string(&lock).expect("read"), "t1-dead");
3071 age_file(&t0);
3073 let c = TurnLock::take(&lease).expect("take").expect("generation 1");
3074 assert_eq!(std::fs::read_to_string(&lock).expect("read"), c.token);
3075 assert!(lock.with_extension("lock.t1-dead.break.1").exists());
3076 assert!(TurnLock::take(&lease).expect("take").is_none());
3077 }
3078
3079 #[test]
3080 fn exhausted_ticket_generations_recover_once_the_sweep_ages_them_out() {
3081 let (_tmp, a, _b) = lease_store();
3082 let lease = a.turn_path("t1");
3083 let lock = lease.with_extension("turn.lock");
3084 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
3085 std::fs::write(&lock, "t1-dead").expect("dead lock");
3086 age_file(&lock);
3087 let tickets: Vec<_> = (0..TICKET_GENERATIONS)
3088 .map(|n| lock.with_extension(format!("lock.t1-dead.break.{n}")))
3089 .collect();
3090 for t in &tickets {
3091 assert!(create_exclusive(t, "").expect("ticket"));
3092 age_file(t);
3093 }
3094 assert!(TurnLock::take(&lease).expect("take").is_none());
3096 assert!(tickets.iter().all(|t| t.exists()));
3097 for t in &tickets {
3099 let f = std::fs::OpenOptions::new()
3100 .write(true)
3101 .open(t)
3102 .expect("open");
3103 f.set_modified(std::time::SystemTime::now() - TICKET_SWEEP_AGE * 2)
3104 .expect("age");
3105 }
3106 assert!(TurnLock::take(&lease).expect("take").is_none());
3107 assert!(TurnLock::take(&lease).expect("take").is_some());
3108 }
3109
3110 #[test]
3111 fn an_aged_empty_lock_is_broken() {
3112 let (_tmp, a, _b) = lease_store();
3113 let lease = a.turn_path("t1");
3114 let lock = lease.with_extension("turn.lock");
3115 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
3116 std::fs::write(&lock, "").expect("empty lock");
3117 age_file(&lock);
3118 assert!(TurnLock::take(&lease).expect("take").is_some());
3119 }
3120
3121 #[test]
3122 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
3123 let (tmp, talks) = store();
3124 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
3125 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3126
3127 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
3128 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
3129 let saved = talks.get(&talk.id).expect("reload queued talk");
3130 assert_eq!(saved.pending, "first\n\nsecond");
3131 assert_eq!(saved.pending_breaks, Some(vec!["first\n\n".len()]));
3132 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
3133
3134 let drained = drain(&mut talk, &talks).expect("drain");
3135 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
3136 let saved = talks.get(&talk.id).expect("reload drained talk");
3137 assert!(saved.pending.is_empty());
3138 assert_eq!(saved.turns.len(), 1);
3139 assert_eq!(saved.turns[0].body, "first\n\nsecond");
3140 assert_eq!(saved.turns[0].breaks, Some(vec!["first\n\n".len()]));
3141 assert_eq!(saved.pending_breaks, None);
3142 }
3143
3144 #[test]
3145 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
3146 let (tmp, talks) = store();
3147 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
3148 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3149 let attachment = Attachment {
3150 id: "a".repeat(32),
3151 name: "shot.png".to_owned(),
3152 mime: "image/png".to_owned(),
3153 bytes: 3,
3154 };
3155
3156 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
3157 assert!(
3158 edit_pending_text(
3159 &mut talk,
3160 &talks,
3161 "corrected",
3162 "first",
3163 std::slice::from_ref(&attachment.id),
3164 )
3165 .expect("edit")
3166 );
3167 let saved = talks.get(&talk.id).expect("reload edited draft");
3168 assert_eq!(saved.pending, "corrected");
3169 assert_eq!(saved.pending_attachments, vec![attachment]);
3170
3171 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
3172 assert!(
3173 !edit_pending_text(
3174 &mut talk,
3175 &talks,
3176 "stale edit",
3177 "corrected",
3178 &["a".repeat(32)],
3179 )
3180 .expect("stale edit is a conflict")
3181 );
3182 assert_eq!(
3183 talks.get(&talk.id).expect("reload after conflict").pending,
3184 "corrected\n\nlater"
3185 );
3186 assert!(
3187 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
3188 .expect("stale clear is a conflict")
3189 );
3190 assert_eq!(
3191 talks
3192 .get(&talk.id)
3193 .expect("reload after stale clear")
3194 .pending,
3195 "corrected\n\nlater"
3196 );
3197 }
3198
3199 #[tokio::test]
3200 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
3201 let (tmp, talks) = store();
3202 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
3203 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
3204 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3205 let id = running.id.clone();
3206 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
3207
3208 let response_talks = talks.clone();
3209 let response_cfg = cfg.clone();
3210 let reply = tokio::spawn(async move {
3211 respond(
3212 &response_talks.claim_turn(&running.id).unwrap().unwrap(),
3213 &mut running,
3214 &response_talks,
3215 &response_cfg,
3216 &first,
3217 )
3218 .await
3219 });
3220 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
3221
3222 let mut queued = talks.get(&id).expect("queued handle");
3223 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
3224 reply.await.expect("join").expect("reply");
3225
3226 let saved = talks.get(&id).expect("reload");
3227 assert_eq!(saved.pending, "next");
3228 assert_eq!(saved.pending_breaks, Some(Vec::new()));
3229 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
3230 }
3231
3232 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
3235 let calls = dir.join(format!("{id}.calls"));
3236 let script = format!(
3237 "#!/bin/sh\necho x >> '{}'\n{body}\n",
3238 calls.to_string_lossy()
3239 );
3240 let path = dir.join(format!("mock-{id}.sh"));
3241 std::fs::write(&path, script).expect("write mock");
3242 AgentSpec {
3243 id: id.to_owned(),
3244 kind: AgentKind::Command,
3245 model: None,
3246 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
3247 extra_args: Vec::new(),
3248 env: BTreeMap::new(),
3249 prompt_delivery: None,
3250 }
3251 }
3252
3253 fn calls(dir: &Path, id: &str) -> usize {
3254 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
3255 }
3256
3257 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
3258 let mut cfg = config(specs[0].clone());
3259 cfg.agents = specs;
3260 cfg.roles.chatter = Some(AgentChoice::Chain(
3261 ids.iter().map(|s| (*s).to_owned()).collect(),
3262 ));
3263 cfg
3264 }
3265
3266 #[tokio::test]
3267 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
3268 let (tmp, talks) = store();
3269 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3270 let b = counting_agent(tmp.path(), "b", "cat");
3271 let cfg = chain_config(vec![a, b], &["a", "b"]);
3272 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3273 assert_eq!(talk.agent, "a");
3274
3275 say(
3276 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3277 &mut talk,
3278 &talks,
3279 &cfg,
3280 "hello there",
3281 Vec::new(),
3282 )
3283 .await
3284 .expect("turn");
3285 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
3286 assert_eq!(calls(tmp.path(), "b"), 1);
3287 assert_eq!(talk.agent, "b", "the switch persists");
3288 assert!(talks.get(&talk.id).unwrap().agent == "b");
3289 let reply = talk.turns.last().unwrap();
3290 assert!(reply.body.contains("hello there"));
3291 assert!(
3292 reply.body.contains("magi task add --solo"),
3293 "a fresh seat gets the full briefing"
3294 );
3295 assert!(
3296 talk.turns
3297 .iter()
3298 .any(|t| t.body.contains("agent changed from a to b")),
3299 "the switch is noted"
3300 );
3301 }
3302
3303 #[tokio::test]
3304 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
3305 let (tmp, talks) = store();
3306 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3307 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
3308 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
3309 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3310
3311 let err = say(
3312 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3313 &mut talk,
3314 &talks,
3315 &cfg,
3316 "hi",
3317 Vec::new(),
3318 )
3319 .await
3320 .expect_err("every agent failed");
3321 assert!(err.to_string().contains("`a`"), "{err:#}");
3322 assert_eq!(calls(tmp.path(), "a"), 1);
3323 assert_eq!(calls(tmp.path(), "b"), 1);
3324 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
3325 }
3326
3327 #[test]
3328 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
3329 let (tmp, talks) = store();
3330 let b = counting_agent(tmp.path(), "b", "cat");
3331 let cfg = chain_config(vec![b], &["ghost", "b"]);
3332 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3333 assert_eq!(talk.agent, "b");
3334 }
3335
3336 #[tokio::test]
3337 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
3338 let (tmp, talks) = store();
3339 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3340 let b = counting_agent(tmp.path(), "b", "cat");
3341 let cfg = chain_config(vec![a, b], &["a", "b"]);
3342 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
3343 say(
3344 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3345 &mut talk,
3346 &talks,
3347 &cfg,
3348 "hi",
3349 Vec::new(),
3350 )
3351 .await
3352 .expect_err("a alone, and it fails");
3353 assert_eq!(calls(tmp.path(), "b"), 0);
3354 assert_eq!(talk.agent, "a");
3355 }
3356
3357 #[tokio::test]
3358 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
3359 let (tmp, talks) = store();
3360 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
3361 let b = counting_agent(tmp.path(), "b", "cat");
3362 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
3363 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
3364 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
3365 say(
3366 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3367 &mut talk,
3368 &talks,
3369 &cfg,
3370 "hi",
3371 Vec::new(),
3372 )
3373 .await
3374 .expect_err("c alone, and it fails");
3375 assert_eq!(calls(tmp.path(), "b"), 0);
3376 }
3377
3378 #[test]
3379 fn briefing_names_the_operator_with_or_without_a_persona() {
3380 let plain = briefing(Path::new("/repo"), "en", false);
3381 let named = briefing_with(Path::new("/repo"), "en", false, None, Some("Commander"));
3382 assert!(named.starts_with(&plain));
3383 assert!(named.contains("# Addressing the operator"));
3384 assert!(named.contains("\"Commander\""));
3385 let rei = crate::persona::builtin_catalog()
3386 .into_iter()
3387 .find(|p| p.id == "rei")
3388 .expect("rei");
3389 let with = briefing_with(
3390 Path::new("/repo"),
3391 "en",
3392 false,
3393 Some(&rei),
3394 Some("Commander"),
3395 );
3396 assert!(with.contains("\"Commander\""));
3397 assert!(!with.contains("# Addressing the operator"));
3398 }
3399
3400 #[test]
3401 fn briefing_carries_a_persona_section_only_when_one_is_chosen() {
3402 let plain = briefing(Path::new("/repo"), "en", false);
3403 assert_eq!(
3404 plain,
3405 briefing_with(Path::new("/repo"), "en", false, None, None)
3406 );
3407 assert!(!plain.contains("Persona"));
3408 let rei = crate::persona::builtin_catalog()
3409 .into_iter()
3410 .find(|p| p.id == "rei")
3411 .expect("rei");
3412 let with = briefing_with(Path::new("/repo"), "en", false, Some(&rei), None);
3413 assert!(with.starts_with(&plain), "the plain briefing is untouched");
3414 assert!(with.contains("# Persona (tone only)"));
3415 assert!(with.contains("TONE ONLY"));
3416 assert!(with.contains("task ids"));
3417 assert!(with.contains("write policy"));
3418 assert!(with.contains("`magi task add`"));
3419 assert!(with.contains("Rei Ayanami"));
3420 }
3421
3422 #[test]
3423 fn a_talk_written_before_personas_still_loads() {
3424 let (tmp, talks) = store();
3425 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3426 let cfg = config(spec);
3427 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3428 let path = talks.root.join(format!("{}.json", talk.id));
3429 let mut v: serde_json::Value =
3430 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3431 v.as_object_mut().unwrap().remove("persona");
3432 v.as_object_mut().unwrap().remove("persona_dirty");
3433 std::fs::write(&path, v.to_string()).unwrap();
3434 let loaded = talks.get(&talk.id).expect("old record loads");
3435 assert_eq!(loaded.persona, "");
3436 assert!(!loaded.persona_dirty);
3437 }
3438
3439 #[test]
3440 fn the_briefing_files_with_solo_or_the_chosen_implementer_count() {
3441 let repo = Path::new("/repo");
3442 let solo = briefing_for(repo, "", false, None, None, 1);
3443 assert_eq!(solo, briefing(repo, "", false));
3444 assert!(solo.contains("magi task add --solo --repo"));
3445 assert!(solo.contains("magi task add --solo --attach"));
3446 for n in [2u8, 3] {
3447 let b = briefing_for(repo, "", false, None, None, n);
3448 assert!(
3449 b.contains(&format!("magi task add --implementers {n} --repo")),
3450 "{b}"
3451 );
3452 assert!(b.contains(&format!("magi task add --implementers {n} --attach")));
3453 assert!(!b.contains("task add --solo"), "{b}");
3454 assert!(!b.contains("Use --solo"), "{b}");
3455 }
3456 }
3457
3458 #[test]
3459 fn implementers_are_validated_and_old_records_load_as_solo() {
3460 let (tmp, talks) = store();
3461 let cfg = config(mock_agent(tmp.path(), ECHO, BTreeMap::new()));
3462 assert_eq!(check_implementers(1, &cfg), Ok(1));
3463 assert_eq!(check_implementers(3, &cfg), Ok(3));
3464 assert!(check_implementers(0, &cfg).is_err());
3465 assert!(check_implementers(4, &cfg).is_err());
3466 let mut empty = cfg.clone();
3467 empty.agents.clear();
3468 assert!(check_implementers(2, &empty).is_err());
3469
3470 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3471 let path = talks.path_of(&talk.id);
3472 let mut v: serde_json::Value =
3473 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3474 v.as_object_mut().unwrap().remove("implementers");
3475 v.as_object_mut().unwrap().remove("implementers_dirty");
3476 std::fs::write(&path, v.to_string()).unwrap();
3477 let loaded = talks.get(&talk.id).expect("load");
3478 assert_eq!(loaded.implementers, 1);
3479 assert!(!loaded.implementers_dirty);
3480 }
3481
3482 #[tokio::test]
3483 async fn an_implementers_switch_is_noted_once_and_kept_across_a_failed_turn() {
3484 let (tmp, talks) = store();
3485 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3486 let cfg = config(spec);
3487 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3488 let lease = talks.claim_turn(&talk.id).unwrap().unwrap();
3489 say(&lease, &mut talk, &talks, &cfg, "hello", Vec::new())
3490 .await
3491 .expect("first turn");
3492
3493 assert!(switch_implementers(&mut talk, &talks, 2).expect("switch"));
3494 assert!(!switch_implementers(&mut talk, &talks, 2).expect("same"));
3495 assert!(talk.implementers_dirty);
3496 say(&lease, &mut talk, &talks, &cfg, "next", Vec::new())
3497 .await
3498 .expect("turn");
3499 let prompt = &talk.turns.last().unwrap().body;
3500 assert!(prompt.contains("# Task filing update"), "{prompt}");
3501 assert!(prompt.contains("--implementers 2"));
3502 assert!(!talk.implementers_dirty, "cleared after a successful turn");
3503
3504 say(&lease, &mut talk, &talks, &cfg, "again", Vec::new())
3505 .await
3506 .expect("turn");
3507 assert!(
3508 !talk
3509 .turns
3510 .last()
3511 .unwrap()
3512 .body
3513 .contains("# Task filing update")
3514 );
3515
3516 let mut no_sessions = cfg.clone();
3518 no_sessions.graph.sessions = false;
3519 say(&lease, &mut talk, &talks, &no_sessions, "one", Vec::new())
3520 .await
3521 .expect("turn");
3522 say(&lease, &mut talk, &talks, &no_sessions, "two", Vec::new())
3523 .await
3524 .expect("turn");
3525 assert!(talk.turns.last().unwrap().body.contains("--implementers 2"));
3526 }
3527
3528 #[tokio::test]
3529 async fn a_persona_switch_notes_marks_dirty_and_updates_the_next_turn_once() {
3530 let (tmp, talks) = store();
3531 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3532 let cfg = config(spec);
3533 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3534 let lease = talks.claim_turn(&talk.id).unwrap().unwrap();
3535 say(&lease, &mut talk, &talks, &cfg, "hello", Vec::new())
3536 .await
3537 .expect("first turn");
3538 assert!(!talk.turns[1].body.contains("Persona"));
3539
3540 assert!(switch_persona(&mut talk, &talks, "misato").expect("switch"));
3541 assert!(!switch_persona(&mut talk, &talks, "misato").expect("same"));
3542 assert!(talk.persona_dirty);
3543 assert!(talk.turns.last().unwrap().body.contains("persona changed"));
3544
3545 say(&lease, &mut talk, &talks, &cfg, "next", Vec::new())
3546 .await
3547 .expect("turn");
3548 let prompt = &talk.turns.last().unwrap().body;
3549 assert!(prompt.contains("# Persona update"), "{prompt}");
3550 assert!(prompt.contains("Misato Katsuragi"));
3551 assert!(!talk.persona_dirty, "cleared after a successful turn");
3552
3553 say(&lease, &mut talk, &talks, &cfg, "again", Vec::new())
3554 .await
3555 .expect("turn");
3556 assert!(!talk.turns.last().unwrap().body.contains("# Persona update"));
3557
3558 assert!(switch_persona(&mut talk, &talks, "default").expect("back"));
3559 assert_eq!(talk.persona, "");
3560 say(&lease, &mut talk, &talks, &cfg, "plain", Vec::new())
3561 .await
3562 .expect("turn");
3563 assert!(
3564 talk.turns
3565 .last()
3566 .unwrap()
3567 .body
3568 .contains("turned the persona off")
3569 );
3570
3571 assert!(switch_persona(&mut talk, &talks, "rei").expect("rei"));
3573 let mut no_sessions = cfg.clone();
3574 no_sessions.graph.sessions = false;
3575 say(&lease, &mut talk, &talks, &no_sessions, "one", Vec::new())
3576 .await
3577 .expect("turn");
3578 say(&lease, &mut talk, &talks, &no_sessions, "two", Vec::new())
3579 .await
3580 .expect("turn");
3581 let last = &talk.turns.last().unwrap().body;
3582 assert!(!talk.persona_dirty);
3583 assert!(last.contains("# Persona (tone only)"), "{last}");
3584 assert!(last.contains("Rei Ayanami"));
3585 }
3586
3587 #[tokio::test]
3588 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
3589 let (tmp, talks) = store();
3590 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3591 let cfg = config(spec);
3592 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3593
3594 say(
3595 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3596 &mut talk,
3597 &talks,
3598 &cfg,
3599 "what does the queue module do?",
3600 Vec::new(),
3601 )
3602 .await
3603 .expect("first turn");
3604 let first_prompt = &talk.turns[1].body;
3605 assert!(first_prompt.contains("magi task add --solo"));
3606 assert!(first_prompt.contains("what does the queue module do?"));
3607
3608 say(
3609 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3610 &mut talk,
3611 &talks,
3612 &cfg,
3613 "and how is it locked?",
3614 Vec::new(),
3615 )
3616 .await
3617 .expect("second turn");
3618 let second_prompt = &talk.turns[3].body;
3619 assert!(
3620 !second_prompt.contains("magi task add --solo"),
3621 "the briefing is sent once, not on every turn: {second_prompt}"
3622 );
3623 assert!(second_prompt.contains("and how is it locked?"));
3624 }
3625
3626 #[tokio::test]
3627 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
3628 let (tmp, talks) = store();
3629 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3630 let mut b = a.clone();
3631 b.id = "other".to_owned();
3632 let mut cfg = config(a.clone());
3633 cfg.agents.push(b.clone());
3634 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
3635 say(
3636 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3637 &mut talk,
3638 &talks,
3639 &cfg,
3640 "remember the walrus",
3641 Vec::new(),
3642 )
3643 .await
3644 .expect("first turn");
3645 let old_session = talk.seat.claude_session.clone();
3646 assert_eq!(talk.seat.turns, 1);
3647
3648 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
3649 assert_eq!(talk.agent, "other");
3650 assert_eq!(talk.seat.turns, 0);
3651 assert_eq!(talk.seat.agent, "other");
3652 assert_ne!(talk.seat.claude_session, old_session);
3653 let note = talk.turns.last().expect("note");
3654 assert_eq!(note.who, Who::Agent);
3655 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3656 assert!(note.body.contains("changed from"), "{}", note.body);
3657 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
3658
3659 let before = talk.turns.len();
3660 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
3661 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
3662
3663 say(
3664 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3665 &mut talk,
3666 &talks,
3667 &cfg,
3668 "what did I say?",
3669 Vec::new(),
3670 )
3671 .await
3672 .expect("turn after switch");
3673 let prompt = &talk.turns.last().expect("reply").body;
3674 assert!(prompt.contains("remember the walrus"), "{prompt}");
3675 assert!(prompt.contains("## magi"), "{prompt}");
3676 assert!(prompt.contains("what did I say?"), "{prompt}");
3677 }
3678
3679 #[tokio::test]
3680 async fn say_appends_the_operator_turn_then_the_agent_turn() {
3681 let (tmp, talks) = store();
3682 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3683 let cfg = config(spec);
3684 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3685
3686 say(
3687 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3688 &mut talk,
3689 &talks,
3690 &cfg,
3691 "can I rename this function?",
3692 Vec::new(),
3693 )
3694 .await
3695 .expect("say");
3696
3697 assert_eq!(talk.turns.len(), 2);
3698 assert_eq!(talk.turns[0].who, Who::Operator);
3699 assert_eq!(talk.turns[0].body, "can I rename this function?");
3700 assert_eq!(talk.turns[1].who, Who::Agent);
3701 assert_eq!(talk.turns[1].body, "go ahead");
3702 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
3703 }
3704
3705 #[tokio::test]
3706 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
3707 let (tmp, talks) = store();
3708 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
3709 let cfg = config(spec);
3710 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3711
3712 let err = say(
3713 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3714 &mut talk,
3715 &talks,
3716 &cfg,
3717 "check the tests",
3718 Vec::new(),
3719 )
3720 .await
3721 .expect_err("a turn with no answer is an error");
3722 assert!(err.to_string().contains("no answer"), "{err}");
3723
3724 let on_disk = talks.get(&talk.id).expect("get");
3725 assert_eq!(on_disk.turns.len(), 2);
3726 assert_eq!(on_disk.turns[0].body, "check the tests");
3727 let note = &on_disk.turns[1];
3728 assert_eq!(note.who, Who::Agent);
3729 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3730 assert!(note.body.contains("your message is saved"));
3731 }
3732
3733 #[tokio::test]
3739 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
3740 let (tmp, talks) = store();
3741 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3742 let cfg = config(spec);
3743 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3744
3745 let text =
3746 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
3747 failpoint::force_put_failures(PUT_RETRIES - 1);
3750 respond(
3751 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3752 &mut talk,
3753 &talks,
3754 &cfg,
3755 &text,
3756 )
3757 .await
3758 .expect("respond must survive a write failure its own retries can outlast");
3759
3760 assert_eq!(talk.turns.len(), 2);
3761 assert_eq!(talk.turns[1].who, Who::Agent);
3762 assert_eq!(talk.turns[1].body, "go ahead");
3763 let on_disk = talks.get(&talk.id).expect("get");
3764 assert_eq!(
3765 on_disk.turns, talk.turns,
3766 "the reply must reach disk despite the early write failures"
3767 );
3768 }
3769
3770 #[tokio::test]
3776 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
3777 let (tmp, talks) = store();
3778 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3779 let cfg = config(spec);
3780 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3781
3782 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3783 failpoint::force_put_failures(PUT_RETRIES);
3788 let err = respond(
3789 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3790 &mut talk,
3791 &talks,
3792 &cfg,
3793 &text,
3794 )
3795 .await
3796 .expect_err("a reply that cannot be saved must be reported, not swallowed");
3797 assert!(err.to_string().contains("could not be saved"), "{err}");
3798
3799 let on_disk = talks.get(&talk.id).expect("get");
3800 assert_eq!(
3801 on_disk.turns.len(),
3802 2,
3803 "the operator turn plus a visible note"
3804 );
3805 assert_eq!(on_disk.turns[0].body, "check the tests");
3806 let note = &on_disk.turns[1];
3807 assert_eq!(note.who, Who::Agent);
3808 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3809 assert!(
3810 note.body.contains("could not be saved"),
3811 "the operator must be told the reply is missing, not left staring \
3812 at a gap with no explanation: {}",
3813 note.body
3814 );
3815 assert_eq!(
3816 talk.turns, on_disk.turns,
3817 "the in-memory talk must match what actually landed on disk"
3818 );
3819
3820 let artifacts = talks.artifacts_of(&talk.id);
3823 let stash = std::fs::read_dir(&artifacts)
3824 .expect("artifacts dir")
3825 .filter_map(|e| e.ok())
3826 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
3827 .expect("a stash file for the lost reply");
3828 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
3829 assert_eq!(stashed, "go ahead");
3830
3831 assert_eq!(
3839 on_disk.seat.turns, 1,
3840 "the note's write must carry the turn the CLI actually took"
3841 );
3842 assert_eq!(
3843 on_disk.seat.claude_session, talk.seat.claude_session,
3844 "the session id handed to the CLI must survive the failed reply"
3845 );
3846 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
3847 assert!(
3848 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
3849 "the next turn must resume, not open the same session id twice"
3850 );
3851 }
3852
3853 #[tokio::test]
3858 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
3859 let (tmp, talks) = store();
3860 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3861 let cfg = config(spec);
3862 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3863
3864 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3865 failpoint::force_put_failures(PUT_RETRIES * 2);
3868 let err = respond(
3869 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3870 &mut talk,
3871 &talks,
3872 &cfg,
3873 &text,
3874 )
3875 .await
3876 .expect_err("neither the reply nor the note could be saved");
3877 assert!(err.to_string().contains("could not be saved"), "{err}");
3878
3879 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
3880 let on_disk = talks.get(&talk.id).expect("get");
3881 assert_eq!(on_disk.turns.len(), 1);
3882
3883 assert_eq!(
3893 on_disk.seat.turns, 0,
3894 "an unwritable file cannot record the turn the CLI took"
3895 );
3896 assert_eq!(
3897 talk.seat.turns, 1,
3898 "the in-memory seat still reports the turn the CLI actually took"
3899 );
3900 assert_eq!(
3901 on_disk.seat.claude_session, talk.seat.claude_session,
3902 "the session id was minted at `begin` and never changes here"
3903 );
3904 }
3905
3906 #[tokio::test]
3910 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
3911 let (tmp, talks) = store();
3912 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3913 let cfg = config(spec);
3914 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3915
3916 let att = talks
3917 .put_attachment(
3918 &talk.id,
3919 "image/png",
3920 "screenshot.png",
3921 b"pretend-png-bytes",
3922 )
3923 .expect("put attachment");
3924
3925 say(
3926 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3927 &mut talk,
3928 &talks,
3929 &cfg,
3930 "",
3931 vec![att.clone()],
3932 )
3933 .await
3934 .expect("an empty body with an attachment is still a turn");
3935
3936 let operator_turn = &talk.turns[0];
3937 assert_eq!(operator_turn.who, Who::Operator);
3938 assert_eq!(operator_turn.body, "");
3939 assert_eq!(operator_turn.attachments, vec![att.clone()]);
3940
3941 let prompt = &talk.turns[1].body;
3942 let expected_path = talks
3943 .attachments_dir(&talk.id)
3944 .join(format!("{}.png", att.id));
3945 assert!(
3946 prompt.contains(&expected_path.display().to_string()),
3947 "the agent must be told the attachment's absolute path: {prompt}"
3948 );
3949 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
3950 }
3951
3952 #[test]
3961 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
3962 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
3963 let att = Attachment {
3964 id: "0".repeat(32),
3965 name: "shot.png".to_owned(),
3966 mime: "image/png".to_owned(),
3967 bytes: 3,
3968 };
3969 let path = talks
3970 .attachment_path("some-talk-id", &att)
3971 .expect("a supported mime always yields a path");
3972 assert!(
3973 path.is_absolute(),
3974 "must be absolute even off a relative store root: {}",
3975 path.display()
3976 );
3977 }
3978
3979 #[tokio::test]
3980 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
3981 let (tmp, talks) = store();
3986 let slow = mock_agent(
3987 tmp.path(),
3988 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
3989 BTreeMap::new(),
3990 );
3991 let mut cfg = config(slow);
3992 cfg.graph.timeout_talk = 1;
3993 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3994
3995 let err = say(
3996 &talks.claim_turn(&talk.id).unwrap().unwrap(),
3997 &mut talk,
3998 &talks,
3999 &cfg,
4000 "check the tests",
4001 Vec::new(),
4002 )
4003 .await
4004 .expect_err("a turn that never answers is an error");
4005 assert!(
4006 err.to_string().contains("did not answer within 1s"),
4007 "{err}"
4008 );
4009
4010 let on_disk = talks.get(&talk.id).expect("get");
4011 let note = on_disk.turns.last().expect("a note turn was recorded");
4012 assert!(
4013 note.body.contains("did not answer within 1s"),
4014 "the transcript must show the configured timeout: {}",
4015 note.body
4016 );
4017 }
4018
4019 #[test]
4020 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
4021 let (tmp, talks) = store();
4022 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4023 let cfg = config(spec);
4024 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4025
4026 close(&mut talk, &talks).expect("close");
4027 assert_eq!(talk.status, TalkStatus::Closed);
4028 close(&mut talk, &talks).expect("closing twice is not an error");
4029
4030 let err =
4031 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
4032 assert!(err.to_string().contains("closed"));
4033 let _ = &cfg; }
4035
4036 #[tokio::test]
4037 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
4038 let (tmp, talks) = store();
4039 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
4040 let cfg = config(spec);
4041 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4044
4045 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
4049 close(&mut closed_elsewhere, &talks).expect("close");
4050 assert_eq!(
4051 talks.get(&in_flight.id).expect("reread").status,
4052 TalkStatus::Closed,
4053 "the close landed on disk before the turn finished"
4054 );
4055
4056 assert_eq!(in_flight.status, TalkStatus::Open);
4060 respond(
4061 &talks.claim_turn(&in_flight.id).unwrap().unwrap(),
4062 &mut in_flight,
4063 &talks,
4064 &cfg,
4065 "one more question",
4066 )
4067 .await
4068 .expect("the turn itself still completes");
4069
4070 let on_disk = talks.get(&in_flight.id).expect("reread");
4071 assert_eq!(
4072 on_disk.status,
4073 TalkStatus::Closed,
4074 "a close must stick even when a turn that started before it finishes after it"
4075 );
4076 assert!(
4079 on_disk.turns.iter().any(|t| t.body == "here you go"),
4080 "the in-flight turn's own reply is still recorded: {:?}",
4081 on_disk.turns
4082 );
4083 }
4084
4085 #[test]
4086 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
4087 let (tmp, talks) = store();
4088 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4089 let cfg = config(spec);
4090 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4093
4094 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
4097 close(&mut closed_elsewhere, &talks).expect("close");
4098 assert_eq!(
4099 talks.get(&stale.id).expect("reread").status,
4100 TalkStatus::Closed,
4101 "the close landed on disk before record was called"
4102 );
4103
4104 assert_eq!(stale.status, TalkStatus::Open);
4108 let err = record(&mut stale, &talks, "still there?", Vec::new())
4109 .expect_err("a close that landed first must be honored, not overwritten");
4110 assert!(err.to_string().contains("closed"));
4111
4112 let on_disk = talks.get(&stale.id).expect("reread");
4113 assert_eq!(
4114 on_disk.status,
4115 TalkStatus::Closed,
4116 "record must not resurrect a conversation closed while its snapshot was stale"
4117 );
4118 assert!(
4119 on_disk.turns.is_empty(),
4120 "the rejected turn must not have been appended: {:?}",
4121 on_disk.turns
4122 );
4123 let _ = &cfg; }
4125
4126 #[test]
4127 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
4128 let (tmp, talks) = store();
4129 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4130 let cfg = config(spec);
4131 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4132
4133 let held = talks.guard().unwrap();
4137
4138 let talks2 = talks.clone();
4139 let id = talk.id.clone();
4140 let closing = std::thread::spawn(move || {
4141 let mut talk = talks2.get(&id).expect("get");
4142 close(&mut talk, &talks2).expect("close");
4143 });
4144
4145 std::thread::sleep(Duration::from_millis(50));
4146 assert!(
4147 !closing.is_finished(),
4148 "close must wait for the guard, not read and write while it is held - \
4149 a re-read alone narrows this window without closing it"
4150 );
4151
4152 drop(held);
4153 closing.join().expect("close thread panicked");
4154
4155 assert_eq!(
4156 talks.get(&talk.id).expect("reread").status,
4157 TalkStatus::Closed,
4158 "once the guard is free, close still lands"
4159 );
4160 let _ = &cfg; }
4162
4163 #[test]
4164 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
4165 let (tmp, talks) = store();
4166 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4167 let cfg = config(spec);
4168 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4169
4170 close(&mut talk, &talks).expect("close");
4171 assert_eq!(talk.status, TalkStatus::Closed);
4172
4173 reopen(&mut talk, &talks).expect("reopen");
4174 assert_eq!(talk.status, TalkStatus::Open);
4175 assert_eq!(
4176 talks.get(&talk.id).expect("reread").status,
4177 TalkStatus::Open
4178 );
4179
4180 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
4182 assert_eq!(talk.status, TalkStatus::Open);
4183
4184 record(&mut talk, &talks, "one more thing", Vec::new())
4185 .expect("a reopened talk takes turns again");
4186 let _ = &cfg; }
4188
4189 #[test]
4190 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
4191 let (tmp, talks) = store();
4192 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4193 let cfg = config(spec);
4194 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4195
4196 let artifacts = talks.artifacts_of(&talk.id);
4197 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
4198 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
4199
4200 talks.remove(&talk.id).expect("remove");
4201 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
4202 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
4203 assert!(
4204 talks.get(&talk.id).is_err(),
4205 "a removed talk cannot be read back"
4206 );
4207
4208 let err = talks
4209 .remove("nonexistent-id")
4210 .expect_err("unknown id refused");
4211 assert!(err.to_string().contains("no talk matches"), "{err}");
4212 let _ = &cfg; }
4214
4215 #[tokio::test]
4216 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
4217 let (tmp, talks) = store();
4218 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
4219 let cfg = config(spec);
4220 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4223
4224 talks.remove(&in_flight.id).expect("remove");
4225 assert!(
4226 talks.get(&in_flight.id).is_err(),
4227 "the delete landed on disk before the turn finished"
4228 );
4229
4230 respond(
4233 &talks.claim_turn(&in_flight.id).unwrap().unwrap(),
4234 &mut in_flight,
4235 &talks,
4236 &cfg,
4237 "one more question",
4238 )
4239 .await
4240 .expect("the turn itself still completes rather than erroring");
4241
4242 assert!(
4243 talks.get(&in_flight.id).is_err(),
4244 "a delete must stick even when a turn that started before it finishes after it"
4245 );
4246 }
4247
4248 #[test]
4249 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
4250 let (tmp, talks) = store();
4251 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4252 let cfg = config(spec);
4253 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4256
4257 talks.remove(&stale.id).expect("remove");
4258
4259 let err = record(&mut stale, &talks, "still there?", Vec::new())
4263 .expect_err("a delete that landed first must be honored, not overwritten");
4264 assert!(err.to_string().contains("deleted"), "{err}");
4265
4266 assert!(
4267 talks.get(&stale.id).is_err(),
4268 "record must not resurrect a conversation deleted while its snapshot was stale"
4269 );
4270 let _ = &cfg; }
4272
4273 #[test]
4274 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
4275 let (tmp, talks) = store();
4276 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4277 let cfg = config(spec);
4278 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4281
4282 talks.remove(&stale.id).expect("remove");
4283
4284 let err = close(&mut stale, &talks)
4288 .expect_err("a delete that landed first must be honored, not overwritten");
4289 assert!(err.to_string().contains("deleted"), "{err}");
4290
4291 assert!(
4292 talks.get(&stale.id).is_err(),
4293 "close must not resurrect a conversation deleted while its snapshot was stale"
4294 );
4295 let _ = &cfg; }
4297
4298 #[test]
4299 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
4300 let (tmp, talks) = store();
4301 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
4302 let cfg = config(spec);
4303 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
4306 close(&mut stale, &talks).expect("close");
4307
4308 talks.remove(&stale.id).expect("remove");
4309
4310 let err = reopen(&mut stale, &talks)
4314 .expect_err("a delete that landed first must be honored, not overwritten");
4315 assert!(err.to_string().contains("deleted"), "{err}");
4316
4317 assert!(
4318 talks.get(&stale.id).is_err(),
4319 "reopen must not resurrect a conversation deleted while its snapshot was stale"
4320 );
4321 let _ = &cfg; }
4323
4324 #[test]
4325 fn list_puts_open_talks_before_closed_ones() {
4326 let (tmp, talks) = store();
4327 let make = |id: &str, status: TalkStatus| {
4328 let mut t = Talk {
4329 pending_breaks: None,
4330 schema: SCHEMA,
4331 id: id.to_owned(),
4332 repo: tmp.path().to_owned(),
4333 agent: "mock".to_owned(),
4334 status,
4335 turns: Vec::new(),
4336 pending: String::new(),
4337 pending_attachments: Vec::new(),
4338 fallback: false,
4339 persona: String::new(),
4340 persona_dirty: false,
4341 implementers: 1,
4342 implementers_dirty: false,
4343 created_at: Timestamp::now(),
4344 updated_at: Timestamp::now(),
4345 seat: SeatState::new(SEAT, "mock", 7),
4346 };
4347 talks.put(&mut t).expect("put");
4348 };
4349 make("20260901-000000-0001", TalkStatus::Open);
4350 make("20260902-000000-0002", TalkStatus::Open);
4351 make("20260903-000000-0003", TalkStatus::Closed);
4352
4353 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
4354 assert_eq!(
4355 ids,
4356 [
4357 "20260902-000000-0002",
4358 "20260901-000000-0001",
4359 "20260903-000000-0003"
4360 ]
4361 );
4362 assert_eq!(talks.count_open(), 2);
4363 }
4364
4365 #[test]
4366 fn tasks_of_finds_only_this_talks_own_tasks() {
4367 let dir = tempfile::tempdir().expect("tempdir");
4368 let queue = Queue::at(dir.path().join("queue"));
4369
4370 let mut mine = Task::new(
4371 "rework the loader".to_owned(),
4372 "rework the loader".to_owned(),
4373 PathBuf::from("/repo"),
4374 Source::Agent {
4375 run: "20260904-014455-ab12".to_owned(),
4376 node: "chat".to_owned(),
4377 },
4378 );
4379 queue.put(&mut mine).expect("put mine");
4380
4381 let mut theirs = Task::new(
4382 "unrelated".to_owned(),
4383 "unrelated".to_owned(),
4384 PathBuf::from("/repo"),
4385 Source::Agent {
4386 run: "20260904-090000-zz99".to_owned(),
4387 node: "implement".to_owned(),
4388 },
4389 );
4390 queue.put(&mut theirs).expect("put theirs");
4391
4392 let mut human = Task::new(
4393 "typed by hand".to_owned(),
4394 "typed by hand".to_owned(),
4395 PathBuf::from("/repo"),
4396 Source::Human,
4397 );
4398 queue.put(&mut human).expect("put human");
4399
4400 let found = tasks_of(&queue, "20260904-014455-ab12");
4401 assert_eq!(found.len(), 1);
4402 assert_eq!(found[0].id, mine.id);
4403 }
4404
4405 #[test]
4406 fn the_briefing_names_solo_task_add() {
4407 let brief = briefing(Path::new("/repo"), "en", false);
4408 assert!(brief.contains("magi task add --solo"));
4409 assert!(brief.contains("/repo"));
4410 assert!(!brief.contains("Hold this conversation in"));
4411 }
4412
4413 #[test]
4419 fn the_briefing_explains_targeting_a_different_repository_by_name() {
4420 let brief = briefing(Path::new("/repo"), "en", false);
4421 assert!(brief.contains("--repo does not have to be a full path"));
4422 assert!(brief.contains("owner/repo"));
4423 assert!(brief.contains("magi repos"));
4424 assert!(brief.contains("ask the operator"));
4425 }
4426
4427 #[test]
4428 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
4429 let brief = briefing(Path::new("/repo"), "en", false);
4430 assert!(brief.contains("--attach <path>"), "{brief}");
4431 assert!(brief.contains("deleting this conversation"), "{brief}");
4432 }
4433
4434 #[test]
4435 fn the_briefing_names_the_language_when_it_is_not_english() {
4436 let brief = briefing(Path::new("/repo"), "Japanese", false);
4437 assert!(brief.contains("Hold this conversation in Japanese"));
4438 }
4439
4440 #[test]
4441 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
4442 let read_only = briefing(Path::new("/repo"), "en", false);
4443 assert!(read_only.contains("Do not write files"));
4444 assert!(!read_only.contains("allow_write"));
4445
4446 let writable = briefing(Path::new("/repo"), "en", true);
4447 assert!(!writable.contains("Do not write files"));
4448 assert!(writable.contains("allow_write = true"));
4449 assert!(writable.contains("magi task add --solo"));
4452 assert!(writable.contains("say plainly what you"));
4453 }
4454}