1pub mod chain;
36pub mod claude_channel;
37pub mod codex;
38pub mod hook;
39pub mod inbox;
40pub mod kimi;
41pub mod registry;
42pub mod spawn;
43
44use std::fmt;
45use std::path::{Path, PathBuf};
46
47use crate::{post_decrypted_with_bearer, reply_hint, DecryptedWake, SEND_COMMAND};
48
49pub use chain::{plan_chain, ChainError, HostEnv, Mechanism, Tier, WakeChain, WebVendor};
50pub use claude_channel::ClaudeChannelAdapter;
51pub use codex::{CodexAppServerAdapter, CodexEndpoint};
52pub use hook::{HookFlavor, InboxHookAdapter, InboxQueueAdapter};
53pub use inbox::{InboxEntry, InboxLetter};
54pub use kimi::KimiServerAdapter;
55pub use registry::SessionRecord;
56pub use spawn::ResumeSpawnAdapter;
57
58pub const PROVIDER_ENV: &str = "M4A_PROVIDER";
60
61pub const INBOX_DIR_ENV: &str = "M4A_INBOX_DIR";
64
65#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
68pub enum ProviderKind {
69 Grok,
71 KimiCode,
73 ClaudeCode,
75 Codex,
77 Cursor,
80}
81
82impl ProviderKind {
83 pub const CORE: [ProviderKind; 4] = [
85 ProviderKind::Grok,
86 ProviderKind::KimiCode,
87 ProviderKind::ClaudeCode,
88 ProviderKind::Codex,
89 ];
90
91 pub const ALL: [ProviderKind; 5] = [
93 ProviderKind::Grok,
94 ProviderKind::KimiCode,
95 ProviderKind::ClaudeCode,
96 ProviderKind::Codex,
97 ProviderKind::Cursor,
98 ];
99
100 pub fn id(self) -> &'static str {
102 match self {
103 ProviderKind::Grok => "grok",
104 ProviderKind::KimiCode => "kimi",
105 ProviderKind::ClaudeCode => "claude",
106 ProviderKind::Codex => "codex",
107 ProviderKind::Cursor => "cursor",
108 }
109 }
110
111 pub fn is_optional(self) -> bool {
113 matches!(self, ProviderKind::Cursor)
114 }
115
116 pub fn parse(value: &str) -> Option<Self> {
118 match value.trim().to_ascii_lowercase().as_str() {
119 "grok" => Some(ProviderKind::Grok),
120 "kimi" | "kimi-code" | "kimicode" => Some(ProviderKind::KimiCode),
121 "claude" | "claude-code" | "claudecode" => Some(ProviderKind::ClaudeCode),
122 "codex" => Some(ProviderKind::Codex),
123 "cursor" | "cursor-agent" => Some(ProviderKind::Cursor),
124 _ => None,
125 }
126 }
127
128 pub fn from_env() -> Result<Self, WakeError> {
130 match std::env::var(PROVIDER_ENV) {
131 Ok(value) if !value.trim().is_empty() => {
132 Self::parse(&value).ok_or(WakeError::UnknownProvider)
133 }
134 _ => Ok(ProviderKind::Grok),
135 }
136 }
137}
138
139impl fmt::Display for ProviderKind {
140 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
141 f.write_str(self.id())
142 }
143}
144
145#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
147pub enum Surface {
148 Web,
150 Local,
152}
153
154#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
156pub struct SessionKind {
157 pub provider: ProviderKind,
159 pub surface: Surface,
161}
162
163impl SessionKind {
164 pub fn local(provider: ProviderKind) -> Self {
166 Self {
167 provider,
168 surface: Surface::Local,
169 }
170 }
171
172 pub fn web(provider: ProviderKind) -> Self {
174 Self {
175 provider,
176 surface: Surface::Web,
177 }
178 }
179}
180
181impl fmt::Display for SessionKind {
182 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
183 let surface = match self.surface {
184 Surface::Web => "web",
185 Surface::Local => "local",
186 };
187 write!(f, "{}/{}", self.provider, surface)
188 }
189}
190
191#[derive(Debug, Clone, PartialEq, Eq)]
193pub struct ProviderSession {
194 pub kind: SessionKind,
196 pub session_id: String,
198 pub nick: String,
200 pub cwd: Option<PathBuf>,
202 pub headless: bool,
205}
206
207#[derive(Debug, Clone, Copy)]
209pub struct WakeLetter<'a> {
210 pub body: &'a str,
212 pub from_nick: &'a str,
214 pub event_id: &'a str,
216 pub room: Option<&'a str>,
218}
219
220#[derive(Debug, Clone, PartialEq, Eq)]
222pub enum WakeOutcome {
223 Delivered,
226 Queued(PathBuf),
229}
230
231#[derive(Debug, Clone, PartialEq, Eq)]
233pub enum WakeError {
234 UnknownProvider,
236 NotImplemented(SessionKind),
238 NoInboundTrigger(SessionKind),
240 Unavailable(String),
242 Transport(String),
244}
245
246impl fmt::Display for WakeError {
247 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
248 match self {
249 WakeError::UnknownProvider => write!(f, "unknown provider in {PROVIDER_ENV}"),
250 WakeError::NotImplemented(kind) => {
251 write!(f, "wake adapter for {kind} is not implemented")
252 }
253 WakeError::NoInboundTrigger(kind) => {
254 write!(
255 f,
256 "{kind} exposes no inbound trigger into a running session"
257 )
258 }
259 WakeError::Unavailable(why) => write!(f, "wake endpoint unavailable: {why}"),
260 WakeError::Transport(why) => write!(f, "wake failed: {why}"),
261 }
262 }
263}
264
265impl std::error::Error for WakeError {}
266
267pub trait WakeAdapter: Send {
273 fn kind(&self) -> SessionKind;
275
276 fn probe(&self, session: &ProviderSession) -> Result<(), WakeError>;
279
280 fn wake(
282 &mut self,
283 session: &ProviderSession,
284 letter: &WakeLetter<'_>,
285 ) -> Result<WakeOutcome, WakeError>;
286}
287
288pub fn wake_prompt(session: &ProviderSession, letter: &WakeLetter<'_>) -> String {
293 let reply = reply_hint(&session.nick, letter.from_nick);
294 format!(
295 "[mail4agent] Letter from {from} to {to} (event {event}).\n\
296 Answer every direct letter with {cmd}; at minimum \"принято\" plus what you will do and when.\n\
297 Reply: {reply}\n\
298 ---\n\
299 {body}",
300 from = letter.from_nick,
301 to = session.nick,
302 event = letter.event_id,
303 cmd = SEND_COMMAND,
304 reply = reply,
305 body = letter.body,
306 )
307}
308
309#[derive(Default, Clone)]
313pub struct AdapterConfig {
314 pub leader_sock: Option<PathBuf>,
316 pub inbox_dir: Option<PathBuf>,
318 pub codex: Option<CodexEndpoint>,
320 pub kimi_url: Option<String>,
322 pub kimi_bearer: Option<String>,
324 pub routine_url: Option<String>,
326 pub routine_bearer: Option<String>,
328 pub claude_fire_url: Option<String>,
330 pub claude_fire_bearer: Option<String>,
332 pub codex_cloud_env: Option<String>,
334 pub spawn_program: Option<PathBuf>,
336}
337
338pub const CLAUDE_FIRE_URL_ENV: &str = "M4A_CLAUDE_ROUTINE_FIRE_URL";
340pub const CLAUDE_FIRE_TOKEN_ENV: &str = "M4A_CLAUDE_ROUTINE_TOKEN";
342
343impl AdapterConfig {
344 pub fn from_env(inbox_dir: Option<PathBuf>) -> Self {
350 let get = |key: &str| std::env::var(key).ok().filter(|value| !value.is_empty());
351 let mut kimi = KimiServerAdapter::from_env();
352 if !kimi.is_configured() {
353 if let Some(found) = kimi::discover_local_server(None) {
354 kimi = found;
355 }
356 }
357 let (kimi_url, kimi_bearer) = kimi.into_parts();
358 Self {
359 leader_sock: get(crate::LEADER_SOCK_ENV).map(PathBuf::from),
360 inbox_dir: get(INBOX_DIR_ENV).map(PathBuf::from).or(inbox_dir),
361 codex: CodexEndpoint::from_env(),
362 kimi_url,
363 kimi_bearer,
364 routine_url: None,
365 routine_bearer: None,
366 claude_fire_url: get(CLAUDE_FIRE_URL_ENV),
367 claude_fire_bearer: get(CLAUDE_FIRE_TOKEN_ENV),
368 codex_cloud_env: get(spawn::CODEX_CLOUD_ENV_ENV),
369 spawn_program: None,
370 }
371 }
372}
373
374impl fmt::Debug for AdapterConfig {
375 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
376 f.debug_struct("AdapterConfig")
378 .field("leader_sock", &self.leader_sock.is_some())
379 .field("inbox_dir", &self.inbox_dir.is_some())
380 .field("codex", &self.codex.is_some())
381 .field("kimi_url", &self.kimi_url.is_some())
382 .field("routine_url", &self.routine_url.is_some())
383 .field("claude_fire_url", &self.claude_fire_url.is_some())
384 .field("codex_cloud_env", &self.codex_cloud_env.is_some())
385 .finish()
386 }
387}
388
389pub fn adapter_for(kind: SessionKind, config: &AdapterConfig) -> Box<dyn WakeAdapter> {
391 match (kind.surface, kind.provider) {
392 (Surface::Local, ProviderKind::Grok) => Box::new(GrokLeaderAdapter {
393 leader_sock: config.leader_sock.clone(),
394 }),
395 (Surface::Local, ProviderKind::Codex) => {
396 Box::new(CodexAppServerAdapter::new(config.codex.clone()))
397 }
398 (Surface::Local, ProviderKind::KimiCode) => Box::new(KimiServerAdapter::new(
399 config.kimi_url.clone(),
400 config.kimi_bearer.clone(),
401 )),
402 (Surface::Local, ProviderKind::ClaudeCode) => {
403 Box::new(ClaudeChannelAdapter::new(config.inbox_dir.clone()))
404 }
405 (Surface::Local, ProviderKind::Cursor) => Box::new(CursorAgentAdapter {
406 inbox_dir: config.inbox_dir.clone(),
407 }),
408 (Surface::Web, ProviderKind::Cursor) => Box::new(RoutineWebhookAdapter {
409 provider: ProviderKind::Cursor,
410 url: config.routine_url.clone(),
411 bearer: config.routine_bearer.clone(),
412 }),
413 (Surface::Web, ProviderKind::ClaudeCode) => Box::new(ClaudeRoutineFireAdapter {
414 url: config.claude_fire_url.clone(),
415 bearer: config.claude_fire_bearer.clone(),
416 }),
417 (Surface::Web, ProviderKind::Codex) => Box::new(CodexCloudAdapter),
418 (Surface::Web, provider @ (ProviderKind::KimiCode | ProviderKind::Grok)) => {
419 Box::new(NoInboundAdapter { provider })
420 }
421 }
422}
423
424pub struct GrokLeaderAdapter {
428 pub leader_sock: Option<PathBuf>,
430}
431
432impl WakeAdapter for GrokLeaderAdapter {
433 fn kind(&self) -> SessionKind {
434 SessionKind::local(ProviderKind::Grok)
435 }
436
437 fn probe(&self, _session: &ProviderSession) -> Result<(), WakeError> {
438 match self.leader_sock.as_deref() {
439 Some(path) if path.exists() => Ok(()),
440 Some(_) => Err(WakeError::Unavailable("leader socket missing".into())),
441 None => Err(WakeError::Unavailable("leader socket unset".into())),
442 }
443 }
444
445 fn wake(
446 &mut self,
447 session: &ProviderSession,
448 letter: &WakeLetter<'_>,
449 ) -> Result<WakeOutcome, WakeError> {
450 self.probe(session)?;
451 let Some(sock) = self.leader_sock.as_deref() else {
452 return Err(WakeError::Unavailable("leader socket unset".into()));
453 };
454 let cwd = session_cwd(session)?;
455 mail4agent_grok::wake_decrypted_room_blocking(
456 sock,
457 &session.session_id,
458 &cwd,
459 &wake_prompt(session, letter),
460 )
461 .map(|()| WakeOutcome::Delivered)
462 .map_err(|err| WakeError::Transport(err.to_string()))
463 }
464}
465
466pub struct CursorAgentAdapter {
472 pub inbox_dir: Option<PathBuf>,
474}
475
476impl WakeAdapter for CursorAgentAdapter {
477 fn kind(&self) -> SessionKind {
478 SessionKind::local(ProviderKind::Cursor)
479 }
480
481 fn probe(&self, session: &ProviderSession) -> Result<(), WakeError> {
482 InboxHookAdapter::new(self.kind(), HookFlavor::CursorStop, self.inbox_dir.clone())
483 .probe(session)
484 }
485
486 fn wake(
487 &mut self,
488 session: &ProviderSession,
489 letter: &WakeLetter<'_>,
490 ) -> Result<WakeOutcome, WakeError> {
491 InboxHookAdapter::new(self.kind(), HookFlavor::CursorStop, self.inbox_dir.clone())
492 .wake(session, letter)
493 }
494}
495
496pub struct RoutineWebhookAdapter {
502 pub provider: ProviderKind,
504 pub url: Option<String>,
506 pub bearer: Option<String>,
508}
509
510impl WakeAdapter for RoutineWebhookAdapter {
511 fn kind(&self) -> SessionKind {
512 SessionKind::web(self.provider)
513 }
514
515 fn probe(&self, _session: &ProviderSession) -> Result<(), WakeError> {
516 match self.url.as_deref() {
517 Some(url) if !url.is_empty() => Ok(()),
518 _ => Err(WakeError::Unavailable("routine url unset".into())),
519 }
520 }
521
522 fn wake(
523 &mut self,
524 session: &ProviderSession,
525 letter: &WakeLetter<'_>,
526 ) -> Result<WakeOutcome, WakeError> {
527 self.probe(session)?;
528 let url = self.url.as_deref().unwrap_or_default();
529 let reply = reply_hint(&session.nick, letter.from_nick);
530 let wake = DecryptedWake {
531 body: letter.body,
532 from: letter.from_nick,
533 nick: None,
534 event_id: letter.event_id,
535 room: letter.room,
536 from_nick: Some(letter.from_nick),
537 to: Some(&session.nick),
538 reply: Some(&reply),
539 };
540 post_decrypted_with_bearer(url, &wake, self.bearer.as_deref())
541 .map(|()| WakeOutcome::Delivered)
542 .map_err(|err| WakeError::Transport(err.to_string()))
543 }
544}
545
546pub struct ClaudeRoutineFireAdapter {
554 pub url: Option<String>,
556 pub bearer: Option<String>,
558}
559
560pub const CLAUDE_ROUTINE_BETA: &str = "experimental-cc-routine-2026-04-01";
562
563impl WakeAdapter for ClaudeRoutineFireAdapter {
564 fn kind(&self) -> SessionKind {
565 SessionKind::web(ProviderKind::ClaudeCode)
566 }
567
568 fn probe(&self, session: &ProviderSession) -> Result<(), WakeError> {
569 if self.url.as_deref().is_none_or(str::is_empty) {
570 return Err(WakeError::Unavailable(
571 "claude routine fire url unset".into(),
572 ));
573 }
574 if self.bearer.as_deref().is_none_or(str::is_empty) {
575 return Err(WakeError::Unavailable("claude routine token unset".into()));
576 }
577 if !session.headless {
578 return Err(WakeError::Unavailable(
579 "session is open; routine fire would start a new session".into(),
580 ));
581 }
582 Ok(())
583 }
584
585 fn wake(
586 &mut self,
587 session: &ProviderSession,
588 letter: &WakeLetter<'_>,
589 ) -> Result<WakeOutcome, WakeError> {
590 self.probe(session)?;
591 let url = self.url.as_deref().unwrap_or_default();
592 let token = self.bearer.as_deref().unwrap_or_default();
593 let parsed = reqwest::Url::parse(url)
594 .ok()
595 .filter(|u| u.scheme() == "https" || is_loopback_http(u))
596 .ok_or_else(|| WakeError::Unavailable("claude routine fire url invalid".into()))?;
597 let beta = std::env::var("M4A_CLAUDE_ROUTINE_BETA")
598 .ok()
599 .filter(|v| !v.is_empty())
600 .unwrap_or_else(|| CLAUDE_ROUTINE_BETA.to_string());
601 let client = reqwest::blocking::Client::builder()
602 .timeout(std::time::Duration::from_secs(15))
603 .build()
604 .map_err(|err| WakeError::Transport(err.without_url().to_string()))?;
605 let response = client
606 .post(parsed)
607 .bearer_auth(token)
608 .header("anthropic-version", "2023-06-01")
609 .header("anthropic-beta", beta)
610 .json(&serde_json::json!({"text": wake_prompt(session, letter)}))
611 .send()
612 .map_err(|err| WakeError::Transport(err.without_url().to_string()))?;
613 if response.status().is_success() {
614 Ok(WakeOutcome::Delivered)
615 } else {
616 Err(WakeError::Transport(format!(
617 "routine fire status {}",
618 response.status().as_u16()
619 )))
620 }
621 }
622}
623
624fn is_loopback_http(url: &reqwest::Url) -> bool {
625 url.scheme() == "http"
626 && matches!(
627 url.host_str(),
628 Some("127.0.0.1" | "localhost" | "[::1]" | "::1")
629 )
630}
631
632pub struct CodexCloudAdapter;
637
638pub struct NoInboundAdapter {
640 pub provider: ProviderKind,
642}
643
644macro_rules! refuse_adapter {
645 ($ty:ty, $kind:expr, $err:ident) => {
646 impl WakeAdapter for $ty {
647 fn kind(&self) -> SessionKind {
648 $kind(self)
649 }
650
651 fn probe(&self, _session: &ProviderSession) -> Result<(), WakeError> {
652 Err(WakeError::$err(self.kind()))
653 }
654
655 fn wake(
656 &mut self,
657 _session: &ProviderSession,
658 _letter: &WakeLetter<'_>,
659 ) -> Result<WakeOutcome, WakeError> {
660 Err(WakeError::$err(self.kind()))
661 }
662 }
663 };
664}
665
666refuse_adapter!(
667 CodexCloudAdapter,
668 |_: &CodexCloudAdapter| SessionKind::web(ProviderKind::Codex),
669 NotImplemented
670);
671refuse_adapter!(
672 NoInboundAdapter,
673 |adapter: &NoInboundAdapter| SessionKind::web(adapter.provider),
674 NoInboundTrigger
675);
676
677pub(crate) fn session_cwd(session: &ProviderSession) -> Result<String, WakeError> {
678 session
679 .cwd
680 .as_deref()
681 .map(Path::to_path_buf)
682 .or_else(|| std::env::current_dir().ok())
683 .map(|path| path.display().to_string())
684 .filter(|cwd| !cwd.is_empty())
685 .ok_or_else(|| WakeError::Unavailable("session cwd is empty".into()))
686}
687
688#[cfg(test)]
689pub(crate) mod tests {
690 use super::*;
691
692 pub(crate) fn session(kind: SessionKind) -> ProviderSession {
693 ProviderSession {
694 kind,
695 session_id: "s-1".into(),
696 nick: "alice".into(),
697 cwd: Some(PathBuf::from("/tmp")),
698 headless: false,
699 }
700 }
701
702 pub(crate) fn letter<'a>(body: &'a str) -> WakeLetter<'a> {
703 WakeLetter {
704 body,
705 from_nick: "carol",
706 event_id: "$ev1",
707 room: None,
708 }
709 }
710
711 #[test]
712 fn ids_round_trip() {
713 for kind in ProviderKind::ALL {
714 assert_eq!(ProviderKind::parse(kind.id()), Some(kind));
715 }
716 assert_eq!(
717 ProviderKind::parse("Claude-Code"),
718 Some(ProviderKind::ClaudeCode)
719 );
720 assert_eq!(ProviderKind::parse("gemini"), None);
721 assert!(ProviderKind::CORE.iter().all(|kind| !kind.is_optional()));
722 assert!(ProviderKind::Cursor.is_optional());
723 assert_eq!(
724 SessionKind::web(ProviderKind::Codex).to_string(),
725 "codex/web"
726 );
727 }
728
729 #[test]
730 fn every_kind_gets_an_adapter_of_that_kind() {
731 let config = AdapterConfig::default();
732 for provider in ProviderKind::ALL {
733 for kind in [SessionKind::local(provider), SessionKind::web(provider)] {
734 assert_eq!(adapter_for(kind, &config).kind(), kind);
735 }
736 }
737 }
738
739 #[test]
740 fn stubs_and_unconfigured_adapters_refuse_without_side_effects() {
741 let config = AdapterConfig::default();
742 let cases = [
743 (SessionKind::local(ProviderKind::Cursor), "unavailable"),
744 (SessionKind::web(ProviderKind::ClaudeCode), "unavailable"),
745 (SessionKind::web(ProviderKind::Codex), "not implemented"),
746 (SessionKind::web(ProviderKind::KimiCode), "no inbound"),
747 (SessionKind::web(ProviderKind::Grok), "no inbound"),
748 (SessionKind::web(ProviderKind::Cursor), "unavailable"),
749 (SessionKind::local(ProviderKind::Grok), "unavailable"),
750 (SessionKind::local(ProviderKind::Codex), "unavailable"),
751 (SessionKind::local(ProviderKind::KimiCode), "unavailable"),
752 (SessionKind::local(ProviderKind::ClaudeCode), "unavailable"),
753 ];
754 for (kind, expect) in cases {
755 let mut adapter = adapter_for(kind, &config);
756 let s = session(kind);
757 let err = adapter.wake(&s, &letter("x")).unwrap_err().to_string();
758 assert!(err.contains(expect), "{kind}: {err}");
759 }
760 }
761
762 #[test]
763 fn prompt_carries_reply_command_and_body() {
764 let text = wake_prompt(
765 &session(SessionKind::local(ProviderKind::Codex)),
766 &letter("ping"),
767 );
768 assert!(text.contains("m4a-send --as alice --to carol"));
769 assert!(text.ends_with("ping"));
770 }
771
772 #[test]
773 fn config_debug_hides_values() {
774 let config = AdapterConfig {
775 kimi_bearer: Some("k-secret".into()),
776 routine_url: Some("https://example.invalid/hook".into()),
777 ..AdapterConfig::default()
778 };
779 let shown = format!("{config:?}");
780 assert!(!shown.contains("secret") && !shown.contains("example"));
781 }
782}