Skip to main content

messaging_api/
model.rs

1use conversation_api::ConversationSurface;
2use serde::{Deserialize, Serialize};
3use std::fmt;
4
5pub const NORMALIZED_INBOUND_SCHEMA_VERSION: u16 = 4;
6
7#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
8#[serde(rename_all = "snake_case")]
9pub enum ConversationAudience {
10    Personal,
11    Shared,
12}
13
14/// Whether a provider could determine that a message explicitly targets the bot.
15///
16/// Shared-conversation policy drops `Unaddressed` messages. `Unknown` remains
17/// admissible for providers that cannot expose an equivalent signal.
18#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
19#[serde(rename_all = "snake_case")]
20pub enum MessageAttention {
21    Addressed,
22    Unaddressed,
23    Unknown,
24}
25
26#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
27#[serde(rename_all = "camelCase", deny_unknown_fields)]
28pub struct MessagingAddress {
29    pub provider: String,
30    pub account_id: String,
31    pub conversation_id: String,
32    #[serde(default, skip_serializing_if = "Option::is_none")]
33    pub lane_id: Option<String>,
34    pub audience: ConversationAudience,
35}
36
37#[derive(Debug, Clone, PartialEq, Eq)]
38pub struct MessagingModelError(&'static str);
39
40impl fmt::Display for MessagingModelError {
41    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
42        formatter.write_str(self.0)
43    }
44}
45
46impl std::error::Error for MessagingModelError {}
47
48impl MessagingAddress {
49    pub fn new(
50        provider: impl Into<String>,
51        account_id: impl Into<String>,
52        conversation_id: impl Into<String>,
53        lane_id: Option<String>,
54        audience: ConversationAudience,
55    ) -> Result<Self, MessagingModelError> {
56        Ok(Self {
57            provider: provider_id(provider.into())?,
58            account_id: segment(account_id.into(), "messaging account id is invalid")?,
59            conversation_id: segment(
60                conversation_id.into(),
61                "messaging conversation id is invalid",
62            )?,
63            lane_id: lane_id
64                .map(|value| segment(value, "messaging lane id is invalid"))
65                .transpose()?,
66            audience,
67        })
68    }
69
70    #[must_use]
71    pub fn base_address(&self) -> Self {
72        let mut base = self.clone();
73        base.lane_id = None;
74        base
75    }
76
77    pub fn conversation_surface(&self) -> Result<ConversationSurface, MessagingModelError> {
78        let result = match self.audience {
79            ConversationAudience::Personal => ConversationSurface::messaging_personal(
80                &self.provider,
81                &self.account_id,
82                &self.conversation_id,
83                self.lane_id.clone(),
84            ),
85            ConversationAudience::Shared => ConversationSurface::messaging_group(
86                &self.provider,
87                &self.account_id,
88                &self.conversation_id,
89                self.lane_id.clone(),
90            ),
91        };
92        result.map_err(|_| MessagingModelError("messaging address cannot form a surface"))
93    }
94
95    pub fn validate(self) -> Result<Self, MessagingModelError> {
96        Self::new(
97            self.provider,
98            self.account_id,
99            self.conversation_id,
100            self.lane_id,
101            self.audience,
102        )
103    }
104}
105
106#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
107#[serde(rename_all = "camelCase", deny_unknown_fields)]
108pub struct ExternalActor {
109    pub provider: String,
110    pub account_id: String,
111    pub external_user_id: String,
112    #[serde(default, skip_serializing_if = "Option::is_none")]
113    pub display_name: Option<String>,
114}
115
116impl ExternalActor {
117    pub fn new(
118        provider: impl Into<String>,
119        account_id: impl Into<String>,
120        external_user_id: impl Into<String>,
121        display_name: Option<String>,
122    ) -> Result<Self, MessagingModelError> {
123        Ok(Self {
124            provider: provider_id(provider.into())?,
125            account_id: segment(account_id.into(), "messaging actor account id is invalid")?,
126            external_user_id: segment(
127                external_user_id.into(),
128                "messaging external user id is invalid",
129            )?,
130            display_name: display_name.map(optional_segment).transpose()?,
131        })
132    }
133
134    pub fn validate(self) -> Result<Self, MessagingModelError> {
135        Self::new(
136            self.provider,
137            self.account_id,
138            self.external_user_id,
139            self.display_name,
140        )
141    }
142}
143
144#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
145#[serde(rename_all = "camelCase", deny_unknown_fields)]
146pub struct ActionOption {
147    pub label: String,
148    pub token: String,
149}
150
151impl ActionOption {
152    pub fn new(
153        label: impl Into<String>,
154        token: impl Into<String>,
155    ) -> Result<Self, MessagingModelError> {
156        Ok(Self {
157            label: segment(label.into(), "action label is invalid")?,
158            token: segment(token.into(), "action token is invalid")?,
159        })
160    }
161}
162
163#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
164#[serde(rename_all = "camelCase", deny_unknown_fields)]
165pub struct ActionSet {
166    pub options: Vec<ActionOption>,
167}
168
169impl ActionSet {
170    pub fn new(options: Vec<ActionOption>) -> Result<Self, MessagingModelError> {
171        if options.is_empty() {
172            return Err(MessagingModelError("action set cannot be empty"));
173        }
174        for option in &options {
175            segment_ref(&option.label, "action label is invalid")?;
176            segment_ref(&option.token, "action token is invalid")?;
177        }
178        Ok(Self { options })
179    }
180}
181
182#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
183#[serde(rename_all = "snake_case")]
184pub enum ProviderMediaKind {
185    Image,
186    Audio,
187    Video,
188    File,
189}
190
191/// Provider-scoped, short-lived media handle. It must be materialized before entering Conversation.
192#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
193#[serde(rename_all = "camelCase", deny_unknown_fields)]
194pub struct ProviderMediaRef {
195    pub handle: String,
196    pub kind: ProviderMediaKind,
197    #[serde(default, skip_serializing_if = "Option::is_none")]
198    pub mime_type: Option<String>,
199    #[serde(default, skip_serializing_if = "Option::is_none")]
200    pub size_bytes: Option<u64>,
201    #[serde(default, skip_serializing_if = "Option::is_none")]
202    pub file_name: Option<String>,
203    #[serde(default, skip_serializing_if = "Option::is_none")]
204    pub duration_ms: Option<u64>,
205    #[serde(default, skip_serializing_if = "Option::is_none")]
206    pub width_px: Option<u32>,
207    #[serde(default, skip_serializing_if = "Option::is_none")]
208    pub height_px: Option<u32>,
209    #[serde(default, skip_serializing_if = "Option::is_none")]
210    pub caption: Option<String>,
211    #[serde(default, skip_serializing_if = "Option::is_none")]
212    pub transcript: Option<String>,
213}
214
215#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
216#[serde(
217    tag = "type",
218    rename_all = "snake_case",
219    rename_all_fields = "camelCase",
220    deny_unknown_fields
221)]
222pub enum InboundMessagePart {
223    Text { text: String },
224    Media { reference: ProviderMediaRef },
225}
226
227#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
228#[serde(
229    tag = "type",
230    rename_all = "snake_case",
231    rename_all_fields = "camelCase",
232    deny_unknown_fields
233)]
234pub enum NormalizedInboundContent {
235    Message {
236        #[serde(default, skip_serializing_if = "Option::is_none")]
237        provider_message_id: Option<String>,
238        attention: MessageAttention,
239        parts: Vec<InboundMessagePart>,
240    },
241    ActionSelected {
242        token: String,
243    },
244}
245
246#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
247#[serde(rename_all = "camelCase", deny_unknown_fields)]
248pub struct NormalizedInbound {
249    pub schema_version: u16,
250    pub event_id: String,
251    pub address: MessagingAddress,
252    pub actor: ExternalActor,
253    pub content: NormalizedInboundContent,
254    #[serde(default, skip_serializing_if = "Option::is_none")]
255    pub conversation_display_name: Option<String>,
256    #[serde(default, skip_serializing_if = "Option::is_none")]
257    pub occurred_at_ms: Option<i64>,
258}
259
260impl NormalizedInbound {
261    pub fn validate(self) -> Result<Self, MessagingModelError> {
262        if self.schema_version != NORMALIZED_INBOUND_SCHEMA_VERSION {
263            return Err(MessagingModelError(
264                "unsupported normalized messaging schema version",
265            ));
266        }
267        let event_id = segment(self.event_id, "messaging event id is invalid")?;
268        let address = self.address.validate()?;
269        let actor = self.actor.validate()?;
270        if actor.provider != address.provider || actor.account_id != address.account_id {
271            return Err(MessagingModelError(
272                "messaging actor and address account do not match",
273            ));
274        }
275        if self.occurred_at_ms.is_some_and(|value| value < 0) {
276            return Err(MessagingModelError(
277                "messaging occurrence time cannot be negative",
278            ));
279        }
280        validate_content(&self.content)?;
281        optional_ref(
282            &self.conversation_display_name,
283            "messaging conversation display name is invalid",
284        )?;
285        Ok(Self {
286            schema_version: self.schema_version,
287            event_id,
288            address,
289            actor,
290            content: self.content,
291            conversation_display_name: self.conversation_display_name,
292            occurred_at_ms: self.occurred_at_ms,
293        })
294    }
295}
296
297fn validate_content(content: &NormalizedInboundContent) -> Result<(), MessagingModelError> {
298    match content {
299        NormalizedInboundContent::Message {
300            provider_message_id,
301            attention: _,
302            parts,
303        } => {
304            if parts.is_empty() {
305                return Err(MessagingModelError("messaging message has no parts"));
306            }
307            optional_ref(provider_message_id, "provider message id is invalid")?;
308            for part in parts {
309                match part {
310                    InboundMessagePart::Text { text } if text.is_empty() => {
311                        return Err(MessagingModelError("messaging text is invalid"));
312                    }
313                    InboundMessagePart::Text { .. } => {}
314                    InboundMessagePart::Media { reference } => validate_media(reference)?,
315                }
316            }
317        }
318        NormalizedInboundContent::ActionSelected { token } => {
319            segment_ref(token, "action token is invalid")?;
320        }
321    }
322    Ok(())
323}
324
325fn validate_media(reference: &ProviderMediaRef) -> Result<(), MessagingModelError> {
326    segment_ref(&reference.handle, "provider media handle is invalid")?;
327    optional_ref(&reference.mime_type, "provider media MIME type is invalid")?;
328    optional_ref(&reference.file_name, "provider media file name is invalid")?;
329    optional_ref(&reference.caption, "provider media caption is invalid")?;
330    optional_ref(
331        &reference.transcript,
332        "provider media transcript is invalid",
333    )?;
334    if reference.size_bytes == Some(0)
335        || reference.duration_ms == Some(0)
336        || reference.width_px == Some(0)
337        || reference.height_px == Some(0)
338    {
339        return Err(MessagingModelError("provider media metadata is invalid"));
340    }
341    if reference.width_px.is_some() != reference.height_px.is_some() {
342        return Err(MessagingModelError(
343            "provider media dimensions must be complete",
344        ));
345    }
346    Ok(())
347}
348
349fn provider_id(value: String) -> Result<String, MessagingModelError> {
350    let normalized = value.trim().to_ascii_lowercase();
351    if normalized.is_empty()
352        || !normalized
353            .as_bytes()
354            .first()
355            .is_some_and(u8::is_ascii_lowercase)
356        || !normalized
357            .bytes()
358            .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
359    {
360        return Err(MessagingModelError("messaging provider id is invalid"));
361    }
362    Ok(normalized)
363}
364
365fn segment(value: String, error: &'static str) -> Result<String, MessagingModelError> {
366    segment_ref(&value, error)?;
367    Ok(value)
368}
369
370fn optional_segment(value: String) -> Result<String, MessagingModelError> {
371    segment(value, "messaging display name is invalid")
372}
373
374fn segment_ref(value: &str, error: &'static str) -> Result<(), MessagingModelError> {
375    if value.is_empty() || value.trim() != value {
376        return Err(MessagingModelError(error));
377    }
378    Ok(())
379}
380
381fn optional_ref(value: &Option<String>, error: &'static str) -> Result<(), MessagingModelError> {
382    if let Some(value) = value {
383        segment_ref(value, error)?;
384    }
385    Ok(())
386}
387
388#[cfg(test)]
389mod tests {
390    use super::*;
391
392    #[test]
393    fn address_owns_lane_and_canonical_surface_mapping() {
394        let address = MessagingAddress::new(
395            "Telegram",
396            "bot:1",
397            "chat:2",
398            Some("topic:3".to_string()),
399            ConversationAudience::Shared,
400        )
401        .unwrap();
402        assert_eq!(address.provider, "telegram");
403        assert_eq!(address.base_address().lane_id, None);
404        let surface = address.conversation_surface().unwrap();
405        let route = surface.messaging_route().unwrap();
406        assert_eq!(route.lane_id, Some("topic:3"));
407        assert!(route.group);
408    }
409
410    #[test]
411    fn inbound_rejects_actor_from_another_provider_account() {
412        let inbound = NormalizedInbound {
413            schema_version: NORMALIZED_INBOUND_SCHEMA_VERSION,
414            event_id: "event".to_string(),
415            address: MessagingAddress::new(
416                "telegram",
417                "bot-a",
418                "chat",
419                None,
420                ConversationAudience::Personal,
421            )
422            .unwrap(),
423            actor: ExternalActor::new("telegram", "bot-b", "user", None).unwrap(),
424            content: NormalizedInboundContent::Message {
425                provider_message_id: None,
426                attention: MessageAttention::Unknown,
427                parts: vec![InboundMessagePart::Text {
428                    text: "hello".to_string(),
429                }],
430            },
431            conversation_display_name: None,
432            occurred_at_ms: None,
433        };
434        assert!(inbound.validate().is_err());
435    }
436
437    #[test]
438    fn actions_only_expose_labels_and_opaque_tokens() {
439        let actions = ActionSet::new(vec![
440            ActionOption::new("Approve", "route-approve").unwrap(),
441            ActionOption::new("Reject", "route-reject").unwrap(),
442        ])
443        .unwrap();
444        assert_eq!(actions.options.len(), 2);
445        assert!(ActionSet::new(Vec::new()).is_err());
446        assert!(ActionOption::new("Approve", " ").is_err());
447    }
448
449    #[test]
450    fn inbound_text_preserves_whitespace_allowed_by_the_schema() {
451        let inbound = NormalizedInbound {
452            schema_version: NORMALIZED_INBOUND_SCHEMA_VERSION,
453            event_id: "event".to_string(),
454            address: MessagingAddress::new(
455                "telegram",
456                "bot",
457                "chat",
458                None,
459                ConversationAudience::Personal,
460            )
461            .unwrap(),
462            actor: ExternalActor::new("telegram", "bot", "user", None).unwrap(),
463            content: NormalizedInboundContent::Message {
464                provider_message_id: None,
465                attention: MessageAttention::Unknown,
466                parts: vec![InboundMessagePart::Text {
467                    text: " message with spacing ".to_string(),
468                }],
469            },
470            conversation_display_name: None,
471            occurred_at_ms: None,
472        };
473        assert!(inbound.validate().is_ok());
474    }
475}