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