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