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`](systemprompt_identifiers::ContextId) (derived from the
15//! conversation) ties multi-turn history 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 conversation;
22pub mod identity;
23
24use std::sync::LazyLock;
25
26use serde_json::json;
27use systemprompt_identifiers::{Actor, AgentName, SessionId, TraceId};
28use systemprompt_oauth::OauthError;
29use systemprompt_runtime::AppContext;
30use systemprompt_security::authz::{AuthzContext, AuthzDecision, AuthzRequest, EntityRef};
31use systemprompt_traits::SenderIdentity;
32use systemprompt_users::UserError;
33
34use crate::services::proxy::ProxyError;
35use a2a::{authenticated_user, build_a2a_request, mint_a2a_token, run_agent};
36pub use conversation::MessagingConversation;
37use identity::resolve_or_link_user;
38
39static GUARDED_CLIENT: LazyLock<Option<reqwest::Client>> = LazyLock::new(|| {
40    systemprompt_client::guarded_client(&systemprompt_client::GuardedClientConfig::default())
41        .inspect_err(|e| tracing::error!(error = %e, "Guarded outbound http client unavailable"))
42        .ok()
43});
44
45// Why: Slack replies target a caller-supplied `response_url`, so they must go
46// through the connect-time SSRF guard rather than the plain client the
47// operator-configured Teams endpoints use.
48#[must_use]
49pub fn guarded_http_client() -> Option<reqwest::Client> {
50    GUARDED_CLIENT.clone()
51}
52
53#[derive(Debug, Clone)]
54pub enum ReplyTarget {
55    Channel { id: String },
56    Url { url: String },
57}
58
59/// A surface-agnostic inbound message ready for dispatch. Per-platform routes
60/// build this from their normalized payload; the dispatch core never sees a
61/// Slack- or Teams-specific type.
62#[derive(Debug, Clone)]
63pub struct MessagingInbound {
64    pub issuer: String,
65    pub conversation: MessagingConversation,
66    pub text: String,
67    pub agent_name: AgentName,
68    pub entity: EntityRef,
69    pub reply: ReplyTarget,
70    pub sender: SenderIdentity,
71}
72
73#[derive(Debug, Clone)]
74pub enum DispatchOutcome {
75    Replied(String),
76    Denied(String),
77}
78
79/// Failures along the dispatch pipeline. This is an internal system surface;
80/// messages are deliberately descriptive for operator debugging.
81#[derive(Debug, thiserror::Error)]
82pub enum MessagingError {
83    #[error("identity resolution failed")]
84    Identity(#[source] UserError),
85    #[error("token minting failed")]
86    Token(#[source] OauthError),
87    #[error("could not encode the agent request")]
88    Encode(#[source] serde_json::Error),
89    #[error("could not build the agent request")]
90    Request(#[source] http::Error),
91    #[error("agent dispatch failed")]
92    Dispatch(#[source] ProxyError),
93    #[error("agent returned JSON-RPC error {code}: {message}")]
94    AgentRejected { code: i32, message: String },
95    #[error("agent response body could not be read")]
96    ResponseBody(#[source] axum::Error),
97    #[error("malformed agent response")]
98    Response(#[source] serde_json::Error),
99}
100
101impl MessagingError {
102    #[must_use]
103    #[expect(
104        clippy::unused_self,
105        reason = "opaque by contract: no part of the error may reach the caller"
106    )]
107    pub fn user_message(&self) -> String {
108        "Sorry — something went wrong handling that.".to_owned()
109    }
110}
111
112pub async fn dispatch_messaging(
113    ctx: &AppContext,
114    inbound: MessagingInbound,
115) -> Result<DispatchOutcome, MessagingError> {
116    let user = resolve_or_link_user(
117        ctx,
118        &inbound.issuer,
119        inbound.conversation.sender_wire_id(),
120        &inbound.sender.claims(),
121    )
122    .await?;
123    let authed = authenticated_user(&user);
124
125    let context_id = inbound.conversation.context_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.conversation.platform()),
139            json!({ "channel": inbound.conversation.channel_key() }),
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, request).await?;
154    Ok(DispatchOutcome::Replied(reply))
155}