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