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 = 2;
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#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
15#[serde(rename_all = "camelCase", deny_unknown_fields)]
16pub struct MessagingAddress {
17    pub provider: String,
18    pub account_id: String,
19    pub conversation_id: String,
20    #[serde(default, skip_serializing_if = "Option::is_none")]
21    pub lane_id: Option<String>,
22    pub audience: ConversationAudience,
23}
24
25#[derive(Debug, Clone, PartialEq, Eq)]
26pub struct MessagingModelError(&'static str);
27
28impl fmt::Display for MessagingModelError {
29    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
30        formatter.write_str(self.0)
31    }
32}
33
34impl std::error::Error for MessagingModelError {}
35
36impl MessagingAddress {
37    pub fn new(
38        provider: impl Into<String>,
39        account_id: impl Into<String>,
40        conversation_id: impl Into<String>,
41        lane_id: Option<String>,
42        audience: ConversationAudience,
43    ) -> Result<Self, MessagingModelError> {
44        Ok(Self {
45            provider: provider_id(provider.into())?,
46            account_id: segment(account_id.into(), "messaging account id is invalid")?,
47            conversation_id: segment(
48                conversation_id.into(),
49                "messaging conversation id is invalid",
50            )?,
51            lane_id: lane_id
52                .map(|value| segment(value, "messaging lane id is invalid"))
53                .transpose()?,
54            audience,
55        })
56    }
57
58    #[must_use]
59    pub fn base_address(&self) -> Self {
60        let mut base = self.clone();
61        base.lane_id = None;
62        base
63    }
64
65    pub fn conversation_surface(&self) -> Result<ConversationSurface, MessagingModelError> {
66        let result = match self.audience {
67            ConversationAudience::Personal => ConversationSurface::messaging_personal(
68                &self.provider,
69                &self.account_id,
70                &self.conversation_id,
71                self.lane_id.clone(),
72            ),
73            ConversationAudience::Shared => ConversationSurface::messaging_group(
74                &self.provider,
75                &self.account_id,
76                &self.conversation_id,
77                self.lane_id.clone(),
78            ),
79        };
80        result.map_err(|_| MessagingModelError("messaging address cannot form a surface"))
81    }
82
83    pub fn validate(self) -> Result<Self, MessagingModelError> {
84        Self::new(
85            self.provider,
86            self.account_id,
87            self.conversation_id,
88            self.lane_id,
89            self.audience,
90        )
91    }
92}
93
94#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
95#[serde(rename_all = "camelCase", deny_unknown_fields)]
96pub struct ExternalActor {
97    pub provider: String,
98    pub account_id: String,
99    pub external_user_id: String,
100    #[serde(default, skip_serializing_if = "Option::is_none")]
101    pub display_name: Option<String>,
102}
103
104impl ExternalActor {
105    pub fn new(
106        provider: impl Into<String>,
107        account_id: impl Into<String>,
108        external_user_id: impl Into<String>,
109        display_name: Option<String>,
110    ) -> Result<Self, MessagingModelError> {
111        Ok(Self {
112            provider: provider_id(provider.into())?,
113            account_id: segment(account_id.into(), "messaging actor account id is invalid")?,
114            external_user_id: segment(
115                external_user_id.into(),
116                "messaging external user id is invalid",
117            )?,
118            display_name: display_name.map(optional_segment).transpose()?,
119        })
120    }
121
122    pub fn validate(self) -> Result<Self, MessagingModelError> {
123        Self::new(
124            self.provider,
125            self.account_id,
126            self.external_user_id,
127            self.display_name,
128        )
129    }
130}
131
132#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
133#[serde(rename_all = "camelCase", deny_unknown_fields)]
134pub struct ActionOption {
135    pub label: String,
136    pub token: String,
137}
138
139impl ActionOption {
140    pub fn new(
141        label: impl Into<String>,
142        token: impl Into<String>,
143    ) -> Result<Self, MessagingModelError> {
144        Ok(Self {
145            label: segment(label.into(), "action label is invalid")?,
146            token: segment(token.into(), "action token is invalid")?,
147        })
148    }
149}
150
151#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
152#[serde(rename_all = "camelCase", deny_unknown_fields)]
153pub struct ActionSet {
154    pub options: Vec<ActionOption>,
155}
156
157impl ActionSet {
158    pub fn new(options: Vec<ActionOption>) -> Result<Self, MessagingModelError> {
159        if options.is_empty() {
160            return Err(MessagingModelError("action set cannot be empty"));
161        }
162        for option in &options {
163            segment_ref(&option.label, "action label is invalid")?;
164            segment_ref(&option.token, "action token is invalid")?;
165        }
166        Ok(Self { options })
167    }
168}
169
170#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
171#[serde(
172    tag = "type",
173    rename_all = "snake_case",
174    rename_all_fields = "camelCase",
175    deny_unknown_fields
176)]
177pub enum NormalizedInboundContent {
178    Text {
179        text: String,
180        #[serde(default, skip_serializing_if = "Option::is_none")]
181        provider_message_id: Option<String>,
182    },
183    Audio {
184        provider_message_id: String,
185        provider_file_id: String,
186        #[serde(default, skip_serializing_if = "Option::is_none")]
187        duration_seconds: Option<u32>,
188    },
189    ActionSelected {
190        token: String,
191    },
192}
193
194#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
195#[serde(rename_all = "camelCase", deny_unknown_fields)]
196pub struct NormalizedInbound {
197    pub schema_version: u16,
198    pub event_id: String,
199    pub address: MessagingAddress,
200    pub actor: ExternalActor,
201    pub content: NormalizedInboundContent,
202    #[serde(default, skip_serializing_if = "Option::is_none")]
203    pub conversation_display_name: Option<String>,
204    #[serde(default, skip_serializing_if = "Option::is_none")]
205    pub occurred_at_ms: Option<i64>,
206}
207
208impl NormalizedInbound {
209    pub fn validate(self) -> Result<Self, MessagingModelError> {
210        if self.schema_version != NORMALIZED_INBOUND_SCHEMA_VERSION {
211            return Err(MessagingModelError(
212                "unsupported normalized messaging schema version",
213            ));
214        }
215        let event_id = segment(self.event_id, "messaging event id is invalid")?;
216        let address = self.address.validate()?;
217        let actor = self.actor.validate()?;
218        if actor.provider != address.provider || actor.account_id != address.account_id {
219            return Err(MessagingModelError(
220                "messaging actor and address account do not match",
221            ));
222        }
223        if self.occurred_at_ms.is_some_and(|value| value < 0) {
224            return Err(MessagingModelError(
225                "messaging occurrence time cannot be negative",
226            ));
227        }
228        validate_content(&self.content)?;
229        optional_ref(
230            &self.conversation_display_name,
231            "messaging conversation display name is invalid",
232        )?;
233        Ok(Self {
234            schema_version: self.schema_version,
235            event_id,
236            address,
237            actor,
238            content: self.content,
239            conversation_display_name: self.conversation_display_name,
240            occurred_at_ms: self.occurred_at_ms,
241        })
242    }
243}
244
245fn validate_content(content: &NormalizedInboundContent) -> Result<(), MessagingModelError> {
246    match content {
247        NormalizedInboundContent::Text {
248            text,
249            provider_message_id,
250        } => {
251            if text.is_empty() {
252                return Err(MessagingModelError("messaging text is invalid"));
253            }
254            optional_ref(provider_message_id, "provider message id is invalid")?;
255        }
256        NormalizedInboundContent::Audio {
257            provider_message_id,
258            provider_file_id,
259            ..
260        } => {
261            segment_ref(provider_message_id, "provider message id is invalid")?;
262            segment_ref(provider_file_id, "provider file id is invalid")?;
263        }
264        NormalizedInboundContent::ActionSelected { token } => {
265            segment_ref(token, "action token is invalid")?;
266        }
267    }
268    Ok(())
269}
270
271fn provider_id(value: String) -> Result<String, MessagingModelError> {
272    let normalized = value.trim().to_ascii_lowercase();
273    if normalized.is_empty()
274        || !normalized
275            .as_bytes()
276            .first()
277            .is_some_and(u8::is_ascii_lowercase)
278        || !normalized
279            .bytes()
280            .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
281    {
282        return Err(MessagingModelError("messaging provider id is invalid"));
283    }
284    Ok(normalized)
285}
286
287fn segment(value: String, error: &'static str) -> Result<String, MessagingModelError> {
288    segment_ref(&value, error)?;
289    Ok(value)
290}
291
292fn optional_segment(value: String) -> Result<String, MessagingModelError> {
293    segment(value, "messaging display name is invalid")
294}
295
296fn segment_ref(value: &str, error: &'static str) -> Result<(), MessagingModelError> {
297    if value.is_empty() || value.trim() != value {
298        return Err(MessagingModelError(error));
299    }
300    Ok(())
301}
302
303fn optional_ref(value: &Option<String>, error: &'static str) -> Result<(), MessagingModelError> {
304    if let Some(value) = value {
305        segment_ref(value, error)?;
306    }
307    Ok(())
308}
309
310#[cfg(test)]
311mod tests {
312    use super::*;
313
314    #[test]
315    fn address_owns_lane_and_canonical_surface_mapping() {
316        let address = MessagingAddress::new(
317            "Telegram",
318            "bot:1",
319            "chat:2",
320            Some("topic:3".to_string()),
321            ConversationAudience::Shared,
322        )
323        .unwrap();
324        assert_eq!(address.provider, "telegram");
325        assert_eq!(address.base_address().lane_id, None);
326        let surface = address.conversation_surface().unwrap();
327        let route = surface.messaging_route().unwrap();
328        assert_eq!(route.lane_id, Some("topic:3"));
329        assert!(route.group);
330    }
331
332    #[test]
333    fn inbound_rejects_actor_from_another_provider_account() {
334        let inbound = NormalizedInbound {
335            schema_version: NORMALIZED_INBOUND_SCHEMA_VERSION,
336            event_id: "event".to_string(),
337            address: MessagingAddress::new(
338                "telegram",
339                "bot-a",
340                "chat",
341                None,
342                ConversationAudience::Personal,
343            )
344            .unwrap(),
345            actor: ExternalActor::new("telegram", "bot-b", "user", None).unwrap(),
346            content: NormalizedInboundContent::Text {
347                text: "hello".to_string(),
348                provider_message_id: None,
349            },
350            conversation_display_name: None,
351            occurred_at_ms: None,
352        };
353        assert!(inbound.validate().is_err());
354    }
355
356    #[test]
357    fn actions_only_expose_labels_and_opaque_tokens() {
358        let actions = ActionSet::new(vec![
359            ActionOption::new("Approve", "route-approve").unwrap(),
360            ActionOption::new("Reject", "route-reject").unwrap(),
361        ])
362        .unwrap();
363        assert_eq!(actions.options.len(), 2);
364        assert!(ActionSet::new(Vec::new()).is_err());
365        assert!(ActionOption::new("Approve", " ").is_err());
366    }
367
368    #[test]
369    fn inbound_text_preserves_whitespace_allowed_by_the_schema() {
370        let inbound = NormalizedInbound {
371            schema_version: NORMALIZED_INBOUND_SCHEMA_VERSION,
372            event_id: "event".to_string(),
373            address: MessagingAddress::new(
374                "telegram",
375                "bot",
376                "chat",
377                None,
378                ConversationAudience::Personal,
379            )
380            .unwrap(),
381            actor: ExternalActor::new("telegram", "bot", "user", None).unwrap(),
382            content: NormalizedInboundContent::Text {
383                text: " message with spacing ".to_string(),
384                provider_message_id: None,
385            },
386            conversation_display_name: None,
387            occurred_at_ms: None,
388        };
389        assert!(inbound.validate().is_ok());
390    }
391}