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_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// Why: Slack replies target a caller-supplied `response_url`, so they must go
48// through the connect-time SSRF guard rather than the plain client the
49// operator-configured Teams endpoints use.
50#[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/// A surface-agnostic inbound message ready for dispatch. Per-platform routes
62/// build this from their normalized payload; the dispatch core never sees a
63/// Slack- or Teams-specific type.
64#[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/// Failures along the dispatch pipeline. This is an internal system surface;
85/// messages are deliberately descriptive for operator debugging.
86#[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}