1use std::path::{Path, PathBuf};
52use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
53use std::time::Duration;
54
55use anyhow::{Context, Result, bail};
56use jiff::Timestamp;
57use serde::{Deserialize, Serialize};
58
59use crate::agent::{self, Invocation, SeatState};
60use crate::config::{AgentSpec, Config};
61use crate::queue::{Queue, Source, Task};
62
63pub const SCHEMA: u32 = 1;
65
66fn turn_timeout(cfg: &Config) -> Duration {
77 Duration::from_secs(cfg.graph.timeout_talk)
78}
79
80const SEAT: &str = "talk";
83
84const MAGI_NOTE: &str = "magi: ";
86
87#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
89#[serde(rename_all = "lowercase")]
90pub enum Who {
91 Operator,
93 Agent,
96}
97
98#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
106#[serde(deny_unknown_fields)]
107pub struct Attachment {
108 pub id: String,
110 pub name: String,
112 pub mime: String,
115 pub bytes: u64,
117}
118
119#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
121#[serde(deny_unknown_fields)]
122pub struct Turn {
123 pub who: Who,
125 pub body: String,
127 pub at: Timestamp,
129 #[serde(default)]
132 pub attachments: Vec<Attachment>,
133 #[serde(default, skip_serializing_if = "Option::is_none")]
138 pub usage: Option<TurnUsage>,
139}
140
141#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
149pub struct TurnUsage {
150 pub context_tokens: u64,
152 pub agent: String,
154 #[serde(default)]
156 pub model: Option<String>,
157}
158
159#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)]
164pub struct ContextUsage {
165 pub tokens: Option<u64>,
167 pub window: Option<u64>,
169 pub percent: Option<u64>,
171 pub warn: bool,
173 pub since_switch: bool,
177 pub model: Option<String>,
179 #[serde(default)]
182 pub estimated: bool,
183}
184
185const CONTEXT_WARN_PERCENT: u64 = 80;
187
188const STANDING_PROMPT_FALLBACK_CHARS: u64 = 7000;
191
192pub fn estimate_context_tokens(talk: &Talk, standing_chars: u64) -> Option<u64> {
205 let mut counted = false;
206 let mut chars = standing_chars;
207 for t in talk.turns.iter().filter(|t| !t.body.starts_with(MAGI_NOTE)) {
208 counted = true;
209 chars += t.body.chars().count() as u64;
210 }
211 counted.then(|| (chars * 2).div_ceil(7))
213}
214
215pub fn context_usage(talk: &Talk, cfg: Option<&Config>) -> ContextUsage {
228 let current = cfg.and_then(|c| c.agents.iter().find(|a| a.id == talk.agent));
229 let model = current.and_then(|a| a.model.clone());
230 let window = cfg
231 .zip(model.as_deref())
232 .and_then(|(c, m)| c.context_window(m))
233 .filter(|w| *w > 0);
234 let usage = talk
235 .turns
236 .iter()
237 .rev()
238 .find(|t| t.who == Who::Agent && !t.body.starts_with(MAGI_NOTE))
239 .and_then(|t| t.usage.as_ref());
240 let measured = usage.map(|u| u.context_tokens);
241 let tokens = measured.or_else(|| {
242 let standing = cfg.map_or(STANDING_PROMPT_FALLBACK_CHARS, |c| {
243 briefing_with(
244 &talk.repo,
245 &c.graph.language,
246 c.talk.allow_write,
247 crate::persona::active(&c.talk.personas, &talk.persona).as_ref(),
248 c.talk.operator_name(),
249 )
250 .chars()
251 .count() as u64
252 });
253 estimate_context_tokens(talk, standing)
254 });
255 let estimated = measured.is_none() && tokens.is_some();
256 let since_switch =
257 usage.is_some_and(|u| u.agent != talk.agent || (current.is_some() && u.model != model));
258 let (percent, warn) = match (tokens, window) {
259 (Some(t), Some(w)) => (
260 Some(t.saturating_mul(100) / w),
261 t.saturating_mul(100) >= w.saturating_mul(CONTEXT_WARN_PERCENT),
262 ),
263 _ => (None, false),
264 };
265 ContextUsage {
266 tokens,
267 window,
268 percent,
269 warn,
270 since_switch,
271 model,
272 estimated,
273 }
274}
275
276#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
280#[serde(rename_all = "lowercase")]
281pub enum TalkStatus {
282 Open,
285 Closed,
287}
288
289impl TalkStatus {
290 pub fn open(self) -> bool {
292 matches!(self, Self::Open)
293 }
294
295 pub fn as_str(self) -> &'static str {
297 match self {
298 Self::Open => "open",
299 Self::Closed => "closed",
300 }
301 }
302}
303
304#[derive(Debug, Clone, Serialize, Deserialize)]
306#[serde(deny_unknown_fields)]
307pub struct Talk {
308 pub schema: u32,
310 pub id: String,
312 pub repo: PathBuf,
314 pub agent: String,
316 pub status: TalkStatus,
318 pub turns: Vec<Turn>,
320 #[serde(default)]
323 pub pending: String,
324 #[serde(default)]
326 pub pending_attachments: Vec<Attachment>,
327 #[serde(default)]
332 pub fallback: bool,
333 #[serde(default)]
336 pub persona: String,
337 #[serde(default)]
341 pub persona_dirty: bool,
342 pub created_at: Timestamp,
344 pub updated_at: Timestamp,
346 seat: SeatState,
351}
352
353impl Talk {
354 pub fn short(&self) -> &str {
356 short(&self.id)
357 }
358}
359
360#[derive(Debug, Clone)]
362pub struct Talks {
363 root: PathBuf,
364 lock: Arc<Mutex<()>>,
373}
374
375impl Talks {
376 pub fn open() -> Self {
378 Self::at(crate::run::home().join("talks"))
379 }
380
381 pub fn at(root: PathBuf) -> Self {
384 Self {
385 root,
386 lock: Arc::new(Mutex::new(())),
387 }
388 }
389
390 fn guard(&self) -> MutexGuard<'_, ()> {
399 self.lock.lock().unwrap_or_else(PoisonError::into_inner)
400 }
401
402 pub fn root(&self) -> &Path {
404 &self.root
405 }
406
407 pub fn path_of(&self, id: &str) -> PathBuf {
409 self.root.join(format!("{id}.json"))
410 }
411
412 pub fn artifacts_of(&self, id: &str) -> PathBuf {
415 self.root.join(format!("{id}.artifacts"))
416 }
417
418 pub fn attachments_dir(&self, id: &str) -> PathBuf {
422 self.artifacts_of(id).join("attachments")
423 }
424
425 pub fn put_attachment(
434 &self,
435 id: &str,
436 mime: &str,
437 name: &str,
438 data: &[u8],
439 ) -> Result<Attachment> {
440 let dir = self.attachments_dir(id);
441 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
442 let ext = attachment_ext(mime).with_context(|| format!("unsupported mime `{mime}`"))?;
443 let att = Attachment {
444 id: new_attachment_id(),
445 name: name.to_owned(),
446 mime: mime.to_owned(),
447 bytes: data.len() as u64,
448 };
449 std::fs::write(dir.join(format!("{}.{ext}", att.id)), data)
450 .with_context(|| format!("write attachment {}", att.id))?;
451 std::fs::write(
452 dir.join(format!("{}.json", att.id)),
453 serde_json::to_string(&att).context("serialize attachment")?,
454 )
455 .with_context(|| format!("write attachment metadata {}", att.id))?;
456 Ok(att)
457 }
458
459 pub fn attachment_meta(&self, id: &str, att_id: &str) -> Result<Option<Attachment>> {
468 if !valid_attachment_id(att_id) {
469 return Ok(None);
470 }
471 let meta_path = self.attachments_dir(id).join(format!("{att_id}.json"));
472 if !meta_path.is_file() {
473 return Ok(None);
474 }
475 let att = serde_json::from_str(
476 &std::fs::read_to_string(&meta_path)
477 .with_context(|| format!("read {}", meta_path.display()))?,
478 )
479 .with_context(|| format!("parse {}", meta_path.display()))?;
480 Ok(Some(att))
481 }
482
483 pub fn read_attachment(&self, id: &str, att_id: &str) -> Result<Option<(Attachment, Vec<u8>)>> {
487 let Some(att) = self.attachment_meta(id, att_id)? else {
488 return Ok(None);
489 };
490 let ext = attachment_ext(&att.mime).with_context(|| {
491 format!("attachment {att_id} has an unsupported mime `{}`", att.mime)
492 })?;
493 let data_path = self.attachments_dir(id).join(format!("{att_id}.{ext}"));
494 let data =
495 std::fs::read(&data_path).with_context(|| format!("read {}", data_path.display()))?;
496 Ok(Some((att, data)))
497 }
498
499 fn attachment_path(&self, id: &str, att: &Attachment) -> Option<PathBuf> {
516 let ext = attachment_ext(&att.mime)?;
517 let path = self.attachments_dir(id).join(format!("{}.{ext}", att.id));
518 std::path::absolute(&path).ok()
519 }
520
521 pub fn put(&self, t: &mut Talk) -> Result<()> {
530 std::fs::create_dir_all(&self.root)
531 .with_context(|| format!("create {}", self.root.display()))?;
532 t.updated_at = Timestamp::now();
533 let body = serde_json::to_string_pretty(t).context("serialize talk")?;
534 let path = self.path_of(&t.id);
535 let tmp = path.with_extension("json.tmp");
536 write_atomic(&tmp, &path, &body)
537 }
538
539 pub fn get(&self, id: &str) -> Result<Talk> {
541 let resolved = self.resolve_id(id)?;
542 read_path(&self.path_of(&resolved))
543 }
544
545 pub fn list(&self) -> Vec<Talk> {
548 self.list_counting_unreadable().0
549 }
550
551 pub fn list_counting_unreadable(&self) -> (Vec<Talk>, usize) {
554 let mut unreadable = 0;
555 let mut all: Vec<Talk> = std::fs::read_dir(&self.root)
556 .into_iter()
557 .flatten()
558 .flatten()
559 .map(|e| e.path())
560 .filter(|p| p.extension().is_some_and(|x| x == "json"))
561 .filter_map(|p| {
562 let talk = read_path(&p).ok();
563 if talk.is_none() {
564 unreadable += 1;
565 }
566 talk
567 })
568 .collect();
569 all.sort_unstable_by(|a, b| {
570 let rank = |t: &Talk| u8::from(!t.status.open());
571 rank(a).cmp(&rank(b)).then_with(|| b.id.cmp(&a.id))
572 });
573 (all, unreadable)
574 }
575
576 pub fn resolve_id(&self, prefix: &str) -> Result<String> {
578 if self.path_of(prefix).is_file() {
579 return Ok(prefix.to_owned());
580 }
581 let hits: Vec<String> = self
582 .list()
583 .into_iter()
584 .map(|t| t.id)
585 .filter(|id| id.starts_with(prefix) || id.ends_with(prefix))
586 .collect();
587 match hits.len() {
588 1 => Ok(hits.into_iter().next().expect("exactly one hit")),
589 0 => bail!("no talk matches `{prefix}`"),
590 _ => bail!(
591 "`{prefix}` matches {} talks: {}",
592 hits.len(),
593 hits.join(", ")
594 ),
595 }
596 }
597
598 pub fn revision(&self) -> u64 {
601 std::fs::read_dir(&self.root)
602 .into_iter()
603 .flatten()
604 .flatten()
605 .filter(|e| e.path().extension().is_none_or(|x| x != "turn"))
606 .filter_map(|e| e.metadata().ok())
607 .filter_map(|m| m.modified().ok())
608 .filter_map(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
609 .map(|d| d.as_millis() as u64)
610 .max()
611 .unwrap_or(0)
612 }
613
614 pub fn count_open(&self) -> usize {
616 self.list().iter().filter(|t| t.status.open()).count()
617 }
618
619 pub fn turn_path(&self, id: &str) -> PathBuf {
622 self.root.join(format!("{id}.turn"))
623 }
624
625 pub fn claim_turn(&self, id: &str) -> Result<Option<TurnLease>> {
634 self.claim_turn_at(id, Timestamp::now())
635 }
636
637 fn claim_turn_at(&self, id: &str, now: Timestamp) -> Result<Option<TurnLease>> {
638 std::fs::create_dir_all(&self.root)
639 .with_context(|| format!("create {}", self.root.display()))?;
640 let path = self.turn_path(id);
641 let token = fresh_token();
642 if create_turn(&path, &token, now)? {
643 return Ok(Some(TurnLease { path, token }));
644 }
645 if read_turn(&path).is_some_and(|r| r.fresh(now)) {
646 return Ok(None);
647 }
648 let Some(_lock) = TurnLock::take(&path)? else {
653 return Ok(None);
654 };
655 if read_turn(&path).is_some_and(|r| r.fresh(now)) {
656 return Ok(None);
657 }
658 let _ = std::fs::remove_file(&path);
659 Ok(create_turn(&path, &token, now)?.then_some(TurnLease { path, token }))
660 }
661
662 pub fn turn_held(&self, id: &str) -> bool {
664 read_turn(&self.turn_path(id)).is_some_and(|r| r.fresh(Timestamp::now()))
665 }
666
667 pub fn remove(&self, id: &str) -> Result<()> {
680 let _guard = self.guard();
681 let resolved = self.resolve_id(id)?;
682 let path = self.path_of(&resolved);
683 std::fs::remove_file(&path).with_context(|| format!("remove {}", path.display()))?;
684 let artifacts = self.artifacts_of(&resolved);
685 if artifacts.is_dir() {
686 std::fs::remove_dir_all(&artifacts)
687 .with_context(|| format!("remove {}", artifacts.display()))?;
688 }
689 let _ = std::fs::remove_file(self.turn_path(&resolved));
690 Ok(())
691 }
692}
693
694#[derive(Debug, Serialize, Deserialize)]
697struct TurnRecord {
698 token: String,
699 pid: u32,
700 beat_at: Timestamp,
701}
702
703impl TurnRecord {
704 fn fresh(&self, now: Timestamp) -> bool {
705 now.as_second() - self.beat_at.as_second() <= crate::ask::LEASE_TTL.as_secs() as i64
706 }
707}
708
709fn fresh_token() -> String {
713 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
714 let n = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
715 let seed = crate::rng::entropy() ^ n.wrapping_mul(0x9E37_79B9_7F4A_7C15);
716 crate::rng::SplitMix64::new(seed).uuid_v4()
717}
718
719fn read_turn(path: &Path) -> Option<TurnRecord> {
720 serde_json::from_str(&std::fs::read_to_string(path).ok()?).ok()
721}
722
723const TAKEOVER_LOCK_TTL: Duration = Duration::from_secs(10);
725
726const TICKET_TTL: Duration = Duration::from_secs(10);
729
730const TICKET_GENERATIONS: u32 = 16;
732
733const TICKET_SWEEP_AGE: Duration = Duration::from_secs(3600);
735
736fn create_exclusive(path: &Path, body: &str) -> Result<bool> {
738 use std::io::Write as _;
739 match std::fs::OpenOptions::new()
740 .write(true)
741 .create_new(true)
742 .open(path)
743 {
744 Ok(mut f) => {
745 if let Err(e) = f.write_all(body.as_bytes()) {
746 drop(f);
747 let _ = std::fs::remove_file(path);
748 return Err(e).with_context(|| format!("write {}", path.display()));
749 }
750 Ok(true)
751 }
752 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
753 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
754 }
755}
756
757fn create_turn(path: &Path, token: &str, now: Timestamp) -> Result<bool> {
758 let record = TurnRecord {
759 token: token.to_owned(),
760 pid: std::process::id(),
761 beat_at: now,
762 };
763 let body = serde_json::to_string(&record).context("serialize turn lease")?;
764 let tmp = path.with_extension(format!("turn.{token}.new"));
767 std::fs::write(&tmp, body).with_context(|| format!("write {}", tmp.display()))?;
768 let linked = std::fs::hard_link(&tmp, path);
769 let _ = std::fs::remove_file(&tmp);
770 match linked {
771 Ok(()) => Ok(true),
772 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
773 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
774 }
775}
776
777struct TurnLock {
781 path: PathBuf,
782 token: String,
783}
784
785impl TurnLock {
786 fn token() -> String {
787 fresh_token()
788 }
789
790 fn publish(path: &Path, token: &str) -> Result<bool> {
793 let tmp = path.with_extension(format!("lock.{token}.new"));
794 std::fs::write(&tmp, token).with_context(|| format!("write {}", tmp.display()))?;
795 let linked = std::fs::hard_link(&tmp, path);
796 let _ = std::fs::remove_file(&tmp);
797 match linked {
798 Ok(()) => Ok(true),
799 Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
800 Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
801 }
802 }
803
804 fn take(lease: &Path) -> Result<Option<Self>> {
805 let path = lease.with_extension("turn.lock");
806 let token = Self::token();
807 if Self::publish(&path, &token)? {
808 return Ok(Some(Self { path, token }));
809 }
810 let Ok(seen) = std::fs::read_to_string(&path) else {
811 return Ok(None);
812 };
813 let key = if !seen.is_empty() && seen.chars().all(|c| c.is_ascii_alphanumeric() || c == '-')
816 {
817 seen.as_str()
818 } else {
819 "invalid"
820 };
821 let aged = std::fs::metadata(&path)
822 .and_then(|m| m.modified())
823 .ok()
824 .and_then(|t| t.elapsed().ok())
825 .is_some_and(|age| age > TAKEOVER_LOCK_TTL);
826 if !aged {
827 return Ok(None);
828 }
829 let mut won = false;
841 for n in 0..TICKET_GENERATIONS {
842 let ticket = path.with_extension(format!("lock.{key}.break.{n}"));
843 if create_exclusive(&ticket, "")? {
844 won = true;
845 break;
846 }
847 let stale = std::fs::metadata(&ticket)
848 .and_then(|m| m.modified())
849 .ok()
850 .and_then(|t| t.elapsed().ok())
851 .is_some_and(|age| age > TICKET_TTL);
852 if !stale {
853 return Ok(None);
854 }
855 }
856 if !won {
857 Self::sweep_tickets(&path);
862 return Ok(None);
863 }
864 Self::sweep_tickets(&path);
865 if std::fs::read_to_string(&path).ok().as_deref() != Some(seen.as_str()) {
868 return Ok(None);
869 }
870 let _ = std::fs::remove_file(&path);
871 if Self::publish(&path, &token)? {
872 return Ok(Some(Self { path, token }));
873 }
874 Ok(None)
875 }
876
877 fn sweep_tickets(path: &Path) {
880 let (Some(dir), Some(name)) = (path.parent(), path.file_name().and_then(|n| n.to_str()))
881 else {
882 return;
883 };
884 let prefix = format!("{name}.");
885 let Ok(entries) = std::fs::read_dir(dir) else {
886 return;
887 };
888 for entry in entries.flatten() {
889 let file = entry.file_name();
890 let Some(file) = file.to_str() else { continue };
891 if !(file.starts_with(&prefix) && file.contains(".break.")) {
892 continue;
893 }
894 let old = entry
895 .metadata()
896 .and_then(|m| m.modified())
897 .ok()
898 .and_then(|t| t.elapsed().ok())
899 .is_some_and(|age| age > TICKET_SWEEP_AGE);
900 if old {
901 let _ = std::fs::remove_file(entry.path());
902 }
903 }
904 }
905
906 fn take_patiently(lease: &Path) -> Option<Self> {
908 for _ in 0..50 {
909 match Self::take(lease) {
910 Ok(Some(lock)) => return Some(lock),
911 Ok(None) => std::thread::sleep(Duration::from_millis(10)),
912 Err(_) => return None,
913 }
914 }
915 None
916 }
917}
918
919impl Drop for TurnLock {
920 fn drop(&mut self) {
921 if std::fs::read_to_string(&self.path).is_ok_and(|t| t == self.token) {
923 let _ = std::fs::remove_file(&self.path);
924 }
925 }
926}
927
928pub const TURN_BEAT: Duration = Duration::from_secs(20);
931
932#[derive(Debug)]
937pub struct TurnLease {
938 path: PathBuf,
939 token: String,
940}
941
942impl TurnLease {
943 pub fn beat(&self) -> Result<bool> {
947 let _lock = TurnLock::take_patiently(&self.path)
948 .with_context(|| format!("lock {} to renew it", self.path.display()))?;
949 let Some(mut record) = read_turn(&self.path).filter(|r| r.token == self.token) else {
950 return Ok(false);
951 };
952 record.beat_at = Timestamp::now();
953 let body = serde_json::to_string(&record).context("serialize turn lease")?;
954 let tmp = self.path.with_extension(format!("turn.{}.tmp", self.token));
955 write_atomic(&tmp, &self.path, &body)?;
956 Ok(true)
957 }
958
959 pub async fn beating<T>(&self, fut: impl std::future::Future<Output = T>) -> Result<T> {
964 self.beating_every(TURN_BEAT, fut).await
965 }
966
967 async fn beating_every<T>(
968 &self,
969 period: Duration,
970 fut: impl std::future::Future<Output = T>,
971 ) -> Result<T> {
972 tokio::pin!(fut);
973 loop {
974 match tokio::time::timeout(period, &mut fut).await {
975 Ok(out) => return Ok(out),
976 Err(_) => match self.beat() {
977 Ok(true) => {}
978 Ok(false) => bail!(
979 "the turn lease {} was taken over; this turn is stopped",
980 self.path.display()
981 ),
982 Err(e) => tracing::warn!("{e:#}"),
983 },
984 }
985 }
986 }
987}
988
989impl Drop for TurnLease {
990 fn drop(&mut self) {
991 if let Some(_lock) = TurnLock::take_patiently(&self.path) {
994 if read_turn(&self.path).is_some_and(|r| r.token == self.token) {
995 let _ = std::fs::remove_file(&self.path);
996 }
997 }
998 }
999}
1000
1001pub fn begin(store: &Talks, cfg: &Config, repo: PathBuf, agent: Option<&str>) -> Result<Talk> {
1011 let repo = repo.canonicalize().unwrap_or(repo);
1014 let spec = match agent {
1017 Some(id) => agent::pick(&cfg.agents, Some(id), &agent::installed)?,
1018 None => agent::pick_chain(
1019 &cfg.agents,
1020 cfg.roles.chatter.as_ref(),
1021 &agent::installed,
1022 "chatter",
1023 )?
1024 .remove(0),
1025 };
1026
1027 let now = Timestamp::now();
1028 let mut talk = Talk {
1029 schema: SCHEMA,
1030 id: new_id(),
1031 repo,
1032 agent: spec.id.clone(),
1033 status: TalkStatus::Open,
1034 turns: Vec::new(),
1035 pending: String::new(),
1036 pending_attachments: Vec::new(),
1037 fallback: agent.is_none(),
1038 persona: String::new(),
1039 persona_dirty: false,
1040 created_at: now,
1041 updated_at: now,
1042 seat: SeatState::new(SEAT, &spec.id, crate::rng::entropy()),
1043 };
1044 store.put(&mut talk)?;
1045 Ok(talk)
1046}
1047
1048pub fn record(
1055 talk: &mut Talk,
1056 store: &Talks,
1057 text: &str,
1058 attachments: Vec<Attachment>,
1059) -> Result<String> {
1060 let _guard = store.guard();
1068 let Ok(fresh) = store.get(&talk.id) else {
1073 bail!("talk {} was deleted", talk.short());
1074 };
1075 talk.status = fresh.status;
1076 talk.pending = fresh.pending;
1079 talk.pending_attachments = fresh.pending_attachments;
1080 if !talk.status.open() {
1081 bail!(
1082 "talk {} is {} and takes no more turns",
1083 talk.short(),
1084 talk.status.as_str()
1085 );
1086 }
1087 let text = text.trim();
1088 if text.is_empty() && attachments.is_empty() {
1089 bail!("nothing to say");
1090 }
1091 talk.turns.push(Turn {
1092 who: Who::Operator,
1093 body: text.to_owned(),
1094 at: Timestamp::now(),
1095 attachments,
1096 usage: None,
1097 });
1098 store.put(talk)?;
1099 Ok(text.to_owned())
1100}
1101
1102pub fn queue(
1104 talk: &mut Talk,
1105 store: &Talks,
1106 text: &str,
1107 attachments: Vec<Attachment>,
1108) -> Result<()> {
1109 let text = text.trim();
1110 if text.is_empty() && attachments.is_empty() {
1111 bail!("nothing to say");
1112 }
1113 let _guard = store.guard();
1114 let mut fresh = store
1115 .get(&talk.id)
1116 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1117 if !fresh.status.open() {
1118 bail!(
1119 "talk {} is {} and takes no more turns",
1120 fresh.short(),
1121 fresh.status.as_str()
1122 );
1123 }
1124 if !text.is_empty() {
1125 if fresh.pending.is_empty() {
1126 fresh.pending = text.to_owned();
1127 } else {
1128 fresh.pending.push_str("\n\n");
1129 fresh.pending.push_str(text);
1130 }
1131 }
1132 fresh.pending_attachments.extend(attachments);
1133 store.put(&mut fresh)?;
1134 *talk = fresh;
1135 Ok(())
1136}
1137
1138pub fn drain(talk: &mut Talk, store: &Talks) -> Result<Option<String>> {
1140 let _guard = store.guard();
1141 let mut fresh = store
1142 .get(&talk.id)
1143 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1144 if !fresh.status.open() || (fresh.pending.is_empty() && fresh.pending_attachments.is_empty()) {
1145 *talk = fresh;
1146 return Ok(None);
1147 }
1148 let text = std::mem::take(&mut fresh.pending);
1149 let attachments = std::mem::take(&mut fresh.pending_attachments);
1150 fresh.turns.push(Turn {
1151 who: Who::Operator,
1152 body: text.clone(),
1153 at: Timestamp::now(),
1154 attachments,
1155 usage: None,
1156 });
1157 store.put(&mut fresh)?;
1158 *talk = fresh;
1159 Ok(Some(text))
1160}
1161
1162pub async fn say(
1165 talk: &mut Talk,
1166 store: &Talks,
1167 cfg: &Config,
1168 text: &str,
1169 attachments: Vec<Attachment>,
1170) -> Result<()> {
1171 let text = record(talk, store, text, attachments)?;
1172 turn(talk, store, cfg, &text).await
1173}
1174
1175pub async fn respond(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1177 turn(talk, store, cfg, text).await
1178}
1179
1180pub fn close(talk: &mut Talk, store: &Talks) -> Result<()> {
1198 let _guard = store.guard();
1199 let mut fresh = store
1200 .get(&talk.id)
1201 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1202 fresh.status = TalkStatus::Closed;
1203 fresh.pending.clear();
1205 fresh.pending_attachments.clear();
1206 store.put(&mut fresh)?;
1207 *talk = fresh;
1208 Ok(())
1209}
1210
1211pub fn reopen(talk: &mut Talk, store: &Talks) -> Result<()> {
1222 let _guard = store.guard();
1223 let mut fresh = store
1224 .get(&talk.id)
1225 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1226 fresh.status = TalkStatus::Open;
1227 store.put(&mut fresh)?;
1228 *talk = fresh;
1229 Ok(())
1230}
1231
1232pub fn switch_agent(talk: &mut Talk, store: &Talks, spec: &AgentSpec) -> Result<bool> {
1245 let _guard = store.guard();
1246 let mut fresh = store
1247 .get(&talk.id)
1248 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1249 if fresh.agent == spec.id {
1250 *talk = fresh;
1251 return Ok(false);
1252 }
1253 let from = std::mem::replace(&mut fresh.agent, spec.id.clone());
1254 fresh.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1255 fresh.fallback = false;
1257 fresh.turns.push(Turn {
1258 who: Who::Agent,
1259 body: format!("{MAGI_NOTE}agent changed from {from} to {}", spec.id),
1260 at: Timestamp::now(),
1261 attachments: Vec::new(),
1262 usage: None,
1263 });
1264 store.put(&mut fresh)?;
1265 *talk = fresh;
1266 Ok(true)
1267}
1268
1269pub fn switch_persona(talk: &mut Talk, store: &Talks, id: &str) -> Result<bool> {
1277 let id = if id.trim() == crate::persona::DEFAULT_ID {
1278 ""
1279 } else {
1280 id.trim()
1281 };
1282 let _guard = store.guard();
1283 let mut fresh = store
1284 .get(&talk.id)
1285 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1286 if fresh.persona == id {
1287 *talk = fresh;
1288 return Ok(false);
1289 }
1290 let from = std::mem::replace(&mut fresh.persona, id.to_owned());
1291 fresh.persona_dirty = true;
1292 let label = |p: &str| {
1293 if p.is_empty() {
1294 crate::persona::DEFAULT_ID.to_owned()
1295 } else {
1296 p.to_owned()
1297 }
1298 };
1299 fresh.turns.push(Turn {
1300 who: Who::Agent,
1301 body: format!(
1302 "{MAGI_NOTE}persona changed from {} to {}",
1303 label(&from),
1304 label(id)
1305 ),
1306 at: Timestamp::now(),
1307 attachments: Vec::new(),
1308 usage: None,
1309 });
1310 store.put(&mut fresh)?;
1311 *talk = fresh;
1312 Ok(true)
1313}
1314
1315pub fn clear_pending(talk: &mut Talk, store: &Talks) -> Result<()> {
1317 let _guard = store.guard();
1318 let mut fresh = store
1319 .get(&talk.id)
1320 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1321 fresh.pending.clear();
1322 fresh.pending_attachments.clear();
1323 store.put(&mut fresh)?;
1324 *talk = fresh;
1325 Ok(())
1326}
1327
1328pub fn clear_pending_if_matches(
1330 talk: &mut Talk,
1331 store: &Talks,
1332 expected_text: &str,
1333 expected_attachments: &[String],
1334) -> Result<bool> {
1335 let _guard = store.guard();
1336 let mut fresh = store
1337 .get(&talk.id)
1338 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1339 if !pending_matches(&fresh, expected_text, expected_attachments) {
1340 *talk = fresh;
1341 return Ok(false);
1342 }
1343 fresh.pending.clear();
1344 fresh.pending_attachments.clear();
1345 store.put(&mut fresh)?;
1346 *talk = fresh;
1347 Ok(true)
1348}
1349
1350pub fn edit_pending_text(
1354 talk: &mut Talk,
1355 store: &Talks,
1356 text: &str,
1357 expected_text: &str,
1358 expected_attachments: &[String],
1359) -> Result<bool> {
1360 let _guard = store.guard();
1361 let mut fresh = store
1362 .get(&talk.id)
1363 .with_context(|| format!("talk {} was deleted", talk.short()))?;
1364 if !pending_matches(&fresh, expected_text, expected_attachments) {
1365 *talk = fresh;
1366 return Ok(false);
1367 }
1368 fresh.pending = text.trim().to_owned();
1369 store.put(&mut fresh)?;
1370 *talk = fresh;
1371 Ok(true)
1372}
1373
1374fn pending_matches(talk: &Talk, expected_text: &str, expected_attachments: &[String]) -> bool {
1375 talk.pending == expected_text
1376 && talk
1377 .pending_attachments
1378 .iter()
1379 .map(|attachment| &attachment.id)
1380 .eq(expected_attachments.iter())
1381}
1382
1383async fn turn(talk: &mut Talk, store: &Talks, cfg: &Config, text: &str) -> Result<()> {
1390 let spec = cfg
1391 .agents
1392 .iter()
1393 .find(|a| a.id == talk.agent)
1394 .with_context(|| {
1395 format!(
1396 "talk {} was opened with agent `{}`, which is no longer in \
1397 the roster; restore it in magi.toml or start a new \
1398 conversation",
1399 talk.short(),
1400 talk.agent
1401 )
1402 })?;
1403
1404 let last_note = attachment_note(
1408 store,
1409 &talk.id,
1410 talk.turns
1411 .last()
1412 .map_or(&[][..], |t| t.attachments.as_slice()),
1413 );
1414
1415 let attachment_paths: Vec<PathBuf> = talk
1421 .turns
1422 .iter()
1423 .flat_map(|t| t.attachments.iter())
1424 .filter_map(|a| store.attachment_path(&talk.id, a))
1425 .collect();
1426
1427 let questions = crate::ask::Questions::open();
1435 let consulted = crate::consult::pending_consults(&questions, &talk.id);
1436 let consult_roots: Vec<PathBuf> = if consulted {
1437 vec![questions.root().to_path_buf()]
1438 } else {
1439 Vec::new()
1440 };
1441
1442 let (allow_write, unsandboxed) = turn_access(cfg.talk.allow_write, consulted);
1443
1444 let artifacts = store.artifacts_of(&talk.id);
1445 let operator_turns = talk.turns.iter().filter(|t| t.who == Who::Operator).count();
1448 let stem = format!("turn-{}", operator_turns.max(1));
1449 let cache_dir = cfg.cache_dir();
1452
1453 let mut chain = vec![spec.clone()];
1456 if let Some(choice) = cfg.roles.chatter.as_ref()
1457 && talk.fallback
1458 {
1459 for id in choice.ids() {
1460 if id == talk.agent || chain.iter().any(|s| s.id == id) {
1461 continue;
1462 }
1463 match agent::pick(&cfg.agents, Some(id), &agent::installed) {
1464 Ok(s) => chain.push(s),
1465 Err(e) => tracing::warn!("[roles] chatter: skipping `{id}`: {e:#}"),
1466 }
1467 }
1468 }
1469
1470 let persona = crate::persona::active(&cfg.talk.personas, &talk.persona);
1471 let operator_name = cfg.talk.operator_name();
1472 let persona_update = if talk.persona_dirty {
1473 format!(
1474 "{}\n\n",
1475 crate::persona::update_block_for(persona.as_ref(), operator_name)
1476 )
1477 } else {
1478 String::new()
1479 };
1480
1481 let mut outcome = None;
1482 let mut fell_back_from: Option<String> = None;
1483 let mut first_try: Option<(String, SeatState)> = None;
1486 for (n, spec) in chain.iter().enumerate() {
1487 if n > 0 {
1488 if first_try.is_none() {
1489 first_try = Some((talk.agent.clone(), talk.seat.clone()));
1490 }
1491 tracing::warn!("chat: falling back from `{}` to `{}`", talk.agent, spec.id);
1492 fell_back_from.get_or_insert_with(|| talk.agent.clone());
1495 talk.agent = spec.id.clone();
1496 talk.seat = SeatState::new(SEAT, &spec.id, crate::rng::entropy());
1497 }
1498 let resuming = agent::has_session(spec.kind, &talk.seat, cfg.graph.sessions);
1499 let first_ever = talk.turns.len() <= 1;
1500 let body = if talk.seat.turns == 0 && first_ever {
1501 format!(
1502 "{}\n\n# Operator\n\n{text}{last_note}",
1503 briefing_with(
1504 &talk.repo,
1505 &cfg.graph.language,
1506 cfg.talk.allow_write,
1507 persona.as_ref(),
1508 operator_name
1509 )
1510 )
1511 } else if talk.seat.turns == 0 {
1512 format!(
1515 "{}\n\n{}\n\n# Operator\n\n{text}{last_note}",
1516 briefing_with(
1517 &talk.repo,
1518 &cfg.graph.language,
1519 cfg.talk.allow_write,
1520 persona.as_ref(),
1521 operator_name
1522 ),
1523 transcript(talk, store)
1524 )
1525 } else if resuming {
1526 format!("{persona_update}{text}{last_note}")
1527 } else {
1528 let mut standing = match (&persona, talk.persona_dirty) {
1532 (Some(p), false) => format!("{}\n", crate::persona::section_for(p, operator_name)),
1533 _ => persona_update.clone(),
1534 };
1535 if let (None, Some(n)) = (&persona, operator_name) {
1538 standing.push_str(&format!(
1539 "# Addressing the operator\n{}\n",
1540 crate::persona::addressing(n)
1541 ));
1542 }
1543 format!("{}\n\n{standing}{text}{last_note}", transcript(talk, store))
1544 };
1545 let attempt_stem = if n == 0 {
1546 stem.clone()
1547 } else {
1548 format!("{stem}-{}", spec.id)
1549 };
1550 let inv = Invocation {
1551 cwd: &talk.repo,
1552 prompt: &body,
1553 timeout: turn_timeout(cfg),
1554 allow_write,
1559 unsandboxed,
1560 sessions: cfg.graph.sessions,
1561 artifacts: &artifacts,
1562 stem: &attempt_stem,
1563 run: &talk.id,
1566 node: crate::queue::CHAT_NODE,
1567 cache_dir: cache_dir.as_deref(),
1568 attachments: &attachment_paths,
1569 writable: &consult_roots,
1570 };
1571 let result = agent::invoke(spec, &mut talk.seat, &inv).await;
1572 let advance = agent::chain_advances(&result);
1573 if n == 0 || !advance {
1574 outcome = Some(result);
1575 } else {
1576 tracing::warn!("chat: fallback agent `{}` also failed", spec.id);
1578 }
1579 if !advance {
1580 break;
1581 }
1582 }
1583 if outcome.as_ref().is_some_and(agent::chain_advances) {
1584 if let Some((id, seat)) = first_try {
1587 talk.agent = id;
1588 talk.seat = seat;
1589 fell_back_from = None;
1590 }
1591 }
1592 let outcome = outcome.expect("a chain holds at least one agent");
1593 let note = |why: String| Turn {
1594 who: Who::Agent,
1595 body: format!("{MAGI_NOTE}{why}"),
1596 at: Timestamp::now(),
1597 attachments: Vec::new(),
1598 usage: None,
1599 };
1600 let (reply, failure) = match outcome {
1601 Err(e) => (
1602 note(format!("could not run agent `{}`: {e}", talk.agent)),
1603 Some(format!("could not run agent `{}`: {e}", talk.agent)),
1604 ),
1605 Ok(out) if out.quota_exhausted() => {
1606 let reset = out
1607 .quota
1608 .as_ref()
1609 .and_then(|q| q.reset.clone())
1610 .map_or_else(String::new, |r| format!(" (resets {r})"));
1611 let why = format!(
1612 "agent `{}` is out of quota{reset}; your message is saved, so \
1613 say it again when the window reopens",
1614 talk.agent
1615 );
1616 (note(why.clone()), Some(why))
1617 }
1618 Ok(out) if out.timed_out => {
1619 let why = format!(
1620 "agent `{}` did not answer within {}s; your message is saved",
1621 talk.agent,
1622 turn_timeout(cfg).as_secs()
1623 );
1624 (note(why.clone()), Some(why))
1625 }
1626 Ok(out) if !out.usable() => {
1627 let why = format!(
1628 "agent `{}` produced no answer (exit {}); your message is saved",
1629 talk.agent,
1630 out.exit_code
1631 .map_or_else(|| "unknown".to_owned(), |c| c.to_string())
1632 );
1633 (note(why.clone()), Some(why))
1634 }
1635 Ok(out) => (
1636 Turn {
1637 who: Who::Agent,
1638 body: out.text.trim().to_owned(),
1639 at: Timestamp::now(),
1640 attachments: Vec::new(),
1641 usage: out.context_tokens.map(|context_tokens| TurnUsage {
1644 context_tokens,
1645 agent: talk.agent.clone(),
1646 model: cfg
1647 .agents
1648 .iter()
1649 .find(|a| a.id == talk.agent)
1650 .and_then(|a| a.model.clone()),
1651 }),
1652 },
1653 None,
1654 ),
1655 };
1656
1657 let _guard = store.guard();
1669 let Ok(fresh) = store.get(&talk.id) else {
1675 return Ok(());
1676 };
1677 talk.status = fresh.status;
1678 talk.pending = fresh.pending;
1682 talk.pending_attachments = fresh.pending_attachments;
1683 if let Some(from) = fell_back_from.filter(|_| failure.is_none()) {
1684 talk.turns.push(note(format!(
1687 "agent changed from {from} to {} (fallback)",
1688 talk.agent
1689 )));
1690 }
1691 if failure.is_none() {
1694 talk.persona_dirty = false;
1695 }
1696 talk.turns.push(reply);
1697 if let Err(put_err) = store.put(talk) {
1698 let lost = talk.turns.pop().expect("just pushed above");
1707 let stash = stash_lost_turn(store, &talk.id, &stem, &lost);
1708 let why = match &stash {
1709 Ok(path) => format!(
1710 "agent `{}` answered, but the reply could not be saved to \
1711 this conversation ({put_err:#}); the raw text was kept at \
1712 {} - your message is saved, ask again",
1713 talk.agent,
1714 path.display()
1715 ),
1716 Err(stash_err) => format!(
1717 "agent `{}` answered, but the reply could not be saved to \
1718 this conversation ({put_err:#}), and it could not be kept \
1719 anywhere else either ({stash_err:#}); your message is \
1720 saved, ask again",
1721 talk.agent
1722 ),
1723 };
1724 talk.turns.push(note(why.clone()));
1725 return match store.put(talk) {
1732 Ok(()) => bail!("{why}"),
1733 Err(note_err) => {
1734 talk.turns.pop();
1754 Err(note_err).context(why)
1755 }
1756 };
1757 }
1758
1759 match failure {
1760 Some(why) => bail!("{why}"),
1761 None => Ok(()),
1762 }
1763}
1764
1765fn transcript(talk: &Talk, store: &Talks) -> String {
1768 let mut out = String::from(
1769 "This conversation cannot resume on the CLI's side, so here is \
1770 everything said so far; answer only the last message.\n",
1771 );
1772 for t in &talk.turns {
1773 let who = match t.who {
1774 Who::Operator => "operator",
1775 Who::Agent if t.body.starts_with(MAGI_NOTE) => "magi",
1776 Who::Agent => "you",
1777 };
1778 out.push_str(&format!("\n## {who}\n\n{}\n", t.body.trim()));
1779 out.push_str(&attachment_note(store, &talk.id, &t.attachments));
1780 }
1781 out
1782}
1783
1784fn attachment_note(store: &Talks, talk_id: &str, attachments: &[Attachment]) -> String {
1789 if attachments.is_empty() {
1790 return String::new();
1791 }
1792 let mut out = String::from(
1793 "\n\nThe operator attached the image(s) below to this message. Open \
1794 and look at each one before you answer.\n",
1795 );
1796 for att in attachments {
1797 if let Some(path) = store.attachment_path(talk_id, att) {
1798 out.push_str(&format!("\n- {} ({})", path.display(), att.mime));
1799 }
1800 }
1801 out.push('\n');
1802 out
1803}
1804
1805pub(crate) fn turn_access(talk_allow_write: bool, consulted: bool) -> (bool, bool) {
1816 (talk_allow_write || consulted, talk_allow_write)
1817}
1818
1819pub fn briefing(repo: &Path, language: &str, allow_write: bool) -> String {
1835 briefing_with(repo, language, allow_write, None, None)
1836}
1837
1838pub fn briefing_with(
1843 repo: &Path,
1844 language: &str,
1845 allow_write: bool,
1846 persona: Option<&crate::persona::Persona>,
1847 operator_name: Option<&str>,
1848) -> String {
1849 let write_policy = if allow_write {
1850 "Write access is enabled for this conversation (`allow_write = \
1851 true`), so you may write files - but only a small, \
1852 already-decided edit the operator names outright in this \
1853 conversation, not an implementation. This is a permission on the \
1854 conversation as a whole, not a property of whichever repository \
1855 it happened to start in: if the operator names a different \
1856 repository for that small edit, the policy allows it there too. \
1857 Your own tool may still confine writes to the repository this \
1858 conversation started in regardless - if a write elsewhere is \
1859 refused, say so plainly rather than working around it. Once you \
1860 have made an edit, say plainly what you edited. Anything bigger, \
1861 or anything still open-ended, still goes through the queue below \
1862 rather than being done here."
1863 } else {
1864 "Do not write files. Implementing a change is not this \
1865 conversation's job; a separate, blind competition of agents does \
1866 that, and a repository this conversation has already edited would \
1867 make their diffs unjudgeable."
1868 };
1869 let mut out = format!(
1870 "You are magi's standing conversation partner for its operator, who \
1871 usually has this open on a phone. Keep replies short: no preamble, \
1872 no restating what they just said.\n\n\
1873 # Repository\n\n{repo}\n\n\
1874 You may look around: read files, run shell commands, search history, \
1875 run tests - whatever answers the question. {write_policy}\n\n\
1876 A short, command-shaped message (\"list\", \"info <id>\", \"show \
1877 3cbf\") is almost always the operator asking you to look something \
1878 up, not an instruction to file - answer it yourself with `magi \
1879 list`, `magi show <id>`, `magi task list`, or the like, the same way \
1880 you would answer any other question in this conversation.\n\n\
1881 # When the operator wants something done\n\n\
1882 Run:\n\n\
1883 magi task add --solo --repo {repo} <instruction>\n\n\
1884 and tell the operator the task id it prints, so they can follow it \
1885 from the Queue. If it refuses with a duplicate warning (the \
1886 instruction names a branch, commit or pull request that an \
1887 unfinished task, run or PR already owns), do not repeat it with \
1888 --force yourself: tell the operator what it matched and let them \
1889 decide. Write <instruction> so that an implementer who has \
1890 never seen this conversation can act on it alone - it is everything \
1891 they get. Use --solo: it runs the task through one implementer \
1892 straight into review instead of the usual multi-agent competition, \
1893 which is the right shape for a change this conversation has already \
1894 settled, rather than one still worth several independent takes.\n\n\
1895 If the operator asks for something in a different repository, \
1896 --repo does not have to be a full path: --repo owner/repo (or just \
1897 repo, when that is unambiguous) is resolved against local checkouts \
1898 the same way `magi repos` lists them. If the command fails because \
1899 nothing matches or more than one checkout shares that name, ask the \
1900 operator which repository they mean (or run `magi repos` yourself \
1901 to see the candidates) rather than guessing.\n\n\
1902 The current state of the code is whatever origin/main holds, not \
1903 whatever a working tree shows: a primary checkout often lags \
1904 upstream, sits on a detached HEAD and carries uncommitted changes. \
1905 Before answering about code, run `git fetch origin` in that \
1906 repository if it is cheap, then read through \
1907 `git show origin/main:<path>` or `git grep <pattern> origin/main`. \
1908 If the working tree differs, say so; if the fetch fails, say that \
1909 too, so the operator knows the answer may be stale.\n\n\
1910 If the operator attached an image (a screenshot, say) that the task \
1911 is about, pass it with `--attach <path>`, using the absolute path \
1912 the turn's attachment note gives; repeat the flag for several. \
1913 `magi task add --solo --attach <path> <instruction>` copies the \
1914 file into the task, so the implementer receives it. Do not paste the \
1915 path into <instruction> instead: deleting this conversation deletes \
1916 its attachments, and then that path reaches no one.\n",
1917 repo = repo.display(),
1918 );
1919 out.push_str(&language_note(language));
1920 match (persona, operator_name) {
1921 (Some(p), n) => out.push_str(&crate::persona::section_for(p, n)),
1922 (None, Some(n)) => {
1923 out.push_str("\n# Addressing the operator\n");
1924 out.push_str(&crate::persona::addressing(n));
1925 }
1926 (None, None) => {}
1927 }
1928 out
1929}
1930
1931fn language_note(language: &str) -> String {
1934 if language.trim().is_empty() || language.eq_ignore_ascii_case("en") {
1935 String::new()
1936 } else {
1937 format!("\nHold this conversation in {language}.\n")
1938 }
1939}
1940
1941pub fn tasks_of(queue: &Queue, talk_id: &str) -> Vec<Task> {
1948 let mut tasks: Vec<Task> = queue
1949 .list()
1950 .into_iter()
1951 .filter(|t| matches!(&t.source, Source::Agent { run, .. } if run == talk_id))
1952 .collect();
1953 tasks.sort_unstable_by(|a, b| a.id.cmp(&b.id));
1954 tasks
1955}
1956
1957fn read_path(path: &Path) -> Result<Talk> {
1958 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
1959 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))
1960}
1961
1962const PUT_RETRIES: u32 = 5;
1965
1966fn write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1978 let mut last_err = None;
1979 for attempt in 0..PUT_RETRIES {
1980 if attempt > 0 {
1981 std::thread::sleep(Duration::from_millis(20 * u64::from(attempt)));
1982 }
1983 match try_write_atomic(tmp, path, body) {
1984 Ok(()) => return Ok(()),
1985 Err(e) => last_err = Some(e),
1986 }
1987 }
1988 Err(last_err.expect("the loop above always runs at least once"))
1989}
1990
1991fn try_write_atomic(tmp: &Path, path: &Path, body: &str) -> Result<()> {
1992 #[cfg(test)]
1993 if failpoint::take_forced_put_failure() {
1994 bail!("simulated write failure (test)");
1995 }
1996 std::fs::write(tmp, body).with_context(|| format!("write {}", tmp.display()))?;
1997 std::fs::rename(tmp, path).with_context(|| format!("replace {}", path.display()))?;
1998 Ok(())
1999}
2000
2001fn stash_lost_turn(store: &Talks, id: &str, stem: &str, reply: &Turn) -> Result<PathBuf> {
2006 let dir = store.artifacts_of(id);
2007 std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
2008 let path = dir.join(format!("{stem}-lost.txt"));
2009 std::fs::write(&path, &reply.body).with_context(|| format!("write {}", path.display()))?;
2010 Ok(path)
2011}
2012
2013#[cfg(test)]
2020mod failpoint {
2021 use std::cell::Cell;
2022
2023 thread_local! {
2024 static FORCE_PUT_FAILURES: Cell<u32> = const { Cell::new(0) };
2025 }
2026
2027 pub(super) fn force_put_failures(count: u32) {
2030 FORCE_PUT_FAILURES.with(|c| c.set(count));
2031 }
2032
2033 pub(super) fn take_forced_put_failure() -> bool {
2036 FORCE_PUT_FAILURES.with(|c| {
2037 let n = c.get();
2038 if n == 0 {
2039 false
2040 } else {
2041 c.set(n - 1);
2042 true
2043 }
2044 })
2045 }
2046}
2047
2048fn short(id: &str) -> &str {
2049 id.split('-').next_back().unwrap_or(id)
2050}
2051
2052fn new_id() -> String {
2053 let stamp = jiff::Zoned::now().strftime("%Y%m%d-%H%M%S");
2054 let seed = crate::rng::entropy();
2055 format!("{stamp}-{:04x}", (seed ^ (seed >> 32)) & 0xffff)
2056}
2057
2058fn attachment_ext(mime: &str) -> Option<&'static str> {
2063 match mime {
2064 "image/png" => Some("png"),
2065 "image/jpeg" => Some("jpg"),
2066 "image/gif" => Some("gif"),
2067 "image/webp" => Some("webp"),
2068 _ => None,
2069 }
2070}
2071
2072pub fn valid_attachment_id(id: &str) -> bool {
2077 id.len() == 32
2078 && id
2079 .bytes()
2080 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
2081}
2082
2083fn new_attachment_id() -> String {
2087 let mut r = crate::rng::SplitMix64::new(crate::rng::entropy());
2088 format!("{:016x}{:016x}", r.next_u64(), r.next_u64())
2089}
2090
2091#[cfg(test)]
2092mod tests {
2093 #[test]
2094 fn turn_access_lifts_the_sandbox_only_for_the_talk_opt_in() {
2095 use super::turn_access;
2096 assert_eq!(turn_access(false, false), (false, false));
2097 assert_eq!(turn_access(true, false), (true, true));
2098 assert_eq!(turn_access(false, true), (true, false));
2100 assert_eq!(turn_access(true, true), (true, true));
2101 }
2102
2103 #[test]
2104 fn the_briefing_points_at_origin_main_not_the_working_tree() {
2105 let b = briefing(Path::new("/r"), "en", false);
2106 assert!(b.contains("origin/main"));
2107 assert!(b.contains("git show origin/main:"));
2108 }
2109 use std::collections::BTreeMap;
2110
2111 use crate::config::{AgentChoice, AgentKind, AgentSpec, Graph};
2112 use crate::queue::{Queue, Source, Task};
2113
2114 use super::*;
2115
2116 fn ctx_agent(id: &str, model: Option<&str>) -> AgentSpec {
2117 AgentSpec {
2118 id: id.to_owned(),
2119 kind: AgentKind::Command,
2120 model: model.map(str::to_owned),
2121 command: Vec::new(),
2122 extra_args: Vec::new(),
2123 env: BTreeMap::new(),
2124 prompt_delivery: None,
2125 }
2126 }
2127
2128 fn ctx_talk(agent: &str, turns: Vec<Turn>) -> Talk {
2129 Talk {
2130 schema: SCHEMA,
2131 id: "20260904-014455-ab12".to_owned(),
2132 repo: PathBuf::from("."),
2133 agent: agent.to_owned(),
2134 status: TalkStatus::Open,
2135 turns,
2136 pending: String::new(),
2137 pending_attachments: Vec::new(),
2138 fallback: false,
2139 persona: String::new(),
2140 persona_dirty: false,
2141 created_at: Timestamp::now(),
2142 updated_at: Timestamp::now(),
2143 seat: SeatState::new(SEAT, agent, 1),
2144 }
2145 }
2146
2147 fn reply(body: &str, usage: Option<(u64, &str, Option<&str>)>) -> Turn {
2148 Turn {
2149 who: Who::Agent,
2150 body: body.to_owned(),
2151 at: Timestamp::now(),
2152 attachments: Vec::new(),
2153 usage: usage.map(|(t, a, m)| TurnUsage {
2154 context_tokens: t,
2155 agent: a.to_owned(),
2156 model: m.map(str::to_owned),
2157 }),
2158 }
2159 }
2160
2161 fn ctx_config(windows: &[(&str, u64)]) -> Config {
2162 Config {
2163 agents: vec![
2164 ctx_agent("small", Some("small-model")),
2165 ctx_agent("big", Some("big-model")),
2166 ctx_agent("plain", None),
2167 ],
2168 context_windows: windows.iter().map(|(k, v)| ((*k).to_owned(), *v)).collect(),
2169 ..Config::default()
2170 }
2171 }
2172
2173 #[test]
2174 fn context_usage_computes_percent_and_warns_at_eighty() {
2175 let cfg = ctx_config(&[("small-model", 1000)]);
2176 let at = |tokens| {
2177 let t = ctx_talk(
2178 "small",
2179 vec![reply("hi", Some((tokens, "small", Some("small-model"))))],
2180 );
2181 context_usage(&t, Some(&cfg))
2182 };
2183 let u = at(799);
2184 assert_eq!((u.percent, u.warn, u.window), (Some(79), false, Some(1000)));
2185 let u = at(800);
2186 assert_eq!((u.percent, u.warn), (Some(80), true));
2187 let u = at(1500);
2188 assert_eq!((u.percent, u.warn), (Some(150), true));
2189 assert!(!u.since_switch);
2190 }
2191
2192 #[test]
2193 fn context_usage_is_unknown_without_usage_and_never_looks_back() {
2194 let cfg = ctx_config(&[("small-model", 1000)]);
2195 let t = ctx_talk(
2196 "small",
2197 vec![
2198 reply("old", Some((900, "small", Some("small-model")))),
2199 reply("new", None),
2200 ],
2201 );
2202 let u = context_usage(&t, Some(&cfg));
2203 assert!(u.estimated);
2205 assert_ne!(u.tokens, Some(900));
2206 assert!(u.tokens.is_some());
2207 let t = ctx_talk(
2209 "small",
2210 vec![
2211 reply("old", Some((900, "small", Some("small-model")))),
2212 reply("magi: could not run agent", None),
2213 ],
2214 );
2215 assert_eq!(context_usage(&t, Some(&cfg)).tokens, Some(900));
2216 assert_eq!(
2217 context_usage(&ctx_talk("small", Vec::new()), Some(&cfg)).tokens,
2218 None
2219 );
2220 }
2221
2222 #[test]
2223 fn estimate_counts_chars_both_sides_and_standing_prompt() {
2224 let mut t = ctx_talk("small", vec![reply("abcdefg", None)]);
2225 assert_eq!(estimate_context_tokens(&t, 0), Some(2)); let op = Turn {
2227 who: Who::Operator,
2228 ..reply("abcdefg", None)
2229 };
2230 t.turns.push(op);
2231 assert_eq!(estimate_context_tokens(&t, 0), Some(4));
2232 assert!(
2233 estimate_context_tokens(&t, 700).unwrap() > estimate_context_tokens(&t, 0).unwrap()
2234 );
2235 let ja = ctx_talk("small", vec![reply("日本語日本語日", None)]);
2237 assert_eq!(estimate_context_tokens(&ja, 0), Some(2));
2238 let note = ctx_talk("small", vec![reply("magi: could not run agent", None)]);
2240 assert_eq!(estimate_context_tokens(¬e, 1000), None);
2241 assert_eq!(
2242 estimate_context_tokens(&ctx_talk("small", Vec::new()), 1000),
2243 None
2244 );
2245 }
2246
2247 #[test]
2248 fn context_usage_measured_wins_and_estimate_gets_percent_and_warn() {
2249 let cfg = ctx_config(&[("small-model", 1000)]);
2250 let t = ctx_talk(
2251 "small",
2252 vec![reply(
2253 &"x".repeat(5000),
2254 Some((10, "small", Some("small-model"))),
2255 )],
2256 );
2257 let u = context_usage(&t, Some(&cfg));
2258 assert_eq!((u.tokens, u.estimated), (Some(10), false));
2259 let t = ctx_talk("small", vec![reply(&"x".repeat(5000), None)]);
2260 let u = context_usage(&t, Some(&cfg));
2261 assert!(u.estimated && !u.since_switch);
2262 assert_eq!(u.window, Some(1000));
2263 assert!(u.warn && u.percent.unwrap() >= 80);
2264 let t = ctx_talk("small", vec![reply("hi", None)]);
2265 let u = context_usage(&t, Some(&cfg));
2266 assert!(u.estimated && u.percent.is_some());
2267 }
2268
2269 #[test]
2270 fn context_usage_without_a_window_shows_tokens_only() {
2271 let cfg = ctx_config(&[]);
2272 let t = ctx_talk("plain", vec![reply("hi", Some((5000, "plain", None)))]);
2274 let u = context_usage(&t, Some(&cfg));
2275 assert_eq!(
2276 (u.tokens, u.window, u.percent, u.warn),
2277 (Some(5000), None, None, false)
2278 );
2279 let t = ctx_talk(
2280 "small",
2281 vec![reply("hi", Some((5000, "small", Some("small-model"))))],
2282 );
2283 assert_eq!(context_usage(&t, Some(&cfg)).percent, None);
2284 assert_eq!(context_usage(&t, None).window, None);
2286 }
2287
2288 #[test]
2289 fn context_usage_switching_model_changes_the_denominator() {
2290 let cfg = ctx_config(&[("small-model", 1000), ("big-model", 10_000)]);
2291 let used = reply("hi", Some((900, "small", Some("small-model"))));
2292 let before = context_usage(&ctx_talk("small", vec![used.clone()]), Some(&cfg));
2293 assert_eq!(
2294 (before.percent, before.warn, before.since_switch),
2295 (Some(90), true, false)
2296 );
2297 let after = context_usage(&ctx_talk("big", vec![used]), Some(&cfg));
2300 assert_eq!(after.window, Some(10_000));
2301 assert_eq!(
2302 (after.percent, after.warn, after.since_switch),
2303 (Some(9), false, true)
2304 );
2305 assert_eq!(after.model.as_deref(), Some("big-model"));
2306 }
2307
2308 #[test]
2309 fn a_turn_recorded_before_usage_existed_still_reads() {
2310 let old = r#"{"who":"agent","body":"hi","at":"2026-09-04T01:44:55Z"}"#;
2311 let turn: Turn = serde_json::from_str(old).expect("old turn reads");
2312 assert!(turn.usage.is_none());
2313 let json = serde_json::to_string(&turn).expect("serialize");
2314 assert!(
2315 !json.contains("usage"),
2316 "absent usage is not written: {json}"
2317 );
2318 }
2319
2320 fn store() -> (tempfile::TempDir, Talks) {
2322 let tmp = tempfile::tempdir().expect("tempdir");
2323 let talks = Talks::at(tmp.path().join("talks"));
2324 (tmp, talks)
2325 }
2326
2327 fn mock_agent(dir: &Path, script: &str, env: BTreeMap<String, String>) -> AgentSpec {
2331 let path = dir.join("mock-talk-agent.sh");
2332 std::fs::write(&path, script).expect("write mock");
2333 AgentSpec {
2334 id: "mock".to_owned(),
2335 kind: AgentKind::Command,
2336 model: None,
2337 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2338 extra_args: Vec::new(),
2339 env,
2340 prompt_delivery: None,
2341 }
2342 }
2343
2344 fn config(spec: AgentSpec) -> Config {
2345 Config {
2346 agents: vec![spec],
2347 graph: Graph {
2348 language: "en".to_owned(),
2349 ..Graph::default()
2350 },
2351 ..Config::default()
2352 }
2353 }
2354
2355 const REPLY: &str = "#!/bin/sh\ncat >/dev/null\nprintf '%s\\n' \"$MOCK_REPLY\"\n";
2357
2358 const BROKEN: &str = "#!/bin/sh\ncat >/dev/null\nexit 3\n";
2360
2361 const ECHO: &str = "#!/bin/sh\ncat\n";
2364
2365 fn env(reply: &str) -> BTreeMap<String, String> {
2366 BTreeMap::from([("MOCK_REPLY".to_owned(), reply.to_owned())])
2367 }
2368
2369 #[test]
2370 fn the_frozen_json_field_names_round_trip_through_disk() {
2371 let (tmp, talks) = store();
2372 let mut talk = Talk {
2373 schema: SCHEMA,
2374 id: "20260904-014455-ab12".to_owned(),
2375 repo: tmp.path().to_owned(),
2376 agent: "sonnet".to_owned(),
2377 status: TalkStatus::Open,
2378 turns: Vec::new(),
2379 pending: String::new(),
2380 pending_attachments: Vec::new(),
2381 fallback: false,
2382 persona: String::new(),
2383 persona_dirty: false,
2384 created_at: Timestamp::now(),
2385 updated_at: Timestamp::now(),
2386 seat: SeatState::new(SEAT, "sonnet", 7),
2387 };
2388 talks.put(&mut talk).expect("put");
2389
2390 let raw = std::fs::read_to_string(talks.path_of(&talk.id)).expect("read back");
2391 let v: serde_json::Value = serde_json::from_str(&raw).expect("parse");
2392 for field in [
2393 "schema",
2394 "id",
2395 "repo",
2396 "agent",
2397 "status",
2398 "turns",
2399 "created_at",
2400 "updated_at",
2401 ] {
2402 assert!(v.get(field).is_some(), "missing field `{field}`");
2403 }
2404 assert_eq!(v["schema"], 1);
2405 assert_eq!(v["status"], "open");
2406
2407 let back = talks.get(&talk.id).expect("get");
2408 assert_eq!(back.id, talk.id);
2409 assert_eq!(back.status, TalkStatus::Open);
2410 }
2411
2412 #[test]
2413 fn opening_a_talk_takes_no_agent_turn() {
2414 let (tmp, talks) = store();
2415 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2419 let cfg = config(spec);
2420
2421 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2422 assert_eq!(talk.status, TalkStatus::Open);
2423 assert!(talk.turns.is_empty(), "nothing has been said yet");
2424
2425 let on_disk = talks.get(&talk.id).expect("get");
2426 assert_eq!(on_disk.turns.len(), 0);
2427 }
2428
2429 #[test]
2437 fn chatter_wins_when_set_and_falls_back_to_pick_s_default_order_otherwise() {
2438 let (tmp, talks) = store();
2439 let first_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2440 let mut chatter_spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
2441 chatter_spec.id = "chatter-mock".to_owned();
2442
2443 let mut cfg = Config {
2444 agents: vec![first_spec.clone(), chatter_spec.clone()],
2445 graph: Graph {
2446 language: "en".to_owned(),
2447 ..Graph::default()
2448 },
2449 ..Config::default()
2450 };
2451 cfg.roles.chatter = Some(chatter_spec.id.as_str().into());
2452
2453 let talk =
2454 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter set");
2455 assert_eq!(talk.agent, chatter_spec.id, "an explicit chatter must win");
2456
2457 cfg.roles.chatter = None;
2458 let fallback =
2459 begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin with chatter unset");
2460 assert_eq!(
2461 fallback.agent, first_spec.id,
2462 "unset chatter must fall back to agent::pick's own default order"
2463 );
2464 }
2465
2466 #[test]
2469 fn a_talk_recorded_without_attachments_still_reads() {
2470 let (tmp, talks) = store();
2471 let path = talks.path_of("20260904-014455-ab12");
2472 std::fs::create_dir_all(talks.root()).expect("talks dir");
2473 std::fs::write(
2474 &path,
2475 serde_json::json!({
2476 "schema": 1,
2477 "id": "20260904-014455-ab12",
2478 "repo": tmp.path(),
2479 "agent": "sonnet",
2480 "status": "open",
2481 "turns": [
2482 { "who": "operator", "body": "still there?",
2483 "at": Timestamp::now().to_string() },
2484 ],
2485 "created_at": Timestamp::now().to_string(),
2486 "updated_at": Timestamp::now().to_string(),
2487 "seat": SeatState::new(SEAT, "sonnet", 7),
2488 })
2489 .to_string(),
2490 )
2491 .expect("write pre-attachments talk");
2492
2493 let talk = talks.get("20260904-014455-ab12").expect("must still read");
2494 assert!(talk.turns[0].attachments.is_empty());
2495 }
2496
2497 fn lease_store() -> (tempfile::TempDir, Talks, Talks) {
2498 let tmp = tempfile::TempDir::new().expect("tmp");
2499 let root = tmp.path().join("talks");
2500 (tmp, Talks::at(root.clone()), Talks::at(root))
2501 }
2502
2503 #[test]
2504 fn two_starters_on_one_talk_one_wins_and_the_other_is_refused() {
2505 let (_tmp, a, b) = lease_store();
2506 let won = a.claim_turn("t1").expect("claim").expect("first wins");
2507 assert!(
2508 b.claim_turn("t1").expect("claim").is_none(),
2509 "second is refused"
2510 );
2511 assert!(b.turn_held("t1"));
2512 assert!(
2513 b.claim_turn("t2").expect("claim").is_some(),
2514 "other talks are free"
2515 );
2516 drop(won);
2517 }
2518
2519 #[test]
2520 fn a_stale_lease_is_taken_over_and_the_old_guard_cannot_release_it() {
2521 let (_tmp, a, b) = lease_store();
2522 let old = a.claim_turn("t1").expect("claim").expect("held");
2523 let later = Timestamp::now()
2524 .checked_add(jiff::SignedDuration::from_secs(
2525 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2526 ))
2527 .expect("later");
2528 let new = b
2529 .claim_turn_at("t1", later)
2530 .expect("claim")
2531 .expect("a stale lease is taken over");
2532 drop(old);
2533 assert!(a.turn_held("t1"), "the old guard left the new lease alone");
2534 assert!(new.beat().expect("beat"), "the new owner still beats");
2535 drop(new);
2536 assert!(!a.turn_held("t1"));
2537 }
2538
2539 #[tokio::test]
2540 async fn a_turn_whose_lease_was_taken_over_is_stopped() {
2541 let (_tmp, a, b) = lease_store();
2542 let old = a.claim_turn("t1").expect("claim").expect("held");
2543 let later = Timestamp::now()
2544 .checked_add(jiff::SignedDuration::from_secs(
2545 crate::ask::LEASE_TTL.as_secs() as i64 + 5,
2546 ))
2547 .expect("later");
2548 let _new = b
2549 .claim_turn_at("t1", later)
2550 .expect("claim")
2551 .expect("taken over");
2552 let out = old
2553 .beating_every(Duration::from_millis(10), std::future::pending::<()>())
2554 .await;
2555 assert!(out.is_err(), "the displaced turn must stop, not run on");
2556 }
2557
2558 #[tokio::test]
2559 async fn a_turn_that_finishes_is_returned_and_keeps_its_lease_beating() {
2560 let (_tmp, a, _b) = lease_store();
2561 let lease = a.claim_turn("t1").expect("claim").expect("held");
2562 let out = lease
2563 .beating_every(Duration::from_millis(5), async {
2564 tokio::time::sleep(Duration::from_millis(40)).await;
2565 7
2566 })
2567 .await
2568 .expect("still ours");
2569 assert_eq!(out, 7);
2570 assert!(a.turn_held("t1"));
2571 }
2572
2573 #[test]
2574 fn an_unreadable_lease_counts_as_stale() {
2575 let (_tmp, a, b) = lease_store();
2576 std::fs::create_dir_all(a.root()).expect("dir");
2577 std::fs::write(a.turn_path("t1"), "not json").expect("write");
2578 assert!(!a.turn_held("t1"));
2579 assert!(b.claim_turn("t1").expect("claim").is_some());
2580 }
2581
2582 #[test]
2583 fn a_lease_is_released_when_the_turn_ends_or_fails() {
2584 let (_tmp, a, b) = lease_store();
2585 let lease = a.claim_turn("t1").expect("claim").expect("held");
2586 let failed: Result<()> = (|| {
2587 let _held = &lease;
2588 bail!("turn failed")
2589 })();
2590 assert!(failed.is_err());
2591 assert!(
2592 b.claim_turn("t1").expect("claim").is_none(),
2593 "held mid-turn"
2594 );
2595 drop(lease);
2596 assert!(
2597 b.claim_turn("t1").expect("claim").is_some(),
2598 "free after the turn"
2599 );
2600 }
2601
2602 #[test]
2603 fn concurrent_takeovers_of_a_stale_lease_have_one_winner() {
2604 let (_tmp, a, _b) = lease_store();
2605 drop(a.claim_turn("t1").expect("claim").expect("held"));
2606 std::fs::write(
2607 a.turn_path("t1"),
2608 serde_json::to_string(&TurnRecord {
2609 token: "gone".into(),
2610 pid: 1,
2611 beat_at: Timestamp::from_second(1).expect("ts"),
2612 })
2613 .expect("json"),
2614 )
2615 .expect("write");
2616 let wins: Vec<_> = std::thread::scope(|sc| {
2617 let hs: Vec<_> = (0..8)
2618 .map(|_| {
2619 let s = a.clone();
2620 sc.spawn(move || s.claim_turn("t1").expect("claim"))
2621 })
2622 .collect();
2623 hs.into_iter().map(|h| h.join().expect("join")).collect()
2624 });
2625 assert_eq!(wins.iter().filter(|w| w.is_some()).count(), 1);
2626 }
2627
2628 fn age_file(path: &Path) {
2629 let f = std::fs::OpenOptions::new()
2630 .write(true)
2631 .open(path)
2632 .expect("open");
2633 f.set_modified(std::time::SystemTime::now() - std::time::Duration::from_secs(60))
2634 .expect("age");
2635 }
2636
2637 #[test]
2638 fn a_late_taker_of_a_broken_lock_cannot_disturb_its_replacement() {
2639 let (_tmp, a, _b) = lease_store();
2640 let lease = a.turn_path("t1");
2641 let lock = lease.with_extension("turn.lock");
2642 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2643 std::fs::write(&lock, "t1-dead").expect("dead lock");
2644 age_file(&lock);
2645 let b = TurnLock::take(&lease).expect("take").expect("b wins");
2647 let fresh = std::fs::read_to_string(&lock).expect("read");
2648 assert_eq!(fresh, b.token);
2649 let ticket = std::fs::read_dir(lock.parent().expect("dir"))
2652 .expect("dir")
2653 .flatten()
2654 .map(|e| e.path())
2655 .find(|p| p.to_string_lossy().ends_with(".break.0"))
2656 .expect("ticket");
2657 assert!(!create_exclusive(&ticket, "").expect("ticket"));
2658 assert!(TurnLock::take(&lease).expect("take").is_none());
2659 assert_eq!(std::fs::read_to_string(&lock).expect("read"), fresh);
2660 age_file(&lock);
2662 age_file(&ticket);
2663 std::mem::forget(b);
2664 let c = TurnLock::take(&lease)
2665 .expect("take")
2666 .expect("next generation");
2667 assert_ne!(c.token, fresh);
2668 }
2669
2670 #[test]
2671 fn a_live_ticket_blocks_and_a_stale_one_hands_over_to_the_next_generation() {
2672 let (_tmp, a, _b) = lease_store();
2673 let lease = a.turn_path("t1");
2674 let lock = lease.with_extension("turn.lock");
2675 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2676 std::fs::write(&lock, "t1-dead").expect("dead lock");
2677 age_file(&lock);
2678 let t0 = lock.with_extension("lock.t1-dead.break.0");
2679 assert!(create_exclusive(&t0, "").expect("ticket"));
2680 assert!(TurnLock::take(&lease).expect("take").is_none());
2682 assert_eq!(std::fs::read_to_string(&lock).expect("read"), "t1-dead");
2683 age_file(&t0);
2685 let c = TurnLock::take(&lease).expect("take").expect("generation 1");
2686 assert_eq!(std::fs::read_to_string(&lock).expect("read"), c.token);
2687 assert!(lock.with_extension("lock.t1-dead.break.1").exists());
2688 assert!(TurnLock::take(&lease).expect("take").is_none());
2689 }
2690
2691 #[test]
2692 fn exhausted_ticket_generations_recover_once_the_sweep_ages_them_out() {
2693 let (_tmp, a, _b) = lease_store();
2694 let lease = a.turn_path("t1");
2695 let lock = lease.with_extension("turn.lock");
2696 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2697 std::fs::write(&lock, "t1-dead").expect("dead lock");
2698 age_file(&lock);
2699 let tickets: Vec<_> = (0..TICKET_GENERATIONS)
2700 .map(|n| lock.with_extension(format!("lock.t1-dead.break.{n}")))
2701 .collect();
2702 for t in &tickets {
2703 assert!(create_exclusive(t, "").expect("ticket"));
2704 age_file(t);
2705 }
2706 assert!(TurnLock::take(&lease).expect("take").is_none());
2708 assert!(tickets.iter().all(|t| t.exists()));
2709 for t in &tickets {
2711 let f = std::fs::OpenOptions::new()
2712 .write(true)
2713 .open(t)
2714 .expect("open");
2715 f.set_modified(std::time::SystemTime::now() - TICKET_SWEEP_AGE * 2)
2716 .expect("age");
2717 }
2718 assert!(TurnLock::take(&lease).expect("take").is_none());
2719 assert!(TurnLock::take(&lease).expect("take").is_some());
2720 }
2721
2722 #[test]
2723 fn an_aged_empty_lock_is_broken() {
2724 let (_tmp, a, _b) = lease_store();
2725 let lease = a.turn_path("t1");
2726 let lock = lease.with_extension("turn.lock");
2727 std::fs::create_dir_all(lock.parent().expect("dir")).expect("dir");
2728 std::fs::write(&lock, "").expect("empty lock");
2729 age_file(&lock);
2730 assert!(TurnLock::take(&lease).expect("take").is_some());
2731 }
2732
2733 #[test]
2734 fn queued_text_is_durable_combined_and_drained_as_one_operator_turn() {
2735 let (tmp, talks) = store();
2736 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2737 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2738
2739 queue(&mut talk, &talks, "first", Vec::new()).expect("queue first");
2740 queue(&mut talk, &talks, "second", Vec::new()).expect("queue second");
2741 let saved = talks.get(&talk.id).expect("reload queued talk");
2742 assert_eq!(saved.pending, "first\n\nsecond");
2743 assert!(saved.turns.is_empty(), "a draft is not a transcript turn");
2744
2745 let drained = drain(&mut talk, &talks).expect("drain");
2746 assert_eq!(drained.as_deref(), Some("first\n\nsecond"));
2747 let saved = talks.get(&talk.id).expect("reload drained talk");
2748 assert!(saved.pending.is_empty());
2749 assert_eq!(saved.turns.len(), 1);
2750 assert_eq!(saved.turns[0].body, "first\n\nsecond");
2751 }
2752
2753 #[test]
2754 fn editing_a_queued_draft_preserves_its_attachments_and_rejects_a_stale_snapshot() {
2755 let (tmp, talks) = store();
2756 let cfg = config(mock_agent(tmp.path(), REPLY, env("reply")));
2757 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2758 let attachment = Attachment {
2759 id: "a".repeat(32),
2760 name: "shot.png".to_owned(),
2761 mime: "image/png".to_owned(),
2762 bytes: 3,
2763 };
2764
2765 queue(&mut talk, &talks, "first", vec![attachment.clone()]).expect("queue");
2766 assert!(
2767 edit_pending_text(
2768 &mut talk,
2769 &talks,
2770 "corrected",
2771 "first",
2772 std::slice::from_ref(&attachment.id),
2773 )
2774 .expect("edit")
2775 );
2776 let saved = talks.get(&talk.id).expect("reload edited draft");
2777 assert_eq!(saved.pending, "corrected");
2778 assert_eq!(saved.pending_attachments, vec![attachment]);
2779
2780 queue(&mut talk, &talks, "later", Vec::new()).expect("queue concurrent draft");
2781 assert!(
2782 !edit_pending_text(
2783 &mut talk,
2784 &talks,
2785 "stale edit",
2786 "corrected",
2787 &["a".repeat(32)],
2788 )
2789 .expect("stale edit is a conflict")
2790 );
2791 assert_eq!(
2792 talks.get(&talk.id).expect("reload after conflict").pending,
2793 "corrected\n\nlater"
2794 );
2795 assert!(
2796 !clear_pending_if_matches(&mut talk, &talks, "corrected", &["a".repeat(32)])
2797 .expect("stale clear is a conflict")
2798 );
2799 assert_eq!(
2800 talks
2801 .get(&talk.id)
2802 .expect("reload after stale clear")
2803 .pending,
2804 "corrected\n\nlater"
2805 );
2806 }
2807
2808 #[tokio::test]
2809 async fn a_reply_save_preserves_pending_accepted_while_the_cli_runs() {
2810 let (tmp, talks) = store();
2811 let slow = "#!/bin/sh\ncat >/dev/null\nsleep 0.1\nprintf reply\n";
2812 let cfg = config(mock_agent(tmp.path(), slow, BTreeMap::new()));
2813 let mut running = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2814 let id = running.id.clone();
2815 let first = record(&mut running, &talks, "first", Vec::new()).expect("record");
2816
2817 let response_talks = talks.clone();
2818 let response_cfg = cfg.clone();
2819 let reply = tokio::spawn(async move {
2820 respond(&mut running, &response_talks, &response_cfg, &first).await
2821 });
2822 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
2823
2824 let mut queued = talks.get(&id).expect("queued handle");
2825 queue(&mut queued, &talks, "next", Vec::new()).expect("queue");
2826 reply.await.expect("join").expect("reply");
2827
2828 let saved = talks.get(&id).expect("reload");
2829 assert_eq!(saved.pending, "next");
2830 assert_eq!(saved.turns.len(), 2, "operator message and reply remain");
2831 }
2832
2833 fn counting_agent(dir: &Path, id: &str, body: &str) -> AgentSpec {
2836 let calls = dir.join(format!("{id}.calls"));
2837 let script = format!(
2838 "#!/bin/sh\necho x >> '{}'\n{body}\n",
2839 calls.to_string_lossy()
2840 );
2841 let path = dir.join(format!("mock-{id}.sh"));
2842 std::fs::write(&path, script).expect("write mock");
2843 AgentSpec {
2844 id: id.to_owned(),
2845 kind: AgentKind::Command,
2846 model: None,
2847 command: vec!["sh".to_owned(), path.to_string_lossy().into_owned()],
2848 extra_args: Vec::new(),
2849 env: BTreeMap::new(),
2850 prompt_delivery: None,
2851 }
2852 }
2853
2854 fn calls(dir: &Path, id: &str) -> usize {
2855 std::fs::read_to_string(dir.join(format!("{id}.calls"))).map_or(0, |s| s.lines().count())
2856 }
2857
2858 fn chain_config(specs: Vec<AgentSpec>, ids: &[&str]) -> Config {
2859 let mut cfg = config(specs[0].clone());
2860 cfg.agents = specs;
2861 cfg.roles.chatter = Some(AgentChoice::Chain(
2862 ids.iter().map(|s| (*s).to_owned()).collect(),
2863 ));
2864 cfg
2865 }
2866
2867 #[tokio::test]
2868 async fn a_chatter_chain_falls_back_resends_the_transcript_and_sticks() {
2869 let (tmp, talks) = store();
2870 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2871 let b = counting_agent(tmp.path(), "b", "cat");
2872 let cfg = chain_config(vec![a, b], &["a", "b"]);
2873 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2874 assert_eq!(talk.agent, "a");
2875
2876 say(&mut talk, &talks, &cfg, "hello there", Vec::new())
2877 .await
2878 .expect("turn");
2879 assert_eq!(calls(tmp.path(), "a"), 1, "each id is tried once");
2880 assert_eq!(calls(tmp.path(), "b"), 1);
2881 assert_eq!(talk.agent, "b", "the switch persists");
2882 assert!(talks.get(&talk.id).unwrap().agent == "b");
2883 let reply = talk.turns.last().unwrap();
2884 assert!(reply.body.contains("hello there"));
2885 assert!(
2886 reply.body.contains("magi task add --solo"),
2887 "a fresh seat gets the full briefing"
2888 );
2889 assert!(
2890 talk.turns
2891 .iter()
2892 .any(|t| t.body.contains("agent changed from a to b")),
2893 "the switch is noted"
2894 );
2895 }
2896
2897 #[tokio::test]
2898 async fn an_exhausted_chatter_chain_fails_like_a_single_seat_and_stays_put() {
2899 let (tmp, talks) = store();
2900 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2901 let b = counting_agent(tmp.path(), "b", "cat >/dev/null\nexit 4");
2902 let cfg = chain_config(vec![a, b], &["a", "b", "a"]);
2903 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2904
2905 let err = say(&mut talk, &talks, &cfg, "hi", Vec::new())
2906 .await
2907 .expect_err("every agent failed");
2908 assert!(err.to_string().contains("`a`"), "{err:#}");
2909 assert_eq!(calls(tmp.path(), "a"), 1);
2910 assert_eq!(calls(tmp.path(), "b"), 1);
2911 assert_eq!(talk.agent, "a", "an exhausted chain leaves the agent alone");
2912 }
2913
2914 #[test]
2915 fn a_chatter_chain_skips_an_unknown_id_at_begin() {
2916 let (tmp, talks) = store();
2917 let b = counting_agent(tmp.path(), "b", "cat");
2918 let cfg = chain_config(vec![b], &["ghost", "b"]);
2919 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
2920 assert_eq!(talk.agent, "b");
2921 }
2922
2923 #[tokio::test]
2924 async fn an_explicit_agent_inside_the_chatter_chain_stays_pinned() {
2925 let (tmp, talks) = store();
2926 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2927 let b = counting_agent(tmp.path(), "b", "cat");
2928 let cfg = chain_config(vec![a, b], &["a", "b"]);
2929 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("a")).expect("begin");
2930 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2931 .await
2932 .expect_err("a alone, and it fails");
2933 assert_eq!(calls(tmp.path(), "b"), 0);
2934 assert_eq!(talk.agent, "a");
2935 }
2936
2937 #[tokio::test]
2938 async fn an_explicit_agent_does_not_borrow_the_chatter_chain() {
2939 let (tmp, talks) = store();
2940 let a = counting_agent(tmp.path(), "a", "cat >/dev/null\nexit 3");
2941 let b = counting_agent(tmp.path(), "b", "cat");
2942 let c = counting_agent(tmp.path(), "c", "cat >/dev/null\nexit 3");
2943 let cfg = chain_config(vec![a, b, c], &["a", "b"]);
2944 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some("c")).expect("begin");
2945 say(&mut talk, &talks, &cfg, "hi", Vec::new())
2946 .await
2947 .expect_err("c alone, and it fails");
2948 assert_eq!(calls(tmp.path(), "b"), 0);
2949 }
2950
2951 #[test]
2952 fn briefing_names_the_operator_with_or_without_a_persona() {
2953 let plain = briefing(Path::new("/repo"), "en", false);
2954 let named = briefing_with(Path::new("/repo"), "en", false, None, Some("Commander"));
2955 assert!(named.starts_with(&plain));
2956 assert!(named.contains("# Addressing the operator"));
2957 assert!(named.contains("\"Commander\""));
2958 let rei = crate::persona::builtin_catalog()
2959 .into_iter()
2960 .find(|p| p.id == "rei")
2961 .expect("rei");
2962 let with = briefing_with(
2963 Path::new("/repo"),
2964 "en",
2965 false,
2966 Some(&rei),
2967 Some("Commander"),
2968 );
2969 assert!(with.contains("\"Commander\""));
2970 assert!(!with.contains("# Addressing the operator"));
2971 }
2972
2973 #[test]
2974 fn briefing_carries_a_persona_section_only_when_one_is_chosen() {
2975 let plain = briefing(Path::new("/repo"), "en", false);
2976 assert_eq!(
2977 plain,
2978 briefing_with(Path::new("/repo"), "en", false, None, None)
2979 );
2980 assert!(!plain.contains("Persona"));
2981 let rei = crate::persona::builtin_catalog()
2982 .into_iter()
2983 .find(|p| p.id == "rei")
2984 .expect("rei");
2985 let with = briefing_with(Path::new("/repo"), "en", false, Some(&rei), None);
2986 assert!(with.starts_with(&plain), "the plain briefing is untouched");
2987 assert!(with.contains("# Persona (tone only)"));
2988 assert!(with.contains("TONE ONLY"));
2989 assert!(with.contains("task ids"));
2990 assert!(with.contains("write policy"));
2991 assert!(with.contains("`magi task add`"));
2992 assert!(with.contains("Rei Ayanami"));
2993 }
2994
2995 #[test]
2996 fn a_talk_written_before_personas_still_loads() {
2997 let (tmp, talks) = store();
2998 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
2999 let cfg = config(spec);
3000 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3001 let path = talks.root.join(format!("{}.json", talk.id));
3002 let mut v: serde_json::Value =
3003 serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
3004 v.as_object_mut().unwrap().remove("persona");
3005 v.as_object_mut().unwrap().remove("persona_dirty");
3006 std::fs::write(&path, v.to_string()).unwrap();
3007 let loaded = talks.get(&talk.id).expect("old record loads");
3008 assert_eq!(loaded.persona, "");
3009 assert!(!loaded.persona_dirty);
3010 }
3011
3012 #[tokio::test]
3013 async fn a_persona_switch_notes_marks_dirty_and_updates_the_next_turn_once() {
3014 let (tmp, talks) = store();
3015 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3016 let cfg = config(spec);
3017 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3018 say(&mut talk, &talks, &cfg, "hello", Vec::new())
3019 .await
3020 .expect("first turn");
3021 assert!(!talk.turns[1].body.contains("Persona"));
3022
3023 assert!(switch_persona(&mut talk, &talks, "misato").expect("switch"));
3024 assert!(!switch_persona(&mut talk, &talks, "misato").expect("same"));
3025 assert!(talk.persona_dirty);
3026 assert!(talk.turns.last().unwrap().body.contains("persona changed"));
3027
3028 say(&mut talk, &talks, &cfg, "next", Vec::new())
3029 .await
3030 .expect("turn");
3031 let prompt = &talk.turns.last().unwrap().body;
3032 assert!(prompt.contains("# Persona update"), "{prompt}");
3033 assert!(prompt.contains("Misato Katsuragi"));
3034 assert!(!talk.persona_dirty, "cleared after a successful turn");
3035
3036 say(&mut talk, &talks, &cfg, "again", Vec::new())
3037 .await
3038 .expect("turn");
3039 assert!(!talk.turns.last().unwrap().body.contains("# Persona update"));
3040
3041 assert!(switch_persona(&mut talk, &talks, "default").expect("back"));
3042 assert_eq!(talk.persona, "");
3043 say(&mut talk, &talks, &cfg, "plain", Vec::new())
3044 .await
3045 .expect("turn");
3046 assert!(
3047 talk.turns
3048 .last()
3049 .unwrap()
3050 .body
3051 .contains("turned the persona off")
3052 );
3053
3054 assert!(switch_persona(&mut talk, &talks, "rei").expect("rei"));
3056 let mut no_sessions = cfg.clone();
3057 no_sessions.graph.sessions = false;
3058 say(&mut talk, &talks, &no_sessions, "one", Vec::new())
3059 .await
3060 .expect("turn");
3061 say(&mut talk, &talks, &no_sessions, "two", Vec::new())
3062 .await
3063 .expect("turn");
3064 let last = &talk.turns.last().unwrap().body;
3065 assert!(!talk.persona_dirty);
3066 assert!(last.contains("# Persona (tone only)"), "{last}");
3067 assert!(last.contains("Rei Ayanami"));
3068 }
3069
3070 #[tokio::test]
3071 async fn the_first_turn_carries_the_briefing_and_later_turns_do_not() {
3072 let (tmp, talks) = store();
3073 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3074 let cfg = config(spec);
3075 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3076
3077 say(
3078 &mut talk,
3079 &talks,
3080 &cfg,
3081 "what does the queue module do?",
3082 Vec::new(),
3083 )
3084 .await
3085 .expect("first turn");
3086 let first_prompt = &talk.turns[1].body;
3087 assert!(first_prompt.contains("magi task add --solo"));
3088 assert!(first_prompt.contains("what does the queue module do?"));
3089
3090 say(&mut talk, &talks, &cfg, "and how is it locked?", Vec::new())
3091 .await
3092 .expect("second turn");
3093 let second_prompt = &talk.turns[3].body;
3094 assert!(
3095 !second_prompt.contains("magi task add --solo"),
3096 "the briefing is sent once, not on every turn: {second_prompt}"
3097 );
3098 assert!(second_prompt.contains("and how is it locked?"));
3099 }
3100
3101 #[tokio::test]
3102 async fn switching_agent_resets_the_seat_notes_it_and_resends_the_transcript() {
3103 let (tmp, talks) = store();
3104 let a = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3105 let mut b = a.clone();
3106 b.id = "other".to_owned();
3107 let mut cfg = config(a.clone());
3108 cfg.agents.push(b.clone());
3109 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), Some(&a.id)).expect("begin");
3110 say(&mut talk, &talks, &cfg, "remember the walrus", Vec::new())
3111 .await
3112 .expect("first turn");
3113 let old_session = talk.seat.claude_session.clone();
3114 assert_eq!(talk.seat.turns, 1);
3115
3116 assert!(switch_agent(&mut talk, &talks, &b).expect("switch"));
3117 assert_eq!(talk.agent, "other");
3118 assert_eq!(talk.seat.turns, 0);
3119 assert_eq!(talk.seat.agent, "other");
3120 assert_ne!(talk.seat.claude_session, old_session);
3121 let note = talk.turns.last().expect("note");
3122 assert_eq!(note.who, Who::Agent);
3123 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3124 assert!(note.body.contains("changed from"), "{}", note.body);
3125 assert_eq!(talks.get(&talk.id).expect("reload").agent, "other");
3126
3127 let before = talk.turns.len();
3128 assert!(!switch_agent(&mut talk, &talks, &b).expect("same agent"));
3129 assert_eq!(talk.turns.len(), before, "a no-op writes no note");
3130
3131 say(&mut talk, &talks, &cfg, "what did I say?", Vec::new())
3132 .await
3133 .expect("turn after switch");
3134 let prompt = &talk.turns.last().expect("reply").body;
3135 assert!(prompt.contains("remember the walrus"), "{prompt}");
3136 assert!(prompt.contains("## magi"), "{prompt}");
3137 assert!(prompt.contains("what did I say?"), "{prompt}");
3138 }
3139
3140 #[tokio::test]
3141 async fn say_appends_the_operator_turn_then_the_agent_turn() {
3142 let (tmp, talks) = store();
3143 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3144 let cfg = config(spec);
3145 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3146
3147 say(
3148 &mut talk,
3149 &talks,
3150 &cfg,
3151 "can I rename this function?",
3152 Vec::new(),
3153 )
3154 .await
3155 .expect("say");
3156
3157 assert_eq!(talk.turns.len(), 2);
3158 assert_eq!(talk.turns[0].who, Who::Operator);
3159 assert_eq!(talk.turns[0].body, "can I rename this function?");
3160 assert_eq!(talk.turns[1].who, Who::Agent);
3161 assert_eq!(talk.turns[1].body, "go ahead");
3162 assert_eq!(talks.get(&talk.id).expect("get").turns, talk.turns);
3163 }
3164
3165 #[tokio::test]
3166 async fn a_failed_turn_keeps_the_operator_message_and_says_what_happened() {
3167 let (tmp, talks) = store();
3168 let spec = mock_agent(tmp.path(), BROKEN, BTreeMap::new());
3169 let cfg = config(spec);
3170 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3171
3172 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
3173 .await
3174 .expect_err("a turn with no answer is an error");
3175 assert!(err.to_string().contains("no answer"), "{err}");
3176
3177 let on_disk = talks.get(&talk.id).expect("get");
3178 assert_eq!(on_disk.turns.len(), 2);
3179 assert_eq!(on_disk.turns[0].body, "check the tests");
3180 let note = &on_disk.turns[1];
3181 assert_eq!(note.who, Who::Agent);
3182 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3183 assert!(note.body.contains("your message is saved"));
3184 }
3185
3186 #[tokio::test]
3192 async fn a_passing_write_failure_while_saving_the_reply_does_not_lose_it() {
3193 let (tmp, talks) = store();
3194 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3195 let cfg = config(spec);
3196 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3197
3198 let text =
3199 record(&mut talk, &talks, "can I rename this function?", Vec::new()).expect("record");
3200 failpoint::force_put_failures(PUT_RETRIES - 1);
3203 respond(&mut talk, &talks, &cfg, &text)
3204 .await
3205 .expect("respond must survive a write failure its own retries can outlast");
3206
3207 assert_eq!(talk.turns.len(), 2);
3208 assert_eq!(talk.turns[1].who, Who::Agent);
3209 assert_eq!(talk.turns[1].body, "go ahead");
3210 let on_disk = talks.get(&talk.id).expect("get");
3211 assert_eq!(
3212 on_disk.turns, talk.turns,
3213 "the reply must reach disk despite the early write failures"
3214 );
3215 }
3216
3217 #[tokio::test]
3223 async fn a_persistent_write_failure_while_saving_the_reply_is_never_silent() {
3224 let (tmp, talks) = store();
3225 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3226 let cfg = config(spec);
3227 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3228
3229 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3230 failpoint::force_put_failures(PUT_RETRIES);
3235 let err = respond(&mut talk, &talks, &cfg, &text)
3236 .await
3237 .expect_err("a reply that cannot be saved must be reported, not swallowed");
3238 assert!(err.to_string().contains("could not be saved"), "{err}");
3239
3240 let on_disk = talks.get(&talk.id).expect("get");
3241 assert_eq!(
3242 on_disk.turns.len(),
3243 2,
3244 "the operator turn plus a visible note"
3245 );
3246 assert_eq!(on_disk.turns[0].body, "check the tests");
3247 let note = &on_disk.turns[1];
3248 assert_eq!(note.who, Who::Agent);
3249 assert!(note.body.starts_with(MAGI_NOTE), "{}", note.body);
3250 assert!(
3251 note.body.contains("could not be saved"),
3252 "the operator must be told the reply is missing, not left staring \
3253 at a gap with no explanation: {}",
3254 note.body
3255 );
3256 assert_eq!(
3257 talk.turns, on_disk.turns,
3258 "the in-memory talk must match what actually landed on disk"
3259 );
3260
3261 let artifacts = talks.artifacts_of(&talk.id);
3264 let stash = std::fs::read_dir(&artifacts)
3265 .expect("artifacts dir")
3266 .filter_map(|e| e.ok())
3267 .find(|e| e.file_name().to_string_lossy().ends_with("-lost.txt"))
3268 .expect("a stash file for the lost reply");
3269 let stashed = std::fs::read_to_string(stash.path()).expect("read stash");
3270 assert_eq!(stashed, "go ahead");
3271
3272 assert_eq!(
3280 on_disk.seat.turns, 1,
3281 "the note's write must carry the turn the CLI actually took"
3282 );
3283 assert_eq!(
3284 on_disk.seat.claude_session, talk.seat.claude_session,
3285 "the session id handed to the CLI must survive the failed reply"
3286 );
3287 assert_eq!(on_disk.seat.captured_session, talk.seat.captured_session);
3288 assert!(
3289 agent::has_session(AgentKind::Command, &on_disk.seat, cfg.graph.sessions),
3290 "the next turn must resume, not open the same session id twice"
3291 );
3292 }
3293
3294 #[tokio::test]
3299 async fn a_write_failure_that_also_loses_the_note_still_reports_it() {
3300 let (tmp, talks) = store();
3301 let spec = mock_agent(tmp.path(), REPLY, env("go ahead"));
3302 let cfg = config(spec);
3303 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3304
3305 let text = record(&mut talk, &talks, "check the tests", Vec::new()).expect("record");
3306 failpoint::force_put_failures(PUT_RETRIES * 2);
3309 let err = respond(&mut talk, &talks, &cfg, &text)
3310 .await
3311 .expect_err("neither the reply nor the note could be saved");
3312 assert!(err.to_string().contains("could not be saved"), "{err}");
3313
3314 assert_eq!(talk.turns.len(), 1, "only the operator's own turn");
3315 let on_disk = talks.get(&talk.id).expect("get");
3316 assert_eq!(on_disk.turns.len(), 1);
3317
3318 assert_eq!(
3328 on_disk.seat.turns, 0,
3329 "an unwritable file cannot record the turn the CLI took"
3330 );
3331 assert_eq!(
3332 talk.seat.turns, 1,
3333 "the in-memory seat still reports the turn the CLI actually took"
3334 );
3335 assert_eq!(
3336 on_disk.seat.claude_session, talk.seat.claude_session,
3337 "the session id was minted at `begin` and never changes here"
3338 );
3339 }
3340
3341 #[tokio::test]
3345 async fn attachments_reach_the_prompt_and_an_empty_body_is_still_a_turn() {
3346 let (tmp, talks) = store();
3347 let spec = mock_agent(tmp.path(), ECHO, BTreeMap::new());
3348 let cfg = config(spec);
3349 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3350
3351 let att = talks
3352 .put_attachment(
3353 &talk.id,
3354 "image/png",
3355 "screenshot.png",
3356 b"pretend-png-bytes",
3357 )
3358 .expect("put attachment");
3359
3360 say(&mut talk, &talks, &cfg, "", vec![att.clone()])
3361 .await
3362 .expect("an empty body with an attachment is still a turn");
3363
3364 let operator_turn = &talk.turns[0];
3365 assert_eq!(operator_turn.who, Who::Operator);
3366 assert_eq!(operator_turn.body, "");
3367 assert_eq!(operator_turn.attachments, vec![att.clone()]);
3368
3369 let prompt = &talk.turns[1].body;
3370 let expected_path = talks
3371 .attachments_dir(&talk.id)
3372 .join(format!("{}.png", att.id));
3373 assert!(
3374 prompt.contains(&expected_path.display().to_string()),
3375 "the agent must be told the attachment's absolute path: {prompt}"
3376 );
3377 assert!(prompt.contains("image/png"), "and its mime: {prompt}");
3378 }
3379
3380 #[test]
3389 fn attachment_path_is_absolute_even_when_the_store_root_is_relative() {
3390 let talks = Talks::at(PathBuf::from("relative-talks-root-for-this-test"));
3391 let att = Attachment {
3392 id: "0".repeat(32),
3393 name: "shot.png".to_owned(),
3394 mime: "image/png".to_owned(),
3395 bytes: 3,
3396 };
3397 let path = talks
3398 .attachment_path("some-talk-id", &att)
3399 .expect("a supported mime always yields a path");
3400 assert!(
3401 path.is_absolute(),
3402 "must be absolute even off a relative store root: {}",
3403 path.display()
3404 );
3405 }
3406
3407 #[tokio::test]
3408 async fn a_turn_past_the_configured_talk_timeout_is_reported_with_that_timeout() {
3409 let (tmp, talks) = store();
3414 let slow = mock_agent(
3415 tmp.path(),
3416 "#!/bin/sh\ncat >/dev/null\nsleep 2\n",
3417 BTreeMap::new(),
3418 );
3419 let mut cfg = config(slow);
3420 cfg.graph.timeout_talk = 1;
3421 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3422
3423 let err = say(&mut talk, &talks, &cfg, "check the tests", Vec::new())
3424 .await
3425 .expect_err("a turn that never answers is an error");
3426 assert!(
3427 err.to_string().contains("did not answer within 1s"),
3428 "{err}"
3429 );
3430
3431 let on_disk = talks.get(&talk.id).expect("get");
3432 let note = on_disk.turns.last().expect("a note turn was recorded");
3433 assert!(
3434 note.body.contains("did not answer within 1s"),
3435 "the transcript must show the configured timeout: {}",
3436 note.body
3437 );
3438 }
3439
3440 #[test]
3441 fn closing_is_idempotent_and_a_closed_talk_takes_no_more_turns() {
3442 let (tmp, talks) = store();
3443 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3444 let cfg = config(spec);
3445 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3446
3447 close(&mut talk, &talks).expect("close");
3448 assert_eq!(talk.status, TalkStatus::Closed);
3449 close(&mut talk, &talks).expect("closing twice is not an error");
3450
3451 let err =
3452 record(&mut talk, &talks, "still there?", Vec::new()).expect_err("closed talks refuse");
3453 assert!(err.to_string().contains("closed"));
3454 let _ = &cfg; }
3456
3457 #[tokio::test]
3458 async fn a_close_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3459 let (tmp, talks) = store();
3460 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3461 let cfg = config(spec);
3462 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3465
3466 let mut closed_elsewhere = talks.get(&in_flight.id).expect("reread");
3470 close(&mut closed_elsewhere, &talks).expect("close");
3471 assert_eq!(
3472 talks.get(&in_flight.id).expect("reread").status,
3473 TalkStatus::Closed,
3474 "the close landed on disk before the turn finished"
3475 );
3476
3477 assert_eq!(in_flight.status, TalkStatus::Open);
3481 respond(&mut in_flight, &talks, &cfg, "one more question")
3482 .await
3483 .expect("the turn itself still completes");
3484
3485 let on_disk = talks.get(&in_flight.id).expect("reread");
3486 assert_eq!(
3487 on_disk.status,
3488 TalkStatus::Closed,
3489 "a close must stick even when a turn that started before it finishes after it"
3490 );
3491 assert!(
3494 on_disk.turns.iter().any(|t| t.body == "here you go"),
3495 "the in-flight turn's own reply is still recorded: {:?}",
3496 on_disk.turns
3497 );
3498 }
3499
3500 #[test]
3501 fn a_close_that_lands_before_record_is_called_is_not_undone_by_it() {
3502 let (tmp, talks) = store();
3503 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3504 let cfg = config(spec);
3505 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3508
3509 let mut closed_elsewhere = talks.get(&stale.id).expect("reread");
3512 close(&mut closed_elsewhere, &talks).expect("close");
3513 assert_eq!(
3514 talks.get(&stale.id).expect("reread").status,
3515 TalkStatus::Closed,
3516 "the close landed on disk before record was called"
3517 );
3518
3519 assert_eq!(stale.status, TalkStatus::Open);
3523 let err = record(&mut stale, &talks, "still there?", Vec::new())
3524 .expect_err("a close that landed first must be honored, not overwritten");
3525 assert!(err.to_string().contains("closed"));
3526
3527 let on_disk = talks.get(&stale.id).expect("reread");
3528 assert_eq!(
3529 on_disk.status,
3530 TalkStatus::Closed,
3531 "record must not resurrect a conversation closed while its snapshot was stale"
3532 );
3533 assert!(
3534 on_disk.turns.is_empty(),
3535 "the rejected turn must not have been appended: {:?}",
3536 on_disk.turns
3537 );
3538 let _ = &cfg; }
3540
3541 #[test]
3542 fn close_blocks_on_records_guard_rather_than_interleaving_with_it() {
3543 let (tmp, talks) = store();
3544 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3545 let cfg = config(spec);
3546 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3547
3548 let held = talks.guard();
3552
3553 let talks2 = talks.clone();
3554 let id = talk.id.clone();
3555 let closing = std::thread::spawn(move || {
3556 let mut talk = talks2.get(&id).expect("get");
3557 close(&mut talk, &talks2).expect("close");
3558 });
3559
3560 std::thread::sleep(Duration::from_millis(50));
3561 assert!(
3562 !closing.is_finished(),
3563 "close must wait for the guard, not read and write while it is held - \
3564 a re-read alone narrows this window without closing it"
3565 );
3566
3567 drop(held);
3568 closing.join().expect("close thread panicked");
3569
3570 assert_eq!(
3571 talks.get(&talk.id).expect("reread").status,
3572 TalkStatus::Closed,
3573 "once the guard is free, close still lands"
3574 );
3575 let _ = &cfg; }
3577
3578 #[test]
3579 fn reopening_a_closed_talk_lets_it_take_turns_again_and_reopening_twice_is_not_an_error() {
3580 let (tmp, talks) = store();
3581 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3582 let cfg = config(spec);
3583 let mut talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3584
3585 close(&mut talk, &talks).expect("close");
3586 assert_eq!(talk.status, TalkStatus::Closed);
3587
3588 reopen(&mut talk, &talks).expect("reopen");
3589 assert_eq!(talk.status, TalkStatus::Open);
3590 assert_eq!(
3591 talks.get(&talk.id).expect("reread").status,
3592 TalkStatus::Open
3593 );
3594
3595 reopen(&mut talk, &talks).expect("reopening an open talk is not an error");
3597 assert_eq!(talk.status, TalkStatus::Open);
3598
3599 record(&mut talk, &talks, "one more thing", Vec::new())
3600 .expect("a reopened talk takes turns again");
3601 let _ = &cfg; }
3603
3604 #[test]
3605 fn removing_a_talk_deletes_its_record_and_artifacts_and_refuses_an_unknown_id() {
3606 let (tmp, talks) = store();
3607 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3608 let cfg = config(spec);
3609 let talk = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3610
3611 let artifacts = talks.artifacts_of(&talk.id);
3612 std::fs::create_dir_all(&artifacts).expect("create artifacts dir");
3613 std::fs::write(artifacts.join("turn-1.txt"), "hello").expect("write artifact");
3614
3615 talks.remove(&talk.id).expect("remove");
3616 assert!(!talks.path_of(&talk.id).is_file(), "the record is gone");
3617 assert!(!artifacts.is_dir(), "the artifacts directory is gone");
3618 assert!(
3619 talks.get(&talk.id).is_err(),
3620 "a removed talk cannot be read back"
3621 );
3622
3623 let err = talks
3624 .remove("nonexistent-id")
3625 .expect_err("unknown id refused");
3626 assert!(err.to_string().contains("no talk matches"), "{err}");
3627 let _ = &cfg; }
3629
3630 #[tokio::test]
3631 async fn a_delete_that_lands_while_a_turn_is_in_flight_is_not_undone_by_the_reply() {
3632 let (tmp, talks) = store();
3633 let spec = mock_agent(tmp.path(), REPLY, env("here you go"));
3634 let cfg = config(spec);
3635 let mut in_flight = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3638
3639 talks.remove(&in_flight.id).expect("remove");
3640 assert!(
3641 talks.get(&in_flight.id).is_err(),
3642 "the delete landed on disk before the turn finished"
3643 );
3644
3645 respond(&mut in_flight, &talks, &cfg, "one more question")
3648 .await
3649 .expect("the turn itself still completes rather than erroring");
3650
3651 assert!(
3652 talks.get(&in_flight.id).is_err(),
3653 "a delete must stick even when a turn that started before it finishes after it"
3654 );
3655 }
3656
3657 #[test]
3658 fn a_delete_that_lands_before_record_is_called_is_not_undone_by_it() {
3659 let (tmp, talks) = store();
3660 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3661 let cfg = config(spec);
3662 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3665
3666 talks.remove(&stale.id).expect("remove");
3667
3668 let err = record(&mut stale, &talks, "still there?", Vec::new())
3672 .expect_err("a delete that landed first must be honored, not overwritten");
3673 assert!(err.to_string().contains("deleted"), "{err}");
3674
3675 assert!(
3676 talks.get(&stale.id).is_err(),
3677 "record must not resurrect a conversation deleted while its snapshot was stale"
3678 );
3679 let _ = &cfg; }
3681
3682 #[test]
3683 fn a_delete_that_lands_before_close_is_called_is_not_undone_by_it() {
3684 let (tmp, talks) = store();
3685 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3686 let cfg = config(spec);
3687 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3690
3691 talks.remove(&stale.id).expect("remove");
3692
3693 let err = close(&mut stale, &talks)
3697 .expect_err("a delete that landed first must be honored, not overwritten");
3698 assert!(err.to_string().contains("deleted"), "{err}");
3699
3700 assert!(
3701 talks.get(&stale.id).is_err(),
3702 "close must not resurrect a conversation deleted while its snapshot was stale"
3703 );
3704 let _ = &cfg; }
3706
3707 #[test]
3708 fn a_delete_that_lands_before_reopen_is_called_is_not_undone_by_it() {
3709 let (tmp, talks) = store();
3710 let spec = mock_agent(tmp.path(), REPLY, env("hi"));
3711 let cfg = config(spec);
3712 let mut stale = begin(&talks, &cfg, tmp.path().to_owned(), None).expect("begin");
3715 close(&mut stale, &talks).expect("close");
3716
3717 talks.remove(&stale.id).expect("remove");
3718
3719 let err = reopen(&mut stale, &talks)
3723 .expect_err("a delete that landed first must be honored, not overwritten");
3724 assert!(err.to_string().contains("deleted"), "{err}");
3725
3726 assert!(
3727 talks.get(&stale.id).is_err(),
3728 "reopen must not resurrect a conversation deleted while its snapshot was stale"
3729 );
3730 let _ = &cfg; }
3732
3733 #[test]
3734 fn list_puts_open_talks_before_closed_ones() {
3735 let (tmp, talks) = store();
3736 let make = |id: &str, status: TalkStatus| {
3737 let mut t = Talk {
3738 schema: SCHEMA,
3739 id: id.to_owned(),
3740 repo: tmp.path().to_owned(),
3741 agent: "mock".to_owned(),
3742 status,
3743 turns: Vec::new(),
3744 pending: String::new(),
3745 pending_attachments: Vec::new(),
3746 fallback: false,
3747 persona: String::new(),
3748 persona_dirty: false,
3749 created_at: Timestamp::now(),
3750 updated_at: Timestamp::now(),
3751 seat: SeatState::new(SEAT, "mock", 7),
3752 };
3753 talks.put(&mut t).expect("put");
3754 };
3755 make("20260901-000000-0001", TalkStatus::Open);
3756 make("20260902-000000-0002", TalkStatus::Open);
3757 make("20260903-000000-0003", TalkStatus::Closed);
3758
3759 let ids: Vec<String> = talks.list().into_iter().map(|t| t.id).collect();
3760 assert_eq!(
3761 ids,
3762 [
3763 "20260902-000000-0002",
3764 "20260901-000000-0001",
3765 "20260903-000000-0003"
3766 ]
3767 );
3768 assert_eq!(talks.count_open(), 2);
3769 }
3770
3771 #[test]
3772 fn tasks_of_finds_only_this_talks_own_tasks() {
3773 let dir = tempfile::tempdir().expect("tempdir");
3774 let queue = Queue::at(dir.path().join("queue"));
3775
3776 let mut mine = Task::new(
3777 "rework the loader".to_owned(),
3778 "rework the loader".to_owned(),
3779 PathBuf::from("/repo"),
3780 Source::Agent {
3781 run: "20260904-014455-ab12".to_owned(),
3782 node: "chat".to_owned(),
3783 },
3784 );
3785 queue.put(&mut mine).expect("put mine");
3786
3787 let mut theirs = Task::new(
3788 "unrelated".to_owned(),
3789 "unrelated".to_owned(),
3790 PathBuf::from("/repo"),
3791 Source::Agent {
3792 run: "20260904-090000-zz99".to_owned(),
3793 node: "implement".to_owned(),
3794 },
3795 );
3796 queue.put(&mut theirs).expect("put theirs");
3797
3798 let mut human = Task::new(
3799 "typed by hand".to_owned(),
3800 "typed by hand".to_owned(),
3801 PathBuf::from("/repo"),
3802 Source::Human,
3803 );
3804 queue.put(&mut human).expect("put human");
3805
3806 let found = tasks_of(&queue, "20260904-014455-ab12");
3807 assert_eq!(found.len(), 1);
3808 assert_eq!(found[0].id, mine.id);
3809 }
3810
3811 #[test]
3812 fn the_briefing_names_solo_task_add() {
3813 let brief = briefing(Path::new("/repo"), "en", false);
3814 assert!(brief.contains("magi task add --solo"));
3815 assert!(brief.contains("/repo"));
3816 assert!(!brief.contains("Hold this conversation in"));
3817 }
3818
3819 #[test]
3825 fn the_briefing_explains_targeting_a_different_repository_by_name() {
3826 let brief = briefing(Path::new("/repo"), "en", false);
3827 assert!(brief.contains("--repo does not have to be a full path"));
3828 assert!(brief.contains("owner/repo"));
3829 assert!(brief.contains("magi repos"));
3830 assert!(brief.contains("ask the operator"));
3831 }
3832
3833 #[test]
3834 fn the_briefing_tells_the_assistant_to_pass_images_with_attach() {
3835 let brief = briefing(Path::new("/repo"), "en", false);
3836 assert!(brief.contains("--attach <path>"), "{brief}");
3837 assert!(brief.contains("deleting this conversation"), "{brief}");
3838 }
3839
3840 #[test]
3841 fn the_briefing_names_the_language_when_it_is_not_english() {
3842 let brief = briefing(Path::new("/repo"), "Japanese", false);
3843 assert!(brief.contains("Hold this conversation in Japanese"));
3844 }
3845
3846 #[test]
3847 fn the_briefing_forbids_writes_unless_the_repository_opted_in() {
3848 let read_only = briefing(Path::new("/repo"), "en", false);
3849 assert!(read_only.contains("Do not write files"));
3850 assert!(!read_only.contains("allow_write"));
3851
3852 let writable = briefing(Path::new("/repo"), "en", true);
3853 assert!(!writable.contains("Do not write files"));
3854 assert!(writable.contains("allow_write = true"));
3855 assert!(writable.contains("magi task add --solo"));
3858 assert!(writable.contains("say plainly what you"));
3859 }
3860}