Skip to main content

systemprompt_api/routes/messaging/
mod.rs

1//! Platform-agnostic dispatch for chat-platform inbound messages.
2//!
3//! Slack and Teams differ only at their edges — request verification, payload
4//! shape, and reply rendering. Everything between (identity, authorization,
5//! deterministic conversation context, per-user A2A token minting, the blocking
6//! `message/send` through the proxy, and reply extraction) is identical and
7//! lives here once. A per-platform route normalizes its wire payload into a
8//! [`MessagingInbound`] and calls [`dispatch_messaging`]; the returned
9//! [`DispatchOutcome`] is rendered back into the platform's UI by the route.
10//!
11//! The pipeline is **synchronous, spawned**: the route acks the platform within
12//! its timeout, then a spawned task runs this blocking dispatch and posts the
13//! reply. There is no responder job and no dispatch-state table — a stable
14//! [`ContextId`] (derived from the conversation) ties multi-turn history
15//! together instead.
16//!
17//! Copyright (c) systemprompt.io — Business Source License 1.1.
18//! See <https://systemprompt.io> for licensing details.
19
20pub 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// Why: Slack replies target a caller-supplied `response_url`, so they must go
50// through the connect-time SSRF guard rather than the plain client the
51// operator-configured Teams endpoints use.
52#[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/// A surface-agnostic inbound message ready for dispatch. Per-platform routes
64/// build this from their normalized payload; the dispatch core never sees a
65/// Slack- or Teams-specific type.
66#[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/// Failures along the dispatch pipeline. This is an internal system surface;
87/// messages are deliberately descriptive for operator debugging.
88#[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}