Skip to main content

mobius_gateway/
swarm.rs

1//! Durable gateway-owned swarm membership and message-board state.
2
3use std::collections::{BTreeMap, BTreeSet, VecDeque};
4use std::sync::Arc;
5use std::time::{SystemTime, UNIX_EPOCH};
6
7use mobius::backend::checkpoint::CheckpointStore;
8use mobius::middleware::swarm::SwarmBackend;
9use serde::{Deserialize, Serialize};
10use tokio::sync::{Mutex, MutexGuard, mpsc};
11use uuid::Uuid;
12
13use crate::wire::{SwarmMemberRecord, SwarmMessageRecord, SwarmRecord};
14use crate::{Error, Result};
15
16const STATE_SCOPE: &str = "gateway";
17const STATE_KEY: &str = "swarms.v1";
18const MAX_HANDLE_BYTES: usize = 64;
19const MAX_ID_BYTES: usize = 512;
20const MAX_TITLE_BYTES: usize = 256;
21const MAX_MESSAGE_BYTES: usize = 24_000;
22const MAX_ACKNOWLEDGED_ENTRIES: usize = 256;
23const MAX_PENDING_DELIVERIES_PER_RECIPIENT: usize = 256;
24const MAX_PAGE_ENTRIES: usize = 256;
25const MAX_SWARM_MEMBERS: usize = 100;
26// ponytail: one 8 MiB catalog; split board storage and wire pages if real usage reaches it.
27const MAX_CATALOG_BYTES: usize = 8 * 1024 * 1024;
28const MAX_TOOL_READ_BYTES: usize = 32_000;
29const MAX_TOOL_READ_BODY_BYTES: usize = 4_000;
30
31/// One conversation participating in a swarm.
32#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
33#[serde(deny_unknown_fields)]
34pub struct SwarmMember {
35    /// Visible gateway session identifier.
36    pub session_id: String,
37    /// Unique, mentionable handle within this swarm.
38    pub handle: String,
39    /// Time this membership was created, in Unix milliseconds.
40    pub joined_at_ms: i64,
41}
42
43/// Compact durable swarm metadata and its current roster.
44#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
45#[serde(deny_unknown_fields)]
46pub struct SwarmSummary {
47    /// Stable generated swarm identifier.
48    pub id: String,
49    /// User-visible swarm title.
50    pub title: String,
51    /// Session identifier of the swarm leader.
52    pub leader_session_id: String,
53    /// Current members ordered by handle.
54    pub members: Vec<SwarmMember>,
55    /// Latest board sequence assigned in this swarm.
56    pub latest_sequence: u64,
57    /// Time the swarm was created, in Unix milliseconds.
58    pub created_at_ms: i64,
59    /// Time its roster or board was last changed, in Unix milliseconds.
60    pub updated_at_ms: i64,
61}
62
63/// A swarm snapshot resolved for one participating session.
64#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65#[serde(deny_unknown_fields)]
66pub struct SwarmSnapshot {
67    /// The swarm and current roster.
68    pub swarm: SwarmSummary,
69    /// This session's immutable handle in the swarm.
70    pub handle: String,
71}
72
73/// One durable swarm board entry.
74#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
75#[serde(deny_unknown_fields)]
76pub struct BoardEntry {
77    /// Stable generated message identifier.
78    pub id: String,
79    /// Monotonic sequence within the swarm board.
80    pub sequence: u64,
81    /// Time the entry was created, in Unix milliseconds.
82    pub created_at_ms: i64,
83    /// Author identity as it existed when the message was posted.
84    pub author: SwarmMember,
85    /// Message text.
86    pub text: String,
87    /// Session identifiers resolved from the entry's `@handle` mentions.
88    pub mentioned_recipient_session_ids: Vec<String>,
89    /// Mentioned sessions which have not acknowledged delivery.
90    pub pending_recipient_session_ids: Vec<String>,
91}
92
93/// The durable entry created by a post and the conversations to notify.
94#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
95#[serde(deny_unknown_fields)]
96pub struct SwarmPost {
97    /// Newly-created durable board entry.
98    pub entry: BoardEntry,
99    /// Resolved recipient sessions, ready for the gateway delivery bridge.
100    pub resolved_recipient_session_ids: Vec<String>,
101}
102
103/// One newest-first page of swarm board entries.
104#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
105#[serde(deny_unknown_fields)]
106pub struct BoardPage {
107    /// Newest-first entries in this page.
108    pub entries: Vec<BoardEntry>,
109    /// Sequence to pass as `before_sequence` for the next older page.
110    pub next_before_sequence: Option<u64>,
111}
112
113/// One message awaiting delivery to a particular conversation.
114#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
115#[serde(deny_unknown_fields)]
116pub struct PendingDelivery {
117    /// Swarm containing the message.
118    pub swarm_id: String,
119    /// User-visible title snapshotted from the current swarm.
120    pub swarm_title: String,
121    /// Durable message awaiting acknowledgement.
122    pub entry: BoardEntry,
123}
124
125/// Wake-up signal consumed by the gateway's swarm delivery worker.
126#[derive(Debug, Clone, PartialEq, Eq)]
127pub(crate) enum SwarmDelivery {
128    /// Durable board contents changed and connected clients need a fresh catalog.
129    Changed,
130    /// Gateway capacity changed and durable pending recipients should be retried.
131    RetryPending,
132    /// A target has at least one durable board message awaiting delivery.
133    Pending { target_session_id: String },
134    /// A target durably recorded one delivered peer message.
135    Acknowledged {
136        target_session_id: String,
137        message_id: String,
138    },
139}
140
141/// Cloneable access to the gateway's durable swarm catalog.
142#[derive(Clone)]
143pub struct SwarmStore {
144    checkpoints: Arc<dyn CheckpointStore>,
145    state: Arc<Mutex<Option<Catalog>>>,
146    deliveries: mpsc::UnboundedSender<SwarmDelivery>,
147}
148
149#[derive(Clone, Default, Serialize, Deserialize)]
150#[serde(deny_unknown_fields)]
151struct Catalog {
152    swarms: BTreeMap<String, StoredSwarm>,
153}
154
155#[derive(Clone, Serialize, Deserialize)]
156#[serde(deny_unknown_fields)]
157struct StoredSwarm {
158    title: String,
159    leader_session_id: String,
160    members: BTreeMap<String, SwarmMember>,
161    latest_sequence: u64,
162    board: VecDeque<BoardEntry>,
163    created_at_ms: i64,
164    updated_at_ms: i64,
165}
166
167impl SwarmStore {
168    /// Creates a lazily-loaded store and its single gateway delivery queue.
169    #[must_use]
170    pub(crate) fn new(
171        checkpoints: Arc<dyn CheckpointStore>,
172    ) -> (Self, mpsc::UnboundedReceiver<SwarmDelivery>) {
173        let (deliveries, receiver) = mpsc::unbounded_channel();
174        (
175            Self {
176                checkpoints,
177                state: Arc::new(Mutex::new(None)),
178                deliveries,
179            },
180            receiver,
181        )
182    }
183
184    /// Returns every swarm ordered by its stable identifier.
185    pub async fn summaries(&self) -> Result<Vec<SwarmSummary>> {
186        let state = self.lock_loaded().await?;
187        Ok(state
188            .as_ref()
189            .expect("swarm catalog loaded")
190            .swarms
191            .iter()
192            .map(|(id, swarm)| swarm.summary(id))
193            .collect())
194    }
195
196    pub(crate) async fn records(&self) -> Result<Vec<SwarmRecord>> {
197        let state = self.lock_loaded().await?;
198        Ok(state
199            .as_ref()
200            .expect("swarm catalog loaded")
201            .swarms
202            .iter()
203            .map(|(id, swarm)| {
204                let mut messages = swarm
205                    .board
206                    .iter()
207                    .rev()
208                    .take(MAX_PAGE_ENTRIES)
209                    .map(|entry| SwarmMessageRecord {
210                        id: entry.id.clone(),
211                        sequence: entry.sequence,
212                        author_session_id: entry.author.session_id.clone(),
213                        author_handle: entry.author.handle.clone(),
214                        body: entry.text.clone(),
215                        created_at_ms: entry.created_at_ms,
216                    })
217                    .collect::<Vec<_>>();
218                messages.reverse();
219                SwarmRecord {
220                    id: id.clone(),
221                    title: swarm.title.clone(),
222                    leader_session_id: swarm.leader_session_id.clone(),
223                    members: swarm
224                        .members
225                        .values()
226                        .map(|member| SwarmMemberRecord {
227                            session_id: member.session_id.clone(),
228                            handle: member.handle.clone(),
229                        })
230                        .collect(),
231                    messages,
232                    updated_at_ms: swarm.updated_at_ms,
233                }
234            })
235            .collect())
236    }
237
238    /// Creates a neutrally named swarm with gateway-generated member handles.
239    pub(crate) async fn create(
240        &self,
241        leader_session_id: String,
242        member_session_ids: Vec<String>,
243    ) -> Result<SwarmSummary> {
244        validate_swarm_members(&leader_session_id, &member_session_ids)?;
245
246        self.mutate(move |catalog| {
247            for session_id in &member_session_ids {
248                ensure_session_available(catalog, session_id)?;
249            }
250            let id = Uuid::new_v4();
251            let title = format!(
252                "Swarm {}",
253                id.simple().to_string().chars().take(8).collect::<String>()
254            );
255            let mut members = BTreeMap::new();
256            for session_id in member_session_ids {
257                let member = generated_member(session_id, &members);
258                members.insert(member.handle.clone(), member);
259            }
260            let now = unix_ms();
261            let swarm = StoredSwarm {
262                title,
263                leader_session_id,
264                members,
265                latest_sequence: 0,
266                board: VecDeque::new(),
267                created_at_ms: now,
268                updated_at_ms: now,
269            };
270            let id = id.to_string();
271            let summary = swarm.summary(&id);
272            catalog.swarms.insert(id, swarm);
273            Ok(summary)
274        })
275        .await
276    }
277
278    /// Adds one session under a new gateway-generated handle.
279    pub(crate) async fn join(&self, swarm_id: &str, session_id: String) -> Result<SwarmSummary> {
280        validate_swarm_id(swarm_id)?;
281        validate_session_id(&session_id)?;
282        let swarm_id = swarm_id.to_owned();
283        self.mutate(move |catalog| {
284            ensure_session_available(catalog, &session_id)?;
285            let swarm = catalog
286                .swarms
287                .get_mut(&swarm_id)
288                .ok_or_else(|| config(format!("unknown swarm `{swarm_id}`")))?;
289            if swarm.members.len() >= MAX_SWARM_MEMBERS {
290                return Err(config(format!(
291                    "a swarm supports at most {MAX_SWARM_MEMBERS} members"
292                )));
293            }
294            let member = generated_member(session_id, &swarm.members);
295            swarm.members.insert(member.handle.clone(), member);
296            swarm.updated_at_ms = unix_ms();
297            Ok(swarm.summary(&swarm_id))
298        })
299        .await
300    }
301
302    /// Removes a non-leader session from a swarm.
303    ///
304    /// A leader must explicitly disband its swarm. A one-member roster remains valid.
305    pub(crate) async fn leave(&self, swarm_id: &str, session_id: &str) -> Result<SwarmSummary> {
306        validate_swarm_id(swarm_id)?;
307        validate_session_id(session_id)?;
308        let swarm_id = swarm_id.to_owned();
309        let session_id = session_id.to_owned();
310        let acknowledged_target = session_id.clone();
311        let (summary, acknowledged_messages) = self
312            .mutate(move |catalog| {
313                let swarm = catalog
314                    .swarms
315                    .get_mut(&swarm_id)
316                    .ok_or_else(|| config(format!("unknown swarm `{swarm_id}`")))?;
317                if swarm.leader_session_id == session_id {
318                    return Err(config("swarm leader must disband instead of leaving"));
319                }
320                let handle = swarm
321                    .members
322                    .iter()
323                    .find_map(|(handle, member)| {
324                        (member.session_id == session_id).then(|| handle.clone())
325                    })
326                    .ok_or_else(|| {
327                        config(format!("session `{session_id}` is not in this swarm"))
328                    })?;
329                swarm.members.remove(&handle);
330                let mut acknowledged_messages = Vec::new();
331                for entry in &mut swarm.board {
332                    if entry
333                        .pending_recipient_session_ids
334                        .iter()
335                        .any(|pending| pending == &session_id)
336                    {
337                        acknowledged_messages.push(entry.id.clone());
338                    }
339                    entry
340                        .pending_recipient_session_ids
341                        .retain(|pending| pending != &session_id);
342                }
343                trim_acknowledged(&mut swarm.board);
344                swarm.updated_at_ms = unix_ms();
345                Ok((swarm.summary(&swarm_id), acknowledged_messages))
346            })
347            .await?;
348        for message_id in acknowledged_messages {
349            self.notify_acknowledged(&message_id, &acknowledged_target);
350        }
351        Ok(summary)
352    }
353
354    /// Permanently removes a swarm and its board.
355    pub(crate) async fn disband(&self, swarm_id: &str) -> Result<SwarmSummary> {
356        validate_swarm_id(swarm_id)?;
357        let swarm_id = swarm_id.to_owned();
358        let (summary, pending_deliveries) = self
359            .mutate(move |catalog| {
360                let swarm = catalog
361                    .swarms
362                    .remove(&swarm_id)
363                    .ok_or_else(|| config(format!("unknown swarm `{swarm_id}`")))?;
364                let pending_deliveries = swarm
365                    .board
366                    .iter()
367                    .flat_map(|entry| {
368                        entry
369                            .pending_recipient_session_ids
370                            .iter()
371                            .map(|target| (entry.id.clone(), target.clone()))
372                    })
373                    .collect::<BTreeSet<_>>();
374                Ok((swarm.summary(&swarm_id), pending_deliveries))
375            })
376            .await?;
377        for (message_id, target_session_id) in pending_deliveries {
378            self.notify_acknowledged(&message_id, &target_session_id);
379        }
380        Ok(summary)
381    }
382
383    /// Resolves the swarm and handle for one participating session.
384    pub async fn snapshot_for_session(&self, session_id: &str) -> Result<Option<SwarmSnapshot>> {
385        validate_session_id(session_id)?;
386        let state = self.lock_loaded().await?;
387        let catalog = state.as_ref().expect("swarm catalog loaded");
388        Ok(catalog.swarms.iter().find_map(|(id, swarm)| {
389            swarm
390                .members
391                .values()
392                .find(|member| member.session_id == session_id)
393                .map(|member| SwarmSnapshot {
394                    swarm: swarm.summary(id),
395                    handle: member.handle.clone(),
396                })
397        }))
398    }
399
400    /// Posts one durable board entry and resolves all `@handle` recipients.
401    pub async fn post(&self, sender_session_id: &str, text: String) -> Result<SwarmPost> {
402        validate_session_id(sender_session_id)?;
403        validate_message(&text)?;
404        let sender_session_id = sender_session_id.to_owned();
405        let post = self
406            .mutate(move |catalog| {
407                let swarm_id =
408                    swarm_id_for_session(catalog, &sender_session_id).ok_or_else(|| {
409                        config(format!("session `{sender_session_id}` is not in a swarm"))
410                    })?;
411                let swarm = catalog
412                    .swarms
413                    .get_mut(&swarm_id)
414                    .expect("resolved swarm exists");
415                let author = swarm
416                    .members
417                    .values()
418                    .find(|member| member.session_id == sender_session_id)
419                    .cloned()
420                    .expect("resolved swarm contains sender");
421
422                let handles = mentioned_handles(&text);
423                let unknown = handles
424                    .iter()
425                    .filter(|handle| !swarm.members.contains_key(*handle))
426                    .cloned()
427                    .collect::<Vec<_>>();
428                if !unknown.is_empty() {
429                    return Err(config(format!(
430                        "unknown swarm mention{}: {}",
431                        if unknown.len() == 1 { "" } else { "s" },
432                        unknown
433                            .iter()
434                            .map(|handle| format!("@{handle}"))
435                            .collect::<Vec<_>>()
436                            .join(", ")
437                    )));
438                }
439
440                let recipients = handles
441                    .iter()
442                    .map(|handle| {
443                        swarm
444                            .members
445                            .get(handle)
446                            .expect("mentions validated against roster")
447                            .session_id
448                            .clone()
449                    })
450                    .collect::<BTreeSet<_>>()
451                    .into_iter()
452                    .collect::<Vec<_>>();
453                if recipients
454                    .iter()
455                    .any(|session| session == &sender_session_id)
456                {
457                    return Err(config("a swarm message cannot mention its author"));
458                }
459                if let Some(recipient) = recipients.iter().find(|recipient| {
460                    swarm
461                        .board
462                        .iter()
463                        .filter(|entry| {
464                            entry
465                                .pending_recipient_session_ids
466                                .iter()
467                                .any(|pending| pending == *recipient)
468                        })
469                        .count()
470                        >= MAX_PENDING_DELIVERIES_PER_RECIPIENT
471                }) {
472                    return Err(config(format!(
473                        "session `{recipient}` has {MAX_PENDING_DELIVERIES_PER_RECIPIENT} pending swarm messages"
474                    )));
475                }
476                let sequence = swarm
477                    .latest_sequence
478                    .checked_add(1)
479                    .ok_or_else(|| config("swarm board sequence exhausted"))?;
480                let entry = BoardEntry {
481                    id: Uuid::new_v4().to_string(),
482                    sequence,
483                    created_at_ms: unix_ms(),
484                    author,
485                    text,
486                    mentioned_recipient_session_ids: recipients.clone(),
487                    pending_recipient_session_ids: recipients.clone(),
488                };
489                swarm.latest_sequence = sequence;
490                swarm.updated_at_ms = entry.created_at_ms;
491                swarm.board.push_back(entry.clone());
492                trim_acknowledged(&mut swarm.board);
493                Ok(SwarmPost {
494                    entry,
495                    resolved_recipient_session_ids: recipients,
496                })
497            })
498            .await?;
499        let _ = self.deliveries.send(SwarmDelivery::Changed);
500        for target_session_id in &post.resolved_recipient_session_ids {
501            self.notify_pending(target_session_id);
502        }
503        Ok(post)
504    }
505
506    /// Loads one newest-first board page before an optional sequence cursor.
507    pub async fn board_page(
508        &self,
509        swarm_id: &str,
510        before_sequence: Option<u64>,
511        limit: usize,
512    ) -> Result<BoardPage> {
513        validate_swarm_id(swarm_id)?;
514        if limit == 0 || limit > MAX_PAGE_ENTRIES {
515            return Err(config(format!(
516                "board page limit must be 1-{MAX_PAGE_ENTRIES}"
517            )));
518        }
519        let state = self.lock_loaded().await?;
520        let swarm = state
521            .as_ref()
522            .expect("swarm catalog loaded")
523            .swarms
524            .get(swarm_id)
525            .ok_or_else(|| config(format!("unknown swarm `{swarm_id}`")))?;
526        let mut matches = swarm
527            .board
528            .iter()
529            .rev()
530            .filter(|entry| before_sequence.is_none_or(|before| entry.sequence < before));
531        let entries = matches.by_ref().take(limit).cloned().collect::<Vec<_>>();
532        let next_before_sequence = matches
533            .next()
534            .and_then(|_| entries.last().map(|entry| entry.sequence));
535        Ok(BoardPage {
536            entries,
537            next_before_sequence,
538        })
539    }
540
541    /// Returns messages still awaiting one target session's acknowledgement.
542    pub async fn pending_deliveries(
543        &self,
544        target_session_id: &str,
545    ) -> Result<Vec<PendingDelivery>> {
546        validate_session_id(target_session_id)?;
547        let state = self.lock_loaded().await?;
548        let catalog = state.as_ref().expect("swarm catalog loaded");
549        Ok(catalog
550            .swarms
551            .iter()
552            .flat_map(|(swarm_id, swarm)| {
553                swarm
554                    .board
555                    .iter()
556                    .filter(|entry| {
557                        entry
558                            .pending_recipient_session_ids
559                            .iter()
560                            .any(|pending| pending == target_session_id)
561                    })
562                    .map(|entry| PendingDelivery {
563                        swarm_id: swarm_id.clone(),
564                        swarm_title: swarm.title.clone(),
565                        entry: entry.clone(),
566                    })
567            })
568            .collect())
569    }
570
571    /// Returns each conversation with at least one durable pending mention.
572    pub(crate) async fn pending_recipient_session_ids(&self) -> Result<Vec<String>> {
573        let state = self.lock_loaded().await?;
574        Ok(state
575            .as_ref()
576            .expect("swarm catalog loaded")
577            .swarms
578            .values()
579            .flat_map(|swarm| &swarm.board)
580            .flat_map(|entry| &entry.pending_recipient_session_ids)
581            .cloned()
582            .collect::<BTreeSet<_>>()
583            .into_iter()
584            .collect())
585    }
586
587    pub(crate) fn notify_pending(&self, target_session_id: &str) {
588        let _ = self.deliveries.send(SwarmDelivery::Pending {
589            target_session_id: target_session_id.to_owned(),
590        });
591    }
592
593    pub(crate) fn retry_pending(&self) {
594        let _ = self.deliveries.send(SwarmDelivery::RetryPending);
595    }
596
597    pub(crate) fn notify_acknowledged(&self, message_id: &str, target_session_id: &str) {
598        let _ = self.deliveries.send(SwarmDelivery::Acknowledged {
599            target_session_id: target_session_id.to_owned(),
600            message_id: message_id.to_owned(),
601        });
602    }
603
604    /// Acknowledges one target's message delivery.
605    ///
606    /// Returns `false` when that target already acknowledged the entry.
607    pub async fn acknowledge(&self, message_id: &str, target_session_id: &str) -> Result<bool> {
608        validate_message_id(message_id)?;
609        validate_session_id(target_session_id)?;
610        let message_id = message_id.to_owned();
611        let target_session_id = target_session_id.to_owned();
612        let acknowledged_message = message_id.clone();
613        let acknowledged_target = target_session_id.clone();
614        let was_pending = self
615            .mutate(move |catalog| {
616                let entry = catalog
617                    .swarms
618                    .values_mut()
619                    .find_map(|swarm| swarm.board.iter_mut().find(|entry| entry.id == message_id))
620                    .ok_or_else(|| config(format!("unknown swarm message `{message_id}`")))?;
621                if !entry
622                    .mentioned_recipient_session_ids
623                    .iter()
624                    .any(|recipient| recipient == &target_session_id)
625                {
626                    return Err(config(format!(
627                        "session `{target_session_id}` is not a recipient of message `{message_id}`"
628                    )));
629                }
630                let was_pending = entry
631                    .pending_recipient_session_ids
632                    .iter()
633                    .any(|pending| pending == &target_session_id);
634                entry
635                    .pending_recipient_session_ids
636                    .retain(|pending| pending != &target_session_id);
637                for swarm in catalog.swarms.values_mut() {
638                    trim_acknowledged(&mut swarm.board);
639                }
640                Ok(was_pending)
641            })
642            .await?;
643        self.notify_acknowledged(&acknowledged_message, &acknowledged_target);
644        Ok(was_pending)
645    }
646
647    async fn mutate<T>(&self, mutation: impl FnOnce(&mut Catalog) -> Result<T>) -> Result<T> {
648        let mut state = self.lock_loaded().await?;
649        let mut candidate = state.as_ref().expect("swarm catalog loaded").clone();
650        let output = mutation(&mut candidate)?;
651        validate_catalog(&candidate)?;
652        self.checkpoints
653            .save_state(STATE_SCOPE, STATE_KEY, &serde_json::to_value(&candidate)?)
654            .await?;
655        *state = Some(candidate);
656        Ok(output)
657    }
658
659    async fn lock_loaded(&self) -> Result<MutexGuard<'_, Option<Catalog>>> {
660        let mut state = self.state.lock().await;
661        if state.is_none() {
662            let catalog = self
663                .checkpoints
664                .load_state(STATE_SCOPE, STATE_KEY)
665                .await?
666                .map(serde_json::from_value)
667                .transpose()?
668                .unwrap_or_default();
669            validate_catalog(&catalog)?;
670            *state = Some(catalog);
671        }
672        Ok(state)
673    }
674}
675
676pub(crate) fn validate_swarm_members(
677    leader_session_id: &str,
678    member_session_ids: &[String],
679) -> Result<()> {
680    validate_session_id(leader_session_id)?;
681    if member_session_ids.len() < 2 {
682        return Err(config("a swarm requires at least two members"));
683    }
684    if member_session_ids.len() > MAX_SWARM_MEMBERS {
685        return Err(config(format!(
686            "a swarm supports at most {MAX_SWARM_MEMBERS} members"
687        )));
688    }
689    let mut unique_sessions = BTreeSet::new();
690    for session_id in member_session_ids {
691        validate_session_id(session_id)?;
692        if !unique_sessions.insert(session_id.as_str()) {
693            return Err(config(format!(
694                "session `{session_id}` appears more than once"
695            )));
696        }
697    }
698    if !unique_sessions.contains(leader_session_id) {
699        return Err(config("swarm leader must be a member"));
700    }
701    Ok(())
702}
703
704impl SwarmBackend for SwarmStore {
705    fn active<'a>(&'a self, session_id: &'a str) -> mobius::BoxFuture<'a, mobius::Result<bool>> {
706        Box::pin(async move {
707            self.snapshot_for_session(session_id)
708                .await
709                .map(|snapshot| snapshot.is_some())
710                .map_err(mobius_error)
711        })
712    }
713
714    fn roster<'a>(&'a self, session_id: &'a str) -> mobius::BoxFuture<'a, mobius::Result<String>> {
715        Box::pin(async move {
716            let snapshot = self
717                .snapshot_for_session(session_id)
718                .await
719                .map_err(mobius_error)?
720                .ok_or_else(|| mobius::Error::Config("this chat is not in a swarm".into()))?;
721            serde_json::to_string(&snapshot).map_err(Into::into)
722        })
723    }
724
725    fn read<'a>(&'a self, session_id: &'a str) -> mobius::BoxFuture<'a, mobius::Result<String>> {
726        Box::pin(async move {
727            let snapshot = self
728                .snapshot_for_session(session_id)
729                .await
730                .map_err(mobius_error)?
731                .ok_or_else(|| mobius::Error::Config("this chat is not in a swarm".into()))?;
732            let page = self
733                .board_page(&snapshot.swarm.id, None, MAX_PAGE_ENTRIES)
734                .await
735                .map_err(mobius_error)?;
736            let mut entries = Vec::new();
737            let mut has_older = page.next_before_sequence.is_some();
738            for entry in page.entries {
739                let end = entry
740                    .text
741                    .floor_char_boundary(MAX_TOOL_READ_BODY_BYTES.min(entry.text.len()));
742                let body = &entry.text[..end];
743                entries.push(serde_json::json!({
744                    "id": entry.id,
745                    "sequence": entry.sequence,
746                    "created_at_ms": entry.created_at_ms,
747                    "author_session_id": entry.author.session_id,
748                    "author_handle": entry.author.handle,
749                    "body": body,
750                    "body_truncated": body.len() != entry.text.len(),
751                }));
752                if serde_json::to_vec(&serde_json::json!({
753                    "entries": &entries,
754                    "has_older": has_older,
755                }))?
756                .len()
757                    > MAX_TOOL_READ_BYTES
758                {
759                    entries.pop();
760                    has_older = true;
761                    break;
762                }
763            }
764            serde_json::to_string(&serde_json::json!({
765                "entries": entries,
766                "has_older": has_older,
767            }))
768            .map_err(Into::into)
769        })
770    }
771
772    fn post<'a>(
773        &'a self,
774        session_id: &'a str,
775        message: String,
776    ) -> mobius::BoxFuture<'a, mobius::Result<String>> {
777        Box::pin(async move {
778            let post = SwarmStore::post(self, session_id, message)
779                .await
780                .map_err(mobius_error)?;
781            serde_json::to_string(&post).map_err(Into::into)
782        })
783    }
784}
785
786impl StoredSwarm {
787    fn summary(&self, id: &str) -> SwarmSummary {
788        SwarmSummary {
789            id: id.to_owned(),
790            title: self.title.clone(),
791            leader_session_id: self.leader_session_id.clone(),
792            members: self.members.values().cloned().collect(),
793            latest_sequence: self.latest_sequence,
794            created_at_ms: self.created_at_ms,
795            updated_at_ms: self.updated_at_ms,
796        }
797    }
798}
799
800fn generated_member(session_id: String, members: &BTreeMap<String, SwarmMember>) -> SwarmMember {
801    let handle = loop {
802        let suffix = Uuid::new_v4()
803            .simple()
804            .to_string()
805            .chars()
806            .take(8)
807            .collect::<String>();
808        let handle = format!("agent_{suffix}");
809        if !members.contains_key(&handle) {
810            break handle;
811        }
812    };
813    SwarmMember {
814        session_id,
815        handle,
816        joined_at_ms: unix_ms(),
817    }
818}
819
820fn ensure_session_available(catalog: &Catalog, session_id: &str) -> Result<()> {
821    if catalog.swarms.values().any(|swarm| {
822        swarm
823            .members
824            .values()
825            .any(|member| member.session_id == session_id)
826    }) {
827        return Err(config(format!(
828            "session `{session_id}` already belongs to a swarm"
829        )));
830    }
831    Ok(())
832}
833
834fn swarm_id_for_session(catalog: &Catalog, session_id: &str) -> Option<String> {
835    catalog.swarms.iter().find_map(|(id, swarm)| {
836        swarm
837            .members
838            .values()
839            .any(|member| member.session_id == session_id)
840            .then(|| id.clone())
841    })
842}
843
844fn trim_acknowledged(board: &mut VecDeque<BoardEntry>) {
845    while board
846        .iter()
847        .filter(|entry| entry.pending_recipient_session_ids.is_empty())
848        .count()
849        > MAX_ACKNOWLEDGED_ENTRIES
850    {
851        let index = board
852            .iter()
853            .position(|entry| entry.pending_recipient_session_ids.is_empty())
854            .expect("acknowledged entry count is positive");
855        board.remove(index);
856    }
857}
858
859fn mentioned_handles(text: &str) -> BTreeSet<String> {
860    let bytes = text.as_bytes();
861    let mut handles = BTreeSet::new();
862    let mut index = 0;
863    while index < bytes.len() {
864        if bytes[index] != b'@'
865            || index
866                .checked_sub(1)
867                .is_some_and(|previous| is_mention_byte(bytes[previous]))
868        {
869            index += 1;
870            continue;
871        }
872        let start = index + 1;
873        let mut end = start;
874        while end < bytes.len() && is_mention_byte(bytes[end]) {
875            end += 1;
876        }
877        if start < end {
878            handles.insert(text[start..end].to_owned());
879        }
880        index = end.max(index + 1);
881    }
882    handles
883}
884
885const fn is_mention_byte(byte: u8) -> bool {
886    byte.is_ascii_alphanumeric() || byte == b'_'
887}
888
889fn validate_catalog(catalog: &Catalog) -> Result<()> {
890    let mut sessions = BTreeSet::new();
891    let mut message_ids = BTreeSet::new();
892    for (id, swarm) in &catalog.swarms {
893        let mut pending_counts = BTreeMap::<&str, usize>::new();
894        validate_swarm_id(id)?;
895        validate_title(&swarm.title)?;
896        validate_session_id(&swarm.leader_session_id)?;
897        if swarm.members.len() > MAX_SWARM_MEMBERS {
898            return Err(config(format!(
899                "swarm `{id}` has more than {MAX_SWARM_MEMBERS} members"
900            )));
901        }
902        if swarm.created_at_ms < 0 || swarm.updated_at_ms < swarm.created_at_ms {
903            return Err(config(format!("swarm `{id}` has invalid timestamps")));
904        }
905        if !swarm
906            .members
907            .values()
908            .any(|member| member.session_id == swarm.leader_session_id)
909        {
910            return Err(config(format!("swarm `{id}` leader is not a member")));
911        }
912        for (handle, member) in &swarm.members {
913            validate_member(member)?;
914            if handle != &member.handle {
915                return Err(config(format!(
916                    "swarm `{id}` member handle does not match its key"
917                )));
918            }
919            if !sessions.insert(member.session_id.clone()) {
920                return Err(config(format!(
921                    "session `{}` belongs to more than one swarm",
922                    member.session_id
923                )));
924            }
925        }
926        let mut previous_sequence = 0;
927        for entry in &swarm.board {
928            validate_message_id(&entry.id)?;
929            validate_member(&entry.author)?;
930            validate_message(&entry.text)?;
931            if entry.created_at_ms < swarm.created_at_ms
932                || entry.created_at_ms > swarm.updated_at_ms
933            {
934                return Err(config(format!(
935                    "swarm message `{}` has an invalid timestamp",
936                    entry.id
937                )));
938            }
939            if !message_ids.insert(entry.id.clone()) {
940                return Err(config(format!(
941                    "swarm message `{}` appears more than once",
942                    entry.id
943                )));
944            }
945            if entry.sequence <= previous_sequence || entry.sequence > swarm.latest_sequence {
946                return Err(config(format!("swarm `{id}` has invalid board sequences")));
947            }
948            previous_sequence = entry.sequence;
949            validate_recipient_ids(&entry.mentioned_recipient_session_ids)?;
950            validate_recipient_ids(&entry.pending_recipient_session_ids)?;
951            for recipient in &entry.pending_recipient_session_ids {
952                let count = pending_counts.entry(recipient).or_default();
953                *count += 1;
954                if *count > MAX_PENDING_DELIVERIES_PER_RECIPIENT {
955                    return Err(config(format!(
956                        "session `{recipient}` has too many pending swarm messages"
957                    )));
958                }
959            }
960            if entry.pending_recipient_session_ids.iter().any(|pending| {
961                !entry
962                    .mentioned_recipient_session_ids
963                    .iter()
964                    .any(|mentioned| mentioned == pending)
965            }) {
966                return Err(config(format!(
967                    "swarm message `{}` has a pending non-recipient",
968                    entry.id
969                )));
970            }
971        }
972        if swarm
973            .board
974            .back()
975            .is_some_and(|entry| entry.sequence != swarm.latest_sequence)
976            || (swarm.board.is_empty() && swarm.latest_sequence != 0)
977            || swarm
978                .board
979                .iter()
980                .filter(|entry| entry.pending_recipient_session_ids.is_empty())
981                .count()
982                > MAX_ACKNOWLEDGED_ENTRIES
983        {
984            return Err(config(format!("swarm `{id}` board state is invalid")));
985        }
986    }
987    if serde_json::to_vec(catalog)?.len() > MAX_CATALOG_BYTES {
988        return Err(config(format!(
989            "swarm catalog exceeds {MAX_CATALOG_BYTES} encoded bytes"
990        )));
991    }
992    Ok(())
993}
994
995fn validate_recipient_ids(recipients: &[String]) -> Result<()> {
996    let mut unique = BTreeSet::new();
997    for recipient in recipients {
998        validate_session_id(recipient)?;
999        if !unique.insert(recipient) {
1000            return Err(config("swarm message recipient appears more than once"));
1001        }
1002    }
1003    Ok(())
1004}
1005
1006fn validate_member(member: &SwarmMember) -> Result<()> {
1007    validate_session_id(&member.session_id)?;
1008    validate_handle(&member.handle)?;
1009    if member.joined_at_ms < 0 {
1010        return Err(config("swarm member join time cannot be negative"));
1011    }
1012    Ok(())
1013}
1014
1015fn validate_handle(handle: &str) -> Result<()> {
1016    if handle.is_empty()
1017        || handle.len() > MAX_HANDLE_BYTES
1018        || !handle
1019            .bytes()
1020            .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
1021    {
1022        return Err(config(
1023            "swarm handle must contain 1-64 lowercase letters, digits, or underscores",
1024        ));
1025    }
1026    Ok(())
1027}
1028
1029fn validate_title(title: &str) -> Result<()> {
1030    if title.trim().is_empty() || title.len() > MAX_TITLE_BYTES {
1031        return Err(config(format!(
1032            "swarm title must contain 1-{MAX_TITLE_BYTES} UTF-8 bytes"
1033        )));
1034    }
1035    Ok(())
1036}
1037
1038fn validate_message(text: &str) -> Result<()> {
1039    if text.trim().is_empty() || text.len() > MAX_MESSAGE_BYTES {
1040        return Err(config(format!(
1041            "swarm message must contain 1-{MAX_MESSAGE_BYTES} UTF-8 bytes"
1042        )));
1043    }
1044    Ok(())
1045}
1046
1047fn validate_session_id(session_id: &str) -> Result<()> {
1048    if session_id.is_empty() || session_id.len() > MAX_ID_BYTES {
1049        return Err(config(format!(
1050            "session id must contain 1-{MAX_ID_BYTES} bytes"
1051        )));
1052    }
1053    Ok(())
1054}
1055
1056fn validate_swarm_id(id: &str) -> Result<()> {
1057    Uuid::parse_str(id)
1058        .map(|_| ())
1059        .map_err(|_| config("invalid swarm id"))
1060}
1061
1062fn validate_message_id(id: &str) -> Result<()> {
1063    Uuid::parse_str(id)
1064        .map(|_| ())
1065        .map_err(|_| config("invalid swarm message id"))
1066}
1067
1068fn config(message: impl Into<String>) -> Error {
1069    Error::Config(message.into())
1070}
1071
1072fn mobius_error(error: Error) -> mobius::Error {
1073    mobius::Error::Config(error.to_string())
1074}
1075
1076fn unix_ms() -> i64 {
1077    SystemTime::now()
1078        .duration_since(UNIX_EPOCH)
1079        .map(|duration| i64::try_from(duration.as_millis()).unwrap_or(i64::MAX))
1080        .unwrap_or_default()
1081}
1082
1083#[cfg(test)]
1084mod tests {
1085    use mobius::backend::checkpoint::sqlite::SqliteCheckpoint;
1086
1087    use super::*;
1088
1089    fn store() -> (
1090        tempfile::TempDir,
1091        Arc<dyn CheckpointStore>,
1092        SwarmStore,
1093        mpsc::UnboundedReceiver<SwarmDelivery>,
1094    ) {
1095        let directory = tempfile::tempdir().expect("workspace");
1096        let checkpoints: Arc<dyn CheckpointStore> = Arc::new(
1097            SqliteCheckpoint::new(directory.path().join("checkpoints.sqlite3"))
1098                .expect("checkpoints"),
1099        );
1100        let (store, deliveries) = SwarmStore::new(Arc::clone(&checkpoints));
1101        (directory, checkpoints, store, deliveries)
1102    }
1103
1104    async fn create_swarm(store: &SwarmStore) -> SwarmSummary {
1105        store
1106            .create("leader".into(), vec!["leader".into(), "reviewer".into()])
1107            .await
1108            .expect("create swarm")
1109    }
1110
1111    #[tokio::test]
1112    async fn membership_is_unique_and_lazily_reloads() {
1113        let (_directory, checkpoints, store, _deliveries) = store();
1114        let created = create_swarm(&store).await;
1115        let created_title = created.title.clone();
1116        assert!(created.created_at_ms > 0);
1117        assert_eq!(created.created_at_ms, created.updated_at_ms);
1118
1119        let duplicate = store
1120            .create("leader".into(), vec!["leader".into(), "third".into()])
1121            .await
1122            .expect_err("session cannot join two swarms");
1123        assert!(duplicate.to_string().contains("already belongs"));
1124
1125        let (reloaded, _deliveries) = SwarmStore::new(checkpoints);
1126        assert_eq!(reloaded.summaries().await.expect("reload"), vec![created]);
1127        let snapshot = reloaded
1128            .snapshot_for_session("reviewer")
1129            .await
1130            .expect("snapshot")
1131            .expect("membership");
1132        assert!(snapshot.handle.starts_with("agent_"));
1133        assert_eq!(snapshot.swarm.title, created_title);
1134    }
1135
1136    #[tokio::test]
1137    async fn joining_cannot_exceed_the_member_limit() {
1138        let (_directory, _checkpoints, store, _deliveries) = store();
1139        let members = (0..MAX_SWARM_MEMBERS)
1140            .map(|index| format!("member-{index}"))
1141            .collect::<Vec<_>>();
1142        let swarm = store
1143            .create(members[0].clone(), members)
1144            .await
1145            .expect("full swarm");
1146
1147        let error = store
1148            .join(&swarm.id, "overflow".into())
1149            .await
1150            .expect_err("member limit");
1151
1152        assert!(error.to_string().contains("at most 100"));
1153    }
1154
1155    #[tokio::test]
1156    async fn mentions_are_resolved_delivered_and_acknowledged() {
1157        let (_directory, checkpoints, store, mut deliveries) = store();
1158        let swarm = create_swarm(&store).await;
1159        let leader_handle = swarm
1160            .members
1161            .iter()
1162            .find(|member| member.session_id == "leader")
1163            .expect("leader")
1164            .handle
1165            .clone();
1166        let reviewer_handle = swarm
1167            .members
1168            .iter()
1169            .find(|member| member.session_id == "reviewer")
1170            .expect("reviewer")
1171            .handle
1172            .clone();
1173
1174        let unknown = store
1175            .post("leader", "Can @missing check this?".into())
1176            .await
1177            .expect_err("unknown mention");
1178        assert!(unknown.to_string().contains("@missing"));
1179
1180        let body = format!("Can @{reviewer_handle} check this?");
1181        let post = store.post("leader", body.clone()).await.expect("post");
1182        assert_eq!(deliveries.recv().await, Some(SwarmDelivery::Changed));
1183        assert_eq!(
1184            deliveries.recv().await,
1185            Some(SwarmDelivery::Pending {
1186                target_session_id: "reviewer".into()
1187            })
1188        );
1189        assert_eq!(post.entry.sequence, 1);
1190        assert!(post.entry.created_at_ms >= swarm.created_at_ms);
1191        assert_eq!(post.entry.author.handle, leader_handle);
1192        assert_eq!(post.resolved_recipient_session_ids, ["reviewer"]);
1193        let records = store.records().await.expect("wire records");
1194        assert_eq!(records[0].messages[0].body, body);
1195        assert!(
1196            SwarmBackend::active(&store, "reviewer")
1197                .await
1198                .expect("active")
1199        );
1200        assert_eq!(
1201            store
1202                .pending_deliveries("reviewer")
1203                .await
1204                .expect("pending")
1205                .len(),
1206            1
1207        );
1208
1209        assert!(
1210            store
1211                .acknowledge(&post.entry.id, "reviewer")
1212                .await
1213                .expect("acknowledge")
1214        );
1215        assert_eq!(
1216            deliveries.recv().await,
1217            Some(SwarmDelivery::Acknowledged {
1218                target_session_id: "reviewer".into(),
1219                message_id: post.entry.id.clone(),
1220            })
1221        );
1222        assert!(
1223            !store
1224                .acknowledge(&post.entry.id, "reviewer")
1225                .await
1226                .expect("idempotent acknowledge")
1227        );
1228        assert_eq!(
1229            deliveries.recv().await,
1230            Some(SwarmDelivery::Acknowledged {
1231                target_session_id: "reviewer".into(),
1232                message_id: post.entry.id.clone(),
1233            })
1234        );
1235        assert!(
1236            store
1237                .pending_deliveries("reviewer")
1238                .await
1239                .expect("pending")
1240                .is_empty()
1241        );
1242
1243        let (reloaded, _deliveries) = SwarmStore::new(checkpoints);
1244        let page = reloaded
1245            .board_page(&swarm.id, None, 10)
1246            .await
1247            .expect("board page");
1248        assert_eq!(page.entries[0].id, post.entry.id);
1249        assert!(page.entries[0].pending_recipient_session_ids.is_empty());
1250    }
1251
1252    #[tokio::test]
1253    async fn leave_and_disband_settle_removed_pending_deliveries() {
1254        let (_directory, _checkpoints, store, mut deliveries) = store();
1255        let swarm = store
1256            .create(
1257                "leader".into(),
1258                vec!["leader".into(), "observer".into(), "reviewer".into()],
1259            )
1260            .await
1261            .expect("create swarm");
1262        let handle = |session_id: &str| {
1263            swarm
1264                .members
1265                .iter()
1266                .find(|member| member.session_id == session_id)
1267                .expect("member")
1268                .handle
1269                .clone()
1270        };
1271        let post = store
1272            .post(
1273                "leader",
1274                format!(
1275                    "@{} @{} please review",
1276                    handle("observer"),
1277                    handle("reviewer")
1278                ),
1279            )
1280            .await
1281            .expect("post");
1282        for _ in 0..3 {
1283            deliveries.recv().await.expect("initial delivery signal");
1284        }
1285
1286        store
1287            .leave(&swarm.id, "reviewer")
1288            .await
1289            .expect("leave swarm");
1290        store.disband(&swarm.id).await.expect("disband swarm");
1291
1292        assert_eq!(
1293            [deliveries.recv().await, deliveries.recv().await],
1294            [
1295                Some(SwarmDelivery::Acknowledged {
1296                    target_session_id: "reviewer".into(),
1297                    message_id: post.entry.id.clone(),
1298                }),
1299                Some(SwarmDelivery::Acknowledged {
1300                    target_session_id: "observer".into(),
1301                    message_id: post.entry.id,
1302                }),
1303            ]
1304        );
1305    }
1306
1307    #[tokio::test]
1308    async fn board_only_posts_signal_catalog_changes_without_peer_delivery() {
1309        let (_directory, _checkpoints, store, mut deliveries) = store();
1310        create_swarm(&store).await;
1311
1312        store
1313            .post("leader", "Shared status update".into())
1314            .await
1315            .expect("board-only post");
1316
1317        assert_eq!(deliveries.recv().await, Some(SwarmDelivery::Changed));
1318        assert!(matches!(
1319            deliveries.try_recv(),
1320            Err(mpsc::error::TryRecvError::Empty)
1321        ));
1322    }
1323
1324    #[tokio::test]
1325    async fn model_board_read_stays_valid_json_below_the_tool_output_cap() {
1326        let (_directory, _checkpoints, store, _deliveries) = store();
1327        create_swarm(&store).await;
1328        store
1329            .post("leader", "\0".repeat(MAX_MESSAGE_BYTES))
1330            .await
1331            .expect("escape-heavy post");
1332
1333        let output = SwarmBackend::read(&store, "leader")
1334            .await
1335            .expect("model board read");
1336        let output: serde_json::Value = serde_json::from_str(&output).expect("valid JSON output");
1337
1338        assert!(serde_json::to_vec(&output).expect("encoded output").len() <= MAX_TOOL_READ_BYTES);
1339        assert_eq!(output["entries"][0]["body_truncated"], true);
1340        assert_eq!(output["has_older"], false);
1341    }
1342
1343    #[tokio::test]
1344    async fn retention_never_evicts_pending_entries() {
1345        let (_directory, _checkpoints, store, _deliveries) = store();
1346        let swarm = create_swarm(&store).await;
1347        let reviewer_handle = swarm
1348            .members
1349            .iter()
1350            .find(|member| member.session_id == "reviewer")
1351            .expect("reviewer")
1352            .handle
1353            .clone();
1354        let pending = store
1355            .post("leader", format!("Please check @{reviewer_handle}"))
1356            .await
1357            .expect("pending post");
1358        for sequence in 0..=MAX_ACKNOWLEDGED_ENTRIES {
1359            store
1360                .post("leader", format!("board update {sequence}"))
1361                .await
1362                .expect("board post");
1363        }
1364
1365        let page = store
1366            .board_page(&swarm.id, None, MAX_PAGE_ENTRIES)
1367            .await
1368            .expect("board page");
1369        assert_eq!(page.entries.len(), MAX_ACKNOWLEDGED_ENTRIES);
1370        assert!(page.next_before_sequence.is_some());
1371
1372        let pending_delivery = store
1373            .pending_deliveries("reviewer")
1374            .await
1375            .expect("pending delivery");
1376        assert_eq!(pending_delivery.len(), 1);
1377        assert_eq!(pending_delivery[0].entry.id, pending.entry.id);
1378    }
1379
1380    #[tokio::test]
1381    async fn posting_backpressures_a_recipient_with_a_full_pending_queue() {
1382        let (_directory, _checkpoints, store, _deliveries) = store();
1383        let swarm = create_swarm(&store).await;
1384        let reviewer_handle = swarm
1385            .members
1386            .iter()
1387            .find(|member| member.session_id == "reviewer")
1388            .expect("reviewer")
1389            .handle
1390            .clone();
1391        for sequence in 0..MAX_PENDING_DELIVERIES_PER_RECIPIENT {
1392            store
1393                .post("leader", format!("@{reviewer_handle} review {sequence}"))
1394                .await
1395                .expect("pending post within limit");
1396        }
1397
1398        let error = store
1399            .post("leader", format!("@{reviewer_handle} one too many"))
1400            .await
1401            .expect_err("pending delivery limit");
1402
1403        assert!(error.to_string().contains("pending swarm messages"));
1404        assert_eq!(
1405            store
1406                .pending_deliveries("reviewer")
1407                .await
1408                .expect("pending deliveries")
1409                .len(),
1410            MAX_PENDING_DELIVERIES_PER_RECIPIENT
1411        );
1412    }
1413
1414    #[test]
1415    fn catalog_encoded_size_is_bounded_before_persistence_or_broadcast() {
1416        let now = unix_ms();
1417        let author = SwarmMember {
1418            session_id: "leader".into(),
1419            handle: "agent_leader".into(),
1420            joined_at_ms: now,
1421        };
1422        let board = (1..=64)
1423            .map(|sequence| BoardEntry {
1424                id: Uuid::new_v4().to_string(),
1425                sequence,
1426                created_at_ms: now,
1427                author: author.clone(),
1428                text: "\0".repeat(MAX_MESSAGE_BYTES),
1429                mentioned_recipient_session_ids: Vec::new(),
1430                pending_recipient_session_ids: Vec::new(),
1431            })
1432            .collect();
1433        let catalog = Catalog {
1434            swarms: BTreeMap::from([(
1435                Uuid::new_v4().to_string(),
1436                StoredSwarm {
1437                    title: "Bounded swarm".into(),
1438                    leader_session_id: author.session_id.clone(),
1439                    members: BTreeMap::from([(author.handle.clone(), author)]),
1440                    latest_sequence: 64,
1441                    board,
1442                    created_at_ms: now,
1443                    updated_at_ms: now,
1444                },
1445            )]),
1446        };
1447
1448        assert!(
1449            validate_catalog(&catalog)
1450                .expect_err("oversized catalog")
1451                .to_string()
1452                .contains("encoded bytes")
1453        );
1454    }
1455
1456    #[test]
1457    fn mention_parser_ignores_email_boundaries_and_deduplicates() {
1458        assert_eq!(
1459            mentioned_handles("@one mail@two @one; @three"),
1460            BTreeSet::from(["one".into(), "three".into()])
1461        );
1462    }
1463}