1use std::fs::File;
70use std::collections::{HashMap, HashSet};
71use std::path::{Path, PathBuf};
72use std::sync::Arc;
73use zeroize::Zeroizing;
74
75use crate::ipc::{SendListener, SendStream};
76pub(crate) use crate::local_bus::{ensure_server_name, LocalBus};
77use crate::{
78 clip_public, register_session, DeviceId, OpenedStore,
79 SessionConfig, SessionWake, ShellError, HOMESERVER_URL_ENV, STORE_ROOT_ENV,
80};
81
82pub const SESSIONS_DIR_ENV: &str = "M4A_SESSIONS_DIR";
87
88pub const AGENTS_DIR_ENV: &str = "M4A_AGENTS_DIR";
92
93pub const DEFAULT_AGENTS_DIR: &str = "agent-data/agents";
97
98pub(crate) fn under_home(rel: &str) -> PathBuf {
100 match std::env::var("HOME").or_else(|_| std::env::var("USERPROFILE")) {
101 Ok(h) if !h.is_empty() => PathBuf::from(h).join(rel),
102 _ => PathBuf::from(rel),
103 }
104}
105
106pub const AGENT_RESCAN_SECS_ENV: &str = "M4A_AGENT_RESCAN_SECS";
109
110pub struct HostSession {
115 pub bot_name: String,
117 pub session_id: String,
120 pub agent_id: Option<String>,
123 pub routine_url: Option<String>,
125 pub routine_bearer: Option<String>,
127 pub invite: Option<String>,
130 pub routine_file: Option<PathBuf>,
133 pub tier: m4a_agent::BackendKind,
135}
136
137impl HostSession {
138 pub fn new(bot_name: impl Into<String>, session_id: impl Into<String>) -> Self {
140 Self {
141 bot_name: bot_name.into(),
142 session_id: session_id.into(),
143 agent_id: None,
144 routine_url: None,
145 routine_bearer: None,
146 invite: None,
147 routine_file: None,
148 tier: m4a_agent::BackendKind::Server,
149 }
150 }
151
152 pub fn with_routine_file(mut self, path: impl Into<PathBuf>) -> Self {
154 self.routine_file = Some(path.into());
155 self
156 }
157
158 pub(crate) fn config(&self, url: &str, store_root: &Path) -> Result<SessionConfig, ShellError> {
159 Ok(SessionConfig::new_identity(url, self.tier, &self.session_id, store_root, self.invite.clone())?.with_nick_request(&self.bot_name))
160 }
161
162 pub fn with_invite(mut self, invite: impl Into<String>) -> Self {
164 self.invite = Some(invite.into());
165 self
166 }
167
168 pub fn with_tier(mut self, tier: m4a_agent::BackendKind) -> Self {
170 self.tier = tier;
171 self
172 }
173
174 pub fn with_routine(mut self, url: impl Into<String>, bearer: Option<String>) -> Self {
176 self.routine_url = Some(url.into());
177 self.routine_bearer = bearer.filter(|token| !token.is_empty());
178 self
179 }
180
181}
182
183impl std::fmt::Debug for HostSession {
184 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
185 f.debug_struct("HostSession")
186 .field("bot_name", &self.bot_name)
187 .field("session_id", &self.session_id)
188 .field("agent_id", &self.agent_id)
189 .field("routine_url", &self.routine_url.as_ref().map(|_| "[set]"))
190 .field(
191 "routine_bearer",
192 &self.routine_bearer.as_ref().map(|_| "[redacted]"),
193 )
194 .finish()
195 }
196}
197
198#[derive(serde::Deserialize)]
199#[serde(deny_unknown_fields)]
200struct SessionFile {
201 bot_name: String,
202 session_id: String,
203 #[serde(default)]
205 agent_id: Option<String>,
206 #[serde(default)]
209 routine_file: Option<String>,
210}
211
212pub fn load_session_records(dir: &Path) -> Result<Vec<HostSession>, ShellError> {
216 if !dir.is_dir() {
217 return Err(ShellError::SessionList(
218 "session directory is missing".to_string(),
219 ));
220 }
221 let mut paths = Vec::new();
222 for entry in std::fs::read_dir(dir)? {
223 let entry = entry?;
224 let path = entry.path();
225 if path.extension().and_then(|ext| ext.to_str()) != Some("json") {
226 continue;
227 }
228 paths.push(path);
229 }
230 paths.sort();
231 if paths.is_empty() {
232 return Err(ShellError::SessionList(
233 "session directory has no records".to_string(),
234 ));
235 }
236 let mut out = Vec::with_capacity(paths.len());
237 for path in paths {
238 let text = std::fs::read_to_string(&path)
239 .map_err(|err| ShellError::SessionList(clip_public(err.to_string())))?;
240 let file: SessionFile = serde_json::from_str(&text)
241 .map_err(|err| ShellError::SessionList(clip_public(err.to_string())))?;
242 let mut session = HostSession::new(file.bot_name, file.session_id);
243 session.agent_id = file
244 .agent_id
245 .map(|id| id.trim().to_string())
246 .filter(|id| !id.is_empty());
247 session.routine_file = file.routine_file.map(|p| PathBuf::from(p.trim())).filter(|p| !p.as_os_str().is_empty());
248 out.push(session);
249 }
250 Ok(out)
251}
252
253#[derive(serde::Deserialize)]
254struct AgentProfileFile {
255 name: String,
256}
257
258pub fn load_agents_dir(dir: &Path) -> Result<Vec<HostSession>, ShellError> {
262 if !dir.is_dir() {
263 return Err(ShellError::SessionList(
264 "agents directory is missing".to_string(),
265 ));
266 }
267 let mut ids = Vec::new();
268 for entry in std::fs::read_dir(dir)? {
269 let entry = entry?;
270 if !entry.file_type()?.is_dir() {
271 continue;
272 }
273 let name = entry.file_name();
274 let Some(id) = name.to_str() else {
275 continue;
276 };
277 if id.is_empty() || id.starts_with('.') {
278 continue;
279 }
280 ids.push(id.to_string());
281 }
282 ids.sort();
283 if ids.is_empty() {
284 return Err(ShellError::SessionList(
285 "agents directory has no agents".to_string(),
286 ));
287 }
288 let mut out = Vec::with_capacity(ids.len());
289 for id in ids {
290 let profile_path = dir.join(&id).join("profile.json");
291 if !profile_path.is_file() {
292 continue;
293 }
294 let text = std::fs::read_to_string(&profile_path)
295 .map_err(|err| ShellError::SessionList(clip_public(err.to_string())))?;
296 let profile: AgentProfileFile = serde_json::from_str(&text)
297 .map_err(|err| ShellError::SessionList(clip_public(err.to_string())))?;
298 let bot_name = profile.name.trim().to_string();
299 if bot_name.is_empty() {
300 continue;
301 }
302 let mut session = HostSession::new(bot_name, id.clone());
303 session.agent_id = Some(id);
304 out.push(session);
305 }
306 if out.is_empty() {
307 return Err(ShellError::SessionList(
308 "agents directory has no profiles".to_string(),
309 ));
310 }
311 Ok(out)
312}
313
314fn explicit_agents_dir(mut get: impl FnMut(&str) -> Option<String>) -> Option<PathBuf> {
319 let path = PathBuf::from(get(AGENTS_DIR_ENV)?);
320 path.is_dir().then_some(path)
321}
322
323fn resolve_agents_dir(mut get: impl FnMut(&str) -> Option<String>) -> Option<PathBuf> {
327 if let Some(path) = get(AGENTS_DIR_ENV).filter(|value| !value.is_empty()) {
328 let path = PathBuf::from(path);
329 if path.is_dir() {
330 return Some(path);
331 }
332 return None;
333 }
334 let default = under_home(DEFAULT_AGENTS_DIR);
335 if default.is_dir() {
336 Some(default)
337 } else {
338 None
339 }
340}
341
342fn load_web_sessions(
350 dir: &Path,
351 mut get: impl FnMut(&str) -> Option<String>,
352) -> Result<Vec<HostSession>, ShellError> {
353 let sessions = load_session_records(dir)?;
354 let mut sessions = drop_skipped(sessions, &skip_nicks(&mut get));
355 attach_routines_from_env(&mut sessions, &mut get)?;
356 Ok(sessions)
357}
358
359fn load_web_agents(
361 dir: &Path,
362 mut get: impl FnMut(&str) -> Option<String>,
363) -> Result<Vec<HostSession>, ShellError> {
364 let sessions = load_agents_dir(dir)?;
365 let mut sessions = drop_skipped(sessions, &skip_nicks(&mut get));
366 apply_session_ids(
367 &mut sessions,
368 &parse_session_ids(get(SESSION_IDS_ENV).as_deref().unwrap_or("")),
369 );
370 attach_routines_from_env(&mut sessions, &mut get)?;
371 Ok(sessions)
372}
373
374fn attach_routines_from_env(
375 sessions: &mut [HostSession],
376 get: &mut impl FnMut(&str) -> Option<String>,
377) -> Result<(), ShellError> {
378 let store_root = get(STORE_ROOT_ENV).map(PathBuf::from);
379 let Some(path) = get(GATEWAY_FILE_ENV).filter(|value| !value.is_empty()) else {
380 return Ok(());
381 };
382 let token = get(GATEWAY_TOKEN_ENV).filter(|value| !value.is_empty());
383 attach_webhook_routines(
384 sessions,
385 Path::new(&path),
386 token.as_deref(),
387 store_root.as_deref(),
388 )?;
389 Ok(())
390}
391
392fn unopened_agent_sessions(
396 found: Vec<HostSession>,
397 skip: &[String],
398 aliases: &[(String, String)],
399 store_root: &Path,
400 held: &[(PathBuf, Option<String>)],
401) -> Vec<(HostSession, String)> {
402 let mut sessions = drop_skipped(found, skip);
403 apply_session_ids(&mut sessions, aliases);
404 sessions
405 .into_iter()
406 .filter_map(|session| {
407 let nick = routine_name_for(&session)?;
408 let dir = crate::session_store_dir(store_root, &session.session_id);
409 let taken = held.iter().any(|(held_dir, held_nick)| {
410 *held_dir == dir
411 || held_nick
412 .as_deref()
413 .is_some_and(|held| held.eq_ignore_ascii_case(&nick))
414 });
415 (!taken).then_some((session, nick))
416 })
417 .collect()
418}
419
420pub const SESSION_IDS_ENV: &str = "M4A_SESSION_IDS";
425
426pub(crate) fn parse_session_ids(raw: &str) -> Vec<(String, String)> {
427 raw.split(',')
428 .filter_map(|pair| {
429 let (agent, session) = pair.split_once('=')?;
430 let (agent, session) = (agent.trim(), session.trim());
431 (!agent.is_empty() && !session.is_empty())
432 .then(|| (agent.to_string(), session.to_string()))
433 })
434 .collect()
435}
436
437pub(crate) fn apply_session_ids(sessions: &mut [HostSession], aliases: &[(String, String)]) {
438 for session in sessions.iter_mut() {
439 let Some(agent_id) = session.agent_id.as_deref() else {
440 continue;
441 };
442 if let Some((_, alias)) = aliases.iter().find(|(agent, _)| agent == agent_id) {
443 session.session_id = alias.clone();
444 }
445 }
446}
447
448pub const SKIP_NICKS_ENV: &str = "M4A_SKIP_NICKS";
452
453fn skip_nicks(get: &mut impl FnMut(&str) -> Option<String>) -> Vec<String> {
454 parse_skip_nicks(get(SKIP_NICKS_ENV).as_deref().unwrap_or(""))
455}
456
457fn parse_skip_nicks(raw: &str) -> Vec<String> {
458 raw.split(',')
459 .map(|item| item.trim().to_ascii_lowercase())
460 .filter(|item| !item.is_empty())
461 .collect()
462}
463
464fn drop_skipped(sessions: Vec<HostSession>, skip: &[String]) -> Vec<HostSession> {
465 if skip.is_empty() {
466 return sessions;
467 }
468 sessions
469 .into_iter()
470 .filter(|session| match routine_name_for(session) {
471 Some(nick) => !skip.iter().any(|item| item.eq_ignore_ascii_case(&nick)),
472 None => true,
473 })
474 .collect()
475}
476
477#[derive(Debug, Clone, PartialEq, Eq)]
480pub enum WakeStatus {
481 Ready,
483 AwaitingBackend,
487 Failed(String),
489}
490
491#[derive(Debug, Clone, PartialEq, Eq)]
493pub struct RoutineReport {
494 pub agent_id: String,
496 pub nick: String,
498 pub folder_id: String,
501 pub status: WakeStatus,
503}
504
505#[derive(Debug, Clone, Default)]
508pub struct WakeOptions {
509 pub skip_nicks: Vec<String>,
511 pub store_root: Option<PathBuf>,
513 pub profile_note: bool,
516 pub session_ids: Vec<(String, String)>,
518}
519
520impl WakeOptions {
521 pub fn from_lookup(mut get: impl FnMut(&str) -> Option<String>) -> Self {
523 Self {
524 skip_nicks: skip_nicks(&mut get),
525 store_root: get(STORE_ROOT_ENV).map(PathBuf::from),
526 profile_note: get(PROFILE_NOTE_ENV).as_deref() == Some("1"),
527 session_ids: parse_session_ids(get(SESSION_IDS_ENV).as_deref().unwrap_or("")),
528 }
529 }
530}
531
532pub fn ensure_agent_webhook_routines(
540 agents_dir: &Path,
541 gateway_file: &Path,
542 token_override: Option<&str>,
543 options: &WakeOptions,
544) -> Result<Vec<RoutineReport>, ShellError> {
545 let store_root = options.store_root.as_deref();
546 let mut sessions = drop_skipped(load_agents_dir(agents_dir)?, &options.skip_nicks);
547 apply_session_ids(&mut sessions, &options.session_ids);
548 let Some(gate) = open_gateway(gateway_file, token_override)? else {
549 return Err(ShellError::Gateway("gateway file is missing".to_string()));
550 };
551 let mut reports = Vec::new();
552 for session in &sessions {
553 let Some(agent_id) = session.agent_id.as_deref().filter(|id| !id.is_empty()) else {
554 continue;
555 };
556 let Some(nick) = routine_name_for(session) else {
557 continue;
558 };
559 let outcome = ensure_wake(&gate, agent_id, &nick);
560 let folder_id = crate::nick::routine_folder_id(&nick).unwrap_or_default();
561 if let (WakeOutcome::Ready { url, key }, Some(root)) = (&outcome, store_root) {
562 save_wake(root, &session.session_id, &folder_id, url, key);
563 }
564 log_outcome(&nick, &folder_id, &outcome);
565 if options.profile_note && matches!(outcome, WakeOutcome::AwaitingBackend) {
566 note_profile(&gate, agents_dir, agent_id, &nick);
567 }
568 reports.push(RoutineReport {
569 agent_id: agent_id.to_string(),
570 nick,
571 folder_id,
572 status: outcome.status(),
573 });
574 }
575 Ok(reports)
576}
577
578pub fn ensure_agent_webhook_routines_from_env() -> Result<Vec<RoutineReport>, ShellError> {
582 let mut get = |key: &str| std::env::var(key).ok().filter(|v| !v.is_empty());
583 let agents_dir = resolve_agents_dir(&mut get)
584 .ok_or_else(|| ShellError::SessionList("agents directory is missing".to_string()))?;
585 let gateway = get(GATEWAY_FILE_ENV).unwrap_or_else(|| DEFAULT_GATEWAY_FILE.to_string());
586 let token = get(GATEWAY_TOKEN_ENV);
587 let options = WakeOptions::from_lookup(&mut get);
588 ensure_agent_webhook_routines(&agents_dir, Path::new(&gateway), token.as_deref(), &options)
589}
590
591pub const PROFILE_NOTE_ENV: &str = "M4A_BOOTSTRAP_PROFILE_NOTE";
597
598const NOTE_OPEN: &str = "<!-- mail4agent:wake -->";
599const NOTE_CLOSE: &str = "<!-- /mail4agent:wake -->";
600const DESCRIPTION_MAX: usize = 20_000;
602
603fn wake_note(routine: &str) -> String {
606 format!(
607 "{NOTE_OPEN}\nmail4agent wake: keep exactly one routine named \"{routine}\" with a webhook trigger, enabled. If it is missing, create it with UpdateRoutine: name \"{routine}\", trigger webhook, prompt \"A mail4agent room message woke this routine. The webhook JSON has body, from, event_id, and nick. Read that message.\" Do not rename it and do not create a second one.\n{NOTE_CLOSE}"
608 )
609}
610
611fn with_wake_note(description: &str, routine: &str) -> Option<String> {
615 let note = wake_note(routine);
616 if description.contains(¬e) {
617 return None;
618 }
619 if let (Some(start), Some(end)) = (description.find(NOTE_OPEN), description.find(NOTE_CLOSE)) {
620 if end > start {
621 let mut out = String::with_capacity(description.len() + note.len());
622 out.push_str(&description[..start]);
623 out.push_str(¬e);
624 out.push_str(&description[end + NOTE_CLOSE.len()..]);
625 return Some(out);
626 }
627 }
628 let base = description.trim_end();
629 Some(if base.is_empty() {
630 note
631 } else {
632 format!("{base}\n\n{note}")
633 })
634}
635
636fn note_profile(gate: &GatewayConn, agents_dir: &Path, agent_id: &str, nick: &str) {
640 match ensure_profile_note(gate, agents_dir, agent_id, nick) {
641 Ok(true) => eprintln!("mail4agent: wake {nick}: bootstrap note written to the profile"),
642 Ok(false) => {}
643 Err(err) => eprintln!("mail4agent: wake {nick}: bootstrap note not written ({err})"),
644 }
645}
646
647fn ensure_profile_note(
648 gate: &GatewayConn,
649 agents_dir: &Path,
650 agent_id: &str,
651 nick: &str,
652) -> Result<bool, ShellError> {
653 let text = std::fs::read_to_string(agents_dir.join(agent_id).join("profile.json"))
654 .map_err(|_| ShellError::Gateway("profile is unreadable".to_string()))?;
655 let profile: serde_json::Value = serde_json::from_str(&text)
656 .map_err(|_| ShellError::Gateway("profile is unreadable".to_string()))?;
657 let field = |key: &str| profile.get(key).and_then(|value| value.as_str());
658 let Some(name) = field("name").filter(|name| !name.trim().is_empty()) else {
659 return Err(ShellError::Gateway("profile has no name".to_string()));
660 };
661 let routine = crate::nick::routine_folder_id(nick)
662 .ok_or_else(|| ShellError::Gateway("nick has no routine name".to_string()))?;
663 let Some(description) = with_wake_note(field("description").unwrap_or(""), &routine) else {
664 return Ok(false);
665 };
666 if description.chars().count() > DESCRIPTION_MAX {
667 return Err(ShellError::Gateway(
668 "description would be too long".to_string(),
669 ));
670 }
671 let mut body = serde_json::json!({ "name": name, "description": description });
672 for key in ["title", "avatarShape", "avatarColor"] {
673 if let Some(value) = field(key) {
674 body[key] = serde_json::Value::String(value.to_string());
675 }
676 }
677 let (status, _) = gateway_post(
678 gate,
679 "updateAgent",
680 &serde_json::json!({ "id": agent_id, "profile": body }),
681 )?;
682 if status != 200 {
683 return Err(ShellError::Gateway(format!(
684 "gateway profile status {status}"
685 )));
686 }
687 Ok(true)
688}
689
690const DEFAULT_GATEWAY_FILE: &str = "agent-data/gateway.json";
693
694pub const GATEWAY_FILE_ENV: &str = "M4A_GATEWAY_FILE";
698
699pub const GATEWAY_TOKEN_ENV: &str = "M4A_GATEWAY_TOKEN";
702
703const MIRROR_ROUTINE_PROMPT: &str = "mail4agent wake mirror. Disabled on purpose: it only lets the local gateway hand out the webhook key for the routine with this folder id, which the bot creates itself with UpdateRoutine.";
707
708pub const WAKE_KEYCHAIN_FILE: &str = "routine-wake.json";
712
713#[derive(serde::Deserialize)]
714struct GatewayFile {
715 port: u16,
716 scheme: String,
717 #[serde(default)]
718 token: Option<String>,
719}
720
721#[derive(serde::Deserialize)]
722struct AutomationCard {
723 id: String,
724 name: String,
725 #[serde(default)]
726 trigger: Option<serde_json::Value>,
727}
728
729impl AutomationCard {
730 fn is_webhook(&self) -> bool {
733 fn has_webhook(value: &serde_json::Value) -> bool {
734 match value {
735 serde_json::Value::Array(items) => items.iter().any(has_webhook),
736 serde_json::Value::Object(map) => {
737 map.get("type").and_then(|kind| kind.as_str()) == Some("webhook")
738 || map
739 .values()
740 .any(|inner| inner.is_array() && has_webhook(inner))
741 }
742 _ => false,
743 }
744 }
745 self.trigger.as_ref().map(has_webhook).unwrap_or(true)
746 }
747}
748
749#[derive(serde::Deserialize)]
750struct WebhookCredential {
751 url: String,
752 key: Option<String>,
753}
754
755#[derive(serde::Serialize, serde::Deserialize)]
756#[serde(deny_unknown_fields)]
757struct StoredWake {
758 folder_id: String,
759 url: String,
760 key: String,
761}
762
763pub(crate) fn routine_name_for(session: &HostSession) -> Option<String> {
768 crate::nick::nick_from_display_name(session.bot_name.trim()).ok()
769}
770
771struct GatewayConn {
772 client: reqwest::blocking::Client,
773 base: String,
774 token: String,
775}
776
777fn open_gateway(
778 gateway_file: &Path,
779 token_override: Option<&str>,
780) -> Result<Option<GatewayConn>, ShellError> {
781 if !gateway_file.is_file() {
782 return Ok(None);
783 }
784 let text = std::fs::read_to_string(gateway_file)
785 .map_err(|_| ShellError::Gateway("gateway file is unreadable".to_string()))?;
786 let file: GatewayFile = serde_json::from_str(&text)
787 .map_err(|_| ShellError::Gateway("gateway file is unreadable".to_string()))?;
788 let scheme = file.scheme.trim();
789 if scheme != "http" && scheme != "https" {
790 return Err(ShellError::Gateway(
791 "gateway scheme is not http or https".to_string(),
792 ));
793 }
794 if file.port == 0 {
795 return Err(ShellError::Gateway("gateway port is unset".to_string()));
796 }
797 let token = token_override
798 .map(str::trim)
799 .filter(|value| !value.is_empty())
800 .map(str::to_string)
801 .or(file.token.filter(|value| !value.trim().is_empty()));
802 let Some(token) = token else {
803 return Err(ShellError::Gateway("gateway token is unset".to_string()));
804 };
805 let base = format!("{scheme}://127.0.0.1:{}", file.port);
807 Ok(Some(GatewayConn {
808 client: gateway_client()?,
809 base,
810 token,
811 }))
812}
813
814enum WakeOutcome {
817 Ready { url: String, key: String },
818 AwaitingBackend,
819 Failed(String),
820}
821
822impl WakeOutcome {
823 fn status(&self) -> WakeStatus {
824 match self {
825 WakeOutcome::Ready { .. } => WakeStatus::Ready,
826 WakeOutcome::AwaitingBackend => WakeStatus::AwaitingBackend,
827 WakeOutcome::Failed(reason) => WakeStatus::Failed(reason.clone()),
828 }
829 }
830}
831
832fn log_outcome(nick: &str, folder_id: &str, outcome: &WakeOutcome) {
833 match outcome {
834 WakeOutcome::Ready { .. } => eprintln!("mail4agent: wake {nick} ({folder_id}): ready"),
835 WakeOutcome::AwaitingBackend => eprintln!(
836 "mail4agent: wake {nick} ({folder_id}): no key yet; the bot has not created its own routine with this folder id; will retry"
837 ),
838 WakeOutcome::Failed(reason) => {
839 eprintln!("mail4agent: wake {nick} ({folder_id}): {reason}")
840 }
841 }
842}
843
844fn ensure_wake(gate: &GatewayConn, agent_id: &str, nick: &str) -> WakeOutcome {
854 let Some(folder) = crate::nick::routine_folder_id(nick) else {
855 return WakeOutcome::Failed("nick has no routine folder id".to_string());
856 };
857 let before = match list_agent_automations(gate, agent_id) {
858 Ok(cards) => cards,
859 Err(err) => return WakeOutcome::Failed(err.to_string()),
860 };
861 match before.iter().find(|card| card.id == folder) {
862 Some(card) if !card.is_webhook() => {
863 return WakeOutcome::Failed(
864 "a routine in this folder is not webhook-triggered; left as is".to_string(),
865 );
866 }
867 Some(_) => {}
868 None => {
869 let after = match create_mirror_routine(gate, agent_id, &folder) {
870 Ok(cards) => cards,
871 Err(err) => return WakeOutcome::Failed(err.to_string()),
872 };
873 if !after.iter().any(|card| card.id == folder) {
874 for stray in after.iter().filter(|card| {
875 card.name == folder && !before.iter().any(|old| old.id == card.id)
876 }) {
877 let _ = delete_agent_automation(gate, agent_id, &stray.id);
878 }
879 return WakeOutcome::Failed(
880 "the mirror did not land in the expected folder; removed it".to_string(),
881 );
882 }
883 }
884 }
885 match read_webhook_credential(gate, agent_id, &folder) {
886 Ok(CredentialRead::Ready { url, key }) => WakeOutcome::Ready { url, key },
887 Ok(CredentialRead::MintFailed) => WakeOutcome::AwaitingBackend,
888 Ok(CredentialRead::Missing) => {
889 WakeOutcome::Failed("the mirror is gone from the gateway".to_string())
890 }
891 Err(err) => WakeOutcome::Failed(err.to_string()),
892 }
893}
894
895fn attach_webhook_routines(
900 sessions: &mut [HostSession],
901 gateway_file: &Path,
902 token_override: Option<&str>,
903 store_root: Option<&Path>,
904) -> Result<(), ShellError> {
905 let gate = open_gateway(gateway_file, token_override)?;
906 for session in sessions.iter_mut() {
907 if session.routine_url.is_some() && session.routine_bearer.is_some() {
908 continue;
909 }
910 let Some(agent_id) = session.agent_id.clone().filter(|id| !id.is_empty()) else {
912 continue;
913 };
914 let Some(nick) = routine_name_for(session) else {
915 continue;
916 };
917 let Some(folder) = crate::nick::routine_folder_id(&nick) else {
918 continue;
919 };
920 let Some(gate) = gate.as_ref() else {
921 if let Some(root) = store_root {
922 if let Some((url, key)) = load_wake(root, &session.session_id, &folder) {
923 session.routine_url = Some(url);
924 session.routine_bearer = Some(key);
925 }
926 }
927 continue;
928 };
929 let outcome = ensure_wake(gate, &agent_id, &nick);
930 log_outcome(&nick, &folder, &outcome);
931 match outcome {
932 WakeOutcome::Ready { url, key } => {
933 if let Some(root) = store_root {
934 save_wake(root, &session.session_id, &folder, &url, &key);
935 crate::wake_policy::note_state(&crate::session_store_dir(root, &session.session_id), "ready", "gateway");
936 }
937 session.routine_url = Some(url);
938 session.routine_bearer = Some(key);
939 }
940 WakeOutcome::AwaitingBackend => {
941 if let Some(root) = store_root {
944 clear_wake(root, &session.session_id);
945 crate::wake_policy::note_state(&crate::session_store_dir(root, &session.session_id), "awaiting", "gateway");
946 }
947 }
948 WakeOutcome::Failed(_) => {
949 if let Some(root) = store_root {
950 crate::wake_policy::note_state(&crate::session_store_dir(root, &session.session_id), "failed", "gateway");
951 }
952 }
953 }
954 }
955 Ok(())
956}
957
958#[derive(serde::Deserialize)]
959#[serde(deny_unknown_fields)]
960struct RoutineFile {
961 url: String,
962 key: String,
963}
964
965pub(crate) fn read_routine_file(path: &Path) -> Result<(String, String), String> {
967 #[cfg(unix)]
968 {
969 use std::os::unix::fs::PermissionsExt;
970 let meta = std::fs::metadata(path).map_err(|_| "the routine file cannot be read".to_string())?;
971 if meta.permissions().mode() & 0o077 != 0 {
972 return Err("the routine file must be mode 0600 (owner only)".to_string());
973 }
974 }
975 let text = std::fs::read_to_string(path).map_err(|_| "the routine file cannot be read".to_string())?;
976 let file: RoutineFile = serde_json::from_str(&text).map_err(|_| "the routine file is not {\"url\",\"key\"} json".to_string())?;
977 let (url, key) = (file.url.trim().to_string(), file.key.trim().to_string());
978 if key.is_empty() || !crate::webhook::url_ok(&url) {
979 return Err("the routine file has no usable url and key".to_string());
980 }
981 Ok((url, key))
982}
983
984fn attach_routine_files(sessions: &mut [HostSession]) {
988 for session in sessions.iter_mut() {
989 let Some(path) = session.routine_file.clone() else { continue };
990 if session.routine_url.is_some() && session.routine_bearer.is_some() {
991 continue;
992 }
993 let label = routine_name_for(session).unwrap_or_else(|| session.session_id.clone());
994 match read_routine_file(&path) {
995 Ok((url, key)) => {
996 session.routine_url = Some(url);
997 session.routine_bearer = Some(key);
998 eprintln!("mail4agent: wake {label}: routine file read");
999 }
1000 Err(reason) => eprintln!("mail4agent: wake {label}: {reason}"),
1001 }
1002 }
1003}
1004
1005pub(crate) fn wake_label(session_id: &str) -> String {
1007 format!("wake/{session_id}")
1008}
1009
1010fn attach_vault_wakes(sessions: &mut [HostSession], store_root: &Path) {
1014 let pending = |s: &HostSession| !(s.routine_url.is_some() && s.routine_bearer.is_some());
1015 if !sessions.iter().any(pending) {
1016 return;
1017 }
1018 let vault = match crate::store_key::vault(store_root) {
1019 Ok(v) => v,
1020 Err(err) => {
1021 eprintln!("mail4agent: wake vault not read: {err}");
1022 return;
1023 }
1024 };
1025 for session in sessions.iter_mut().filter(|s| pending(s)) {
1026 let Ok(Some(raw)) = vault.get(&wake_label(&session.session_id)) else { continue };
1027 #[derive(serde::Deserialize)]
1028 struct Held {
1029 url: String,
1030 key: String,
1031 }
1032 match serde_json::from_slice::<Held>(&raw) {
1033 Ok(h) if !h.key.is_empty() && crate::webhook::url_ok(&h.url) => {
1034 session.routine_url = Some(h.url);
1035 session.routine_bearer = Some(h.key);
1036 crate::wake_policy::note_state(&crate::session_store_dir(store_root, &session.session_id), "ready", "vault");
1037 eprintln!("mail4agent: wake {}: taken from the vault", routine_name_for(session).unwrap_or_else(|| session.session_id.clone()));
1038 }
1039 _ => eprintln!("mail4agent: wake {}: the vault entry is unusable", session.session_id),
1040 }
1041 }
1042}
1043
1044#[cfg(test)]
1045pub(crate) fn attach_vault_wakes_for_test(sessions: &mut [HostSession], store_root: &Path) {
1046 attach_vault_wakes(sessions, store_root)
1047}
1048
1049pub const STORE_LOCK_FILE: &str = ".lock";
1054
1055fn now_ms() -> i64 {
1056 std::time::SystemTime::now()
1057 .duration_since(std::time::UNIX_EPOCH)
1058 .map(|elapsed| elapsed.as_millis() as i64)
1059 .unwrap_or(0)
1060}
1061
1062pub(crate) fn lock_store(dir: &Path) -> Result<File, ShellError> {
1063 std::fs::create_dir_all(dir)?;
1064 let file = std::fs::OpenOptions::new()
1065 .create(true)
1066 .truncate(false)
1067 .write(true)
1068 .open(dir.join(STORE_LOCK_FILE))?;
1069 match file.try_lock() {
1070 Ok(()) => Ok(file),
1071 Err(_) => Err(ShellError::SessionList(
1072 "session store is in use by another process".to_string(),
1073 )),
1074 }
1075}
1076
1077fn write_secret_file(path: &Path, bytes: &[u8]) -> std::io::Result<()> {
1079 let dir = path
1080 .parent()
1081 .ok_or_else(|| std::io::Error::other("no parent"))?;
1082 let mut builder = std::fs::DirBuilder::new();
1083 builder.recursive(true);
1084 #[cfg(unix)]
1085 {
1086 use std::os::unix::fs::DirBuilderExt;
1087 builder.mode(0o700);
1088 }
1089 builder.create(dir)?;
1090 let name = path
1091 .file_name()
1092 .and_then(|name| name.to_str())
1093 .unwrap_or("secret");
1094 let tmp = dir.join(format!("{name}.tmp"));
1095 let mut options = std::fs::OpenOptions::new();
1096 options.write(true).create(true).truncate(true);
1097 #[cfg(unix)]
1098 {
1099 use std::os::unix::fs::OpenOptionsExt;
1100 options.mode(0o600);
1101 }
1102 let mut file = options.open(&tmp)?;
1103 std::io::Write::write_all(&mut file, bytes)?;
1104 file.sync_all()?;
1105 drop(file);
1106 std::fs::rename(&tmp, path)
1107}
1108
1109fn keychain_path(store_root: &Path, session_id: &str) -> PathBuf {
1110 crate::session_store_dir(store_root, session_id).join(WAKE_KEYCHAIN_FILE)
1111}
1112
1113fn save_wake(store_root: &Path, session_id: &str, folder_id: &str, url: &str, key: &str) {
1116 let stored = StoredWake {
1117 folder_id: folder_id.to_string(),
1118 url: url.to_string(),
1119 key: key.to_string(),
1120 };
1121 let result = serde_json::to_vec(&stored)
1122 .map_err(std::io::Error::other)
1123 .and_then(|bytes| write_secret_file(&keychain_path(store_root, session_id), &bytes));
1124 if result.is_err() {
1125 eprintln!("mail4agent: wake keychain write failed for folder {folder_id}");
1126 }
1127}
1128
1129fn load_wake(store_root: &Path, session_id: &str, folder_id: &str) -> Option<(String, String)> {
1130 let text = std::fs::read_to_string(keychain_path(store_root, session_id)).ok()?;
1131 let stored: StoredWake = serde_json::from_str(&text).ok()?;
1132 if stored.folder_id != folder_id || stored.url.is_empty() || stored.key.is_empty() {
1133 return None;
1134 }
1135 Some((stored.url, stored.key))
1136}
1137
1138fn clear_wake(store_root: &Path, session_id: &str) {
1139 let _ = std::fs::remove_file(keychain_path(store_root, session_id));
1140}
1141
1142fn list_agent_automations(
1143 gate: &GatewayConn,
1144 agent_id: &str,
1145) -> Result<Vec<AutomationCard>, ShellError> {
1146 let (status, body) = gateway_post(
1147 gate,
1148 "getAgentAutomations",
1149 &serde_json::json!({ "id": agent_id }),
1150 )?;
1151 if status != 200 {
1152 return Err(ShellError::Gateway(format!("gateway list status {status}")));
1153 }
1154 serde_json::from_slice(&body)
1155 .map_err(|_| ShellError::Gateway("gateway list was not understood".to_string()))
1156}
1157
1158enum CredentialRead {
1159 Ready {
1160 url: String,
1161 key: String,
1162 },
1163 MintFailed,
1165 Missing,
1166}
1167
1168fn gateway_client() -> Result<reqwest::blocking::Client, ShellError> {
1169 reqwest::blocking::Client::builder()
1170 .timeout(std::time::Duration::from_secs(8))
1171 .redirect(reqwest::redirect::Policy::none())
1172 .http1_only()
1173 .build()
1174 .map_err(|_| ShellError::Gateway("gateway client could not start".to_string()))
1175}
1176
1177fn gateway_post(
1178 gate: &GatewayConn,
1179 method: &str,
1180 body: &serde_json::Value,
1181) -> Result<(u16, Vec<u8>), ShellError> {
1182 let url = format!("{}/api/{method}", gate.base);
1183 let authorization = crate::bearer_header(&gate.token).map_err(|_| {
1184 ShellError::Gateway("gateway token is not a single header value".to_string())
1185 })?;
1186 let bytes = serde_json::to_vec(body)
1187 .map_err(|_| ShellError::Gateway("gateway request was not json".to_string()))?;
1188 let response = gate
1189 .client
1190 .post(&url)
1191 .header(reqwest::header::AUTHORIZATION, authorization)
1192 .header(reqwest::header::CONTENT_TYPE, "application/json")
1193 .body(bytes)
1194 .send()
1195 .map_err(|_| ShellError::Gateway("gateway request failed".to_string()))?;
1196 let status = response.status().as_u16();
1197 let body = response
1198 .bytes()
1199 .map_err(|_| ShellError::Gateway("gateway response failed".to_string()))?;
1200 Ok((status, body.to_vec()))
1201}
1202
1203fn read_webhook_credential(
1204 gate: &GatewayConn,
1205 agent_id: &str,
1206 automation_id: &str,
1207) -> Result<CredentialRead, ShellError> {
1208 let (status, body) = gateway_post(
1209 gate,
1210 "getAutomationWebhookCredential",
1211 &serde_json::json!({
1212 "id": agent_id,
1213 "automationId": automation_id,
1214 }),
1215 )?;
1216 if status == 200 {
1217 let parsed: WebhookCredential = serde_json::from_slice(&body).map_err(|_| {
1218 ShellError::Gateway("gateway credential was not understood".to_string())
1219 })?;
1220 let key = parsed.key.filter(|key| !key.is_empty());
1221 if parsed.url.is_empty() {
1222 return Err(ShellError::Gateway(
1223 "gateway credential was not understood".to_string(),
1224 ));
1225 }
1226 return Ok(match key {
1227 Some(key) => CredentialRead::Ready {
1228 url: parsed.url,
1229 key,
1230 },
1231 None => CredentialRead::MintFailed,
1232 });
1233 }
1234 if credential_is_missing(status, &body) {
1235 return Ok(CredentialRead::Missing);
1236 }
1237 Err(ShellError::Gateway(format!(
1238 "gateway credential status {status}"
1239 )))
1240}
1241
1242fn credential_is_missing(status: u16, body: &[u8]) -> bool {
1245 if status == 404 {
1246 return true;
1247 }
1248 if status != 500 {
1249 return false;
1250 }
1251 std::str::from_utf8(body)
1252 .map(|text| text.contains("Automation not found"))
1253 .unwrap_or(false)
1254}
1255
1256fn create_mirror_routine(
1258 gate: &GatewayConn,
1259 agent_id: &str,
1260 routine: &str,
1261) -> Result<Vec<AutomationCard>, ShellError> {
1262 let (status, body) = gateway_post(
1263 gate,
1264 "createAgentAutomation",
1265 &serde_json::json!({
1266 "id": agent_id,
1267 "spec": {
1268 "name": routine,
1269 "prompt": MIRROR_ROUTINE_PROMPT,
1270 "trigger": { "type": "webhook" },
1271 "isEnabled": false,
1272 },
1273 }),
1274 )?;
1275 if status != 200 {
1276 return Err(ShellError::Gateway(format!(
1277 "gateway create status {status}"
1278 )));
1279 }
1280 serde_json::from_slice(&body)
1281 .map_err(|_| ShellError::Gateway("gateway create was not understood".to_string()))
1282}
1283
1284fn delete_agent_automation(
1285 gate: &GatewayConn,
1286 agent_id: &str,
1287 automation_id: &str,
1288) -> Result<(), ShellError> {
1289 let (status, _) = gateway_post(
1290 gate,
1291 "deleteAgentAutomation",
1292 &serde_json::json!({ "id": agent_id, "automationId": automation_id }),
1293 )?;
1294 if status != 200 {
1295 return Err(ShellError::Gateway(format!(
1296 "gateway delete status {status}"
1297 )));
1298 }
1299 Ok(())
1300}
1301
1302pub(crate) struct Prepared {
1303 pub(crate) config: SessionConfig,
1304 pub(crate) nick: String,
1305 pub(crate) user_id: String,
1306 pub(crate) device_id: DeviceId,
1307 pub(crate) bearer: Zeroizing<String>,
1308 pub(crate) backend: Arc<dyn m4a_agent::Backend>,
1309 pub(crate) routine_url: Option<String>,
1310 pub(crate) routine_bearer: Option<String>,
1311}
1312
1313pub struct MachineClient {
1316 product_mode: bool,
1318 sessions: Vec<OpenedStore>,
1319 bus: Arc<LocalBus>,
1320 push: m4a_agent::engine::PushLink,
1321 agents_dir: Option<PathBuf>,
1325 gateway_file: Option<PathBuf>,
1326 gateway_token: Option<String>,
1327 last_agent_poll: std::time::Instant,
1328 agent_rescan_secs: u64,
1329 sessions_dir: Option<PathBuf>,
1331 store_root: Option<PathBuf>,
1333 skip_nicks: Vec<String>,
1335 profile_note: bool,
1337 session_ids: Vec<(String, String)>,
1339 last_full_drive: std::time::Instant,
1341 pending_wakes: Vec<PendingWake>,
1343 ready_agents: HashSet<String>,
1346 _locks: Vec<File>,
1348 send_listener: Option<SendListener>,
1350 send_sock: Option<PathBuf>,
1351 send_queue: Vec<PendingSend>,
1353 homeserver_url: String,
1356 bus_next_user: i64,
1358 failed_opens: HashMap<String, std::time::Instant>,
1361}
1362
1363const LATE_OPEN_RETRY_SECS: u64 = 60;
1366
1367struct PendingSend {
1369 stream: SendStream,
1370 request: crate::SendRequest,
1371 started: std::time::Instant,
1372 room: Option<String>,
1373 peer: Option<String>,
1374}
1375
1376const SEND_JOIN_WAIT_SECS: u64 = 120;
1378
1379impl Drop for MachineClient {
1380 fn drop(&mut self) {
1381 if let Some(path) = self.send_sock.take() {
1382 let _ = std::fs::remove_file(path);
1383 }
1384 }
1385}
1386
1387#[derive(Debug, Default)]
1390pub struct TickReport {
1391 pub pushed: Vec<(String, String)>,
1393 pub joined: Vec<(String, String)>,
1395 pub errors: Vec<(String, String)>,
1397 pub sent: Vec<(String, String, crate::SendReply)>,
1399 pub alerts: Vec<(String, String)>,
1401}
1402
1403struct PendingWake {
1405 store_dir: PathBuf,
1406 agent_id: String,
1407}
1408
1409impl MachineClient {
1410 pub fn open(
1417 homeserver_url: &str,
1418 store_root: &Path,
1419 sessions: Vec<HostSession>,
1420 ) -> Result<Self, ShellError> {
1421 Self::open_with(homeserver_url, store_root, sessions, false)
1422 }
1423
1424 fn open_with(
1428 homeserver_url: &str,
1429 store_root: &Path,
1430 sessions: Vec<HostSession>,
1431 lenient: bool,
1432 ) -> Result<Self, ShellError> {
1433 let mut sessions = sessions;
1434 attach_routine_files(&mut sessions);
1435 attach_vault_wakes(&mut sessions, store_root);
1436 if sessions.is_empty() {
1437 return Err(ShellError::SessionList("session list is empty".to_string()));
1438 }
1439 let mut prepared = Vec::with_capacity(sessions.len());
1440 let mut seen_ids = Vec::new();
1441 let mut seen_nicks = Vec::new();
1442 let mut locks = Vec::new();
1443 for session in sessions {
1444 let config = session.config(homeserver_url, store_root)?;
1445 if seen_ids.iter().any(|id: &String| id == config.session_id()) {
1446 return Err(ShellError::SessionList("duplicate session id".to_string()));
1447 }
1448 seen_ids.push(config.session_id().to_string());
1449 let lock = match lock_store(&config.store_dir()) {
1450 Ok(lock) => lock,
1451 Err(err) if lenient => {
1452 eprintln!("mail4agent: session {} left out: {err}", config.session_id());
1453 continue;
1454 }
1455 Err(err) => return Err(err),
1456 };
1457 let registered = match register_session(&config) {
1458 Ok(registered) => registered,
1459 Err(err) if lenient => {
1460 eprintln!("mail4agent: session {} left out: {err}", config.session_id());
1461 continue;
1462 }
1463 Err(err) => return Err(err),
1464 };
1465 if seen_nicks.iter().any(|nick: &String| nick.eq_ignore_ascii_case(®istered.nick)) {
1466 return Err(ShellError::SessionList("duplicate nick".to_string()));
1467 }
1468 seen_nicks.push(registered.nick.clone());
1469 locks.push(lock);
1470 prepared.push(Prepared {
1471 config,
1472 nick: registered.nick,
1473 user_id: registered.user_id,
1474 device_id: registered.device_id,
1475 bearer: registered.bearer,
1476 backend: registered.backend,
1477 routine_url: session.routine_url,
1478 routine_bearer: session.routine_bearer,
1479 });
1480 }
1481 if prepared.is_empty() {
1482 return Err(ShellError::SessionList(
1483 "no session could be opened".to_string(),
1484 ));
1485 }
1486 let server_name = server_name_of(&prepared[0].user_id)?;
1487 for item in &prepared[1..] {
1488 if server_name_of(&item.user_id)? != server_name {
1489 return Err(ShellError::SessionList(
1490 "sessions disagree on the homeserver name".to_string(),
1491 ));
1492 }
1493 }
1494 ensure_server_name(server_name)?;
1495 let bus = Arc::new(LocalBus::open(&prepared)?);
1496 let peers: Vec<(String, String)> = prepared
1497 .iter()
1498 .map(|item| (item.nick.to_string(), item.user_id.clone()))
1499 .collect();
1500 let mut opened = Vec::with_capacity(prepared.len());
1501 for item in &prepared {
1502 let server_name = server_name_of(&item.user_id)?;
1503 let mut store = OpenedStore::open(
1504 &item.config.store_dir(),
1505 item.config.session_id(),
1506 item.device_id.clone(),
1507 &item.user_id,
1508 server_name,
1509 Arc::clone(&item.backend),
1510 item.bearer.as_str(),
1511 )?;
1512 store.set_registered_nick(item.nick.to_string());
1513 store.attach_bus(Arc::clone(&bus));
1514 store.set_local_peers(
1515 peers
1516 .iter()
1517 .filter(|(nick, _)| *nick != item.nick)
1518 .cloned()
1519 .collect(),
1520 );
1521 store.set_wake(SessionWake {
1522 routine_url: item.routine_url.clone(),
1523 routine_bearer: item.routine_bearer.clone(),
1524 leader_sock: None,
1525 leader_cwd: None,
1526 });
1527 if item.routine_url.is_none() {
1528 attach_detected_chain(&mut store, &item.config, &item.nick);
1529 }
1530 store.drive(1_000, false)?;
1531 store.abandon_inflight_sync(1_000)?;
1532 opened.push(store);
1533 }
1534 let tokens: Vec<String> = prepared
1535 .iter()
1536 .map(|item| item.bearer.as_str().to_string())
1537 .collect();
1538 let product_mode = true;
1539 let push = m4a_agent::engine::PushLink::open(&prepared[0].config.homeserver_url, prepared[0].backend.keep_prefix(), tokens, product_mode)?;
1540 Ok(Self {
1541 product_mode,
1542 sessions: opened,
1543 bus,
1544 push,
1545 agents_dir: None,
1546 gateway_file: None,
1547 gateway_token: None,
1548 last_agent_poll: std::time::Instant::now(),
1549 agent_rescan_secs: 0,
1550 sessions_dir: None,
1551 store_root: None,
1552 skip_nicks: Vec::new(),
1553 profile_note: false,
1554 session_ids: Vec::new(),
1555 last_full_drive: std::time::Instant::now(),
1556 pending_wakes: Vec::new(),
1557 ready_agents: HashSet::new(),
1558 _locks: locks,
1559 send_listener: None,
1560 send_sock: None,
1561 send_queue: Vec::new(),
1562 homeserver_url: homeserver_url.to_string(),
1563 bus_next_user: prepared.len() as i64 + 1,
1564 failed_opens: HashMap::new(),
1565 })
1566 }
1567
1568 pub fn open_session_dir(
1571 homeserver_url: &str,
1572 store_root: &Path,
1573 sessions_dir: &Path,
1574 ) -> Result<Self, ShellError> {
1575 let sessions = load_session_records(sessions_dir)?;
1576 Self::open(homeserver_url, store_root, sessions)
1577 }
1578
1579 pub fn from_env() -> Result<Self, ShellError> {
1590 Self::from_env_filtered(None)
1591 }
1592
1593 pub fn from_env_for(nick: &str) -> Result<Self, ShellError> {
1597 Self::from_env_filtered(Some(nick))
1598 }
1599
1600 fn from_env_filtered(only: Option<&str>) -> Result<Self, ShellError> {
1601 let homeserver_url = std::env::var(HOMESERVER_URL_ENV)
1602 .ok()
1603 .filter(|value| !value.is_empty())
1604 .ok_or(ShellError::HomeserverUrl)?;
1605 let store_root = std::env::var(STORE_ROOT_ENV)
1606 .ok()
1607 .filter(|value| !value.is_empty())
1608 .ok_or(ShellError::StoreRoot)?;
1609 let mut get = |key: &str| -> Option<String> {
1610 std::env::var(key).ok().filter(|value| !value.is_empty())
1611 };
1612 let agents_dir = explicit_agents_dir(&mut get);
1613 let mut sessions_dir_seen: Option<PathBuf> = None;
1614 let (sessions, agents_dir) = if let Some(dir) = agents_dir {
1615 (load_web_agents(&dir, &mut get)?, Some(dir))
1616 } else {
1617 let sessions_dir = get(SESSIONS_DIR_ENV)
1618 .ok_or_else(|| ShellError::SessionList("session directory is unset".to_string()))?;
1619 sessions_dir_seen = Some(PathBuf::from(&sessions_dir));
1620 (load_web_sessions(Path::new(&sessions_dir), &mut get)?, None)
1621 };
1622 let sessions = match only {
1623 Some(nick) => {
1624 let needle = crate::nick::lookup_nick(nick)?;
1625 let kept: Vec<HostSession> = sessions
1626 .into_iter()
1627 .filter(|session| {
1628 routine_name_for(session)
1629 .is_some_and(|name| name.eq_ignore_ascii_case(&needle))
1630 })
1631 .collect();
1632 if kept.is_empty() {
1633 return Err(ShellError::UnknownNick);
1634 }
1635 kept
1636 }
1637 None => sessions,
1638 };
1639 let mut sessions = sessions;
1641 attach_routine_files(&mut sessions);
1642 attach_vault_wakes(&mut sessions, Path::new(&store_root));
1643 let gateway_file = get(GATEWAY_FILE_ENV).map(PathBuf::from);
1644 let gateway_token = get(GATEWAY_TOKEN_ENV);
1645 let agent_rescan_secs = get(AGENT_RESCAN_SECS_ENV)
1646 .and_then(|value| value.parse::<u64>().ok())
1647 .unwrap_or(0);
1648 let mut pending_wakes = Vec::new();
1649 let mut ready_agents = HashSet::new();
1650 for session in &sessions {
1651 let Some(agent_id) = session.agent_id.clone() else {
1652 continue;
1653 };
1654 if session.routine_url.is_some() && session.routine_bearer.is_some() {
1655 ready_agents.insert(agent_id);
1656 } else {
1657 pending_wakes.push(PendingWake {
1658 store_dir: crate::session_store_dir(
1659 Path::new(&store_root),
1660 &session.session_id,
1661 ),
1662 agent_id,
1663 });
1664 }
1665 }
1666 let options = WakeOptions::from_lookup(&mut get);
1667 let mut client = Self::open_with(&homeserver_url, Path::new(&store_root), sessions, true)?;
1668 client.store_root = Some(PathBuf::from(&store_root));
1669 client.skip_nicks = options.skip_nicks;
1670 client.profile_note = options.profile_note;
1671 client.session_ids = options.session_ids;
1672 client.pending_wakes = pending_wakes;
1673 client.ready_agents = ready_agents;
1674 client.agents_dir = agents_dir;
1675 client.gateway_file = gateway_file;
1676 client.gateway_token = gateway_token;
1677 client.agent_rescan_secs = agent_rescan_secs;
1678 client.sessions_dir = sessions_dir_seen;
1679 client.last_agent_poll = std::time::Instant::now();
1680 Ok(client)
1681 }
1682
1683 pub fn poll_agent_directory(&mut self) -> Result<Vec<RoutineReport>, ShellError> {
1694 let Some(agents_dir) = self.agents_dir.clone() else {
1695 return self.poll_session_dir_wakes();
1696 };
1697 if self.agent_rescan_secs > 0 {
1698 let elapsed = self.last_agent_poll.elapsed().as_secs();
1699 if elapsed < self.agent_rescan_secs {
1700 return Ok(Vec::new());
1701 }
1702 }
1703 self.last_agent_poll = std::time::Instant::now();
1704 self.open_new_agent_sessions(&agents_dir);
1705 let gateway = self
1706 .gateway_file
1707 .clone()
1708 .unwrap_or_else(|| under_home(DEFAULT_GATEWAY_FILE));
1709 let Some(gate) = open_gateway(&gateway, self.gateway_token.as_deref())? else {
1710 return Ok(Vec::new());
1711 };
1712 let mut sessions = drop_skipped(load_agents_dir(&agents_dir)?, &self.skip_nicks);
1713 apply_session_ids(&mut sessions, &self.session_ids);
1714 let mut reports = Vec::new();
1715 for session in &sessions {
1716 let Some(agent_id) = session.agent_id.clone() else {
1717 continue;
1718 };
1719 let Some(nick) = routine_name_for(session) else {
1720 continue;
1721 };
1722 let folder_id = crate::nick::routine_folder_id(&nick).unwrap_or_default();
1723 if self.ready_agents.contains(&agent_id) {
1724 reports.push(RoutineReport {
1725 agent_id,
1726 nick,
1727 folder_id,
1728 status: WakeStatus::Ready,
1729 });
1730 continue;
1731 }
1732 let outcome = ensure_wake(&gate, &agent_id, &nick);
1733 log_outcome(&nick, &folder_id, &outcome);
1734 if self.profile_note && matches!(outcome, WakeOutcome::AwaitingBackend) {
1735 note_profile(&gate, &agents_dir, &agent_id, &nick);
1736 }
1737 if let WakeOutcome::Ready { url, key } = &outcome {
1738 if let Some(root) = &self.store_root {
1739 save_wake(root, &session.session_id, &folder_id, url, key);
1740 }
1741 if let Some(index) = self
1742 .pending_wakes
1743 .iter()
1744 .position(|pending| pending.agent_id == agent_id)
1745 {
1746 let pending = self.pending_wakes.remove(index);
1747 if let Some(store) = self
1748 .sessions
1749 .iter_mut()
1750 .find(|store| store.store_dir() == pending.store_dir)
1751 {
1752 store.set_wake(SessionWake {
1753 routine_url: Some(url.clone()),
1754 routine_bearer: Some(key.clone()),
1755 leader_sock: None,
1756 leader_cwd: None,
1757 });
1758 }
1759 }
1760 self.ready_agents.insert(agent_id.clone());
1761 if let Some(user_id) = self
1762 .sessions
1763 .iter()
1764 .find(|store| {
1765 store
1766 .nick()
1767 .map(|n| n.eq_ignore_ascii_case(&nick))
1768 .unwrap_or(false)
1769 })
1770 .map(|store| store.user_id().to_string())
1771 {
1772 self.announce_peer_joined(&nick, &user_id);
1773 }
1774 }
1775 reports.push(RoutineReport {
1776 agent_id,
1777 nick,
1778 folder_id,
1779 status: outcome.status(),
1780 });
1781 }
1782 Ok(reports)
1783 }
1784
1785 fn poll_session_dir_wakes(&mut self) -> Result<Vec<RoutineReport>, ShellError> {
1790 let Some(dir) = self.sessions_dir.clone() else {
1791 return Ok(Vec::new());
1792 };
1793 if self.pending_wakes.is_empty() {
1794 return Ok(Vec::new());
1795 }
1796 let every = if self.agent_rescan_secs > 0 { self.agent_rescan_secs } else { 60 };
1797 if self.last_agent_poll.elapsed().as_secs() < every {
1798 return Ok(Vec::new());
1799 }
1800 self.last_agent_poll = std::time::Instant::now();
1801 let Some(gateway) = self.gateway_file.clone() else {
1802 return Ok(Vec::new());
1803 };
1804 let Some(gate) = open_gateway(&gateway, self.gateway_token.as_deref())? else {
1805 return Ok(Vec::new());
1806 };
1807 let mut sessions = load_session_records(&dir)?;
1808 apply_session_ids(&mut sessions, &self.session_ids);
1809 let mut reports = Vec::new();
1810 for session in &sessions {
1811 let Some(agent_id) = session.agent_id.clone() else { continue };
1812 if !self.pending_wakes.iter().any(|p| p.agent_id == agent_id) {
1813 continue;
1814 }
1815 let Some(nick) = routine_name_for(session) else { continue };
1816 let folder_id = crate::nick::routine_folder_id(&nick).unwrap_or_default();
1817 let outcome = ensure_wake(&gate, &agent_id, &nick);
1818 log_outcome(&nick, &folder_id, &outcome);
1819 let dir_of = self.store_root.as_ref().map(|r| crate::session_store_dir(r, &session.session_id));
1820 match &outcome {
1821 WakeOutcome::Ready { url, key } => {
1822 if let Some(root) = &self.store_root {
1823 save_wake(root, &session.session_id, &folder_id, url, key);
1824 }
1825 if let Some(d) = &dir_of {
1826 crate::wake_policy::note_state(d, "ready", "gateway");
1827 }
1828 if let Some(index) = self.pending_wakes.iter().position(|p| p.agent_id == agent_id) {
1829 let pending = self.pending_wakes.remove(index);
1830 if let Some(store) = self.sessions.iter_mut().find(|s| s.store_dir() == pending.store_dir) {
1831 store.set_wake(SessionWake { routine_url: Some(url.clone()), routine_bearer: Some(key.clone()), leader_sock: None, leader_cwd: None });
1832 }
1833 }
1834 self.ready_agents.insert(agent_id.clone());
1835 }
1836 WakeOutcome::AwaitingBackend => {
1837 if let Some(d) = &dir_of {
1838 crate::wake_policy::note_state(d, "awaiting", "gateway");
1839 }
1840 }
1841 WakeOutcome::Failed(_) => {
1842 if let Some(d) = &dir_of {
1843 crate::wake_policy::note_state(d, "failed", "gateway");
1844 }
1845 }
1846 }
1847 reports.push(RoutineReport { agent_id, nick, folder_id, status: outcome.status() });
1848 }
1849 Ok(reports)
1850 }
1851
1852 fn open_new_agent_sessions(&mut self, agents_dir: &Path) {
1858 let Some(store_root) = self.store_root.clone() else {
1859 return;
1860 };
1861 let found = match load_agents_dir(agents_dir) {
1862 Ok(found) => found,
1863 Err(err) => {
1864 eprintln!("mail4agent: agent rescan: {err}");
1865 return;
1866 }
1867 };
1868 let held: Vec<(PathBuf, Option<String>)> = self
1869 .sessions
1870 .iter()
1871 .map(|store| {
1872 (
1873 store.store_dir().to_path_buf(),
1874 store.nick().map(str::to_string),
1875 )
1876 })
1877 .collect();
1878 let fresh = unopened_agent_sessions(
1879 found,
1880 &self.skip_nicks,
1881 &self.session_ids,
1882 &store_root,
1883 &held,
1884 );
1885 for (session, nick) in fresh {
1886 if let Some(failed_at) = self.failed_opens.get(&session.session_id) {
1887 if failed_at.elapsed().as_secs() < LATE_OPEN_RETRY_SECS {
1888 continue;
1889 }
1890 }
1891 let session_id = session.session_id.clone();
1892 match self.open_late_session(&store_root, session, &nick) {
1893 Ok(()) => {
1894 self.failed_opens.remove(&session_id);
1895 println!("mail4agent: session {nick} registered and open");
1896 }
1897 Err(err) => {
1898 eprintln!("mail4agent: session {nick} not opened: {err}");
1899 self.failed_opens
1900 .insert(session_id, std::time::Instant::now());
1901 }
1902 }
1903 }
1904 }
1905
1906 fn open_late_session(
1907 &mut self,
1908 store_root: &Path,
1909 mut session: HostSession,
1910 nick: &str,
1911 ) -> Result<(), ShellError> {
1912 if session.routine_url.is_none() || session.routine_bearer.is_none() {
1913 if let Some(folder) = crate::nick::routine_folder_id(nick) {
1914 if let Some((url, key)) = load_wake(store_root, &session.session_id, &folder) {
1915 session.routine_url = Some(url);
1916 session.routine_bearer = Some(key);
1917 }
1918 }
1919 }
1920 let config = session.config(&self.homeserver_url, store_root)?;
1921 let lock = lock_store(&config.store_dir())?;
1922 let registered = register_session(&config)?;
1923 if let Some(first) = self.sessions.first() {
1924 if server_name_of(®istered.user_id)? != server_name_of(first.user_id())? {
1925 return Err(ShellError::SessionList(
1926 "sessions disagree on the homeserver name".to_string(),
1927 ));
1928 }
1929 }
1930 let item = Prepared {
1931 config,
1932 nick: registered.nick,
1933 user_id: registered.user_id,
1934 device_id: registered.device_id,
1935 bearer: registered.bearer,
1936 backend: registered.backend,
1937 routine_url: session.routine_url.clone(),
1938 routine_bearer: session.routine_bearer.clone(),
1939 };
1940 self.bus.seed(self.bus_next_user, &item)?;
1941 self.bus_next_user += 1;
1942 let mut store = OpenedStore::open(
1943 &item.config.store_dir(),
1944 item.config.session_id(),
1945 item.device_id.clone(),
1946 &item.user_id,
1947 server_name_of(&item.user_id)?,
1948 Arc::clone(&item.backend),
1949 item.bearer.as_str(),
1950 )?;
1951 store.set_registered_nick(item.nick.to_string());
1952 store.attach_bus(Arc::clone(&self.bus));
1953 store.set_wake(SessionWake {
1954 routine_url: item.routine_url.clone(),
1955 routine_bearer: item.routine_bearer.clone(),
1956 leader_sock: None,
1957 leader_cwd: None,
1958 });
1959 if item.routine_url.is_none() {
1960 attach_detected_chain(&mut store, &item.config, &item.nick);
1961 }
1962 store.drive(1_000, false)?;
1963 store.abandon_inflight_sync(1_000)?;
1964 self.sessions.push(store);
1965 self._locks.push(lock);
1966 self.refresh_local_peers();
1967 let tokens: Vec<String> = self
1968 .sessions
1969 .iter()
1970 .map(|store| store.device_bearer().to_string())
1971 .collect();
1972 match m4a_agent::engine::PushLink::open(&self.homeserver_url, self.sessions.first().is_some_and(|s| s.keep_prefix()), tokens, self.product_mode) {
1973 Ok(push) => {
1974 let old = std::mem::replace(&mut self.push, push);
1975 for (recipient, event) in old.drain() {
1976 if let Some(store) = self
1977 .sessions
1978 .iter_mut()
1979 .find(|store| store.user_id() == recipient)
1980 {
1981 store.record_push(event);
1982 }
1983 }
1984 }
1985 Err(err) => {
1986 eprintln!("mail4agent: push socket not reopened for {nick}: {err}");
1987 }
1988 }
1989 if let Some(agent_id) = session.agent_id.clone() {
1990 if item.routine_url.is_some() && item.routine_bearer.is_some() {
1991 self.ready_agents.insert(agent_id);
1992 self.announce_peer_joined(nick, &item.user_id);
1993 } else {
1994 self.ready_agents.remove(&agent_id);
1996 if !self.pending_wakes.iter().any(|p| p.agent_id == agent_id) {
1997 self.pending_wakes.push(PendingWake {
1998 store_dir: item.config.store_dir(),
1999 agent_id,
2000 });
2001 }
2002 }
2003 }
2004 Ok(())
2005 }
2006
2007 fn refresh_local_peers(&mut self) {
2009 let peers: Vec<(String, String)> = self
2010 .sessions
2011 .iter()
2012 .filter_map(|store| Some((store.nick()?.to_string(), store.user_id().to_string())))
2013 .collect();
2014 for store in self.sessions.iter_mut() {
2015 let own = store.nick().unwrap_or("").to_string();
2016 store.set_local_peers(
2017 peers
2018 .iter()
2019 .filter(|(nick, _)| *nick != own)
2020 .cloned()
2021 .collect(),
2022 );
2023 }
2024 }
2025
2026 fn announce_peer_joined(&self, joined_nick: &str, joined_user_id: &str) {
2030 let body = serde_json::json!({
2031 "kind": "peer_joined",
2032 "nick": joined_nick,
2033 "user_id": joined_user_id,
2034 });
2035 for store in &self.sessions {
2036 let Some(nick) = store.nick() else {
2037 continue;
2038 };
2039 if nick.eq_ignore_ascii_case(joined_nick) {
2040 continue;
2041 }
2042 if !store.has_routine() {
2043 continue;
2044 }
2045 let Some((url, bearer)) = store.routine_target() else {
2046 continue;
2047 };
2048 match crate::post_routine_json(&url, &body, bearer.as_deref()) {
2049 Ok(()) => {
2050 eprintln!("mail4agent: peer_joined {joined_nick} -> {nick} status=200")
2051 }
2052 Err(err) => {
2053 eprintln!("mail4agent: peer_joined {joined_nick} -> {nick}: {err}")
2054 }
2055 }
2056 }
2057 }
2058
2059 pub fn pending_wake_agents(&self) -> Vec<String> {
2061 self.pending_wakes
2062 .iter()
2063 .map(|pending| pending.agent_id.clone())
2064 .collect()
2065 }
2066
2067 pub fn set_local_delivery(&self, enabled: bool) {
2078 self.bus.set_local_only(enabled);
2079 }
2080
2081 pub fn homeserver_hits(&self) -> u64 {
2085 self.bus.hits()
2086 }
2087
2088 pub fn session_mut(&mut self, name_or_nick: &str) -> Result<&mut OpenedStore, ShellError> {
2090 let needle = crate::nick::lookup_nick(name_or_nick)?;
2091 self.sessions
2092 .iter_mut()
2093 .find(|store| {
2094 store
2095 .nick()
2096 .is_some_and(|nick| nick.eq_ignore_ascii_case(&needle))
2097 })
2098 .ok_or(ShellError::UnknownNick)
2099 }
2100
2101 pub fn store_dir(&self, name_or_nick: &str) -> Result<PathBuf, ShellError> {
2103 let needle = crate::nick::lookup_nick(name_or_nick)?;
2104 self.sessions
2105 .iter()
2106 .find(|store| {
2107 store
2108 .nick()
2109 .is_some_and(|nick| nick.eq_ignore_ascii_case(&needle))
2110 })
2111 .map(|store| store.store_dir().to_path_buf())
2112 .ok_or(ShellError::UnknownNick)
2113 }
2114
2115 pub fn deliver_pushed(&mut self) -> usize {
2119 let batch = self.push.drain();
2120 let mut delivered = 0;
2121 for (recipient, event) in batch {
2122 let Some(session) = self
2123 .sessions
2124 .iter_mut()
2125 .find(|store| store.user_id() == recipient)
2126 else {
2127 continue;
2128 };
2129 session.record_push(event);
2130 delivered += 1;
2131 }
2132 delivered
2133 }
2134
2135 pub fn tick(&mut self, now_ms: i64, full_drive_secs: u64) -> TickReport {
2145 let mut report = TickReport::default();
2146 let mut due: Vec<usize> = Vec::new();
2147 for (recipient, event) in self.push.drain() {
2148 let Some(index) = self
2149 .sessions
2150 .iter()
2151 .position(|store| store.user_id() == recipient)
2152 else {
2153 continue;
2154 };
2155 let nick = self.sessions[index].nick().unwrap_or("").to_string();
2156 report.pushed.push((nick, event.event_id.clone()));
2157 self.sessions[index].record_push(event);
2158 if !due.contains(&index) {
2159 due.push(index);
2160 }
2161 }
2162 let full = self.last_full_drive.elapsed().as_secs() >= full_drive_secs;
2163 if full {
2164 self.last_full_drive = std::time::Instant::now();
2165 }
2166 for index in 0..self.sessions.len() {
2167 let pushed = due.contains(&index);
2168 if !pushed && !full {
2169 continue;
2170 }
2171 let store = &mut self.sessions[index];
2172 let nick = store.nick().unwrap_or("").to_string();
2173 let drove = store.drive(now_ms, pushed);
2174 for text in store.take_security_alerts() {
2175 report.alerts.push((nick.clone(), text));
2176 }
2177 if let Err(err) = drove {
2178 report.errors.push((nick.clone(), err.to_string()));
2179 continue;
2180 }
2181 match store.accept_direct_invites(now_ms) {
2182 Ok(joined) => {
2183 for room in joined {
2184 report.joined.push((nick.clone(), room));
2185 }
2186 }
2187 Err(err) => report.errors.push((nick, err.to_string())),
2188 }
2189 }
2190 self.serve_sends(now_ms, &mut report);
2191 report
2192 }
2193
2194 pub fn listen_for_sends(&mut self, path: &Path) -> Result<(), ShellError> {
2202 let listener = SendListener::bind(path).map_err(|err| {
2203 if err.kind() == std::io::ErrorKind::AlreadyExists {
2204 ShellError::SessionList(
2205 "another client already listens on the send socket".to_string(),
2206 )
2207 } else {
2208 ShellError::Io(err)
2209 }
2210 })?;
2211 listener.set_nonblocking(true)?;
2212 self.send_listener = Some(listener);
2213 self.send_sock = Some(path.to_path_buf());
2214 Ok(())
2215 }
2216
2217 pub fn listen_for_sends_from_env(&mut self) -> Result<PathBuf, ShellError> {
2220 let root = self.store_root.clone().ok_or(ShellError::StoreRoot)?;
2221 let path = crate::send_sock_path(
2222 |key| std::env::var(key).ok().filter(|value| !value.is_empty()),
2223 &root,
2224 );
2225 self.listen_for_sends(&path)?;
2226 Ok(path)
2227 }
2228
2229 pub fn send_blocking(
2233 &mut self,
2234 as_nick: &str,
2235 to: &str,
2236 text: &str,
2237 wait: std::time::Duration,
2238 ) -> crate::SendReply {
2239 let started = std::time::Instant::now();
2240 let mut room = None;
2241 let mut peer = None;
2242 loop {
2243 let now = now_ms();
2244 match self.try_send(as_nick, to, text, now, &mut room, &mut peer) {
2245 Some(reply) => return reply,
2246 None if started.elapsed() >= wait => {
2247 return crate::SendReply {
2248 room,
2249 ..crate::SendReply::failed(format!("{to} has not joined the DM yet"))
2250 }
2251 }
2252 None => {
2253 if let Ok(store) = self.session_mut(as_nick) {
2254 let _ = store.drive(now, false);
2255 }
2256 std::thread::sleep(std::time::Duration::from_millis(500));
2257 }
2258 }
2259 }
2260 }
2261
2262 fn try_send(
2264 &mut self,
2265 as_nick: &str,
2266 to: &str,
2267 text: &str,
2268 now_ms: i64,
2269 room: &mut Option<String>,
2270 peer: &mut Option<String>,
2271 ) -> Option<crate::SendReply> {
2272 let store = match self.session_mut(as_nick) {
2273 Ok(store) => store,
2274 Err(_) => {
2275 return Some(crate::SendReply::failed(format!(
2276 "{as_nick} is not a session on this client"
2277 )))
2278 }
2279 };
2280 if peer.is_none() {
2281 let mut last_err = None;
2284 for attempt in 0..3 {
2285 if attempt > 0 {
2286 let _ = store.drive(now_ms, false);
2287 }
2288 match store.find_nick(to, now_ms) {
2289 Ok(found) => {
2290 *peer = Some(found.user_id);
2291 last_err = None;
2292 break;
2293 }
2294 Err(err) => last_err = Some(err),
2295 }
2296 }
2297 if let Some(err) = last_err {
2298 return Some(crate::SendReply::failed(format!("find {to}: {err}")));
2299 }
2300 }
2301 if room.is_none() {
2302 match store.ensure_dm(to, now_ms) {
2303 Ok(room_id) => *room = Some(room_id),
2304 Err(err) => return Some(crate::SendReply::failed(format!("open DM: {err}"))),
2305 }
2306 }
2307 let (room_id, peer_id) = (room.clone()?, peer.clone()?);
2308 if !store.member_joined(&room_id, &peer_id) {
2309 return None;
2310 }
2311 match store.write_to_nick(to, text, now_ms) {
2312 Ok(room_id) => {
2313 let event_id = store
2314 .texts()
2315 .into_iter()
2316 .rev()
2317 .find(|row| row.room_id == room_id && row.body == text)
2318 .and_then(|row| row.event_id);
2319 Some(crate::SendReply {
2320 ok: true,
2321 room: Some(room_id),
2322 event_id,
2323 error: None,
2324 })
2325 }
2326 Err(err) => Some(crate::SendReply {
2327 room: Some(room_id),
2328 ..crate::SendReply::failed(format!("send: {err}"))
2329 }),
2330 }
2331 }
2332
2333 fn serve_sends(&mut self, now_ms: i64, report: &mut TickReport) {
2334 let mut cmds: Vec<(crate::ipc::SendStream, crate::CmdRequest)> = Vec::new();
2335 if let Some(listener) = &self.send_listener {
2336 loop {
2337 match listener.accept() {
2338 Ok(mut stream) => match crate::send::read_incoming(&mut stream) {
2339 Ok(crate::send::Incoming::Send(request)) => self.send_queue.push(PendingSend {
2340 stream,
2341 request,
2342 started: std::time::Instant::now(),
2343 room: None,
2344 peer: None,
2345 }),
2346 Ok(crate::send::Incoming::Cmd(cmd)) => cmds.push((stream, cmd)),
2347 Err(err) => {
2348 crate::send::write_reply(&mut stream, &crate::SendReply::failed(err))
2349 }
2350 },
2351 Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => break,
2352 Err(_) => break,
2353 }
2354 }
2355 }
2356 for (mut stream, cmd) in cmds {
2357 let reply = match self.sessions.iter_mut().find(|store| {
2358 store.nick().is_some_and(|n| n.eq_ignore_ascii_case(cmd.as_nick.trim()))
2359 }) {
2360 Some(store) => store.run_command(&cmd, now_ms),
2361 None => crate::CmdReply::failed(format!("{} is not a session on this client", cmd.as_nick)),
2362 };
2363 crate::send::write_cmd_reply(&mut stream, &reply);
2364 }
2365 let queue = std::mem::take(&mut self.send_queue);
2366 for mut pending in queue {
2367 let (as_nick, to, text) = (
2368 pending.request.as_nick.clone(),
2369 pending.request.to.clone(),
2370 pending.request.text.clone(),
2371 );
2372 let outcome = self.try_send(
2373 &as_nick,
2374 &to,
2375 &text,
2376 now_ms,
2377 &mut pending.room,
2378 &mut pending.peer,
2379 );
2380 let reply = match outcome {
2381 Some(reply) => reply,
2382 None if pending.started.elapsed().as_secs() >= SEND_JOIN_WAIT_SECS => {
2383 crate::SendReply {
2384 room: pending.room.clone(),
2385 ..crate::SendReply::failed(format!("{to} has not joined the DM yet"))
2386 }
2387 }
2388 None => {
2389 self.send_queue.push(pending);
2390 continue;
2391 }
2392 };
2393 crate::send::write_reply(&mut pending.stream, &reply);
2394 report.sent.push((as_nick, to, reply));
2395 }
2396 }
2397
2398 pub fn wake_log(&self) -> Vec<(String, crate::WakeAttempt)> {
2401 self.sessions
2402 .iter()
2403 .flat_map(|store| {
2404 let nick = store.nick().unwrap_or("").to_string();
2405 store
2406 .wake_log()
2407 .iter()
2408 .cloned()
2409 .map(move |attempt| (nick.clone(), attempt))
2410 })
2411 .collect()
2412 }
2413
2414 pub fn holds(&self, name_or_nick: &str) -> bool {
2416 let Ok(needle) = crate::nick::lookup_nick(name_or_nick) else {
2417 return false;
2418 };
2419 self.sessions.iter().any(|store| {
2420 store
2421 .nick()
2422 .is_some_and(|nick| nick.eq_ignore_ascii_case(&needle))
2423 })
2424 }
2425}
2426
2427fn server_name_of(mxid: &str) -> Result<&str, ShellError> {
2428 mxid.split_once(':')
2429 .map(|(_, server)| server)
2430 .filter(|server| !server.is_empty())
2431 .ok_or_else(|| ShellError::Register("user id has no server".to_string()))
2432}
2433
2434
2435fn attach_detected_chain(store: &mut OpenedStore, config: &SessionConfig, nick: &str) {
2441 if let Some((session, chain)) = crate::provider::chain::detected_web_chain(
2442 config.session_id(),
2443 nick,
2444 &config.store_root,
2445 ) {
2446 store.set_wake_chain(session, chain);
2447 }
2448}
2449
2450#[cfg(test)]
2451mod tests {
2452 #[cfg(unix)]
2453 #[test]
2454 fn a_routine_file_gives_a_session_without_an_agent_id_its_wake_and_leaks_nothing() {
2455 use std::os::unix::fs::PermissionsExt;
2456 let dir = tempfile::tempdir().expect("dir");
2457 let file = dir.path().join("wake.json");
2458 let write = |text: &str, mode: u32| {
2459 std::fs::write(&file, text).expect("write");
2460 std::fs::set_permissions(&file, std::fs::Permissions::from_mode(mode)).expect("mode");
2461 };
2462 let mut set = vec![HostSession::new("Courier", "web-courier").with_routine_file(&file)];
2463
2464 write(r#"{"url":"http://127.0.0.1:9/hook","key":"k-secret"}"#, 0o600);
2466 attach_routine_files(&mut set);
2467 assert_eq!(set[0].routine_url.as_deref(), Some("http://127.0.0.1:9/hook"));
2468 assert_eq!(set[0].routine_bearer.as_deref(), Some("k-secret"));
2469 assert!(!format!("{:?}", set[0]).contains("k-secret"), "Debug must not show the key");
2470
2471 for (text, mode) in [(r#"{"url":"http://127.0.0.1:9/hook","key":"k-secret"}"#, 0o644), (r#"{"url":"http://127.0.0.1:9/hook","key":"k-secret","extra":1}"#, 0o600), (r#"{"url":"http://127.0.0.1:9/hook","key":" "}"#, 0o600), (r#"{"url":"not a url","key":"k-secret"}"#, 0o600)] {
2473 write(text, mode);
2474 let err = read_routine_file(&file).expect_err("refused");
2475 assert!(!err.contains("k-secret") && !err.contains("127.0.0.1"), "{err}");
2476 let mut one = vec![HostSession::new("Courier", "web-courier").with_routine_file(&file)];
2477 attach_routine_files(&mut one);
2478 assert!(one[0].routine_url.is_none() && one[0].routine_bearer.is_none());
2479 }
2480 write(r#"{"url":"http://127.0.0.1:9/other","key":"k2"}"#, 0o600);
2482 let mut held = vec![HostSession::new("Courier", "web-courier").with_routine_file(&file).with_routine("http://127.0.0.1:9/mine", Some("mine".into()))];
2483 attach_routine_files(&mut held);
2484 assert_eq!(held[0].routine_url.as_deref(), Some("http://127.0.0.1:9/mine"));
2485 }
2486
2487 #[test]
2488 fn a_session_record_may_name_a_routine_file_but_never_carry_the_secret() {
2489 let dir = tempfile::tempdir().expect("dir");
2490 std::fs::write(dir.path().join("a.json"), r#"{"bot_name":"Courier","session_id":"web-courier","routine_file":"/some/where/wake.json"}"#).expect("write");
2491 let loaded = load_session_records(dir.path()).expect("records");
2492 assert_eq!(loaded[0].routine_file.as_deref(), Some(Path::new("/some/where/wake.json")));
2493 assert!(loaded[0].routine_url.is_none());
2494 std::fs::write(dir.path().join("a.json"), r#"{"bot_name":"Courier","session_id":"web-courier","routine_url":"http://127.0.0.1:9/h"}"#).expect("write");
2495 assert!(load_session_records(dir.path()).is_err(), "a URL in the record is still refused");
2496 }
2497
2498 use super::*;
2499
2500 #[test]
2501 fn web_client_agents_dir_is_only_the_env_path() {
2502 assert!(explicit_agents_dir(|_| None).is_none());
2503 let missing = std::env::temp_dir().join(format!(
2504 "m4a-missing-agents-{}",
2505 std::process::id()
2506 ));
2507 let _ = std::fs::remove_dir_all(&missing);
2508 assert!(explicit_agents_dir(|key| {
2509 (key == AGENTS_DIR_ENV).then(|| missing.display().to_string())
2510 })
2511 .is_none());
2512 let present = std::env::temp_dir().join(format!(
2513 "m4a-present-agents-{}-{}",
2514 std::process::id(),
2515 std::time::SystemTime::now()
2516 .duration_since(std::time::UNIX_EPOCH)
2517 .expect("clock")
2518 .as_nanos()
2519 ));
2520 std::fs::create_dir_all(&present).expect("dir");
2521 let got = explicit_agents_dir(|key| {
2522 (key == AGENTS_DIR_ENV).then(|| present.display().to_string())
2523 });
2524 assert_eq!(got.as_deref(), Some(present.as_path()));
2525 let _ = std::fs::remove_dir_all(&present);
2526 }
2527
2528 #[test]
2529 fn rescan_opens_only_bots_without_a_session() {
2530 let root = Path::new("/tmp/m4a-rescan-test");
2531 let agent = |name: &str, id: &str| {
2532 let mut session = HostSession::new(name, id);
2533 session.agent_id = Some(id.to_string());
2534 session
2535 };
2536 let found = vec![
2537 agent("alice", "a-hatch"),
2538 agent("m4a-proba2", "a-proba2"),
2539 agent("skipme", "a-skip"),
2540 agent("aliased", "a-alias"),
2541 ];
2542 let held = vec![
2543 (crate::session_store_dir(root, "a-hatch"), Some("alice".to_string())),
2544 (crate::session_store_dir(root, "old-alias"), None),
2545 ];
2546 let fresh = unopened_agent_sessions(
2547 found,
2548 &["skipme".to_string()],
2549 &[("a-alias".to_string(), "old-alias".to_string())],
2550 root,
2551 &held,
2552 );
2553 let nicks: Vec<&str> = fresh.iter().map(|(_, nick)| nick.as_str()).collect();
2554 assert_eq!(nicks, vec!["m4a-proba2"]);
2555 assert_eq!(fresh[0].0.session_id, "a-proba2");
2556 }
2557
2558 #[test]
2559 fn session_directory_is_name_and_id_only() {
2560 let dir = std::env::temp_dir().join(format!(
2561 "m4a-sessions-{}-{}",
2562 std::process::id(),
2563 std::time::SystemTime::now()
2564 .duration_since(std::time::UNIX_EPOCH)
2565 .expect("clock")
2566 .as_nanos()
2567 ));
2568 std::fs::create_dir_all(&dir).expect("dir");
2569 std::fs::write(
2570 dir.join("alice.json"),
2571 r#"{"bot_name":"Alice","session_id":"web-alice"}"#,
2572 )
2573 .expect("write");
2574 std::fs::write(
2575 dir.join("chief.json"),
2576 "{\"bot_name\":\"Привет мир\",\"session_id\":\"web-chief\"}",
2577 )
2578 .expect("write");
2579 let loaded = load_session_records(&dir).expect("records");
2580 assert_eq!(loaded.len(), 2);
2581 assert_eq!(loaded[1].bot_name, "Привет мир");
2582 assert_eq!(loaded[0].bot_name, "Alice");
2583 assert!(loaded[1].agent_id.is_none());
2584 assert!(loaded[0].agent_id.is_none());
2585 assert!(loaded[1].routine_url.is_none());
2586 assert!(loaded[1].routine_bearer.is_none());
2587 assert!(loaded[1].invite.is_none());
2588
2589 std::fs::write(
2590 dir.join("leaked.json"),
2591 r#"{"bot_name":"Courier","session_id":"web-courier","routine_url":"http://127.0.0.1/hook","routine_bearer":"not-a-file"}"#,
2592 )
2593 .expect("write");
2594 let refused = load_session_records(&dir).expect_err("webhook file");
2595 let text = refused.to_string();
2596 assert!(!text.contains("not-a-file"));
2597 assert!(!text.contains("127.0.0.1/hook"));
2598 let _ = std::fs::remove_dir_all(&dir);
2599 }
2600
2601 #[test]
2602 fn web_client_does_not_copy_one_routine_onto_every_session_and_node_cli_refuses_it() {
2603 let dir = std::env::temp_dir().join(format!(
2604 "m4a-split-{}-{}",
2605 std::process::id(),
2606 std::time::SystemTime::now()
2607 .duration_since(std::time::UNIX_EPOCH)
2608 .expect("clock")
2609 .as_nanos()
2610 ));
2611 std::fs::create_dir_all(&dir).expect("dir");
2612 let record = dir.join("alice.json");
2613 let body = r#"{"bot_name":"Alice","session_id":"web-alice"}"#;
2614 std::fs::write(&record, body).expect("write");
2615 let routine = "http://127.0.0.1:9/routine";
2616 let bearer = "host-injected-bearer";
2617 let sessions = load_web_sessions(&dir, |key| match key {
2618 crate::ROUTINE_URL_ENV => Some(routine.to_string()),
2619 crate::ROUTINE_BEARER_ENV => Some(bearer.to_string()),
2620 crate::LEADER_SOCK_ENV => Some("/tmp/leader.sock".to_string()),
2621 _ => None,
2622 })
2623 .expect("web sessions");
2624 assert_eq!(sessions.len(), 1);
2625 assert!(sessions[0].routine_url.is_none());
2626 assert!(sessions[0].routine_bearer.is_none());
2627 assert_eq!(std::fs::read_to_string(&record).expect("reread"), body);
2628 let wake = crate::SessionWake::web_from_lookup(|key| match key {
2629 crate::ROUTINE_URL_ENV => Some(routine.to_string()),
2630 crate::LEADER_SOCK_ENV => Some("/tmp/leader.sock".to_string()),
2631 _ => None,
2632 });
2633 assert!(wake.leader_sock.is_none());
2634 assert!(wake.leader_cwd.is_none());
2635
2636 std::fs::write(
2637 dir.join("leaked.json"),
2638 r#"{"bot_name":"Courier","session_id":"web-courier","routine_url":"http://127.0.0.1:9/from-file","routine_bearer":"file-bearer"}"#,
2639 )
2640 .expect("leak");
2641 let refused = load_web_sessions(&dir, |_| None).expect_err("json routine");
2642 let text = refused.to_string();
2643 assert!(!text.contains("from-file"));
2644 assert!(!text.contains("file-bearer"));
2645
2646 let node = crate::OpenedStore::connect_node_from_lookup(
2647 |key| match key {
2648 crate::ROUTINE_URL_ENV => Some(routine.to_string()),
2649 crate::ROUTINE_BEARER_ENV => Some(bearer.to_string()),
2650 crate::LEADER_SOCK_ENV => Some("/tmp/leader.sock".to_string()),
2651 _ => None,
2652 },
2653 None,
2654 );
2655 let err = match node {
2656 Ok(_) => panic!("node cli accepted a routine url"),
2657 Err(err) => err,
2658 };
2659 assert!(matches!(err, crate::ShellError::NodeRoutine));
2660 let text = err.to_string();
2661 assert!(!text.contains(routine));
2662 assert!(!text.contains(bearer));
2663 assert!(!text.contains("leader.sock"));
2664
2665 let node = crate::SessionWake::node_from_lookup(|key| match key {
2666 crate::LEADER_SOCK_ENV => Some("/tmp/node-leader.sock".to_string()),
2667 crate::LEADER_CWD_ENV => Some("/tmp/node".to_string()),
2668 _ => None,
2669 })
2670 .expect("node leader");
2671 assert!(node.routine_url.is_none());
2672 assert!(node.routine_bearer.is_none());
2673 assert_eq!(
2674 node.leader_sock.as_deref(),
2675 Some(std::path::Path::new("/tmp/node-leader.sock"))
2676 );
2677 assert_eq!(node.leader_cwd.as_deref(), Some("/tmp/node"));
2678 let _ = std::fs::remove_dir_all(&dir);
2679 }
2680
2681 #[test]
2682 fn mirrors_follow_the_nick_folder_and_keys_wait_for_the_bot_routine() {
2683 use std::collections::HashSet;
2684 use std::io::{Read, Write};
2685 use std::net::TcpListener;
2686 use std::sync::atomic::{AtomicBool, Ordering};
2687 use std::sync::{Arc, Mutex};
2688 use std::thread;
2689
2690 struct Card {
2691 agent: String,
2692 id: String,
2693 name: String,
2694 trigger: &'static str,
2695 enabled: bool,
2696 }
2697 #[derive(Default)]
2698 struct Gateway {
2699 creates: usize,
2700 deletes: usize,
2701 cards: Vec<Card>,
2702 backend: HashSet<(String, String)>,
2704 clash: HashSet<String>,
2706 agents_seen: HashSet<String>,
2707 profiles: Vec<(String, serde_json::Value)>,
2708 }
2709
2710 fn read_http(sock: &mut std::net::TcpStream) -> Option<(String, Vec<u8>)> {
2711 let _ = sock.set_read_timeout(Some(std::time::Duration::from_secs(2)));
2712 let mut buf = Vec::new();
2713 let mut tmp = [0u8; 2048];
2714 loop {
2715 let n = sock.read(&mut tmp).unwrap_or(0);
2716 if n == 0 {
2717 return None;
2718 }
2719 buf.extend_from_slice(&tmp[..n]);
2720 let Some(end) = buf.windows(4).position(|w| w == b"\r\n\r\n") else {
2721 continue;
2722 };
2723 let headers = String::from_utf8_lossy(&buf[..end]).to_string();
2724 let length = headers
2725 .lines()
2726 .find_map(|line| {
2727 let (name, value) = line.split_once(':')?;
2728 name.eq_ignore_ascii_case("content-length")
2729 .then(|| value.trim().parse::<usize>().ok())
2730 .flatten()
2731 })
2732 .unwrap_or(0);
2733 if buf.len() >= end + 4 + length {
2734 return Some((headers, buf[end + 4..end + 4 + length].to_vec()));
2735 }
2736 }
2737 }
2738 fn reply(status: &str, body: &serde_json::Value) -> Vec<u8> {
2739 let payload = serde_json::to_vec(body).expect("json");
2740 let mut out = format!(
2741 "HTTP/1.1 {status}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
2742 payload.len()
2743 )
2744 .into_bytes();
2745 out.extend(payload);
2746 out
2747 }
2748 fn cards_of(gate: &Gateway, agent: &str) -> serde_json::Value {
2749 serde_json::Value::Array(
2750 gate.cards
2751 .iter()
2752 .filter(|card| card.agent == agent)
2753 .map(|card| {
2754 serde_json::json!({
2755 "id": card.id,
2756 "name": card.name,
2757 "trigger": {"type": card.trigger},
2758 "isEnabled": card.enabled,
2759 })
2760 })
2761 .collect(),
2762 )
2763 }
2764
2765 let listener = TcpListener::bind("127.0.0.1:0").expect("bind");
2766 listener.set_nonblocking(true).expect("nonblocking");
2767 let port = listener.local_addr().expect("addr").port();
2768 let state = Arc::new(Mutex::new(Gateway::default()));
2769 {
2770 let mut gate = state.lock().expect("gate");
2771 gate.backend.insert(("agent-h".into(), "alice".into()));
2772 gate.clash.insert("agent-x".into());
2773 gate.cards.push(Card {
2774 agent: "agent-k".into(),
2775 id: "cron-bot".into(),
2776 name: "cron_bot".into(),
2777 trigger: "cron",
2778 enabled: true,
2779 });
2780 }
2781 let shared = Arc::clone(&state);
2782 let token = "gw-test-token";
2783 let done = Arc::new(AtomicBool::new(false));
2784 let flag = Arc::clone(&done);
2785 let server = thread::spawn(move || {
2786 while !flag.load(Ordering::Relaxed) {
2787 let (mut sock, _) = match listener.accept() {
2788 Ok(pair) => pair,
2789 Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
2790 thread::sleep(std::time::Duration::from_millis(5));
2791 continue;
2792 }
2793 Err(_) => break,
2794 };
2795 let _ = sock.set_nonblocking(false);
2796 let Some((headers, body)) = read_http(&mut sock) else {
2797 continue;
2798 };
2799 let lower = headers.to_ascii_lowercase();
2800 let authorized = lower.contains(&format!("authorization: bearer {token}"));
2801 let request: serde_json::Value = serde_json::from_slice(&body).unwrap_or_default();
2802 let agent = request["id"].as_str().unwrap_or("").to_string();
2803 let mut gate = shared.lock().expect("gate");
2804 gate.agents_seen.insert(agent.clone());
2805 let response = if !authorized {
2806 reply("401 Unauthorized", &serde_json::json!({}))
2807 } else if lower.starts_with("post /api/getagentautomations ") {
2808 reply("200 OK", &cards_of(&gate, &agent))
2809 } else if lower.starts_with("post /api/createagentautomation ") {
2810 let spec = &request["spec"];
2811 assert_eq!(spec["trigger"]["type"], "webhook");
2812 assert_eq!(spec["isEnabled"], false, "mirror must be disabled");
2813 let name = spec["name"].as_str().unwrap_or("").to_string();
2814 let prompt = spec["prompt"].as_str().unwrap_or("");
2815 assert!(!prompt.is_empty());
2816 assert!(!prompt.contains("http"));
2817 gate.creates += 1;
2818 let folder = crate::nick::routine_folder_id(&name).expect("slug");
2819 let id = if gate.clash.contains(&agent) {
2820 format!("{folder}-2")
2821 } else {
2822 folder
2823 };
2824 gate.cards.push(Card {
2825 agent: agent.clone(),
2826 id,
2827 name,
2828 trigger: "webhook",
2829 enabled: false,
2830 });
2831 reply("200 OK", &cards_of(&gate, &agent))
2832 } else if lower.starts_with("post /api/updateagent ") {
2833 gate.profiles
2834 .push((agent.clone(), request["profile"].clone()));
2835 reply("200 OK", &serde_json::json!({"id": agent}))
2836 } else if lower.starts_with("post /api/deleteagentautomation ") {
2837 let id = request["automationId"].as_str().unwrap_or("");
2838 gate.deletes += 1;
2839 gate.cards
2840 .retain(|card| !(card.agent == agent && card.id == id));
2841 reply("200 OK", &cards_of(&gate, &agent))
2842 } else if lower.starts_with("post /api/getautomationwebhookcredential ") {
2843 let id = request["automationId"].as_str().unwrap_or("").to_string();
2844 let local = gate
2845 .cards
2846 .iter()
2847 .find(|card| card.agent == agent && card.id == id);
2848 match local {
2849 None => reply(
2850 "500 Internal Server Error",
2851 &serde_json::json!({"error": format!("Automation not found: {id}")}),
2852 ),
2853 Some(card) if card.trigger != "webhook" => reply(
2854 "500 Internal Server Error",
2855 &serde_json::json!({"error": "Automation is not webhook-triggered"}),
2856 ),
2857 Some(_) => {
2858 let minted = gate.backend.contains(&(agent.clone(), id.clone()));
2859 reply(
2860 "200 OK",
2861 &serde_json::json!({
2862 "url": format!("https://backend.invalid/automations/webhook/{agent}-{id}"),
2863 "key": minted.then(|| format!("key-{agent}-{id}")),
2864 }),
2865 )
2866 }
2867 }
2868 } else {
2869 reply("404 Not Found", &serde_json::json!({}))
2870 };
2871 drop(gate);
2872 let _ = sock.write_all(&response);
2873 }
2874 });
2875
2876 let dir = std::env::temp_dir().join(format!(
2877 "m4a-mirror-{}-{}",
2878 std::process::id(),
2879 std::time::SystemTime::now()
2880 .duration_since(std::time::UNIX_EPOCH)
2881 .expect("clock")
2882 .as_nanos()
2883 ));
2884 let agents = dir.join("agents");
2885 let store_root = dir.join("stores");
2886 let describe = |id: &str| {
2887 if id == "agent-c" {
2888 "about the chief"
2889 } else {
2890 ""
2891 }
2892 };
2893 for (id, name) in [
2894 ("agent-h", "Alice"),
2895 ("agent-c", "Привет мир"),
2896 ("agent-s", "Свой браузер"),
2897 ("agent-x", "Clash Bot"),
2898 ("agent-k", "Cron Bot"),
2899 ] {
2900 std::fs::create_dir_all(agents.join(id)).expect("agent dir");
2901 std::fs::write(
2902 agents.join(id).join("profile.json"),
2903 serde_json::to_vec(&serde_json::json!({"name": name, "description": describe(id)}))
2904 .expect("profile"),
2905 )
2906 .expect("profile");
2907 }
2908 let gateway = dir.join("gateway-file.json");
2909 std::fs::write(
2910 &gateway,
2911 format!(r#"{{"host":"192.0.2.1","port":{port},"scheme":"http","token":"{token}"}}"#),
2912 )
2913 .expect("gateway");
2914 let skip = parse_skip_nicks(" svoi-brauzer , ");
2915 assert_eq!(skip, vec!["svoi-brauzer".to_string()]);
2916 let options = WakeOptions {
2917 skip_nicks: skip.clone(),
2918 store_root: Some(store_root.clone()),
2919 profile_note: false,
2920 session_ids: Vec::new(),
2921 };
2922
2923 let status_of = |reports: &[RoutineReport], agent: &str| {
2924 reports
2925 .iter()
2926 .find(|report| report.agent_id == agent)
2927 .map(|report| report.status.clone())
2928 };
2929 let first =
2930 ensure_agent_webhook_routines(&agents, &gateway, None, &options).expect("first pass");
2931 assert_eq!(status_of(&first, "agent-h"), Some(WakeStatus::Ready));
2932 assert_eq!(
2933 status_of(&first, "agent-c"),
2934 Some(WakeStatus::AwaitingBackend)
2935 );
2936 assert!(matches!(
2937 status_of(&first, "agent-x"),
2938 Some(WakeStatus::Failed(_))
2939 ));
2940 assert!(matches!(
2941 status_of(&first, "agent-k"),
2942 Some(WakeStatus::Failed(_))
2943 ));
2944 assert_eq!(status_of(&first, "agent-s"), None);
2945 let chief = first
2946 .iter()
2947 .find(|r| r.agent_id == "agent-c")
2948 .expect("chief");
2949 assert_eq!(chief.nick, "privet-mir");
2950 assert_eq!(chief.folder_id, "privet-mir");
2951 let shown = format!("{first:?}");
2952 assert!(!shown.contains("key-"));
2953 assert!(!shown.contains("automations/webhook"));
2954 {
2955 let gate = state.lock().expect("gate");
2956 assert_eq!(gate.creates, 3);
2959 assert_eq!(gate.deletes, 1);
2960 assert!(!gate.agents_seen.contains("agent-s"));
2961 assert!(gate.cards.iter().all(|card| card.agent != "agent-x"));
2962 let cron = gate
2963 .cards
2964 .iter()
2965 .find(|card| card.agent == "agent-k")
2966 .expect("cron");
2967 assert_eq!(cron.trigger, "cron");
2968 assert!(cron.enabled);
2969 let mirror = gate
2970 .cards
2971 .iter()
2972 .find(|card| card.agent == "agent-h")
2973 .expect("mirror");
2974 assert_eq!(
2975 (mirror.id.as_str(), mirror.name.as_str()),
2976 ("alice", "alice")
2977 );
2978 let chief_mirror = gate
2979 .cards
2980 .iter()
2981 .find(|card| card.agent == "agent-c")
2982 .expect("chief mirror");
2983 assert_eq!(chief_mirror.name, "privet-mir");
2985 assert_eq!(chief_mirror.id, "privet-mir");
2986 assert!(!mirror.enabled);
2987 }
2988 let alice_file = keychain_path(&store_root, "agent-h");
2990 let stored: serde_json::Value =
2991 serde_json::from_slice(&std::fs::read(&alice_file).expect("keychain"))
2992 .expect("keychain json");
2993 assert_eq!(stored["folder_id"], "alice");
2994 assert_eq!(stored["key"], "key-agent-h-alice");
2995 #[cfg(unix)]
2996 {
2997 use std::os::unix::fs::PermissionsExt;
2998 let mode = std::fs::metadata(&alice_file)
2999 .expect("meta")
3000 .permissions()
3001 .mode();
3002 assert_eq!(mode & 0o777, 0o600);
3003 }
3004 assert!(!keychain_path(&store_root, "agent-c").exists());
3005
3006 let second =
3008 ensure_agent_webhook_routines(&agents, &gateway, None, &options).expect("second pass");
3009 assert_eq!(
3010 status_of(&second, "agent-c"),
3011 Some(WakeStatus::AwaitingBackend)
3012 );
3013 assert_eq!(
3014 state.lock().expect("gate").creates,
3015 4,
3016 "only the clash bot retries"
3017 );
3018 assert!(state.lock().expect("gate").profiles.is_empty());
3019 let noted = WakeOptions {
3022 profile_note: true,
3023 ..options.clone()
3024 };
3025 ensure_agent_webhook_routines(&agents, &gateway, None, ¬ed).expect("note pass");
3026 {
3027 let gate = state.lock().expect("gate");
3028 assert_eq!(gate.profiles.len(), 1);
3029 let (agent, profile) = &gate.profiles[0];
3030 assert_eq!(agent, "agent-c");
3031 assert_eq!(profile["name"], "Привет мир");
3032 let description = profile["description"].as_str().unwrap_or("");
3033 assert!(description.starts_with("about the chief"));
3034 assert!(description.contains("routine named \"privet-mir\""));
3035 assert!(!description.contains("http"));
3036 }
3037 state.lock().expect("gate").creates = 4;
3038 state
3041 .lock()
3042 .expect("gate")
3043 .backend
3044 .insert(("agent-c".into(), "privet-mir".into()));
3045 let third =
3046 ensure_agent_webhook_routines(&agents, &gateway, None, &options).expect("third pass");
3047 assert_eq!(status_of(&third, "agent-c"), Some(WakeStatus::Ready));
3048 assert_eq!(state.lock().expect("gate").creates, 5);
3049 assert!(keychain_path(&store_root, "agent-c").exists());
3050
3051 let gateway_path = gateway.display().to_string();
3053 let root_text = store_root.display().to_string();
3054 let sessions = load_web_agents(&agents, |key| match key {
3055 GATEWAY_FILE_ENV => Some(gateway_path.clone()),
3056 SKIP_NICKS_ENV => Some("svoi-brauzer".to_string()),
3057 crate::STORE_ROOT_ENV => Some(root_text.clone()),
3058 _ => None,
3059 })
3060 .expect("web agents");
3061 assert_eq!(sessions.len(), 4);
3062 let alice = sessions
3063 .iter()
3064 .find(|s| s.session_id == "agent-h")
3065 .expect("h");
3066 assert!(alice
3067 .routine_url
3068 .as_deref()
3069 .unwrap_or("")
3070 .starts_with("https://"));
3071 assert!(alice.routine_bearer.is_some());
3072 assert!(!format!("{alice:?}").contains("key-"));
3073
3074 let mut offline = vec![HostSession::new("Alice", "agent-h")];
3076 offline[0].agent_id = Some("agent-h".to_string());
3077 let mut renamed = HostSession::new("Alice Two", "agent-h");
3078 renamed.agent_id = Some("agent-h".to_string());
3079 offline.push(renamed);
3080 attach_webhook_routines(
3081 &mut offline,
3082 &dir.join("absent.json"),
3083 None,
3084 Some(&store_root),
3085 )
3086 .expect("offline");
3087 assert_eq!(
3088 offline[0].routine_bearer.as_deref(),
3089 Some("key-agent-h-alice")
3090 );
3091 assert!(offline[1].routine_url.is_none());
3092 let mut bare = vec![HostSession::new("Alice", "agent-h")];
3093 attach_webhook_routines(&mut bare, &dir.join("absent.json"), None, None).expect("bare");
3094 assert!(bare[0].routine_url.is_none());
3095
3096 done.store(true, Ordering::Relaxed);
3097 let _ = server.join();
3098 let _ = std::fs::remove_dir_all(&dir);
3099 }
3100
3101 #[test]
3102 fn wake_note_is_appended_once_and_replaced_in_place() {
3103 let first = with_wake_note("Runs the alice.", "alice").expect("added");
3104 assert!(first.starts_with("Runs the alice.\n\n<!-- mail4agent:wake -->"));
3105 assert!(first.ends_with("<!-- /mail4agent:wake -->"));
3106 assert!(with_wake_note(&first, "alice").is_none());
3107 let renamed = with_wake_note(&first, "alice-two").expect("replaced");
3108 assert_eq!(renamed.matches("<!-- mail4agent:wake -->").count(), 1);
3109 assert!(renamed.contains("\"alice-two\""));
3110 assert!(!renamed.contains("\"alice\""));
3111 assert!(renamed.starts_with("Runs the alice."));
3112 assert_eq!(with_wake_note("", "carol").expect("empty"), wake_note("carol"));
3113 }
3114
3115 #[test]
3116 fn agents_directory_becomes_sessions_without_hardcoded_ids() {
3117 let dir = std::env::temp_dir().join(format!(
3118 "m4a-agents-{}-{}",
3119 std::process::id(),
3120 std::time::SystemTime::now()
3121 .duration_since(std::time::UNIX_EPOCH)
3122 .expect("clock")
3123 .as_nanos()
3124 ));
3125 let alice = dir.join("agent-alice");
3126 let chief = dir.join("agent-chief");
3127 std::fs::create_dir_all(&alice).expect("dir");
3128 std::fs::create_dir_all(&chief).expect("dir");
3129 std::fs::write(
3130 alice.join("profile.json"),
3131 r#"{"name":"Alice","description":"x"}"#,
3132 )
3133 .expect("profile");
3134 std::fs::write(
3135 chief.join("profile.json"),
3136 concat!(
3137 "{\"name\":\"",
3138 "Привет мир",
3139 "\",\"description\":\"x\"}"
3140 ),
3141 )
3142 .expect("profile");
3143 std::fs::write(dir.join("active-agent.json"), "{}").expect("skip file");
3144 let loaded = load_agents_dir(&dir).expect("agents");
3145 assert_eq!(loaded.len(), 2);
3146 assert_eq!(loaded[1].session_id, "agent-chief");
3147 assert_eq!(loaded[1].agent_id.as_deref(), Some("agent-chief"));
3148 assert_eq!(loaded[1].bot_name, "Привет мир");
3149 assert_eq!(
3150 routine_name_for(&loaded[1]).as_deref(),
3151 Some("privet-mir")
3152 );
3153 assert_ne!(
3154 routine_name_for(&loaded[1]).as_deref(),
3155 Some(loaded[1].session_id.as_str())
3156 );
3157 assert_eq!(loaded[0].bot_name, "Alice");
3158 assert_eq!(routine_name_for(&loaded[0]).as_deref(), Some("alice"));
3159 assert!(loaded.iter().all(|session| session.routine_url.is_none()));
3160 let _ = std::fs::remove_dir_all(&dir);
3161 }
3162}