1use 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;
26const 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
33#[serde(deny_unknown_fields)]
34pub struct SwarmMember {
35 pub session_id: String,
37 pub handle: String,
39 pub joined_at_ms: i64,
41}
42
43#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
45#[serde(deny_unknown_fields)]
46pub struct SwarmSummary {
47 pub id: String,
49 pub title: String,
51 pub leader_session_id: String,
53 pub members: Vec<SwarmMember>,
55 pub latest_sequence: u64,
57 pub created_at_ms: i64,
59 pub updated_at_ms: i64,
61}
62
63#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65#[serde(deny_unknown_fields)]
66pub struct SwarmSnapshot {
67 pub swarm: SwarmSummary,
69 pub handle: String,
71}
72
73#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
75#[serde(deny_unknown_fields)]
76pub struct BoardEntry {
77 pub id: String,
79 pub sequence: u64,
81 pub created_at_ms: i64,
83 pub author: SwarmMember,
85 pub text: String,
87 pub mentioned_recipient_session_ids: Vec<String>,
89 pub pending_recipient_session_ids: Vec<String>,
91}
92
93#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
95#[serde(deny_unknown_fields)]
96pub struct SwarmPost {
97 pub entry: BoardEntry,
99 pub resolved_recipient_session_ids: Vec<String>,
101}
102
103#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
105#[serde(deny_unknown_fields)]
106pub struct BoardPage {
107 pub entries: Vec<BoardEntry>,
109 pub next_before_sequence: Option<u64>,
111}
112
113#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
115#[serde(deny_unknown_fields)]
116pub struct PendingDelivery {
117 pub swarm_id: String,
119 pub swarm_title: String,
121 pub entry: BoardEntry,
123}
124
125#[derive(Debug, Clone, PartialEq, Eq)]
127pub(crate) enum SwarmDelivery {
128 Changed,
130 RetryPending,
132 Pending { target_session_id: String },
134 Acknowledged {
136 target_session_id: String,
137 message_id: String,
138 },
139}
140
141#[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 #[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 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 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 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 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 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 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 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 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 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 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 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}