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#[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#[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}