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