Skip to main content

mail4agent_messenger_shell/
machine.rs

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