1use std::fs::File;
61use std::collections::{HashMap, HashSet};
62use std::path::{Path, PathBuf};
63use std::sync::Arc;
64use zeroize::Zeroizing;
65
66use crate::ipc::{SendListener, SendStream};
67pub(crate) use crate::local_bus::{ensure_server_name, LocalBus};
68use crate::{
69 clip_public, register_session, DeviceId, OpenedStore,
70 SessionConfig, SessionWake, ShellError, HOMESERVER_URL_ENV, STORE_ROOT_ENV,
71};
72
73pub const SESSIONS_DIR_ENV: &str = "M4A_SESSIONS_DIR";
78
79pub const AGENTS_DIR_ENV: &str = "M4A_AGENTS_DIR";
83
84pub const DEFAULT_AGENTS_DIR: &str = "agent-data/agents";
88
89pub(crate) fn under_home(rel: &str) -> PathBuf {
91 match std::env::var("HOME").or_else(|_| std::env::var("USERPROFILE")) {
92 Ok(h) if !h.is_empty() => PathBuf::from(h).join(rel),
93 _ => PathBuf::from(rel),
94 }
95}
96
97pub const AGENT_RESCAN_SECS_ENV: &str = "M4A_AGENT_RESCAN_SECS";
100
101pub struct HostSession {
106 pub bot_name: String,
108 pub session_id: String,
111 pub agent_id: Option<String>,
114 pub routine_url: Option<String>,
116 pub routine_bearer: Option<String>,
118 pub invite: Option<String>,
121 pub tier: m4a_agent::BackendKind,
123}
124
125impl HostSession {
126 pub fn new(bot_name: impl Into<String>, session_id: impl Into<String>) -> Self {
128 Self {
129 bot_name: bot_name.into(),
130 session_id: session_id.into(),
131 agent_id: None,
132 routine_url: None,
133 routine_bearer: None,
134 invite: None,
135 tier: m4a_agent::BackendKind::Server,
136 }
137 }
138
139 pub(crate) fn config(&self, url: &str, store_root: &Path) -> Result<SessionConfig, ShellError> {
140 SessionConfig::new_identity(url, self.tier, &self.session_id, store_root, self.invite.clone())
141 }
142
143 pub fn with_invite(mut self, invite: impl Into<String>) -> Self {
145 self.invite = Some(invite.into());
146 self
147 }
148
149 pub fn with_tier(mut self, tier: m4a_agent::BackendKind) -> Self {
151 self.tier = tier;
152 self
153 }
154
155 pub fn with_routine(mut self, url: impl Into<String>, bearer: Option<String>) -> Self {
157 self.routine_url = Some(url.into());
158 self.routine_bearer = bearer.filter(|token| !token.is_empty());
159 self
160 }
161
162}
163
164impl std::fmt::Debug for HostSession {
165 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
166 f.debug_struct("HostSession")
167 .field("bot_name", &self.bot_name)
168 .field("session_id", &self.session_id)
169 .field("agent_id", &self.agent_id)
170 .field("routine_url", &self.routine_url.as_ref().map(|_| "[set]"))
171 .field(
172 "routine_bearer",
173 &self.routine_bearer.as_ref().map(|_| "[redacted]"),
174 )
175 .finish()
176 }
177}
178
179#[derive(serde::Deserialize)]
180#[serde(deny_unknown_fields)]
181struct SessionFile {
182 bot_name: String,
183 session_id: String,
184 #[serde(default)]
186 agent_id: Option<String>,
187}
188
189pub fn load_session_records(dir: &Path) -> Result<Vec<HostSession>, ShellError> {
193 if !dir.is_dir() {
194 return Err(ShellError::SessionList(
195 "session directory is missing".to_string(),
196 ));
197 }
198 let mut paths = Vec::new();
199 for entry in std::fs::read_dir(dir)? {
200 let entry = entry?;
201 let path = entry.path();
202 if path.extension().and_then(|ext| ext.to_str()) != Some("json") {
203 continue;
204 }
205 paths.push(path);
206 }
207 paths.sort();
208 if paths.is_empty() {
209 return Err(ShellError::SessionList(
210 "session directory has no records".to_string(),
211 ));
212 }
213 let mut out = Vec::with_capacity(paths.len());
214 for path in paths {
215 let text = std::fs::read_to_string(&path)
216 .map_err(|err| ShellError::SessionList(clip_public(err.to_string())))?;
217 let file: SessionFile = serde_json::from_str(&text)
218 .map_err(|err| ShellError::SessionList(clip_public(err.to_string())))?;
219 let mut session = HostSession::new(file.bot_name, file.session_id);
220 session.agent_id = file
221 .agent_id
222 .map(|id| id.trim().to_string())
223 .filter(|id| !id.is_empty());
224 out.push(session);
225 }
226 Ok(out)
227}
228
229#[derive(serde::Deserialize)]
230struct AgentProfileFile {
231 name: String,
232}
233
234pub fn load_agents_dir(dir: &Path) -> Result<Vec<HostSession>, ShellError> {
238 if !dir.is_dir() {
239 return Err(ShellError::SessionList(
240 "agents directory is missing".to_string(),
241 ));
242 }
243 let mut ids = Vec::new();
244 for entry in std::fs::read_dir(dir)? {
245 let entry = entry?;
246 if !entry.file_type()?.is_dir() {
247 continue;
248 }
249 let name = entry.file_name();
250 let Some(id) = name.to_str() else {
251 continue;
252 };
253 if id.is_empty() || id.starts_with('.') {
254 continue;
255 }
256 ids.push(id.to_string());
257 }
258 ids.sort();
259 if ids.is_empty() {
260 return Err(ShellError::SessionList(
261 "agents directory has no agents".to_string(),
262 ));
263 }
264 let mut out = Vec::with_capacity(ids.len());
265 for id in ids {
266 let profile_path = dir.join(&id).join("profile.json");
267 if !profile_path.is_file() {
268 continue;
269 }
270 let text = std::fs::read_to_string(&profile_path)
271 .map_err(|err| ShellError::SessionList(clip_public(err.to_string())))?;
272 let profile: AgentProfileFile = serde_json::from_str(&text)
273 .map_err(|err| ShellError::SessionList(clip_public(err.to_string())))?;
274 let bot_name = profile.name.trim().to_string();
275 if bot_name.is_empty() {
276 continue;
277 }
278 let mut session = HostSession::new(bot_name, id.clone());
279 session.agent_id = Some(id);
280 out.push(session);
281 }
282 if out.is_empty() {
283 return Err(ShellError::SessionList(
284 "agents directory has no profiles".to_string(),
285 ));
286 }
287 Ok(out)
288}
289
290fn resolve_agents_dir(mut get: impl FnMut(&str) -> Option<String>) -> Option<PathBuf> {
294 if let Some(path) = get(AGENTS_DIR_ENV).filter(|value| !value.is_empty()) {
295 let path = PathBuf::from(path);
296 if path.is_dir() {
297 return Some(path);
298 }
299 return None;
300 }
301 let default = under_home(DEFAULT_AGENTS_DIR);
302 if default.is_dir() {
303 Some(default)
304 } else {
305 None
306 }
307}
308
309fn load_web_sessions(
317 dir: &Path,
318 mut get: impl FnMut(&str) -> Option<String>,
319) -> Result<Vec<HostSession>, ShellError> {
320 let sessions = load_session_records(dir)?;
321 let mut sessions = drop_skipped(sessions, &skip_nicks(&mut get));
322 attach_routines_from_env(&mut sessions, &mut get)?;
323 Ok(sessions)
324}
325
326fn load_web_agents(
328 dir: &Path,
329 mut get: impl FnMut(&str) -> Option<String>,
330) -> Result<Vec<HostSession>, ShellError> {
331 let sessions = load_agents_dir(dir)?;
332 let mut sessions = drop_skipped(sessions, &skip_nicks(&mut get));
333 apply_session_ids(
334 &mut sessions,
335 &parse_session_ids(get(SESSION_IDS_ENV).as_deref().unwrap_or("")),
336 );
337 attach_routines_from_env(&mut sessions, &mut get)?;
338 Ok(sessions)
339}
340
341fn attach_routines_from_env(
342 sessions: &mut [HostSession],
343 get: &mut impl FnMut(&str) -> Option<String>,
344) -> Result<(), ShellError> {
345 let store_root = get(STORE_ROOT_ENV).map(PathBuf::from);
346 let Some(path) = get(GATEWAY_FILE_ENV).filter(|value| !value.is_empty()) else {
347 return Ok(());
348 };
349 let token = get(GATEWAY_TOKEN_ENV).filter(|value| !value.is_empty());
350 attach_webhook_routines(
351 sessions,
352 Path::new(&path),
353 token.as_deref(),
354 store_root.as_deref(),
355 )?;
356 Ok(())
357}
358
359fn unopened_agent_sessions(
363 found: Vec<HostSession>,
364 skip: &[String],
365 aliases: &[(String, String)],
366 store_root: &Path,
367 held: &[(PathBuf, Option<String>)],
368) -> Vec<(HostSession, String)> {
369 let mut sessions = drop_skipped(found, skip);
370 apply_session_ids(&mut sessions, aliases);
371 sessions
372 .into_iter()
373 .filter_map(|session| {
374 let nick = routine_name_for(&session)?;
375 let dir = crate::session_store_dir(store_root, &session.session_id);
376 let taken = held.iter().any(|(held_dir, held_nick)| {
377 *held_dir == dir
378 || held_nick
379 .as_deref()
380 .is_some_and(|held| held.eq_ignore_ascii_case(&nick))
381 });
382 (!taken).then_some((session, nick))
383 })
384 .collect()
385}
386
387pub const SESSION_IDS_ENV: &str = "M4A_SESSION_IDS";
392
393pub(crate) fn parse_session_ids(raw: &str) -> Vec<(String, String)> {
394 raw.split(',')
395 .filter_map(|pair| {
396 let (agent, session) = pair.split_once('=')?;
397 let (agent, session) = (agent.trim(), session.trim());
398 (!agent.is_empty() && !session.is_empty())
399 .then(|| (agent.to_string(), session.to_string()))
400 })
401 .collect()
402}
403
404pub(crate) fn apply_session_ids(sessions: &mut [HostSession], aliases: &[(String, String)]) {
405 for session in sessions.iter_mut() {
406 let Some(agent_id) = session.agent_id.as_deref() else {
407 continue;
408 };
409 if let Some((_, alias)) = aliases.iter().find(|(agent, _)| agent == agent_id) {
410 session.session_id = alias.clone();
411 }
412 }
413}
414
415pub const SKIP_NICKS_ENV: &str = "M4A_SKIP_NICKS";
419
420fn skip_nicks(get: &mut impl FnMut(&str) -> Option<String>) -> Vec<String> {
421 parse_skip_nicks(get(SKIP_NICKS_ENV).as_deref().unwrap_or(""))
422}
423
424fn parse_skip_nicks(raw: &str) -> Vec<String> {
425 raw.split(',')
426 .map(|item| item.trim().to_ascii_lowercase())
427 .filter(|item| !item.is_empty())
428 .collect()
429}
430
431fn drop_skipped(sessions: Vec<HostSession>, skip: &[String]) -> Vec<HostSession> {
432 if skip.is_empty() {
433 return sessions;
434 }
435 sessions
436 .into_iter()
437 .filter(|session| match routine_name_for(session) {
438 Some(nick) => !skip.iter().any(|item| item.eq_ignore_ascii_case(&nick)),
439 None => true,
440 })
441 .collect()
442}
443
444#[derive(Debug, Clone, PartialEq, Eq)]
447pub enum WakeStatus {
448 Ready,
450 AwaitingBackend,
454 Failed(String),
456}
457
458#[derive(Debug, Clone, PartialEq, Eq)]
460pub struct RoutineReport {
461 pub agent_id: String,
463 pub nick: String,
465 pub folder_id: String,
468 pub status: WakeStatus,
470}
471
472#[derive(Debug, Clone, Default)]
475pub struct WakeOptions {
476 pub skip_nicks: Vec<String>,
478 pub store_root: Option<PathBuf>,
480 pub profile_note: bool,
483 pub session_ids: Vec<(String, String)>,
485}
486
487impl WakeOptions {
488 pub fn from_lookup(mut get: impl FnMut(&str) -> Option<String>) -> Self {
490 Self {
491 skip_nicks: skip_nicks(&mut get),
492 store_root: get(STORE_ROOT_ENV).map(PathBuf::from),
493 profile_note: get(PROFILE_NOTE_ENV).as_deref() == Some("1"),
494 session_ids: parse_session_ids(get(SESSION_IDS_ENV).as_deref().unwrap_or("")),
495 }
496 }
497}
498
499pub fn ensure_agent_webhook_routines(
507 agents_dir: &Path,
508 gateway_file: &Path,
509 token_override: Option<&str>,
510 options: &WakeOptions,
511) -> Result<Vec<RoutineReport>, ShellError> {
512 let store_root = options.store_root.as_deref();
513 let mut sessions = drop_skipped(load_agents_dir(agents_dir)?, &options.skip_nicks);
514 apply_session_ids(&mut sessions, &options.session_ids);
515 let Some(gate) = open_gateway(gateway_file, token_override)? else {
516 return Err(ShellError::Gateway("gateway file is missing".to_string()));
517 };
518 let mut reports = Vec::new();
519 for session in &sessions {
520 let Some(agent_id) = session.agent_id.as_deref().filter(|id| !id.is_empty()) else {
521 continue;
522 };
523 let Some(nick) = routine_name_for(session) else {
524 continue;
525 };
526 let outcome = ensure_wake(&gate, agent_id, &nick);
527 let folder_id = crate::nick::routine_folder_id(&nick).unwrap_or_default();
528 if let (WakeOutcome::Ready { url, key }, Some(root)) = (&outcome, store_root) {
529 save_wake(root, &session.session_id, &folder_id, url, key);
530 }
531 log_outcome(&nick, &folder_id, &outcome);
532 if options.profile_note && matches!(outcome, WakeOutcome::AwaitingBackend) {
533 note_profile(&gate, agents_dir, agent_id, &nick);
534 }
535 reports.push(RoutineReport {
536 agent_id: agent_id.to_string(),
537 nick,
538 folder_id,
539 status: outcome.status(),
540 });
541 }
542 Ok(reports)
543}
544
545pub fn ensure_agent_webhook_routines_from_env() -> Result<Vec<RoutineReport>, ShellError> {
549 let mut get = |key: &str| std::env::var(key).ok().filter(|v| !v.is_empty());
550 let agents_dir = resolve_agents_dir(&mut get)
551 .ok_or_else(|| ShellError::SessionList("agents directory is missing".to_string()))?;
552 let gateway = get(GATEWAY_FILE_ENV).unwrap_or_else(|| DEFAULT_GATEWAY_FILE.to_string());
553 let token = get(GATEWAY_TOKEN_ENV);
554 let options = WakeOptions::from_lookup(&mut get);
555 ensure_agent_webhook_routines(&agents_dir, Path::new(&gateway), token.as_deref(), &options)
556}
557
558pub const PROFILE_NOTE_ENV: &str = "M4A_BOOTSTRAP_PROFILE_NOTE";
564
565const NOTE_OPEN: &str = "<!-- mail4agent:wake -->";
566const NOTE_CLOSE: &str = "<!-- /mail4agent:wake -->";
567const DESCRIPTION_MAX: usize = 20_000;
569
570fn wake_note(routine: &str) -> String {
573 format!(
574 "{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}"
575 )
576}
577
578fn with_wake_note(description: &str, routine: &str) -> Option<String> {
582 let note = wake_note(routine);
583 if description.contains(¬e) {
584 return None;
585 }
586 if let (Some(start), Some(end)) = (description.find(NOTE_OPEN), description.find(NOTE_CLOSE)) {
587 if end > start {
588 let mut out = String::with_capacity(description.len() + note.len());
589 out.push_str(&description[..start]);
590 out.push_str(¬e);
591 out.push_str(&description[end + NOTE_CLOSE.len()..]);
592 return Some(out);
593 }
594 }
595 let base = description.trim_end();
596 Some(if base.is_empty() {
597 note
598 } else {
599 format!("{base}\n\n{note}")
600 })
601}
602
603fn note_profile(gate: &GatewayConn, agents_dir: &Path, agent_id: &str, nick: &str) {
607 match ensure_profile_note(gate, agents_dir, agent_id, nick) {
608 Ok(true) => eprintln!("mail4agent: wake {nick}: bootstrap note written to the profile"),
609 Ok(false) => {}
610 Err(err) => eprintln!("mail4agent: wake {nick}: bootstrap note not written ({err})"),
611 }
612}
613
614fn ensure_profile_note(
615 gate: &GatewayConn,
616 agents_dir: &Path,
617 agent_id: &str,
618 nick: &str,
619) -> Result<bool, ShellError> {
620 let text = std::fs::read_to_string(agents_dir.join(agent_id).join("profile.json"))
621 .map_err(|_| ShellError::Gateway("profile is unreadable".to_string()))?;
622 let profile: serde_json::Value = serde_json::from_str(&text)
623 .map_err(|_| ShellError::Gateway("profile is unreadable".to_string()))?;
624 let field = |key: &str| profile.get(key).and_then(|value| value.as_str());
625 let Some(name) = field("name").filter(|name| !name.trim().is_empty()) else {
626 return Err(ShellError::Gateway("profile has no name".to_string()));
627 };
628 let routine = crate::nick::routine_folder_id(nick)
629 .ok_or_else(|| ShellError::Gateway("nick has no routine name".to_string()))?;
630 let Some(description) = with_wake_note(field("description").unwrap_or(""), &routine) else {
631 return Ok(false);
632 };
633 if description.chars().count() > DESCRIPTION_MAX {
634 return Err(ShellError::Gateway(
635 "description would be too long".to_string(),
636 ));
637 }
638 let mut body = serde_json::json!({ "name": name, "description": description });
639 for key in ["title", "avatarShape", "avatarColor"] {
640 if let Some(value) = field(key) {
641 body[key] = serde_json::Value::String(value.to_string());
642 }
643 }
644 let (status, _) = gateway_post(
645 gate,
646 "updateAgent",
647 &serde_json::json!({ "id": agent_id, "profile": body }),
648 )?;
649 if status != 200 {
650 return Err(ShellError::Gateway(format!(
651 "gateway profile status {status}"
652 )));
653 }
654 Ok(true)
655}
656
657const DEFAULT_GATEWAY_FILE: &str = "agent-data/gateway.json";
660
661pub const GATEWAY_FILE_ENV: &str = "M4A_GATEWAY_FILE";
664
665pub const GATEWAY_TOKEN_ENV: &str = "M4A_GATEWAY_TOKEN";
668
669const 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.";
673
674pub const WAKE_KEYCHAIN_FILE: &str = "routine-wake.json";
678
679#[derive(serde::Deserialize)]
680struct GatewayFile {
681 port: u16,
682 scheme: String,
683 #[serde(default)]
684 token: Option<String>,
685}
686
687#[derive(serde::Deserialize)]
688struct AutomationCard {
689 id: String,
690 name: String,
691 #[serde(default)]
692 trigger: Option<serde_json::Value>,
693}
694
695impl AutomationCard {
696 fn is_webhook(&self) -> bool {
699 fn has_webhook(value: &serde_json::Value) -> bool {
700 match value {
701 serde_json::Value::Array(items) => items.iter().any(has_webhook),
702 serde_json::Value::Object(map) => {
703 map.get("type").and_then(|kind| kind.as_str()) == Some("webhook")
704 || map
705 .values()
706 .any(|inner| inner.is_array() && has_webhook(inner))
707 }
708 _ => false,
709 }
710 }
711 self.trigger.as_ref().map(has_webhook).unwrap_or(true)
712 }
713}
714
715#[derive(serde::Deserialize)]
716struct WebhookCredential {
717 url: String,
718 key: Option<String>,
719}
720
721#[derive(serde::Serialize, serde::Deserialize)]
722#[serde(deny_unknown_fields)]
723struct StoredWake {
724 folder_id: String,
725 url: String,
726 key: String,
727}
728
729fn routine_name_for(session: &HostSession) -> Option<String> {
734 crate::nick::nick_from_display_name(session.bot_name.trim()).ok()
735}
736
737struct GatewayConn {
738 client: reqwest::blocking::Client,
739 base: String,
740 token: String,
741}
742
743fn open_gateway(
744 gateway_file: &Path,
745 token_override: Option<&str>,
746) -> Result<Option<GatewayConn>, ShellError> {
747 if !gateway_file.is_file() {
748 return Ok(None);
749 }
750 let text = std::fs::read_to_string(gateway_file)
751 .map_err(|_| ShellError::Gateway("gateway file is unreadable".to_string()))?;
752 let file: GatewayFile = serde_json::from_str(&text)
753 .map_err(|_| ShellError::Gateway("gateway file is unreadable".to_string()))?;
754 let scheme = file.scheme.trim();
755 if scheme != "http" && scheme != "https" {
756 return Err(ShellError::Gateway(
757 "gateway scheme is not http or https".to_string(),
758 ));
759 }
760 if file.port == 0 {
761 return Err(ShellError::Gateway("gateway port is unset".to_string()));
762 }
763 let token = token_override
764 .map(str::trim)
765 .filter(|value| !value.is_empty())
766 .map(str::to_string)
767 .or(file.token.filter(|value| !value.trim().is_empty()));
768 let Some(token) = token else {
769 return Err(ShellError::Gateway("gateway token is unset".to_string()));
770 };
771 let base = format!("{scheme}://127.0.0.1:{}", file.port);
773 Ok(Some(GatewayConn {
774 client: gateway_client()?,
775 base,
776 token,
777 }))
778}
779
780enum WakeOutcome {
783 Ready { url: String, key: String },
784 AwaitingBackend,
785 Failed(String),
786}
787
788impl WakeOutcome {
789 fn status(&self) -> WakeStatus {
790 match self {
791 WakeOutcome::Ready { .. } => WakeStatus::Ready,
792 WakeOutcome::AwaitingBackend => WakeStatus::AwaitingBackend,
793 WakeOutcome::Failed(reason) => WakeStatus::Failed(reason.clone()),
794 }
795 }
796}
797
798fn log_outcome(nick: &str, folder_id: &str, outcome: &WakeOutcome) {
799 match outcome {
800 WakeOutcome::Ready { .. } => eprintln!("mail4agent: wake {nick} ({folder_id}): ready"),
801 WakeOutcome::AwaitingBackend => eprintln!(
802 "mail4agent: wake {nick} ({folder_id}): no key yet; the bot has not created its own routine with this folder id; will retry"
803 ),
804 WakeOutcome::Failed(reason) => {
805 eprintln!("mail4agent: wake {nick} ({folder_id}): {reason}")
806 }
807 }
808}
809
810fn ensure_wake(gate: &GatewayConn, agent_id: &str, nick: &str) -> WakeOutcome {
820 let Some(folder) = crate::nick::routine_folder_id(nick) else {
821 return WakeOutcome::Failed("nick has no routine folder id".to_string());
822 };
823 let before = match list_agent_automations(gate, agent_id) {
824 Ok(cards) => cards,
825 Err(err) => return WakeOutcome::Failed(err.to_string()),
826 };
827 match before.iter().find(|card| card.id == folder) {
828 Some(card) if !card.is_webhook() => {
829 return WakeOutcome::Failed(
830 "a routine in this folder is not webhook-triggered; left as is".to_string(),
831 );
832 }
833 Some(_) => {}
834 None => {
835 let after = match create_mirror_routine(gate, agent_id, &folder) {
836 Ok(cards) => cards,
837 Err(err) => return WakeOutcome::Failed(err.to_string()),
838 };
839 if !after.iter().any(|card| card.id == folder) {
840 for stray in after.iter().filter(|card| {
841 card.name == folder && !before.iter().any(|old| old.id == card.id)
842 }) {
843 let _ = delete_agent_automation(gate, agent_id, &stray.id);
844 }
845 return WakeOutcome::Failed(
846 "the mirror did not land in the expected folder; removed it".to_string(),
847 );
848 }
849 }
850 }
851 match read_webhook_credential(gate, agent_id, &folder) {
852 Ok(CredentialRead::Ready { url, key }) => WakeOutcome::Ready { url, key },
853 Ok(CredentialRead::MintFailed) => WakeOutcome::AwaitingBackend,
854 Ok(CredentialRead::Missing) => {
855 WakeOutcome::Failed("the mirror is gone from the gateway".to_string())
856 }
857 Err(err) => WakeOutcome::Failed(err.to_string()),
858 }
859}
860
861fn attach_webhook_routines(
866 sessions: &mut [HostSession],
867 gateway_file: &Path,
868 token_override: Option<&str>,
869 store_root: Option<&Path>,
870) -> Result<(), ShellError> {
871 let gate = open_gateway(gateway_file, token_override)?;
872 for session in sessions.iter_mut() {
873 if session.routine_url.is_some() && session.routine_bearer.is_some() {
874 continue;
875 }
876 let Some(agent_id) = session.agent_id.clone().filter(|id| !id.is_empty()) else {
878 continue;
879 };
880 let Some(nick) = routine_name_for(session) else {
881 continue;
882 };
883 let Some(folder) = crate::nick::routine_folder_id(&nick) else {
884 continue;
885 };
886 let Some(gate) = gate.as_ref() else {
887 if let Some(root) = store_root {
888 if let Some((url, key)) = load_wake(root, &session.session_id, &folder) {
889 session.routine_url = Some(url);
890 session.routine_bearer = Some(key);
891 }
892 }
893 continue;
894 };
895 let outcome = ensure_wake(gate, &agent_id, &nick);
896 log_outcome(&nick, &folder, &outcome);
897 match outcome {
898 WakeOutcome::Ready { url, key } => {
899 if let Some(root) = store_root {
900 save_wake(root, &session.session_id, &folder, &url, &key);
901 }
902 session.routine_url = Some(url);
903 session.routine_bearer = Some(key);
904 }
905 WakeOutcome::AwaitingBackend => {
906 if let Some(root) = store_root {
909 clear_wake(root, &session.session_id);
910 }
911 }
912 WakeOutcome::Failed(_) => {}
913 }
914 }
915 Ok(())
916}
917
918pub const STORE_LOCK_FILE: &str = ".lock";
923
924fn now_ms() -> i64 {
925 std::time::SystemTime::now()
926 .duration_since(std::time::UNIX_EPOCH)
927 .map(|elapsed| elapsed.as_millis() as i64)
928 .unwrap_or(0)
929}
930
931pub(crate) fn lock_store(dir: &Path) -> Result<File, ShellError> {
932 std::fs::create_dir_all(dir)?;
933 let file = std::fs::OpenOptions::new()
934 .create(true)
935 .truncate(false)
936 .write(true)
937 .open(dir.join(STORE_LOCK_FILE))?;
938 match file.try_lock() {
939 Ok(()) => Ok(file),
940 Err(_) => Err(ShellError::SessionList(
941 "session store is in use by another process".to_string(),
942 )),
943 }
944}
945
946fn write_secret_file(path: &Path, bytes: &[u8]) -> std::io::Result<()> {
948 let dir = path
949 .parent()
950 .ok_or_else(|| std::io::Error::other("no parent"))?;
951 let mut builder = std::fs::DirBuilder::new();
952 builder.recursive(true);
953 #[cfg(unix)]
954 {
955 use std::os::unix::fs::DirBuilderExt;
956 builder.mode(0o700);
957 }
958 builder.create(dir)?;
959 let name = path
960 .file_name()
961 .and_then(|name| name.to_str())
962 .unwrap_or("secret");
963 let tmp = dir.join(format!("{name}.tmp"));
964 let mut options = std::fs::OpenOptions::new();
965 options.write(true).create(true).truncate(true);
966 #[cfg(unix)]
967 {
968 use std::os::unix::fs::OpenOptionsExt;
969 options.mode(0o600);
970 }
971 let mut file = options.open(&tmp)?;
972 std::io::Write::write_all(&mut file, bytes)?;
973 file.sync_all()?;
974 drop(file);
975 std::fs::rename(&tmp, path)
976}
977
978fn keychain_path(store_root: &Path, session_id: &str) -> PathBuf {
979 crate::session_store_dir(store_root, session_id).join(WAKE_KEYCHAIN_FILE)
980}
981
982fn save_wake(store_root: &Path, session_id: &str, folder_id: &str, url: &str, key: &str) {
985 let stored = StoredWake {
986 folder_id: folder_id.to_string(),
987 url: url.to_string(),
988 key: key.to_string(),
989 };
990 let result = serde_json::to_vec(&stored)
991 .map_err(std::io::Error::other)
992 .and_then(|bytes| write_secret_file(&keychain_path(store_root, session_id), &bytes));
993 if result.is_err() {
994 eprintln!("mail4agent: wake keychain write failed for folder {folder_id}");
995 }
996}
997
998fn load_wake(store_root: &Path, session_id: &str, folder_id: &str) -> Option<(String, String)> {
999 let text = std::fs::read_to_string(keychain_path(store_root, session_id)).ok()?;
1000 let stored: StoredWake = serde_json::from_str(&text).ok()?;
1001 if stored.folder_id != folder_id || stored.url.is_empty() || stored.key.is_empty() {
1002 return None;
1003 }
1004 Some((stored.url, stored.key))
1005}
1006
1007fn clear_wake(store_root: &Path, session_id: &str) {
1008 let _ = std::fs::remove_file(keychain_path(store_root, session_id));
1009}
1010
1011fn list_agent_automations(
1012 gate: &GatewayConn,
1013 agent_id: &str,
1014) -> Result<Vec<AutomationCard>, ShellError> {
1015 let (status, body) = gateway_post(
1016 gate,
1017 "getAgentAutomations",
1018 &serde_json::json!({ "id": agent_id }),
1019 )?;
1020 if status != 200 {
1021 return Err(ShellError::Gateway(format!("gateway list status {status}")));
1022 }
1023 serde_json::from_slice(&body)
1024 .map_err(|_| ShellError::Gateway("gateway list was not understood".to_string()))
1025}
1026
1027enum CredentialRead {
1028 Ready {
1029 url: String,
1030 key: String,
1031 },
1032 MintFailed,
1034 Missing,
1035}
1036
1037fn gateway_client() -> Result<reqwest::blocking::Client, ShellError> {
1038 reqwest::blocking::Client::builder()
1039 .timeout(std::time::Duration::from_secs(8))
1040 .redirect(reqwest::redirect::Policy::none())
1041 .http1_only()
1042 .build()
1043 .map_err(|_| ShellError::Gateway("gateway client could not start".to_string()))
1044}
1045
1046fn gateway_post(
1047 gate: &GatewayConn,
1048 method: &str,
1049 body: &serde_json::Value,
1050) -> Result<(u16, Vec<u8>), ShellError> {
1051 let url = format!("{}/api/{method}", gate.base);
1052 let authorization = crate::bearer_header(&gate.token).map_err(|_| {
1053 ShellError::Gateway("gateway token is not a single header value".to_string())
1054 })?;
1055 let bytes = serde_json::to_vec(body)
1056 .map_err(|_| ShellError::Gateway("gateway request was not json".to_string()))?;
1057 let response = gate
1058 .client
1059 .post(&url)
1060 .header(reqwest::header::AUTHORIZATION, authorization)
1061 .header(reqwest::header::CONTENT_TYPE, "application/json")
1062 .body(bytes)
1063 .send()
1064 .map_err(|_| ShellError::Gateway("gateway request failed".to_string()))?;
1065 let status = response.status().as_u16();
1066 let body = response
1067 .bytes()
1068 .map_err(|_| ShellError::Gateway("gateway response failed".to_string()))?;
1069 Ok((status, body.to_vec()))
1070}
1071
1072fn read_webhook_credential(
1073 gate: &GatewayConn,
1074 agent_id: &str,
1075 automation_id: &str,
1076) -> Result<CredentialRead, ShellError> {
1077 let (status, body) = gateway_post(
1078 gate,
1079 "getAutomationWebhookCredential",
1080 &serde_json::json!({
1081 "id": agent_id,
1082 "automationId": automation_id,
1083 }),
1084 )?;
1085 if status == 200 {
1086 let parsed: WebhookCredential = serde_json::from_slice(&body).map_err(|_| {
1087 ShellError::Gateway("gateway credential was not understood".to_string())
1088 })?;
1089 let key = parsed.key.filter(|key| !key.is_empty());
1090 if parsed.url.is_empty() {
1091 return Err(ShellError::Gateway(
1092 "gateway credential was not understood".to_string(),
1093 ));
1094 }
1095 return Ok(match key {
1096 Some(key) => CredentialRead::Ready {
1097 url: parsed.url,
1098 key,
1099 },
1100 None => CredentialRead::MintFailed,
1101 });
1102 }
1103 if credential_is_missing(status, &body) {
1104 return Ok(CredentialRead::Missing);
1105 }
1106 Err(ShellError::Gateway(format!(
1107 "gateway credential status {status}"
1108 )))
1109}
1110
1111fn credential_is_missing(status: u16, body: &[u8]) -> bool {
1114 if status == 404 {
1115 return true;
1116 }
1117 if status != 500 {
1118 return false;
1119 }
1120 std::str::from_utf8(body)
1121 .map(|text| text.contains("Automation not found"))
1122 .unwrap_or(false)
1123}
1124
1125fn create_mirror_routine(
1127 gate: &GatewayConn,
1128 agent_id: &str,
1129 routine: &str,
1130) -> Result<Vec<AutomationCard>, ShellError> {
1131 let (status, body) = gateway_post(
1132 gate,
1133 "createAgentAutomation",
1134 &serde_json::json!({
1135 "id": agent_id,
1136 "spec": {
1137 "name": routine,
1138 "prompt": MIRROR_ROUTINE_PROMPT,
1139 "trigger": { "type": "webhook" },
1140 "isEnabled": false,
1141 },
1142 }),
1143 )?;
1144 if status != 200 {
1145 return Err(ShellError::Gateway(format!(
1146 "gateway create status {status}"
1147 )));
1148 }
1149 serde_json::from_slice(&body)
1150 .map_err(|_| ShellError::Gateway("gateway create was not understood".to_string()))
1151}
1152
1153fn delete_agent_automation(
1154 gate: &GatewayConn,
1155 agent_id: &str,
1156 automation_id: &str,
1157) -> Result<(), ShellError> {
1158 let (status, _) = gateway_post(
1159 gate,
1160 "deleteAgentAutomation",
1161 &serde_json::json!({ "id": agent_id, "automationId": automation_id }),
1162 )?;
1163 if status != 200 {
1164 return Err(ShellError::Gateway(format!(
1165 "gateway delete status {status}"
1166 )));
1167 }
1168 Ok(())
1169}
1170
1171pub(crate) struct Prepared {
1172 pub(crate) config: SessionConfig,
1173 pub(crate) nick: String,
1174 pub(crate) user_id: String,
1175 pub(crate) device_id: DeviceId,
1176 pub(crate) bearer: Zeroizing<String>,
1177 pub(crate) backend: Arc<dyn m4a_agent::Backend>,
1178 pub(crate) routine_url: Option<String>,
1179 pub(crate) routine_bearer: Option<String>,
1180}
1181
1182pub struct MachineClient {
1185 product_mode: bool,
1187 sessions: Vec<OpenedStore>,
1188 bus: Arc<LocalBus>,
1189 push: m4a_agent::engine::PushLink,
1190 agents_dir: Option<PathBuf>,
1194 gateway_file: Option<PathBuf>,
1195 gateway_token: Option<String>,
1196 last_agent_poll: std::time::Instant,
1197 agent_rescan_secs: u64,
1198 store_root: Option<PathBuf>,
1200 skip_nicks: Vec<String>,
1202 profile_note: bool,
1204 session_ids: Vec<(String, String)>,
1206 last_full_drive: std::time::Instant,
1208 pending_wakes: Vec<PendingWake>,
1210 ready_agents: HashSet<String>,
1213 _locks: Vec<File>,
1215 send_listener: Option<SendListener>,
1217 send_sock: Option<PathBuf>,
1218 send_queue: Vec<PendingSend>,
1220 homeserver_url: String,
1223 bus_next_user: i64,
1225 failed_opens: HashMap<String, std::time::Instant>,
1228}
1229
1230const LATE_OPEN_RETRY_SECS: u64 = 60;
1233
1234struct PendingSend {
1236 stream: SendStream,
1237 request: crate::SendRequest,
1238 started: std::time::Instant,
1239 room: Option<String>,
1240 peer: Option<String>,
1241}
1242
1243const SEND_JOIN_WAIT_SECS: u64 = 120;
1245
1246impl Drop for MachineClient {
1247 fn drop(&mut self) {
1248 if let Some(path) = self.send_sock.take() {
1249 let _ = std::fs::remove_file(path);
1250 }
1251 }
1252}
1253
1254#[derive(Debug, Default)]
1257pub struct TickReport {
1258 pub pushed: Vec<(String, String)>,
1260 pub joined: Vec<(String, String)>,
1262 pub errors: Vec<(String, String)>,
1264 pub sent: Vec<(String, String, crate::SendReply)>,
1266 pub alerts: Vec<(String, String)>,
1268}
1269
1270struct PendingWake {
1272 store_dir: PathBuf,
1273 agent_id: String,
1274}
1275
1276impl MachineClient {
1277 pub fn open(
1284 homeserver_url: &str,
1285 store_root: &Path,
1286 sessions: Vec<HostSession>,
1287 ) -> Result<Self, ShellError> {
1288 Self::open_with(homeserver_url, store_root, sessions, false)
1289 }
1290
1291 fn open_with(
1295 homeserver_url: &str,
1296 store_root: &Path,
1297 sessions: Vec<HostSession>,
1298 lenient: bool,
1299 ) -> Result<Self, ShellError> {
1300 if sessions.is_empty() {
1301 return Err(ShellError::SessionList("session list is empty".to_string()));
1302 }
1303 let mut prepared = Vec::with_capacity(sessions.len());
1304 let mut seen_ids = Vec::new();
1305 let mut seen_nicks = Vec::new();
1306 let mut locks = Vec::new();
1307 for session in sessions {
1308 let config = session.config(homeserver_url, store_root)?;
1309 if seen_ids.iter().any(|id: &String| id == config.session_id()) {
1310 return Err(ShellError::SessionList("duplicate session id".to_string()));
1311 }
1312 seen_ids.push(config.session_id().to_string());
1313 let lock = match lock_store(&config.store_dir()) {
1314 Ok(lock) => lock,
1315 Err(err) if lenient => {
1316 eprintln!("mail4agent: session {} left out: {err}", config.session_id());
1317 continue;
1318 }
1319 Err(err) => return Err(err),
1320 };
1321 let registered = match register_session(&config) {
1322 Ok(registered) => registered,
1323 Err(err) if lenient => {
1324 eprintln!("mail4agent: session {} left out: {err}", config.session_id());
1325 continue;
1326 }
1327 Err(err) => return Err(err),
1328 };
1329 if seen_nicks.iter().any(|nick: &String| nick.eq_ignore_ascii_case(®istered.nick)) {
1330 return Err(ShellError::SessionList("duplicate nick".to_string()));
1331 }
1332 seen_nicks.push(registered.nick.clone());
1333 locks.push(lock);
1334 prepared.push(Prepared {
1335 config,
1336 nick: registered.nick,
1337 user_id: registered.user_id,
1338 device_id: registered.device_id,
1339 bearer: registered.bearer,
1340 backend: registered.backend,
1341 routine_url: session.routine_url,
1342 routine_bearer: session.routine_bearer,
1343 });
1344 }
1345 if prepared.is_empty() {
1346 return Err(ShellError::SessionList(
1347 "no session could be opened".to_string(),
1348 ));
1349 }
1350 let server_name = server_name_of(&prepared[0].user_id)?;
1351 for item in &prepared[1..] {
1352 if server_name_of(&item.user_id)? != server_name {
1353 return Err(ShellError::SessionList(
1354 "sessions disagree on the homeserver name".to_string(),
1355 ));
1356 }
1357 }
1358 ensure_server_name(server_name)?;
1359 let bus = Arc::new(LocalBus::open(&prepared)?);
1360 let peers: Vec<(String, String)> = prepared
1361 .iter()
1362 .map(|item| (item.nick.to_string(), item.user_id.clone()))
1363 .collect();
1364 let mut opened = Vec::with_capacity(prepared.len());
1365 for item in &prepared {
1366 let server_name = server_name_of(&item.user_id)?;
1367 let mut store = OpenedStore::open(
1368 &item.config.store_dir(),
1369 item.config.session_id(),
1370 item.device_id.clone(),
1371 &item.user_id,
1372 server_name,
1373 Arc::clone(&item.backend),
1374 item.bearer.as_str(),
1375 )?;
1376 store.set_registered_nick(item.nick.to_string());
1377 store.attach_bus(Arc::clone(&bus));
1378 store.set_local_peers(
1379 peers
1380 .iter()
1381 .filter(|(nick, _)| *nick != item.nick)
1382 .cloned()
1383 .collect(),
1384 );
1385 store.set_wake(SessionWake {
1386 routine_url: item.routine_url.clone(),
1387 routine_bearer: item.routine_bearer.clone(),
1388 leader_sock: None,
1389 leader_cwd: None,
1390 });
1391 if item.routine_url.is_none() {
1392 attach_detected_chain(&mut store, &item.config, &item.nick);
1393 }
1394 store.drive(1_000, false)?;
1395 store.abandon_inflight_sync(1_000)?;
1396 opened.push(store);
1397 }
1398 let tokens: Vec<String> = prepared
1399 .iter()
1400 .map(|item| item.bearer.as_str().to_string())
1401 .collect();
1402 let product_mode = true;
1403 let push = m4a_agent::engine::PushLink::open(&prepared[0].config.homeserver_url, prepared[0].backend.keep_prefix(), tokens, product_mode)?;
1404 Ok(Self {
1405 product_mode,
1406 sessions: opened,
1407 bus,
1408 push,
1409 agents_dir: None,
1410 gateway_file: None,
1411 gateway_token: None,
1412 last_agent_poll: std::time::Instant::now(),
1413 agent_rescan_secs: 0,
1414 store_root: None,
1415 skip_nicks: Vec::new(),
1416 profile_note: false,
1417 session_ids: Vec::new(),
1418 last_full_drive: std::time::Instant::now(),
1419 pending_wakes: Vec::new(),
1420 ready_agents: HashSet::new(),
1421 _locks: locks,
1422 send_listener: None,
1423 send_sock: None,
1424 send_queue: Vec::new(),
1425 homeserver_url: homeserver_url.to_string(),
1426 bus_next_user: prepared.len() as i64 + 1,
1427 failed_opens: HashMap::new(),
1428 })
1429 }
1430
1431 pub fn open_session_dir(
1434 homeserver_url: &str,
1435 store_root: &Path,
1436 sessions_dir: &Path,
1437 ) -> Result<Self, ShellError> {
1438 let sessions = load_session_records(sessions_dir)?;
1439 Self::open(homeserver_url, store_root, sessions)
1440 }
1441
1442 pub fn from_env() -> Result<Self, ShellError> {
1454 Self::from_env_filtered(None)
1455 }
1456
1457 pub fn from_env_for(nick: &str) -> Result<Self, ShellError> {
1461 Self::from_env_filtered(Some(nick))
1462 }
1463
1464 fn from_env_filtered(only: Option<&str>) -> Result<Self, ShellError> {
1465 let homeserver_url = std::env::var(HOMESERVER_URL_ENV)
1466 .ok()
1467 .filter(|value| !value.is_empty())
1468 .ok_or(ShellError::HomeserverUrl)?;
1469 let store_root = std::env::var(STORE_ROOT_ENV)
1470 .ok()
1471 .filter(|value| !value.is_empty())
1472 .ok_or(ShellError::StoreRoot)?;
1473 let mut get = |key: &str| -> Option<String> {
1474 if key == GATEWAY_FILE_ENV {
1475 return std::env::var(GATEWAY_FILE_ENV)
1476 .ok()
1477 .filter(|value| !value.is_empty())
1478 .or_else(|| Some(under_home(DEFAULT_GATEWAY_FILE).to_string_lossy().into_owned()));
1479 }
1480 std::env::var(key).ok().filter(|value| !value.is_empty())
1481 };
1482 let agents_dir = resolve_agents_dir(&mut get);
1483 let (sessions, agents_dir) = if let Some(dir) = agents_dir {
1484 (load_web_agents(&dir, &mut get)?, Some(dir))
1485 } else {
1486 let sessions_dir = get(SESSIONS_DIR_ENV)
1487 .ok_or_else(|| ShellError::SessionList("session directory is unset".to_string()))?;
1488 (load_web_sessions(Path::new(&sessions_dir), &mut get)?, None)
1489 };
1490 let sessions = match only {
1491 Some(nick) => {
1492 let needle = crate::nick::lookup_nick(nick)?;
1493 let kept: Vec<HostSession> = sessions
1494 .into_iter()
1495 .filter(|session| {
1496 routine_name_for(session)
1497 .is_some_and(|name| name.eq_ignore_ascii_case(&needle))
1498 })
1499 .collect();
1500 if kept.is_empty() {
1501 return Err(ShellError::UnknownNick);
1502 }
1503 kept
1504 }
1505 None => sessions,
1506 };
1507 let gateway_file = get(GATEWAY_FILE_ENV).map(PathBuf::from);
1508 let gateway_token = get(GATEWAY_TOKEN_ENV);
1509 let agent_rescan_secs = get(AGENT_RESCAN_SECS_ENV)
1510 .and_then(|value| value.parse::<u64>().ok())
1511 .unwrap_or(0);
1512 let mut pending_wakes = Vec::new();
1513 let mut ready_agents = HashSet::new();
1514 for session in &sessions {
1515 let Some(agent_id) = session.agent_id.clone() else {
1516 continue;
1517 };
1518 if session.routine_url.is_some() && session.routine_bearer.is_some() {
1519 ready_agents.insert(agent_id);
1520 } else {
1521 pending_wakes.push(PendingWake {
1522 store_dir: crate::session_store_dir(
1523 Path::new(&store_root),
1524 &session.session_id,
1525 ),
1526 agent_id,
1527 });
1528 }
1529 }
1530 let options = WakeOptions::from_lookup(&mut get);
1531 let mut client = Self::open_with(&homeserver_url, Path::new(&store_root), sessions, true)?;
1532 client.store_root = Some(PathBuf::from(&store_root));
1533 client.skip_nicks = options.skip_nicks;
1534 client.profile_note = options.profile_note;
1535 client.session_ids = options.session_ids;
1536 client.pending_wakes = pending_wakes;
1537 client.ready_agents = ready_agents;
1538 client.agents_dir = agents_dir;
1539 client.gateway_file = gateway_file;
1540 client.gateway_token = gateway_token;
1541 client.agent_rescan_secs = agent_rescan_secs;
1542 client.last_agent_poll = std::time::Instant::now();
1543 Ok(client)
1544 }
1545
1546 pub fn poll_agent_directory(&mut self) -> Result<Vec<RoutineReport>, ShellError> {
1557 let Some(agents_dir) = self.agents_dir.clone() else {
1558 return Ok(Vec::new());
1559 };
1560 if self.agent_rescan_secs > 0 {
1561 let elapsed = self.last_agent_poll.elapsed().as_secs();
1562 if elapsed < self.agent_rescan_secs {
1563 return Ok(Vec::new());
1564 }
1565 }
1566 self.last_agent_poll = std::time::Instant::now();
1567 self.open_new_agent_sessions(&agents_dir);
1568 let gateway = self
1569 .gateway_file
1570 .clone()
1571 .unwrap_or_else(|| under_home(DEFAULT_GATEWAY_FILE));
1572 let Some(gate) = open_gateway(&gateway, self.gateway_token.as_deref())? else {
1573 return Ok(Vec::new());
1574 };
1575 let mut sessions = drop_skipped(load_agents_dir(&agents_dir)?, &self.skip_nicks);
1576 apply_session_ids(&mut sessions, &self.session_ids);
1577 let mut reports = Vec::new();
1578 for session in &sessions {
1579 let Some(agent_id) = session.agent_id.clone() else {
1580 continue;
1581 };
1582 let Some(nick) = routine_name_for(session) else {
1583 continue;
1584 };
1585 let folder_id = crate::nick::routine_folder_id(&nick).unwrap_or_default();
1586 if self.ready_agents.contains(&agent_id) {
1587 reports.push(RoutineReport {
1588 agent_id,
1589 nick,
1590 folder_id,
1591 status: WakeStatus::Ready,
1592 });
1593 continue;
1594 }
1595 let outcome = ensure_wake(&gate, &agent_id, &nick);
1596 log_outcome(&nick, &folder_id, &outcome);
1597 if self.profile_note && matches!(outcome, WakeOutcome::AwaitingBackend) {
1598 note_profile(&gate, &agents_dir, &agent_id, &nick);
1599 }
1600 if let WakeOutcome::Ready { url, key } = &outcome {
1601 if let Some(root) = &self.store_root {
1602 save_wake(root, &session.session_id, &folder_id, url, key);
1603 }
1604 if let Some(index) = self
1605 .pending_wakes
1606 .iter()
1607 .position(|pending| pending.agent_id == agent_id)
1608 {
1609 let pending = self.pending_wakes.remove(index);
1610 if let Some(store) = self
1611 .sessions
1612 .iter_mut()
1613 .find(|store| store.store_dir() == pending.store_dir)
1614 {
1615 store.set_wake(SessionWake {
1616 routine_url: Some(url.clone()),
1617 routine_bearer: Some(key.clone()),
1618 leader_sock: None,
1619 leader_cwd: None,
1620 });
1621 }
1622 }
1623 self.ready_agents.insert(agent_id.clone());
1624 if let Some(user_id) = self
1625 .sessions
1626 .iter()
1627 .find(|store| {
1628 store
1629 .nick()
1630 .map(|n| n.eq_ignore_ascii_case(&nick))
1631 .unwrap_or(false)
1632 })
1633 .map(|store| store.user_id().to_string())
1634 {
1635 self.announce_peer_joined(&nick, &user_id);
1636 }
1637 }
1638 reports.push(RoutineReport {
1639 agent_id,
1640 nick,
1641 folder_id,
1642 status: outcome.status(),
1643 });
1644 }
1645 Ok(reports)
1646 }
1647
1648 fn open_new_agent_sessions(&mut self, agents_dir: &Path) {
1654 let Some(store_root) = self.store_root.clone() else {
1655 return;
1656 };
1657 let found = match load_agents_dir(agents_dir) {
1658 Ok(found) => found,
1659 Err(err) => {
1660 eprintln!("mail4agent: agent rescan: {err}");
1661 return;
1662 }
1663 };
1664 let held: Vec<(PathBuf, Option<String>)> = self
1665 .sessions
1666 .iter()
1667 .map(|store| {
1668 (
1669 store.store_dir().to_path_buf(),
1670 store.nick().map(str::to_string),
1671 )
1672 })
1673 .collect();
1674 let fresh = unopened_agent_sessions(
1675 found,
1676 &self.skip_nicks,
1677 &self.session_ids,
1678 &store_root,
1679 &held,
1680 );
1681 for (session, nick) in fresh {
1682 if let Some(failed_at) = self.failed_opens.get(&session.session_id) {
1683 if failed_at.elapsed().as_secs() < LATE_OPEN_RETRY_SECS {
1684 continue;
1685 }
1686 }
1687 let session_id = session.session_id.clone();
1688 match self.open_late_session(&store_root, session, &nick) {
1689 Ok(()) => {
1690 self.failed_opens.remove(&session_id);
1691 println!("mail4agent: session {nick} registered and open");
1692 }
1693 Err(err) => {
1694 eprintln!("mail4agent: session {nick} not opened: {err}");
1695 self.failed_opens
1696 .insert(session_id, std::time::Instant::now());
1697 }
1698 }
1699 }
1700 }
1701
1702 fn open_late_session(
1703 &mut self,
1704 store_root: &Path,
1705 mut session: HostSession,
1706 nick: &str,
1707 ) -> Result<(), ShellError> {
1708 if session.routine_url.is_none() || session.routine_bearer.is_none() {
1709 if let Some(folder) = crate::nick::routine_folder_id(nick) {
1710 if let Some((url, key)) = load_wake(store_root, &session.session_id, &folder) {
1711 session.routine_url = Some(url);
1712 session.routine_bearer = Some(key);
1713 }
1714 }
1715 }
1716 let config = session.config(&self.homeserver_url, store_root)?;
1717 let lock = lock_store(&config.store_dir())?;
1718 let registered = register_session(&config)?;
1719 if let Some(first) = self.sessions.first() {
1720 if server_name_of(®istered.user_id)? != server_name_of(first.user_id())? {
1721 return Err(ShellError::SessionList(
1722 "sessions disagree on the homeserver name".to_string(),
1723 ));
1724 }
1725 }
1726 let item = Prepared {
1727 config,
1728 nick: registered.nick,
1729 user_id: registered.user_id,
1730 device_id: registered.device_id,
1731 bearer: registered.bearer,
1732 backend: registered.backend,
1733 routine_url: session.routine_url.clone(),
1734 routine_bearer: session.routine_bearer.clone(),
1735 };
1736 self.bus.seed(self.bus_next_user, &item)?;
1737 self.bus_next_user += 1;
1738 let mut store = OpenedStore::open(
1739 &item.config.store_dir(),
1740 item.config.session_id(),
1741 item.device_id.clone(),
1742 &item.user_id,
1743 server_name_of(&item.user_id)?,
1744 Arc::clone(&item.backend),
1745 item.bearer.as_str(),
1746 )?;
1747 store.set_registered_nick(item.nick.to_string());
1748 store.attach_bus(Arc::clone(&self.bus));
1749 store.set_wake(SessionWake {
1750 routine_url: item.routine_url.clone(),
1751 routine_bearer: item.routine_bearer.clone(),
1752 leader_sock: None,
1753 leader_cwd: None,
1754 });
1755 if item.routine_url.is_none() {
1756 attach_detected_chain(&mut store, &item.config, &item.nick);
1757 }
1758 store.drive(1_000, false)?;
1759 store.abandon_inflight_sync(1_000)?;
1760 self.sessions.push(store);
1761 self._locks.push(lock);
1762 self.refresh_local_peers();
1763 let tokens: Vec<String> = self
1764 .sessions
1765 .iter()
1766 .map(|store| store.device_bearer().to_string())
1767 .collect();
1768 match m4a_agent::engine::PushLink::open(&self.homeserver_url, self.sessions.first().is_some_and(|s| s.keep_prefix()), tokens, self.product_mode) {
1769 Ok(push) => {
1770 let old = std::mem::replace(&mut self.push, push);
1771 for (recipient, event) in old.drain() {
1772 if let Some(store) = self
1773 .sessions
1774 .iter_mut()
1775 .find(|store| store.user_id() == recipient)
1776 {
1777 store.record_push(event);
1778 }
1779 }
1780 }
1781 Err(err) => {
1782 eprintln!("mail4agent: push socket not reopened for {nick}: {err}");
1783 }
1784 }
1785 if let Some(agent_id) = session.agent_id.clone() {
1786 if item.routine_url.is_some() && item.routine_bearer.is_some() {
1787 self.ready_agents.insert(agent_id);
1788 self.announce_peer_joined(nick, &item.user_id);
1789 } else {
1790 self.ready_agents.remove(&agent_id);
1792 if !self.pending_wakes.iter().any(|p| p.agent_id == agent_id) {
1793 self.pending_wakes.push(PendingWake {
1794 store_dir: item.config.store_dir(),
1795 agent_id,
1796 });
1797 }
1798 }
1799 }
1800 Ok(())
1801 }
1802
1803 fn refresh_local_peers(&mut self) {
1805 let peers: Vec<(String, String)> = self
1806 .sessions
1807 .iter()
1808 .filter_map(|store| Some((store.nick()?.to_string(), store.user_id().to_string())))
1809 .collect();
1810 for store in self.sessions.iter_mut() {
1811 let own = store.nick().unwrap_or("").to_string();
1812 store.set_local_peers(
1813 peers
1814 .iter()
1815 .filter(|(nick, _)| *nick != own)
1816 .cloned()
1817 .collect(),
1818 );
1819 }
1820 }
1821
1822 fn announce_peer_joined(&self, joined_nick: &str, joined_user_id: &str) {
1826 let body = serde_json::json!({
1827 "kind": "peer_joined",
1828 "nick": joined_nick,
1829 "user_id": joined_user_id,
1830 });
1831 for store in &self.sessions {
1832 let Some(nick) = store.nick() else {
1833 continue;
1834 };
1835 if nick.eq_ignore_ascii_case(joined_nick) {
1836 continue;
1837 }
1838 if !store.has_routine() {
1839 continue;
1840 }
1841 let Some((url, bearer)) = store.routine_target() else {
1842 continue;
1843 };
1844 match crate::post_routine_json(&url, &body, bearer.as_deref()) {
1845 Ok(()) => {
1846 eprintln!("mail4agent: peer_joined {joined_nick} -> {nick} status=200")
1847 }
1848 Err(err) => {
1849 eprintln!("mail4agent: peer_joined {joined_nick} -> {nick}: {err}")
1850 }
1851 }
1852 }
1853 }
1854
1855 pub fn pending_wake_agents(&self) -> Vec<String> {
1857 self.pending_wakes
1858 .iter()
1859 .map(|pending| pending.agent_id.clone())
1860 .collect()
1861 }
1862
1863 pub fn set_local_delivery(&self, enabled: bool) {
1874 self.bus.set_local_only(enabled);
1875 }
1876
1877 pub fn homeserver_hits(&self) -> u64 {
1881 self.bus.hits()
1882 }
1883
1884 pub fn session_mut(&mut self, name_or_nick: &str) -> Result<&mut OpenedStore, ShellError> {
1886 let needle = crate::nick::lookup_nick(name_or_nick)?;
1887 self.sessions
1888 .iter_mut()
1889 .find(|store| {
1890 store
1891 .nick()
1892 .is_some_and(|nick| nick.eq_ignore_ascii_case(&needle))
1893 })
1894 .ok_or(ShellError::UnknownNick)
1895 }
1896
1897 pub fn store_dir(&self, name_or_nick: &str) -> Result<PathBuf, ShellError> {
1899 let needle = crate::nick::lookup_nick(name_or_nick)?;
1900 self.sessions
1901 .iter()
1902 .find(|store| {
1903 store
1904 .nick()
1905 .is_some_and(|nick| nick.eq_ignore_ascii_case(&needle))
1906 })
1907 .map(|store| store.store_dir().to_path_buf())
1908 .ok_or(ShellError::UnknownNick)
1909 }
1910
1911 pub fn deliver_pushed(&mut self) -> usize {
1915 let batch = self.push.drain();
1916 let mut delivered = 0;
1917 for (recipient, event) in batch {
1918 let Some(session) = self
1919 .sessions
1920 .iter_mut()
1921 .find(|store| store.user_id() == recipient)
1922 else {
1923 continue;
1924 };
1925 session.record_push(event);
1926 delivered += 1;
1927 }
1928 delivered
1929 }
1930
1931 pub fn tick(&mut self, now_ms: i64, full_drive_secs: u64) -> TickReport {
1941 let mut report = TickReport::default();
1942 let mut due: Vec<usize> = Vec::new();
1943 for (recipient, event) in self.push.drain() {
1944 let Some(index) = self
1945 .sessions
1946 .iter()
1947 .position(|store| store.user_id() == recipient)
1948 else {
1949 continue;
1950 };
1951 let nick = self.sessions[index].nick().unwrap_or("").to_string();
1952 report.pushed.push((nick, event.event_id.clone()));
1953 self.sessions[index].record_push(event);
1954 if !due.contains(&index) {
1955 due.push(index);
1956 }
1957 }
1958 let full = self.last_full_drive.elapsed().as_secs() >= full_drive_secs;
1959 if full {
1960 self.last_full_drive = std::time::Instant::now();
1961 }
1962 for index in 0..self.sessions.len() {
1963 let pushed = due.contains(&index);
1964 if !pushed && !full {
1965 continue;
1966 }
1967 let store = &mut self.sessions[index];
1968 let nick = store.nick().unwrap_or("").to_string();
1969 let drove = store.drive(now_ms, pushed);
1970 for text in store.take_security_alerts() {
1971 report.alerts.push((nick.clone(), text));
1972 }
1973 if let Err(err) = drove {
1974 report.errors.push((nick.clone(), err.to_string()));
1975 continue;
1976 }
1977 match store.accept_direct_invites(now_ms) {
1978 Ok(joined) => {
1979 for room in joined {
1980 report.joined.push((nick.clone(), room));
1981 }
1982 }
1983 Err(err) => report.errors.push((nick, err.to_string())),
1984 }
1985 }
1986 self.serve_sends(now_ms, &mut report);
1987 report
1988 }
1989
1990 pub fn listen_for_sends(&mut self, path: &Path) -> Result<(), ShellError> {
1998 let listener = SendListener::bind(path).map_err(|err| {
1999 if err.kind() == std::io::ErrorKind::AlreadyExists {
2000 ShellError::SessionList(
2001 "another client already listens on the send socket".to_string(),
2002 )
2003 } else {
2004 ShellError::Io(err)
2005 }
2006 })?;
2007 listener.set_nonblocking(true)?;
2008 self.send_listener = Some(listener);
2009 self.send_sock = Some(path.to_path_buf());
2010 Ok(())
2011 }
2012
2013 pub fn listen_for_sends_from_env(&mut self) -> Result<PathBuf, ShellError> {
2016 let root = self.store_root.clone().ok_or(ShellError::StoreRoot)?;
2017 let path = crate::send_sock_path(
2018 |key| std::env::var(key).ok().filter(|value| !value.is_empty()),
2019 &root,
2020 );
2021 self.listen_for_sends(&path)?;
2022 Ok(path)
2023 }
2024
2025 pub fn send_blocking(
2029 &mut self,
2030 as_nick: &str,
2031 to: &str,
2032 text: &str,
2033 wait: std::time::Duration,
2034 ) -> crate::SendReply {
2035 let started = std::time::Instant::now();
2036 let mut room = None;
2037 let mut peer = None;
2038 loop {
2039 let now = now_ms();
2040 match self.try_send(as_nick, to, text, now, &mut room, &mut peer) {
2041 Some(reply) => return reply,
2042 None if started.elapsed() >= wait => {
2043 return crate::SendReply {
2044 room,
2045 ..crate::SendReply::failed(format!("{to} has not joined the DM yet"))
2046 }
2047 }
2048 None => {
2049 if let Ok(store) = self.session_mut(as_nick) {
2050 let _ = store.drive(now, false);
2051 }
2052 std::thread::sleep(std::time::Duration::from_millis(500));
2053 }
2054 }
2055 }
2056 }
2057
2058 fn try_send(
2060 &mut self,
2061 as_nick: &str,
2062 to: &str,
2063 text: &str,
2064 now_ms: i64,
2065 room: &mut Option<String>,
2066 peer: &mut Option<String>,
2067 ) -> Option<crate::SendReply> {
2068 let store = match self.session_mut(as_nick) {
2069 Ok(store) => store,
2070 Err(_) => {
2071 return Some(crate::SendReply::failed(format!(
2072 "{as_nick} is not a session on this client"
2073 )))
2074 }
2075 };
2076 if peer.is_none() {
2077 let mut last_err = None;
2080 for attempt in 0..3 {
2081 if attempt > 0 {
2082 let _ = store.drive(now_ms, false);
2083 }
2084 match store.find_nick(to, now_ms) {
2085 Ok(found) => {
2086 *peer = Some(found.user_id);
2087 last_err = None;
2088 break;
2089 }
2090 Err(err) => last_err = Some(err),
2091 }
2092 }
2093 if let Some(err) = last_err {
2094 return Some(crate::SendReply::failed(format!("find {to}: {err}")));
2095 }
2096 }
2097 if room.is_none() {
2098 match store.ensure_dm(to, now_ms) {
2099 Ok(room_id) => *room = Some(room_id),
2100 Err(err) => return Some(crate::SendReply::failed(format!("open DM: {err}"))),
2101 }
2102 }
2103 let (room_id, peer_id) = (room.clone()?, peer.clone()?);
2104 if !store.member_joined(&room_id, &peer_id) {
2105 return None;
2106 }
2107 match store.write_to_nick(to, text, now_ms) {
2108 Ok(room_id) => {
2109 let event_id = store
2110 .texts()
2111 .into_iter()
2112 .rev()
2113 .find(|row| row.room_id == room_id && row.body == text)
2114 .and_then(|row| row.event_id);
2115 Some(crate::SendReply {
2116 ok: true,
2117 room: Some(room_id),
2118 event_id,
2119 error: None,
2120 })
2121 }
2122 Err(err) => Some(crate::SendReply {
2123 room: Some(room_id),
2124 ..crate::SendReply::failed(format!("send: {err}"))
2125 }),
2126 }
2127 }
2128
2129 fn serve_sends(&mut self, now_ms: i64, report: &mut TickReport) {
2130 let mut cmds: Vec<(crate::ipc::SendStream, crate::CmdRequest)> = Vec::new();
2131 if let Some(listener) = &self.send_listener {
2132 loop {
2133 match listener.accept() {
2134 Ok(mut stream) => match crate::send::read_incoming(&mut stream) {
2135 Ok(crate::send::Incoming::Send(request)) => self.send_queue.push(PendingSend {
2136 stream,
2137 request,
2138 started: std::time::Instant::now(),
2139 room: None,
2140 peer: None,
2141 }),
2142 Ok(crate::send::Incoming::Cmd(cmd)) => cmds.push((stream, cmd)),
2143 Err(err) => {
2144 crate::send::write_reply(&mut stream, &crate::SendReply::failed(err))
2145 }
2146 },
2147 Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => break,
2148 Err(_) => break,
2149 }
2150 }
2151 }
2152 for (mut stream, cmd) in cmds {
2153 let reply = match self.sessions.iter_mut().find(|store| {
2154 store.nick().is_some_and(|n| n.eq_ignore_ascii_case(cmd.as_nick.trim()))
2155 }) {
2156 Some(store) => store.run_command(&cmd, now_ms),
2157 None => crate::CmdReply::failed(format!("{} is not a session on this client", cmd.as_nick)),
2158 };
2159 crate::send::write_cmd_reply(&mut stream, &reply);
2160 }
2161 let queue = std::mem::take(&mut self.send_queue);
2162 for mut pending in queue {
2163 let (as_nick, to, text) = (
2164 pending.request.as_nick.clone(),
2165 pending.request.to.clone(),
2166 pending.request.text.clone(),
2167 );
2168 let outcome = self.try_send(
2169 &as_nick,
2170 &to,
2171 &text,
2172 now_ms,
2173 &mut pending.room,
2174 &mut pending.peer,
2175 );
2176 let reply = match outcome {
2177 Some(reply) => reply,
2178 None if pending.started.elapsed().as_secs() >= SEND_JOIN_WAIT_SECS => {
2179 crate::SendReply {
2180 room: pending.room.clone(),
2181 ..crate::SendReply::failed(format!("{to} has not joined the DM yet"))
2182 }
2183 }
2184 None => {
2185 self.send_queue.push(pending);
2186 continue;
2187 }
2188 };
2189 crate::send::write_reply(&mut pending.stream, &reply);
2190 report.sent.push((as_nick, to, reply));
2191 }
2192 }
2193
2194 pub fn wake_log(&self) -> Vec<(String, crate::WakeAttempt)> {
2197 self.sessions
2198 .iter()
2199 .flat_map(|store| {
2200 let nick = store.nick().unwrap_or("").to_string();
2201 store
2202 .wake_log()
2203 .iter()
2204 .cloned()
2205 .map(move |attempt| (nick.clone(), attempt))
2206 })
2207 .collect()
2208 }
2209
2210 pub fn holds(&self, name_or_nick: &str) -> bool {
2212 let Ok(needle) = crate::nick::lookup_nick(name_or_nick) else {
2213 return false;
2214 };
2215 self.sessions.iter().any(|store| {
2216 store
2217 .nick()
2218 .is_some_and(|nick| nick.eq_ignore_ascii_case(&needle))
2219 })
2220 }
2221}
2222
2223fn server_name_of(mxid: &str) -> Result<&str, ShellError> {
2224 mxid.split_once(':')
2225 .map(|(_, server)| server)
2226 .filter(|server| !server.is_empty())
2227 .ok_or_else(|| ShellError::Register("user id has no server".to_string()))
2228}
2229
2230
2231fn attach_detected_chain(store: &mut OpenedStore, config: &SessionConfig, nick: &str) {
2237 if let Some((session, chain)) = crate::provider::chain::detected_web_chain(
2238 config.session_id(),
2239 nick,
2240 &config.store_root,
2241 ) {
2242 store.set_wake_chain(session, chain);
2243 }
2244}
2245
2246#[cfg(test)]
2247mod tests {
2248 use super::*;
2249
2250 #[test]
2251 fn rescan_opens_only_bots_without_a_session() {
2252 let root = Path::new("/tmp/m4a-rescan-test");
2253 let agent = |name: &str, id: &str| {
2254 let mut session = HostSession::new(name, id);
2255 session.agent_id = Some(id.to_string());
2256 session
2257 };
2258 let found = vec![
2259 agent("alice", "a-hatch"),
2260 agent("m4a-proba2", "a-proba2"),
2261 agent("skipme", "a-skip"),
2262 agent("aliased", "a-alias"),
2263 ];
2264 let held = vec![
2265 (crate::session_store_dir(root, "a-hatch"), Some("alice".to_string())),
2266 (crate::session_store_dir(root, "old-alias"), None),
2267 ];
2268 let fresh = unopened_agent_sessions(
2269 found,
2270 &["skipme".to_string()],
2271 &[("a-alias".to_string(), "old-alias".to_string())],
2272 root,
2273 &held,
2274 );
2275 let nicks: Vec<&str> = fresh.iter().map(|(_, nick)| nick.as_str()).collect();
2276 assert_eq!(nicks, vec!["m4a-proba2"]);
2277 assert_eq!(fresh[0].0.session_id, "a-proba2");
2278 }
2279
2280 #[test]
2281 fn session_directory_is_name_and_id_only() {
2282 let dir = std::env::temp_dir().join(format!(
2283 "m4a-sessions-{}-{}",
2284 std::process::id(),
2285 std::time::SystemTime::now()
2286 .duration_since(std::time::UNIX_EPOCH)
2287 .expect("clock")
2288 .as_nanos()
2289 ));
2290 std::fs::create_dir_all(&dir).expect("dir");
2291 std::fs::write(
2292 dir.join("alice.json"),
2293 r#"{"bot_name":"Alice","session_id":"web-alice"}"#,
2294 )
2295 .expect("write");
2296 std::fs::write(
2297 dir.join("chief.json"),
2298 "{\"bot_name\":\"Привет мир\",\"session_id\":\"web-chief\"}",
2299 )
2300 .expect("write");
2301 let loaded = load_session_records(&dir).expect("records");
2302 assert_eq!(loaded.len(), 2);
2303 assert_eq!(loaded[1].bot_name, "Привет мир");
2304 assert_eq!(loaded[0].bot_name, "Alice");
2305 assert!(loaded[1].agent_id.is_none());
2306 assert!(loaded[0].agent_id.is_none());
2307 assert!(loaded[1].routine_url.is_none());
2308 assert!(loaded[1].routine_bearer.is_none());
2309 assert!(loaded[1].invite.is_none());
2310
2311 std::fs::write(
2312 dir.join("leaked.json"),
2313 r#"{"bot_name":"Courier","session_id":"web-courier","routine_url":"http://127.0.0.1/hook","routine_bearer":"not-a-file"}"#,
2314 )
2315 .expect("write");
2316 let refused = load_session_records(&dir).expect_err("webhook file");
2317 let text = refused.to_string();
2318 assert!(!text.contains("not-a-file"));
2319 assert!(!text.contains("127.0.0.1/hook"));
2320 let _ = std::fs::remove_dir_all(&dir);
2321 }
2322
2323 #[test]
2324 fn web_client_does_not_copy_one_routine_onto_every_session_and_node_cli_refuses_it() {
2325 let dir = std::env::temp_dir().join(format!(
2326 "m4a-split-{}-{}",
2327 std::process::id(),
2328 std::time::SystemTime::now()
2329 .duration_since(std::time::UNIX_EPOCH)
2330 .expect("clock")
2331 .as_nanos()
2332 ));
2333 std::fs::create_dir_all(&dir).expect("dir");
2334 let record = dir.join("alice.json");
2335 let body = r#"{"bot_name":"Alice","session_id":"web-alice"}"#;
2336 std::fs::write(&record, body).expect("write");
2337 let routine = "http://127.0.0.1:9/routine";
2338 let bearer = "host-injected-bearer";
2339 let sessions = load_web_sessions(&dir, |key| match key {
2340 crate::ROUTINE_URL_ENV => Some(routine.to_string()),
2341 crate::ROUTINE_BEARER_ENV => Some(bearer.to_string()),
2342 crate::LEADER_SOCK_ENV => Some("/tmp/leader.sock".to_string()),
2343 _ => None,
2344 })
2345 .expect("web sessions");
2346 assert_eq!(sessions.len(), 1);
2347 assert!(sessions[0].routine_url.is_none());
2348 assert!(sessions[0].routine_bearer.is_none());
2349 assert_eq!(std::fs::read_to_string(&record).expect("reread"), body);
2350 let wake = crate::SessionWake::web_from_lookup(|key| match key {
2351 crate::ROUTINE_URL_ENV => Some(routine.to_string()),
2352 crate::LEADER_SOCK_ENV => Some("/tmp/leader.sock".to_string()),
2353 _ => None,
2354 });
2355 assert!(wake.leader_sock.is_none());
2356 assert!(wake.leader_cwd.is_none());
2357
2358 std::fs::write(
2359 dir.join("leaked.json"),
2360 r#"{"bot_name":"Courier","session_id":"web-courier","routine_url":"http://127.0.0.1:9/from-file","routine_bearer":"file-bearer"}"#,
2361 )
2362 .expect("leak");
2363 let refused = load_web_sessions(&dir, |_| None).expect_err("json routine");
2364 let text = refused.to_string();
2365 assert!(!text.contains("from-file"));
2366 assert!(!text.contains("file-bearer"));
2367
2368 let node = crate::OpenedStore::connect_node_from_lookup(
2369 |key| match key {
2370 crate::ROUTINE_URL_ENV => Some(routine.to_string()),
2371 crate::ROUTINE_BEARER_ENV => Some(bearer.to_string()),
2372 crate::LEADER_SOCK_ENV => Some("/tmp/leader.sock".to_string()),
2373 _ => None,
2374 },
2375 None,
2376 );
2377 let err = match node {
2378 Ok(_) => panic!("node cli accepted a routine url"),
2379 Err(err) => err,
2380 };
2381 assert!(matches!(err, crate::ShellError::NodeRoutine));
2382 let text = err.to_string();
2383 assert!(!text.contains(routine));
2384 assert!(!text.contains(bearer));
2385 assert!(!text.contains("leader.sock"));
2386
2387 let node = crate::SessionWake::node_from_lookup(|key| match key {
2388 crate::LEADER_SOCK_ENV => Some("/tmp/node-leader.sock".to_string()),
2389 crate::LEADER_CWD_ENV => Some("/tmp/node".to_string()),
2390 _ => None,
2391 })
2392 .expect("node leader");
2393 assert!(node.routine_url.is_none());
2394 assert!(node.routine_bearer.is_none());
2395 assert_eq!(
2396 node.leader_sock.as_deref(),
2397 Some(std::path::Path::new("/tmp/node-leader.sock"))
2398 );
2399 assert_eq!(node.leader_cwd.as_deref(), Some("/tmp/node"));
2400 let _ = std::fs::remove_dir_all(&dir);
2401 }
2402
2403 #[test]
2404 fn mirrors_follow_the_nick_folder_and_keys_wait_for_the_bot_routine() {
2405 use std::collections::HashSet;
2406 use std::io::{Read, Write};
2407 use std::net::TcpListener;
2408 use std::sync::atomic::{AtomicBool, Ordering};
2409 use std::sync::{Arc, Mutex};
2410 use std::thread;
2411
2412 struct Card {
2413 agent: String,
2414 id: String,
2415 name: String,
2416 trigger: &'static str,
2417 enabled: bool,
2418 }
2419 #[derive(Default)]
2420 struct Gateway {
2421 creates: usize,
2422 deletes: usize,
2423 cards: Vec<Card>,
2424 backend: HashSet<(String, String)>,
2426 clash: HashSet<String>,
2428 agents_seen: HashSet<String>,
2429 profiles: Vec<(String, serde_json::Value)>,
2430 }
2431
2432 fn read_http(sock: &mut std::net::TcpStream) -> Option<(String, Vec<u8>)> {
2433 let _ = sock.set_read_timeout(Some(std::time::Duration::from_secs(2)));
2434 let mut buf = Vec::new();
2435 let mut tmp = [0u8; 2048];
2436 loop {
2437 let n = sock.read(&mut tmp).unwrap_or(0);
2438 if n == 0 {
2439 return None;
2440 }
2441 buf.extend_from_slice(&tmp[..n]);
2442 let Some(end) = buf.windows(4).position(|w| w == b"\r\n\r\n") else {
2443 continue;
2444 };
2445 let headers = String::from_utf8_lossy(&buf[..end]).to_string();
2446 let length = headers
2447 .lines()
2448 .find_map(|line| {
2449 let (name, value) = line.split_once(':')?;
2450 name.eq_ignore_ascii_case("content-length")
2451 .then(|| value.trim().parse::<usize>().ok())
2452 .flatten()
2453 })
2454 .unwrap_or(0);
2455 if buf.len() >= end + 4 + length {
2456 return Some((headers, buf[end + 4..end + 4 + length].to_vec()));
2457 }
2458 }
2459 }
2460 fn reply(status: &str, body: &serde_json::Value) -> Vec<u8> {
2461 let payload = serde_json::to_vec(body).expect("json");
2462 let mut out = format!(
2463 "HTTP/1.1 {status}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
2464 payload.len()
2465 )
2466 .into_bytes();
2467 out.extend(payload);
2468 out
2469 }
2470 fn cards_of(gate: &Gateway, agent: &str) -> serde_json::Value {
2471 serde_json::Value::Array(
2472 gate.cards
2473 .iter()
2474 .filter(|card| card.agent == agent)
2475 .map(|card| {
2476 serde_json::json!({
2477 "id": card.id,
2478 "name": card.name,
2479 "trigger": {"type": card.trigger},
2480 "isEnabled": card.enabled,
2481 })
2482 })
2483 .collect(),
2484 )
2485 }
2486
2487 let listener = TcpListener::bind("127.0.0.1:0").expect("bind");
2488 listener.set_nonblocking(true).expect("nonblocking");
2489 let port = listener.local_addr().expect("addr").port();
2490 let state = Arc::new(Mutex::new(Gateway::default()));
2491 {
2492 let mut gate = state.lock().expect("gate");
2493 gate.backend.insert(("agent-h".into(), "alice".into()));
2494 gate.clash.insert("agent-x".into());
2495 gate.cards.push(Card {
2496 agent: "agent-k".into(),
2497 id: "cron-bot".into(),
2498 name: "cron_bot".into(),
2499 trigger: "cron",
2500 enabled: true,
2501 });
2502 }
2503 let shared = Arc::clone(&state);
2504 let token = "gw-test-token";
2505 let done = Arc::new(AtomicBool::new(false));
2506 let flag = Arc::clone(&done);
2507 let server = thread::spawn(move || {
2508 while !flag.load(Ordering::Relaxed) {
2509 let (mut sock, _) = match listener.accept() {
2510 Ok(pair) => pair,
2511 Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
2512 thread::sleep(std::time::Duration::from_millis(5));
2513 continue;
2514 }
2515 Err(_) => break,
2516 };
2517 let _ = sock.set_nonblocking(false);
2518 let Some((headers, body)) = read_http(&mut sock) else {
2519 continue;
2520 };
2521 let lower = headers.to_ascii_lowercase();
2522 let authorized = lower.contains(&format!("authorization: bearer {token}"));
2523 let request: serde_json::Value = serde_json::from_slice(&body).unwrap_or_default();
2524 let agent = request["id"].as_str().unwrap_or("").to_string();
2525 let mut gate = shared.lock().expect("gate");
2526 gate.agents_seen.insert(agent.clone());
2527 let response = if !authorized {
2528 reply("401 Unauthorized", &serde_json::json!({}))
2529 } else if lower.starts_with("post /api/getagentautomations ") {
2530 reply("200 OK", &cards_of(&gate, &agent))
2531 } else if lower.starts_with("post /api/createagentautomation ") {
2532 let spec = &request["spec"];
2533 assert_eq!(spec["trigger"]["type"], "webhook");
2534 assert_eq!(spec["isEnabled"], false, "mirror must be disabled");
2535 let name = spec["name"].as_str().unwrap_or("").to_string();
2536 let prompt = spec["prompt"].as_str().unwrap_or("");
2537 assert!(!prompt.is_empty());
2538 assert!(!prompt.contains("http"));
2539 gate.creates += 1;
2540 let folder = crate::nick::routine_folder_id(&name).expect("slug");
2541 let id = if gate.clash.contains(&agent) {
2542 format!("{folder}-2")
2543 } else {
2544 folder
2545 };
2546 gate.cards.push(Card {
2547 agent: agent.clone(),
2548 id,
2549 name,
2550 trigger: "webhook",
2551 enabled: false,
2552 });
2553 reply("200 OK", &cards_of(&gate, &agent))
2554 } else if lower.starts_with("post /api/updateagent ") {
2555 gate.profiles
2556 .push((agent.clone(), request["profile"].clone()));
2557 reply("200 OK", &serde_json::json!({"id": agent}))
2558 } else if lower.starts_with("post /api/deleteagentautomation ") {
2559 let id = request["automationId"].as_str().unwrap_or("");
2560 gate.deletes += 1;
2561 gate.cards
2562 .retain(|card| !(card.agent == agent && card.id == id));
2563 reply("200 OK", &cards_of(&gate, &agent))
2564 } else if lower.starts_with("post /api/getautomationwebhookcredential ") {
2565 let id = request["automationId"].as_str().unwrap_or("").to_string();
2566 let local = gate
2567 .cards
2568 .iter()
2569 .find(|card| card.agent == agent && card.id == id);
2570 match local {
2571 None => reply(
2572 "500 Internal Server Error",
2573 &serde_json::json!({"error": format!("Automation not found: {id}")}),
2574 ),
2575 Some(card) if card.trigger != "webhook" => reply(
2576 "500 Internal Server Error",
2577 &serde_json::json!({"error": "Automation is not webhook-triggered"}),
2578 ),
2579 Some(_) => {
2580 let minted = gate.backend.contains(&(agent.clone(), id.clone()));
2581 reply(
2582 "200 OK",
2583 &serde_json::json!({
2584 "url": format!("https://backend.invalid/automations/webhook/{agent}-{id}"),
2585 "key": minted.then(|| format!("key-{agent}-{id}")),
2586 }),
2587 )
2588 }
2589 }
2590 } else {
2591 reply("404 Not Found", &serde_json::json!({}))
2592 };
2593 drop(gate);
2594 let _ = sock.write_all(&response);
2595 }
2596 });
2597
2598 let dir = std::env::temp_dir().join(format!(
2599 "m4a-mirror-{}-{}",
2600 std::process::id(),
2601 std::time::SystemTime::now()
2602 .duration_since(std::time::UNIX_EPOCH)
2603 .expect("clock")
2604 .as_nanos()
2605 ));
2606 let agents = dir.join("agents");
2607 let store_root = dir.join("stores");
2608 let describe = |id: &str| {
2609 if id == "agent-c" {
2610 "about the chief"
2611 } else {
2612 ""
2613 }
2614 };
2615 for (id, name) in [
2616 ("agent-h", "Alice"),
2617 ("agent-c", "Привет мир"),
2618 ("agent-s", "Свой браузер"),
2619 ("agent-x", "Clash Bot"),
2620 ("agent-k", "Cron Bot"),
2621 ] {
2622 std::fs::create_dir_all(agents.join(id)).expect("agent dir");
2623 std::fs::write(
2624 agents.join(id).join("profile.json"),
2625 serde_json::to_vec(&serde_json::json!({"name": name, "description": describe(id)}))
2626 .expect("profile"),
2627 )
2628 .expect("profile");
2629 }
2630 let gateway = dir.join("gateway-file.json");
2631 std::fs::write(
2632 &gateway,
2633 format!(r#"{{"host":"192.0.2.1","port":{port},"scheme":"http","token":"{token}"}}"#),
2634 )
2635 .expect("gateway");
2636 let skip = parse_skip_nicks(" svoi-brauzer , ");
2637 assert_eq!(skip, vec!["svoi-brauzer".to_string()]);
2638 let options = WakeOptions {
2639 skip_nicks: skip.clone(),
2640 store_root: Some(store_root.clone()),
2641 profile_note: false,
2642 session_ids: Vec::new(),
2643 };
2644
2645 let status_of = |reports: &[RoutineReport], agent: &str| {
2646 reports
2647 .iter()
2648 .find(|report| report.agent_id == agent)
2649 .map(|report| report.status.clone())
2650 };
2651 let first =
2652 ensure_agent_webhook_routines(&agents, &gateway, None, &options).expect("first pass");
2653 assert_eq!(status_of(&first, "agent-h"), Some(WakeStatus::Ready));
2654 assert_eq!(
2655 status_of(&first, "agent-c"),
2656 Some(WakeStatus::AwaitingBackend)
2657 );
2658 assert!(matches!(
2659 status_of(&first, "agent-x"),
2660 Some(WakeStatus::Failed(_))
2661 ));
2662 assert!(matches!(
2663 status_of(&first, "agent-k"),
2664 Some(WakeStatus::Failed(_))
2665 ));
2666 assert_eq!(status_of(&first, "agent-s"), None);
2667 let chief = first
2668 .iter()
2669 .find(|r| r.agent_id == "agent-c")
2670 .expect("chief");
2671 assert_eq!(chief.nick, "privet-mir");
2672 assert_eq!(chief.folder_id, "privet-mir");
2673 let shown = format!("{first:?}");
2674 assert!(!shown.contains("key-"));
2675 assert!(!shown.contains("automations/webhook"));
2676 {
2677 let gate = state.lock().expect("gate");
2678 assert_eq!(gate.creates, 3);
2681 assert_eq!(gate.deletes, 1);
2682 assert!(!gate.agents_seen.contains("agent-s"));
2683 assert!(gate.cards.iter().all(|card| card.agent != "agent-x"));
2684 let cron = gate
2685 .cards
2686 .iter()
2687 .find(|card| card.agent == "agent-k")
2688 .expect("cron");
2689 assert_eq!(cron.trigger, "cron");
2690 assert!(cron.enabled);
2691 let mirror = gate
2692 .cards
2693 .iter()
2694 .find(|card| card.agent == "agent-h")
2695 .expect("mirror");
2696 assert_eq!(
2697 (mirror.id.as_str(), mirror.name.as_str()),
2698 ("alice", "alice")
2699 );
2700 let chief_mirror = gate
2701 .cards
2702 .iter()
2703 .find(|card| card.agent == "agent-c")
2704 .expect("chief mirror");
2705 assert_eq!(chief_mirror.name, "privet-mir");
2707 assert_eq!(chief_mirror.id, "privet-mir");
2708 assert!(!mirror.enabled);
2709 }
2710 let alice_file = keychain_path(&store_root, "agent-h");
2712 let stored: serde_json::Value =
2713 serde_json::from_slice(&std::fs::read(&alice_file).expect("keychain"))
2714 .expect("keychain json");
2715 assert_eq!(stored["folder_id"], "alice");
2716 assert_eq!(stored["key"], "key-agent-h-alice");
2717 #[cfg(unix)]
2718 {
2719 use std::os::unix::fs::PermissionsExt;
2720 let mode = std::fs::metadata(&alice_file)
2721 .expect("meta")
2722 .permissions()
2723 .mode();
2724 assert_eq!(mode & 0o777, 0o600);
2725 }
2726 assert!(!keychain_path(&store_root, "agent-c").exists());
2727
2728 let second =
2730 ensure_agent_webhook_routines(&agents, &gateway, None, &options).expect("second pass");
2731 assert_eq!(
2732 status_of(&second, "agent-c"),
2733 Some(WakeStatus::AwaitingBackend)
2734 );
2735 assert_eq!(
2736 state.lock().expect("gate").creates,
2737 4,
2738 "only the clash bot retries"
2739 );
2740 assert!(state.lock().expect("gate").profiles.is_empty());
2741 let noted = WakeOptions {
2744 profile_note: true,
2745 ..options.clone()
2746 };
2747 ensure_agent_webhook_routines(&agents, &gateway, None, ¬ed).expect("note pass");
2748 {
2749 let gate = state.lock().expect("gate");
2750 assert_eq!(gate.profiles.len(), 1);
2751 let (agent, profile) = &gate.profiles[0];
2752 assert_eq!(agent, "agent-c");
2753 assert_eq!(profile["name"], "Привет мир");
2754 let description = profile["description"].as_str().unwrap_or("");
2755 assert!(description.starts_with("about the chief"));
2756 assert!(description.contains("routine named \"privet-mir\""));
2757 assert!(!description.contains("http"));
2758 }
2759 state.lock().expect("gate").creates = 4;
2760 state
2763 .lock()
2764 .expect("gate")
2765 .backend
2766 .insert(("agent-c".into(), "privet-mir".into()));
2767 let third =
2768 ensure_agent_webhook_routines(&agents, &gateway, None, &options).expect("third pass");
2769 assert_eq!(status_of(&third, "agent-c"), Some(WakeStatus::Ready));
2770 assert_eq!(state.lock().expect("gate").creates, 5);
2771 assert!(keychain_path(&store_root, "agent-c").exists());
2772
2773 let gateway_path = gateway.display().to_string();
2775 let root_text = store_root.display().to_string();
2776 let sessions = load_web_agents(&agents, |key| match key {
2777 GATEWAY_FILE_ENV => Some(gateway_path.clone()),
2778 SKIP_NICKS_ENV => Some("svoi-brauzer".to_string()),
2779 crate::STORE_ROOT_ENV => Some(root_text.clone()),
2780 _ => None,
2781 })
2782 .expect("web agents");
2783 assert_eq!(sessions.len(), 4);
2784 let alice = sessions
2785 .iter()
2786 .find(|s| s.session_id == "agent-h")
2787 .expect("h");
2788 assert!(alice
2789 .routine_url
2790 .as_deref()
2791 .unwrap_or("")
2792 .starts_with("https://"));
2793 assert!(alice.routine_bearer.is_some());
2794 assert!(!format!("{alice:?}").contains("key-"));
2795
2796 let mut offline = vec![HostSession::new("Alice", "agent-h")];
2798 offline[0].agent_id = Some("agent-h".to_string());
2799 let mut renamed = HostSession::new("Alice Two", "agent-h");
2800 renamed.agent_id = Some("agent-h".to_string());
2801 offline.push(renamed);
2802 attach_webhook_routines(
2803 &mut offline,
2804 &dir.join("absent.json"),
2805 None,
2806 Some(&store_root),
2807 )
2808 .expect("offline");
2809 assert_eq!(
2810 offline[0].routine_bearer.as_deref(),
2811 Some("key-agent-h-alice")
2812 );
2813 assert!(offline[1].routine_url.is_none());
2814 let mut bare = vec![HostSession::new("Alice", "agent-h")];
2815 attach_webhook_routines(&mut bare, &dir.join("absent.json"), None, None).expect("bare");
2816 assert!(bare[0].routine_url.is_none());
2817
2818 done.store(true, Ordering::Relaxed);
2819 let _ = server.join();
2820 let _ = std::fs::remove_dir_all(&dir);
2821 }
2822
2823 #[test]
2824 fn wake_note_is_appended_once_and_replaced_in_place() {
2825 let first = with_wake_note("Runs the alice.", "alice").expect("added");
2826 assert!(first.starts_with("Runs the alice.\n\n<!-- mail4agent:wake -->"));
2827 assert!(first.ends_with("<!-- /mail4agent:wake -->"));
2828 assert!(with_wake_note(&first, "alice").is_none());
2829 let renamed = with_wake_note(&first, "alice-two").expect("replaced");
2830 assert_eq!(renamed.matches("<!-- mail4agent:wake -->").count(), 1);
2831 assert!(renamed.contains("\"alice-two\""));
2832 assert!(!renamed.contains("\"alice\""));
2833 assert!(renamed.starts_with("Runs the alice."));
2834 assert_eq!(with_wake_note("", "carol").expect("empty"), wake_note("carol"));
2835 }
2836
2837 #[test]
2838 fn agents_directory_becomes_sessions_without_hardcoded_ids() {
2839 let dir = std::env::temp_dir().join(format!(
2840 "m4a-agents-{}-{}",
2841 std::process::id(),
2842 std::time::SystemTime::now()
2843 .duration_since(std::time::UNIX_EPOCH)
2844 .expect("clock")
2845 .as_nanos()
2846 ));
2847 let alice = dir.join("agent-alice");
2848 let chief = dir.join("agent-chief");
2849 std::fs::create_dir_all(&alice).expect("dir");
2850 std::fs::create_dir_all(&chief).expect("dir");
2851 std::fs::write(
2852 alice.join("profile.json"),
2853 r#"{"name":"Alice","description":"x"}"#,
2854 )
2855 .expect("profile");
2856 std::fs::write(
2857 chief.join("profile.json"),
2858 concat!(
2859 "{\"name\":\"",
2860 "Привет мир",
2861 "\",\"description\":\"x\"}"
2862 ),
2863 )
2864 .expect("profile");
2865 std::fs::write(dir.join("active-agent.json"), "{}").expect("skip file");
2866 let loaded = load_agents_dir(&dir).expect("agents");
2867 assert_eq!(loaded.len(), 2);
2868 assert_eq!(loaded[1].session_id, "agent-chief");
2869 assert_eq!(loaded[1].agent_id.as_deref(), Some("agent-chief"));
2870 assert_eq!(loaded[1].bot_name, "Привет мир");
2871 assert_eq!(
2872 routine_name_for(&loaded[1]).as_deref(),
2873 Some("privet-mir")
2874 );
2875 assert_ne!(
2876 routine_name_for(&loaded[1]).as_deref(),
2877 Some(loaded[1].session_id.as_str())
2878 );
2879 assert_eq!(loaded[0].bot_name, "Alice");
2880 assert_eq!(routine_name_for(&loaded[0]).as_deref(), Some("alice"));
2881 assert!(loaded.iter().all(|session| session.routine_url.is_none()));
2882 let _ = std::fs::remove_dir_all(&dir);
2883 }
2884}