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