systemprompt_api/routes/messaging/
mod.rs1pub mod a2a;
21pub mod identity;
22
23use std::sync::LazyLock;
24
25use serde_json::json;
26use systemprompt_identifiers::{Actor, AgentName, ContextId, SessionId, TraceId};
27use systemprompt_runtime::AppContext;
28use systemprompt_security::authz::{AuthzContext, AuthzDecision, AuthzRequest, EntityRef};
29use systemprompt_traits::SenderIdentity;
30
31use a2a::{authenticated_user, build_a2a_request, mint_a2a_token, run_agent};
32use identity::resolve_or_link_user;
33
34static CLIENT: LazyLock<reqwest::Client> = LazyLock::new(reqwest::Client::new);
35
36static GUARDED_CLIENT: LazyLock<Option<reqwest::Client>> = LazyLock::new(|| {
37 systemprompt_models::net::guarded_client(
38 &systemprompt_models::net::GuardedClientConfig::default(),
39 )
40 .inspect_err(|e| tracing::error!(error = %e, "Guarded outbound http client unavailable"))
41 .ok()
42});
43
44#[must_use]
45pub fn http_client() -> reqwest::Client {
46 CLIENT.clone()
47}
48
49#[must_use]
53pub fn guarded_http_client() -> Option<reqwest::Client> {
54 GUARDED_CLIENT.clone()
55}
56
57#[derive(Debug, Clone)]
58pub enum ReplyTarget {
59 Channel { id: String },
60 Url { url: String },
61}
62
63#[derive(Debug, Clone)]
67pub struct MessagingInbound {
68 pub platform: &'static str,
69 pub issuer: String,
70 pub org_id: String,
71 pub channel_id: String,
72 pub external_user_id: String,
73 pub text: String,
74 pub agent_name: AgentName,
75 pub entity: EntityRef,
76 pub reply: ReplyTarget,
77 pub sender: SenderIdentity,
78}
79
80#[derive(Debug, Clone)]
81pub enum DispatchOutcome {
82 Replied(String),
83 Denied(String),
84}
85
86#[derive(Debug, thiserror::Error)]
89pub enum MessagingError {
90 #[error("identity resolution failed: {0}")]
91 Identity(String),
92 #[error("token minting failed: {0}")]
93 Token(String),
94 #[error("agent dispatch failed: {0}")]
95 Dispatch(String),
96 #[error("malformed agent response: {0}")]
97 Response(String),
98}
99
100impl MessagingError {
101 #[must_use]
102 #[expect(
103 clippy::unused_self,
104 reason = "opaque by contract: no part of the error may reach the caller"
105 )]
106 pub fn user_message(&self) -> String {
107 "Sorry — something went wrong handling that.".to_owned()
108 }
109}
110
111pub async fn dispatch_messaging(
112 ctx: &AppContext,
113 inbound: MessagingInbound,
114) -> Result<DispatchOutcome, MessagingError> {
115 let user = resolve_or_link_user(
116 ctx,
117 &inbound.issuer,
118 &inbound.external_user_id,
119 &inbound.sender.claims(),
120 )
121 .await?;
122 let authed = authenticated_user(&user)?;
123
124 let context_id =
125 ContextId::derived_from_messaging(inbound.platform, &inbound.org_id, &inbound.channel_id);
126
127 let authz = AuthzRequest {
128 entity: inbound.entity.clone(),
129 user_id: user.id.clone(),
130 actor: Some(Actor::user(user.id.clone())),
131 client_id: None,
132 access_scope: None,
133 roles: user.roles.clone(),
134 attributes: std::collections::BTreeMap::new(),
135 trace_id: TraceId::generate(),
136 session_id: None,
137 context: AuthzContext::extension(
138 format!("{}.message", inbound.platform),
139 json!({ "channel": inbound.channel_id }),
140 ),
141 context_id: Some(context_id.clone()),
142 task_id: None,
143 act_chain: Vec::new(),
144 };
145 if let AuthzDecision::Deny { reason, policy } = ctx.authz_hook().evaluate(authz).await {
146 return Ok(DispatchOutcome::Denied(format!("{policy}: {reason}")));
147 }
148
149 let session_id = SessionId::new(uuid::Uuid::new_v4().to_string());
150 let token = mint_a2a_token(ctx, &authed, &session_id)?;
151
152 let request = build_a2a_request(&inbound, &authed, &session_id, &token, &context_id)?;
153 let reply = run_agent(ctx, inbound.agent_name.as_str(), request).await?;
154 Ok(DispatchOutcome::Replied(reply))
155}