1use std::collections::BTreeMap;
14use std::path::{Path, PathBuf};
15use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
16
17use serde::{Deserialize, Serialize};
18
19use crate::mailbox::{local_machine_name, mail_root, Envelope, MailAddress, MailKind, Mailbox};
20
21pub const AGENT_HARNESS: &str = "agent";
23
24pub const SESSION_THREAD_PREFIX: &str = "s-";
27
28#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
30pub struct Agent {
31 pub name: String,
33 pub main_session: MailAddress,
35 #[serde(default, skip_serializing_if = "Option::is_none")]
38 pub folder: Option<PathBuf>,
39 #[serde(default, skip_serializing_if = "Option::is_none")]
41 pub harness: Option<String>,
42 #[serde(default, skip_serializing_if = "Option::is_none")]
45 pub idle_minutes: Option<u64>,
46}
47
48impl Agent {
49 pub fn address(&self) -> std::io::Result<MailAddress> {
51 agent_address(&self.name)
52 }
53}
54
55pub fn agent_address(name: &str) -> std::io::Result<MailAddress> {
57 MailAddress::new(local_machine_name(), AGENT_HARNESS, name)
58 .map_err(|error| std::io::Error::other(error.to_string()))
59}
60
61fn agents_dir() -> PathBuf {
62 mail_root().join("agents")
63}
64
65fn threads_dir() -> PathBuf {
66 mail_root().join("threads")
67}
68
69fn valid_name(name: &str) -> bool {
70 !name.is_empty()
71 && name
72 .chars()
73 .all(|character| character.is_ascii_alphanumeric() || "-_.".contains(character))
74}
75
76fn write_atomically(path: &Path, bytes: &[u8]) -> std::io::Result<()> {
77 if let Some(parent) = path.parent() {
78 std::fs::create_dir_all(parent)?;
79 }
80 let temporary = path.with_extension(format!("tmp.{}", std::process::id()));
81 std::fs::write(&temporary, bytes)?;
82 std::fs::rename(temporary, path)
83}
84
85pub fn declare(agent: &Agent) -> std::io::Result<()> {
87 if !valid_name(&agent.name) {
88 return Err(std::io::Error::other(format!(
89 "`{}` is not an agent name: use letters, digits, `-`, `_` and `.`",
90 agent.name
91 )));
92 }
93 let bytes = serde_json::to_vec_pretty(agent).map_err(std::io::Error::other)?;
94 write_atomically(&agents_dir().join(format!("{}.json", agent.name)), &bytes)
95}
96
97pub fn load(name: &str) -> Option<Agent> {
99 if !valid_name(name) {
100 return None;
101 }
102 let bytes = std::fs::read(agents_dir().join(format!("{name}.json"))).ok()?;
103 serde_json::from_slice(&bytes).ok()
104}
105
106pub fn all_agents() -> Vec<Agent> {
108 let Ok(entries) = std::fs::read_dir(agents_dir()) else {
109 return Vec::new();
110 };
111 let mut agents: Vec<Agent> = entries
112 .flatten()
113 .filter(|entry| entry.path().extension().is_some_and(|ext| ext == "json"))
114 .filter_map(|entry| std::fs::read(entry.path()).ok())
115 .filter_map(|bytes| serde_json::from_slice(&bytes).ok())
116 .collect();
117 agents.sort_by(|a, b| a.name.cmp(&b.name));
118 agents
119}
120
121fn owners_account_manager_file() -> PathBuf {
122 agents_dir().join("owners-account-manager")
123}
124
125pub fn set_owners_account_manager(name: &str) -> std::io::Result<()> {
128 write_atomically(
129 &owners_account_manager_file(),
130 format!("{name}\n").as_bytes(),
131 )
132}
133
134pub fn owners_account_manager() -> Option<Agent> {
136 let name = std::fs::read_to_string(owners_account_manager_file()).ok()?;
137 load(name.trim())
138}
139
140#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
142#[serde(rename_all = "snake_case")]
143pub enum Role {
144 To,
146 Cc,
148}
149
150#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
152pub struct Participant {
153 pub address: MailAddress,
155 pub role: Role,
157}
158
159#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
161pub struct Thread {
162 pub id: String,
164 pub agent: String,
166 pub holder: MailAddress,
168 pub participants: Vec<Participant>,
170 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
172 pub markers: BTreeMap<String, String>,
173 #[serde(default, skip_serializing_if = "Option::is_none")]
175 pub surface: Option<String>,
176 pub created_at_ms: u64,
178 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
181 pub holder_launched: bool,
182 #[serde(default, skip_serializing_if = "Option::is_none")]
186 pub parent: Option<String>,
187}
188
189impl Thread {
190 pub fn role_of(&self, address: &MailAddress) -> Option<Role> {
192 self.participants
193 .iter()
194 .find(|participant| &participant.address == address)
195 .map(|participant| participant.role)
196 }
197
198 pub fn join(&mut self, address: &MailAddress, role: Role) {
200 match self
201 .participants
202 .iter_mut()
203 .find(|participant| &participant.address == address)
204 {
205 Some(participant) => participant.role = role,
206 None => self.participants.push(Participant {
207 address: address.clone(),
208 role,
209 }),
210 }
211 }
212
213 pub fn delegated(&self) -> bool {
215 load(&self.agent).is_some_and(|agent| agent.main_session != self.holder)
216 }
217}
218
219fn thread_path(id: &str) -> Option<PathBuf> {
220 valid_name(id).then(|| threads_dir().join(format!("{id}.json")))
221}
222
223pub fn thread(id: &str) -> Option<Thread> {
225 let bytes = std::fs::read(thread_path(id)?).ok()?;
226 serde_json::from_slice(&bytes).ok()
227}
228
229pub fn all_threads() -> Vec<Thread> {
231 let Ok(entries) = std::fs::read_dir(threads_dir()) else {
232 return Vec::new();
233 };
234 let mut threads: Vec<Thread> = entries
235 .flatten()
236 .filter(|entry| entry.path().extension().is_some_and(|ext| ext == "json"))
237 .filter_map(|entry| std::fs::read(entry.path()).ok())
238 .filter_map(|bytes| serde_json::from_slice(&bytes).ok())
239 .collect();
240 threads.sort_by(|a, b| (a.created_at_ms, &a.id).cmp(&(b.created_at_ms, &b.id)));
241 threads
242}
243
244pub fn threads_of(name: &str) -> Vec<Thread> {
246 all_threads()
247 .into_iter()
248 .filter(|thread| thread.agent == name)
249 .collect()
250}
251
252pub fn thread_held_by(address: &MailAddress) -> Option<Thread> {
255 all_threads().into_iter().rev().find(|thread| {
256 &thread.holder == address
257 && thread.parent.is_none()
258 && load(&thread.agent).is_some_and(|agent| &agent.main_session != address)
259 })
260}
261
262pub fn thread_of_marker(marker: &str) -> Option<(Thread, String)> {
264 all_threads().into_iter().find_map(|thread| {
265 let id = thread.markers.get(marker)?.clone();
266 Some((thread, id))
267 })
268}
269
270pub fn agent_of_session(address: &MailAddress) -> Option<Agent> {
272 if let Some(agent) = all_agents()
273 .into_iter()
274 .find(|agent| &agent.main_session == address)
275 {
276 return Some(agent);
277 }
278 thread_held_by(address).and_then(|thread| load(&thread.agent))
279}
280
281pub fn update_thread(
285 id: &str,
286 create: impl FnOnce() -> Option<Thread>,
287 change: impl FnOnce(&mut Thread),
288) -> std::io::Result<Option<Thread>> {
289 let Some(path) = thread_path(id) else {
290 return Ok(None);
291 };
292 std::fs::create_dir_all(threads_dir())?;
293 let lock = path.with_extension("lock");
294 let started = Instant::now();
295 loop {
296 match std::fs::OpenOptions::new()
297 .write(true)
298 .create_new(true)
299 .open(&lock)
300 {
301 Ok(_) => break,
302 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
303 let stale = std::fs::metadata(&lock)
305 .and_then(|meta| meta.modified())
306 .ok()
307 .and_then(|modified| modified.elapsed().ok())
308 .is_some_and(|age| age > Duration::from_secs(10));
309 if stale {
310 std::fs::remove_file(&lock).ok();
311 } else if started.elapsed() > Duration::from_secs(30) {
312 return Err(std::io::Error::new(
314 std::io::ErrorKind::TimedOut,
315 format!("thread {id} is locked by another writer"),
316 ));
317 }
318 std::thread::sleep(Duration::from_millis(20));
319 }
320 Err(error) => return Err(error),
321 }
322 }
323 let result = (|| {
324 let current = std::fs::read(&path)
325 .ok()
326 .and_then(|bytes| serde_json::from_slice::<Thread>(&bytes).ok())
327 .or_else(create);
328 let Some(mut thread) = current else {
329 return Ok(None);
330 };
331 change(&mut thread);
332 let bytes = serde_json::to_vec_pretty(&thread).map_err(std::io::Error::other)?;
333 write_atomically(&path, &bytes)?;
334 Ok(Some(thread))
335 })();
336 std::fs::remove_file(&lock).ok();
337 result
338}
339
340#[derive(Debug, Clone, PartialEq, Eq)]
342pub struct Recipient {
343 pub address: MailAddress,
345 pub wake: bool,
347}
348
349#[derive(Debug, Clone, PartialEq, Eq)]
351pub struct Plan {
352 pub recipients: Vec<Recipient>,
354 pub agent: Option<MailAddress>,
356 pub thread: Option<Thread>,
358}
359
360fn now_ms() -> u64 {
361 SystemTime::now()
362 .duration_since(UNIX_EPOCH)
363 .map(|elapsed| elapsed.as_millis() as u64)
364 .unwrap_or_default()
365}
366
367#[derive(Debug, Clone, Default, PartialEq, Eq)]
370pub struct Channel {
371 pub marker: Option<String>,
373 pub reply_marker: Option<String>,
375 pub surface: Option<String>,
377}
378
379fn record_marker(thread: &mut Thread, channel: &Channel, id: &str) {
380 if let Some(marker) = &channel.marker {
381 thread.markers.insert(marker.clone(), id.to_string());
382 }
383}
384
385pub fn plan(
390 envelope: &mut Envelope,
391 to: &MailAddress,
392 channel: &Channel,
393) -> std::io::Result<Option<Plan>> {
394 let sender = envelope.from.clone();
395 let id = envelope.id.clone();
396 if envelope.in_reply_to.is_none() {
398 if let Some((thread, parent)) = channel.reply_marker.as_deref().and_then(thread_of_marker) {
399 envelope.in_reply_to = Some(parent);
400 envelope.thread = Some(thread.id);
401 }
402 }
403
404 let known = envelope
406 .in_reply_to
407 .as_ref()
408 .and(envelope.thread.as_deref())
409 .and_then(thread);
410 if let Some(known) = known {
411 let receiver = (to.harness != AGENT_HARNESS && to != &sender).then(|| to.clone());
412 let saved = update_thread(
413 &known.id,
414 || None,
415 |thread| {
416 if thread.role_of(&sender).is_none() {
417 thread.join(&sender, Role::To);
418 }
419 if let Some(receiver) = &receiver {
420 if thread.role_of(receiver).is_none() {
421 thread.join(receiver, Role::To);
422 }
423 }
424 record_marker(thread, channel, &id);
425 },
426 )?
427 .unwrap_or(known);
428 envelope.thread = Some(saved.id.clone());
429 let recipients = saved
430 .participants
431 .iter()
432 .filter(|participant| participant.address != sender)
433 .map(|participant| Recipient {
434 address: participant.address.clone(),
435 wake: participant.role == Role::To,
436 })
437 .collect();
438 return Ok(Some(Plan {
439 recipients,
440 agent: None,
441 thread: Some(saved),
442 }));
443 }
444
445 if envelope.in_reply_to.is_none() {
451 if let Some(held) = thread_held_by(&sender) {
452 let agent = load(&held.agent);
453 let own_agent = to.harness == AGENT_HARNESS && to.session_id == held.agent;
454 let to_main = own_agent || agent.as_ref().is_some_and(|a| &a.main_session == to);
455 if to_main {
456 let main = agent.map(|a| a.main_session).unwrap_or_else(|| to.clone());
457 let saved = update_thread(
458 &held.id,
459 || None,
460 |thread| record_marker(thread, channel, &id),
461 )?
462 .unwrap_or(held);
463 envelope.thread = Some(saved.id.clone());
464 return Ok(Some(Plan {
465 recipients: vec![Recipient {
466 address: main,
467 wake: true,
468 }],
469 agent: None,
470 thread: Some(saved),
471 }));
472 }
473 let root = id.clone();
474 let receiver = if to.harness == AGENT_HARNESS {
476 load(&to.session_id)
477 .map(|other| other.main_session)
478 .ok_or_else(|| {
479 std::io::Error::other(format!(
480 "no agent named {} is declared on this machine",
481 to.session_id
482 ))
483 })?
484 } else {
485 to.clone()
486 };
487 let mut participants = vec![
488 Participant {
489 address: sender.clone(),
490 role: Role::To,
491 },
492 Participant {
493 address: receiver,
494 role: Role::To,
495 },
496 ];
497 if let Some(agent) = &agent {
498 participants.push(Participant {
499 address: agent.main_session.clone(),
500 role: Role::Cc,
501 });
502 }
503 let saved = update_thread(
504 &root,
505 || {
506 Some(Thread {
507 id: root.clone(),
508 agent: held.agent.clone(),
509 holder: sender.clone(),
510 participants,
511 markers: BTreeMap::new(),
512 surface: None,
513 created_at_ms: now_ms(),
514 holder_launched: false,
515 parent: Some(held.id.clone()),
516 })
517 },
518 |thread| record_marker(thread, channel, &id),
519 )?;
520 envelope.thread = Some(root);
521 let recipients = saved
522 .as_ref()
523 .map(|thread| {
524 thread
525 .participants
526 .iter()
527 .filter(|p| p.address != sender)
528 .map(|p| Recipient {
529 address: p.address.clone(),
530 wake: p.role == Role::To,
531 })
532 .collect()
533 })
534 .unwrap_or_default();
535 return Ok(Some(Plan {
536 recipients,
537 agent: None,
538 thread: saved,
539 }));
540 }
541 }
542
543 if envelope.in_reply_to.is_none() && envelope.voice_for.as_ref() == Some(to) {
546 if let Some(held) = thread_held_by(to) {
547 envelope.thread = Some(held.id.clone());
548 let mut recipients = vec![Recipient {
549 address: to.clone(),
550 wake: true,
551 }];
552 recipients.extend(
553 held.participants
554 .iter()
555 .filter(|p| p.role == Role::Cc && p.address != sender && &p.address != to)
556 .map(|p| Recipient {
557 address: p.address.clone(),
558 wake: false,
559 }),
560 );
561 return Ok(Some(Plan {
562 recipients,
563 agent: None,
564 thread: Some(held),
565 }));
566 }
567 }
568
569 if to.harness == AGENT_HARNESS {
571 let agent = load(&to.session_id).ok_or_else(|| {
572 std::io::Error::other(format!(
573 "no agent named {} is declared on this machine",
574 to.session_id
575 ))
576 })?;
577 let main = agent.main_session.clone();
578 if main == sender {
579 return Err(std::io::Error::other(format!(
580 "that is your own agent's address ({}); you are its main session",
581 agent.name
582 )));
583 }
584 if agent_of_session(&sender).is_some_and(|own| own.name == agent.name) {
586 return Ok(Some(Plan {
587 recipients: vec![Recipient {
588 address: main,
589 wake: true,
590 }],
591 agent: None,
592 thread: None,
593 }));
594 }
595 envelope.thread = None;
596 let root = id.clone();
597 let surface = channel.surface.clone();
598 let saved = update_thread(
599 &root,
600 || {
601 Some(Thread {
602 id: root.clone(),
603 agent: agent.name.clone(),
604 holder: main.clone(),
605 participants: vec![
606 Participant {
607 address: main.clone(),
608 role: Role::To,
609 },
610 Participant {
611 address: sender.clone(),
612 role: Role::To,
613 },
614 ],
615 markers: BTreeMap::new(),
616 surface,
617 created_at_ms: now_ms(),
618 holder_launched: false,
619 parent: None,
620 })
621 },
622 |thread| record_marker(thread, channel, &id),
623 )?;
624 return Ok(Some(Plan {
625 recipients: vec![Recipient {
626 address: main,
627 wake: true,
628 }],
629 agent: Some(to.clone()),
630 thread: saved,
631 }));
632 }
633
634 if envelope.in_reply_to.is_none() {
637 if let Some(agent) = all_agents()
638 .into_iter()
639 .find(|agent| agent.main_session == sender)
640 {
641 if to != &sender && to.harness != AGENT_HARNESS {
642 let root = id.clone();
643 let saved = update_thread(
644 &root,
645 || {
646 Some(Thread {
647 id: root.clone(),
648 agent: agent.name.clone(),
649 holder: sender.clone(),
650 participants: vec![
651 Participant {
652 address: sender.clone(),
653 role: Role::To,
654 },
655 Participant {
656 address: to.clone(),
657 role: Role::To,
658 },
659 ],
660 markers: BTreeMap::new(),
661 surface: None,
662 created_at_ms: now_ms(),
663 holder_launched: false,
664 parent: None,
665 })
666 },
667 |thread| record_marker(thread, channel, &id),
668 )?;
669 return Ok(Some(Plan {
670 recipients: vec![Recipient {
671 address: to.clone(),
672 wake: true,
673 }],
674 agent: None,
675 thread: saved,
676 }));
677 }
678 }
679 }
680 Ok(None)
681}
682
683pub fn root_of_main(main: &MailAddress, envelope: &Envelope) -> std::io::Result<Option<Thread>> {
690 let Some(agent) = agent_of_session(main).filter(|agent| &agent.main_session == main) else {
691 return Ok(None);
692 };
693 if &envelope.from == main || envelope.thread_id() != envelope.id {
694 return Ok(None);
695 }
696 let sender = envelope.from.clone();
697 update_thread(
698 &envelope.id,
699 || {
700 Some(Thread {
701 id: envelope.id.clone(),
702 agent: agent.name.clone(),
703 holder: main.clone(),
704 participants: vec![
705 Participant {
706 address: main.clone(),
707 role: Role::To,
708 },
709 Participant {
710 address: sender,
711 role: Role::To,
712 },
713 ],
714 markers: BTreeMap::new(),
715 surface: None,
716 created_at_ms: now_ms(),
717 holder_launched: false,
718 parent: None,
719 })
720 },
721 |_| {},
722 )
723}
724
725pub fn delegate_to(
726 thread_id: &str,
727 holder: &MailAddress,
728 launched: bool,
729) -> std::io::Result<Option<Thread>> {
730 let Some(current) = thread(thread_id) else {
731 return Ok(None);
732 };
733 let Some(agent) = load(¤t.agent) else {
734 return Ok(None);
735 };
736 update_thread(
737 thread_id,
738 || None,
739 |thread| {
740 if holder == &agent.main_session && thread.holder_launched && &thread.holder != holder {
742 let previous = thread.holder.clone();
743 thread.participants.retain(|p| p.address != previous);
744 }
745 thread.holder = holder.clone();
746 thread.holder_launched = launched && holder != &agent.main_session;
747 thread.join(holder, Role::To);
748 if holder != &agent.main_session {
749 thread.join(&agent.main_session, Role::Cc);
750 }
751 },
752 )
753}
754
755pub const RESUME_WAIT: Duration = Duration::from_secs(45);
757
758pub fn resumable(address: &MailAddress) -> bool {
762 all_threads()
763 .iter()
764 .any(|thread| &thread.holder == address && thread.holder_launched)
765 && recorded_resume_arguments(address).is_some()
766}
767
768fn recorded_resume_arguments(address: &MailAddress) -> Option<Vec<String>> {
773 let query = crate::DiscoveryQuery {
774 harnesses: vec![crate::HarnessId::new(&address.harness)],
775 query: Some(address.session_id.clone()),
776 limit: Some(50),
777 ..Default::default()
778 };
779 let locator = crate::sdk::discover_sessions(&query)
780 .ok()?
781 .into_iter()
782 .find(|row| row.locator.session_id == address.session_id)?
783 .locator;
784 let config = crate::harness_service::recorded_config(&locator).ok()?;
785 if config.get("session_id").and_then(|v| v.as_str()) != Some(address.session_id.as_str()) {
786 return None;
787 }
788 resume_arguments(&config)
789}
790
791fn resume_arguments(config: &serde_json::Value) -> Option<Vec<String>> {
795 let text = |key: &str| config.get(key).and_then(|v| v.as_str());
796 let mut args = Vec::new();
797 match text("harness")? {
798 "codex" => {
799 let approval = text("approval_policy")?;
800 if !["untrusted", "on-failure", "on-request", "never"].contains(&approval) {
801 return None;
802 }
803 let policy = config.get("sandbox_policy")?.as_object()?;
804 let kind = policy.get("type")?.as_str()?;
805 let allowed: &[&str] = match kind {
806 "read-only" | "danger-full-access" => &["type"],
807 "workspace-write" => &[
808 "type",
809 "writable_roots",
810 "network_access",
811 "exclude_tmpdir_env_var",
812 "exclude_slash_tmp",
813 ],
814 _ => return None,
815 };
816 if policy.keys().any(|key| !allowed.contains(&key.as_str())) {
817 return None;
818 }
819 args.extend([
820 "--ask-for-approval".into(),
821 approval.into(),
822 "--sandbox".into(),
823 kind.into(),
824 ]);
825 if kind == "workspace-write" {
826 let roots = policy
827 .get("writable_roots")
828 .cloned()
829 .unwrap_or(serde_json::json!([]));
830 if !roots
831 .as_array()?
832 .iter()
833 .all(|root| root.as_str().is_some_and(|r| Path::new(r).is_absolute()))
834 {
835 return None;
836 }
837 args.extend([
838 "-c".into(),
839 format!("sandbox_workspace_write.writable_roots={roots}"),
840 ]);
841 for key in [
842 "network_access",
843 "exclude_tmpdir_env_var",
844 "exclude_slash_tmp",
845 ] {
846 let value = match policy.get(key) {
847 Some(value) => value.as_bool()?,
848 None => key != "network_access",
849 };
850 args.extend([
851 "-c".into(),
852 format!("sandbox_workspace_write.{key}={value}"),
853 ]);
854 }
855 }
856 }
857 "claude-code" => {
858 let mode = text("permission_mode")?;
859 if !["default", "acceptEdits", "plan", "bypassPermissions"].contains(&mode) {
860 return None;
861 }
862 args.extend(["--permission-mode".into(), mode.into()]);
863 }
864 _ => return None,
865 }
866 if let Some(model) = text("model").filter(|m| !m.is_empty()) {
867 args.extend(["--model".into(), model.into()]);
868 }
869 Some(args)
870}
871
872fn resumed_path(address: &MailAddress) -> PathBuf {
875 let hash = blake3::hash(address.to_string().as_bytes()).to_hex();
876 agents_dir().join("resumed").join(&hash[..24])
877}
878
879pub fn resume(address: &MailAddress) -> std::io::Result<()> {
883 let mut recorded = recorded_resume_arguments(address).ok_or_else(|| {
884 std::io::Error::other(format!(
885 "{address} has no recorded permissions to resume with; its mail waits"
886 ))
887 })?;
888 if address.harness == "claude-code" {
892 if let Some(thread) = all_threads()
893 .into_iter()
894 .find(|thread| &thread.holder == address)
895 {
896 let id: String = address.session_id.chars().take(6).collect();
897 recorded.extend(["--name".to_string(), format!("{}-{id}", thread.agent)]);
898 }
899 }
900 let program = crate::claude_relay::supercode_program().map_err(std::io::Error::other)?;
901 let marker = resumed_path(address);
902 if let Some(dir) = marker.parent() {
903 std::fs::create_dir_all(dir)?;
904 }
905 write_atomically(&marker, now_ms().to_string().as_bytes())?;
906 let output = std::process::Command::new(program)
907 .args(["open", &address.session_id, "--detach", "--"])
909 .args(&recorded)
910 .stdin(std::process::Stdio::null())
911 .output()?;
912 if output.status.success() {
913 Ok(())
914 } else {
915 Err(std::io::Error::other(crate::mailbox::error_line(
916 &String::from_utf8_lossy(&output.stderr),
917 )))
918 }
919}
920
921pub fn file_unread(address: &MailAddress, envelope: &Envelope) -> std::io::Result<()> {
923 Mailbox::open(&mail_root(), address)?
924 .deliver(envelope)
925 .map(|_| ())
926}
927
928pub fn register_folder_sessions(homes: &crate::HarnessHomes) {
932 let agents: Vec<Agent> = all_agents()
933 .into_iter()
934 .filter(|agent| agent.folder.is_some())
935 .collect();
936 if agents.is_empty() {
937 return;
938 }
939 let holders: Vec<MailAddress> = all_threads()
940 .into_iter()
941 .map(|thread| thread.holder)
942 .collect();
943 for session in crate::mail_route::LiveSessions::read(homes).all() {
944 let Some(cwd) = &session.cwd else { continue };
945 let Some(agent) = agents.iter().find(|agent| {
946 agent.folder.as_deref() == Some(cwd.as_path()) && agent.main_session != session.address
947 }) else {
948 continue;
949 };
950 if holders.contains(&session.address) {
951 continue;
952 }
953 let id = format!("{SESSION_THREAD_PREFIX}{}", session.address.session_id);
954 let holder = session.address.clone();
955 let main = agent.main_session.clone();
956 update_thread(
957 &id,
958 || {
959 Some(Thread {
960 id: id.clone(),
961 agent: agent.name.clone(),
962 holder: holder.clone(),
963 participants: vec![
964 Participant {
965 address: holder.clone(),
966 role: Role::To,
967 },
968 Participant {
969 address: main.clone(),
970 role: Role::Cc,
971 },
972 ],
973 markers: BTreeMap::new(),
974 surface: None,
975 created_at_ms: now_ms(),
976 holder_launched: false,
977 parent: None,
978 })
979 },
980 |_| {},
981 )
982 .ok();
983 }
984}
985
986fn typed_cursor_path(address: &MailAddress) -> PathBuf {
987 let hash = blake3::hash(address.to_string().as_bytes()).to_hex();
988 agents_dir().join("typed").join(&hash[..24])
989}
990
991pub fn copy_typed_lines(homes: &crate::HarnessHomes) {
997 let agents = all_agents();
998 if agents.is_empty() {
999 return;
1000 }
1001 let threads = all_threads();
1002 let account_manager = owners_account_manager();
1003 let mut sessions: Vec<(MailAddress, String)> = agents
1004 .iter()
1005 .map(|agent| (agent.main_session.clone(), agent.name.clone()))
1006 .collect();
1007 for thread in &threads {
1008 if !sessions
1009 .iter()
1010 .any(|(address, _)| address == &thread.holder)
1011 {
1012 sessions.push((thread.holder.clone(), thread.agent.clone()));
1013 }
1014 }
1015 let now = now_ms();
1016 for (session, agent_name) in sessions {
1017 let cursor_path = typed_cursor_path(&session);
1018 let Some(cursor) = std::fs::read_to_string(&cursor_path)
1019 .ok()
1020 .and_then(|text| text.trim().parse::<u64>().ok())
1021 else {
1022 write_atomically(&cursor_path, now.to_string().as_bytes()).ok();
1023 continue;
1024 };
1025 if let Some(path) = crate::mail_transcript::transcript(homes, &session) {
1027 let changed = std::fs::metadata(&path)
1028 .and_then(|meta| meta.modified())
1029 .ok()
1030 .and_then(|modified| modified.duration_since(UNIX_EPOCH).ok())
1031 .map(|since| since.as_millis() as u64);
1032 if changed.is_some_and(|changed| changed <= cursor) {
1033 continue;
1034 }
1035 }
1036 let lines: Vec<(u64, String)> = if session.harness == "codex" {
1039 Mailbox::open(&mail_root(), &session)
1040 .and_then(|mailbox| mailbox.list())
1041 .unwrap_or_default()
1042 .into_iter()
1043 .filter(|stored| {
1044 stored.envelope.kind == MailKind::User
1045 && stored
1046 .envelope
1047 .id
1048 .starts_with(crate::mail_transcript::TYPED_ID_PREFIX)
1049 && stored.envelope.created_at_ms > cursor
1050 })
1051 .map(|stored| (stored.envelope.created_at_ms, stored.envelope.body))
1052 .collect()
1053 } else {
1054 let Some(mail) = crate::mail_transcript::read(homes, &session) else {
1055 continue;
1056 };
1057 let mut typed_by_supercode: std::collections::HashMap<String, usize> =
1060 std::collections::HashMap::new();
1061 for stored in Mailbox::open(&mail_root(), &session)
1062 .and_then(|mailbox| mailbox.list())
1063 .unwrap_or_default()
1064 {
1065 if stored.envelope.kind == MailKind::User
1066 && !stored
1067 .envelope
1068 .id
1069 .starts_with(crate::mail_transcript::TYPED_ID_PREFIX)
1070 {
1071 *typed_by_supercode
1072 .entry(stored.envelope.body.trim().to_string())
1073 .or_default() += 1;
1074 }
1075 }
1076 mail.typed
1077 .into_iter()
1078 .filter(|line| !line.withdrawn && line.sent_at_ms > cursor)
1079 .filter(|line| match typed_by_supercode.get_mut(line.text.trim()) {
1080 Some(left) if *left > 0 => {
1081 *left -= 1;
1082 false
1083 }
1084 _ => true,
1085 })
1086 .map(|line| (line.sent_at_ms, line.text))
1087 .collect()
1088 };
1089 let Some(newest) = lines.iter().map(|(at, _)| *at).max() else {
1090 continue;
1091 };
1092 let held = threads
1093 .iter()
1094 .rev()
1095 .find(|thread| thread.holder == session && thread.delegated());
1096 let mut receivers: Vec<MailAddress> = held
1097 .map(|thread| {
1098 thread
1099 .participants
1100 .iter()
1101 .filter(|p| p.role == Role::Cc && p.address != session)
1102 .map(|p| p.address.clone())
1103 .collect()
1104 })
1105 .unwrap_or_default();
1106 if let Some(account_manager) = &account_manager {
1107 if agent_name != account_manager.name
1108 && !receivers.contains(&account_manager.main_session)
1109 {
1110 receivers.push(account_manager.main_session.clone());
1111 }
1112 }
1113 let Ok(user) = crate::mail_transcript::user_address(&session.machine) else {
1114 continue;
1115 };
1116 for (sent_at_ms, text) in lines {
1117 let copy = Envelope {
1118 id: crate::mail_transcript::typed_line_id(&session, sent_at_ms, &text),
1119 created_at_ms: sent_at_ms,
1120 from: user.clone(),
1121 from_name: format!("user@{} (typed into {session})", session.machine),
1122 sender_identity: None,
1123 kind: MailKind::Typed,
1124 reply_via: crate::mailbox::ReplyVia::None,
1125 in_reply_to: None,
1126 in_reply_to_inferred: false,
1127 thread: held.map(|thread| thread.id.clone()),
1128 native_from: None,
1129 voice_for: None,
1130 subject: None,
1131 body: text,
1132 };
1133 for receiver in &receivers {
1134 if let Err(error) = file_unread(receiver, ©) {
1136 eprintln!("supercode mail: the line typed into {session} was not copied to {receiver}: {error}");
1137 }
1138 }
1139 }
1140 write_atomically(&cursor_path, newest.to_string().as_bytes()).ok();
1141 }
1142}
1143
1144pub fn close_idle_threads(homes: &crate::HarnessHomes) {
1148 let launched: Vec<(Thread, u64)> = all_threads()
1149 .into_iter()
1150 .filter(|thread| thread.holder_launched)
1151 .filter_map(|thread| {
1152 let minutes = load(&thread.agent)?.idle_minutes?;
1153 Some((thread, minutes))
1154 })
1155 .collect();
1156 if launched.is_empty() {
1157 return;
1158 }
1159 let live = crate::mail_route::LiveSessions::read(homes);
1160 let now = now_ms();
1161 for (thread, minutes) in launched {
1162 let Some(session) = live
1163 .all()
1164 .iter()
1165 .find(|session| session.address == thread.holder)
1166 else {
1167 continue;
1168 };
1169 if session.status != "idle" {
1170 continue;
1171 }
1172 let Some(last) = session.last_message_at_ms(homes) else {
1173 continue;
1174 };
1175 let resumed = std::fs::read_to_string(resumed_path(&thread.holder))
1176 .ok()
1177 .and_then(|text| text.trim().parse::<u64>().ok())
1178 .unwrap_or(0);
1179 let last = last.max(resumed);
1180 if now.saturating_sub(last) < minutes.saturating_mul(60_000) {
1181 continue;
1182 }
1183 let Ok(program) = crate::claude_relay::supercode_program() else {
1184 return;
1185 };
1186 std::process::Command::new(program)
1187 .args(["close", &thread.holder.to_string()])
1188 .stdin(std::process::Stdio::null())
1189 .stdout(std::process::Stdio::null())
1190 .stderr(std::process::Stdio::null())
1191 .status()
1192 .ok();
1193 }
1194}