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