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() || parts.len() > 16 {
305                return Err(MessagingModelError(
306                    "messaging message part count is invalid",
307                ));
308            }
309            optional_ref(provider_message_id, "provider message id is invalid")?;
310            for part in parts {
311                match part {
312                    InboundMessagePart::Text { text } if text.is_empty() => {
313                        return Err(MessagingModelError("messaging text is invalid"));
314                    }
315                    InboundMessagePart::Text { .. } => {}
316                    InboundMessagePart::Media { reference } => validate_media(reference)?,
317                }
318            }
319        }
320        NormalizedInboundContent::ActionSelected { token } => {
321            segment_ref(token, "action token is invalid")?;
322        }
323    }
324    Ok(())
325}
326
327fn validate_media(reference: &ProviderMediaRef) -> Result<(), MessagingModelError> {
328    segment_ref(&reference.handle, "provider media handle is invalid")?;
329    optional_ref(&reference.mime_type, "provider media MIME type is invalid")?;
330    optional_ref(&reference.file_name, "provider media file name is invalid")?;
331    optional_ref(&reference.caption, "provider media caption is invalid")?;
332    optional_ref(
333        &reference.transcript,
334        "provider media transcript is invalid",
335    )?;
336    if reference.size_bytes == Some(0)
337        || reference.duration_ms == Some(0)
338        || reference.width_px == Some(0)
339        || reference.height_px == Some(0)
340    {
341        return Err(MessagingModelError("provider media metadata is invalid"));
342    }
343    if reference.width_px.is_some() != reference.height_px.is_some() {
344        return Err(MessagingModelError(
345            "provider media dimensions must be complete",
346        ));
347    }
348    Ok(())
349}
350
351fn provider_id(value: String) -> Result<String, MessagingModelError> {
352    let normalized = value.trim().to_ascii_lowercase();
353    if normalized.is_empty()
354        || !normalized
355            .as_bytes()
356            .first()
357            .is_some_and(u8::is_ascii_lowercase)
358        || !normalized
359            .bytes()
360            .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
361    {
362        return Err(MessagingModelError("messaging provider id is invalid"));
363    }
364    Ok(normalized)
365}
366
367fn segment(value: String, error: &'static str) -> Result<String, MessagingModelError> {
368    segment_ref(&value, error)?;
369    Ok(value)
370}
371
372fn optional_segment(value: String) -> Result<String, MessagingModelError> {
373    segment(value, "messaging display name is invalid")
374}
375
376fn segment_ref(value: &str, error: &'static str) -> Result<(), MessagingModelError> {
377    if value.is_empty() || value.trim() != value {
378        return Err(MessagingModelError(error));
379    }
380    Ok(())
381}
382
383fn optional_ref(value: &Option<String>, error: &'static str) -> Result<(), MessagingModelError> {
384    if let Some(value) = value {
385        segment_ref(value, error)?;
386    }
387    Ok(())
388}
389
390#[cfg(test)]
391mod tests {
392    use super::*;
393
394    #[test]
395    fn address_owns_lane_and_canonical_surface_mapping() {
396        let address = MessagingAddress::new(
397            "Telegram",
398            "bot:1",
399            "chat:2",
400            Some("topic:3".to_string()),
401            ConversationAudience::Shared,
402        )
403        .unwrap();
404        assert_eq!(address.provider, "telegram");
405        assert_eq!(address.base_address().lane_id, None);
406        let surface = address.conversation_surface().unwrap();
407        let route = surface.messaging_route().unwrap();
408        assert_eq!(route.lane_id, Some("topic:3"));
409        assert!(route.group);
410    }
411
412    #[test]
413    fn inbound_rejects_actor_from_another_provider_account() {
414        let inbound = NormalizedInbound {
415            schema_version: NORMALIZED_INBOUND_SCHEMA_VERSION,
416            event_id: "event".to_string(),
417            address: MessagingAddress::new(
418                "telegram",
419                "bot-a",
420                "chat",
421                None,
422                ConversationAudience::Personal,
423            )
424            .unwrap(),
425            actor: ExternalActor::new("telegram", "bot-b", "user", None).unwrap(),
426            content: NormalizedInboundContent::Message {
427                provider_message_id: None,
428                attention: MessageAttention::Unknown,
429                parts: vec![InboundMessagePart::Text {
430                    text: "hello".to_string(),
431                }],
432            },
433            conversation_display_name: None,
434            occurred_at_ms: None,
435        };
436        assert!(inbound.validate().is_err());
437    }
438
439    #[test]
440    fn actions_only_expose_labels_and_opaque_tokens() {
441        let actions = ActionSet::new(vec![
442            ActionOption::new("Approve", "route-approve").unwrap(),
443            ActionOption::new("Reject", "route-reject").unwrap(),
444        ])
445        .unwrap();
446        assert_eq!(actions.options.len(), 2);
447        assert!(ActionSet::new(Vec::new()).is_err());
448        assert!(ActionOption::new("Approve", " ").is_err());
449    }
450
451    #[test]
452    fn inbound_text_preserves_whitespace_allowed_by_the_schema() {
453        let inbound = NormalizedInbound {
454            schema_version: NORMALIZED_INBOUND_SCHEMA_VERSION,
455            event_id: "event".to_string(),
456            address: MessagingAddress::new(
457                "telegram",
458                "bot",
459                "chat",
460                None,
461                ConversationAudience::Personal,
462            )
463            .unwrap(),
464            actor: ExternalActor::new("telegram", "bot", "user", None).unwrap(),
465            content: NormalizedInboundContent::Message {
466                provider_message_id: None,
467                attention: MessageAttention::Unknown,
468                parts: vec![InboundMessagePart::Text {
469                    text: " message with spacing ".to_string(),
470                }],
471            },
472            conversation_display_name: None,
473            occurred_at_ms: None,
474        };
475        assert!(inbound.validate().is_ok());
476    }
477
478    #[test]
479    fn inbound_rejects_more_parts_than_the_schema_allows() {
480        let inbound = NormalizedInbound {
481            schema_version: NORMALIZED_INBOUND_SCHEMA_VERSION,
482            event_id: "event".to_string(),
483            address: MessagingAddress::new(
484                "telegram",
485                "bot",
486                "chat",
487                None,
488                ConversationAudience::Personal,
489            )
490            .unwrap(),
491            actor: ExternalActor::new("telegram", "bot", "user", None).unwrap(),
492            content: NormalizedInboundContent::Message {
493                provider_message_id: None,
494                attention: MessageAttention::Unknown,
495                parts: (0..17)
496                    .map(|index| InboundMessagePart::Text {
497                        text: format!("part-{index}"),
498                    })
499                    .collect(),
500            },
501            conversation_display_name: None,
502            occurred_at_ms: None,
503        };
504        assert!(inbound.validate().is_err());
505    }
506}