Skip to main content

mobius/middleware/
messages.rs

1//! Durable conversation-message delivery.
2
3pub mod voice;
4
5use super::{
6    ActiveCommandContext, MessageRouteContext, Middleware, MiddlewareCommandContext,
7    MiddlewareCommandOutput, SessionStartContext, SessionStartSource, SubmissionResult,
8};
9use crate::backend::checkpoint::QueuedMessageBoundary;
10use crate::protocol::{
11    ActiveMessageDelivery, EventMsg, FrontendBlock, FrontendCommand, FrontendContribution,
12    FrontendEvent, FrontendSlot, FrontendSymbol, FrontendTone, FrontendWidget,
13    MAX_CAPABILITY_INPUT_BYTES, MessageAuthor, MessageDelivery, MessageEvent, Op,
14};
15use crate::{BoxFuture, Error, Result};
16
17mod text {
18    #[derive(serde::Deserialize)]
19    #[serde(deny_unknown_fields)]
20    pub(super) struct Definition {
21        #[serde(deserialize_with = "crate::middleware::manifest::deserialize_settings")]
22        pub(super) settings: Vec<crate::middleware::manifest::MiddlewareSettingManifest>,
23        pub(super) voice_user_label: String,
24        pub(super) voice_speaker_label: String,
25        pub(super) voice_workspace_label: String,
26        pub(super) voice_retained_context: String,
27        pub(super) voice_clarify_request: String,
28        pub(super) voice_delegation_policy: String,
29        pub(super) voice_current_request_heading: String,
30        pub(super) voice_recent_discussion_heading: String,
31        pub(super) voice_latest_request: String,
32        pub(super) voice_tool_started: String,
33        pub(super) voice_tool_result: String,
34        pub(super) voice_tool_failed: String,
35        pub(super) voice_tool_finished: String,
36        pub(super) voice_work_stopped: String,
37        pub(super) voice_request_failed: String,
38        pub(super) voice_request_empty: String,
39        pub(super) voice_instructions: String,
40        pub(super) voice_workspace_heading: String,
41        pub(super) voice_previous_heading: String,
42        pub(super) default_enabled: bool,
43        pub(super) manifest_description: String,
44        pub(super) manifest_label: String,
45    }
46    crate::embedded_config! { pub(super) static DEFINITION: Definition = include_str!("messages.toml"); }
47}
48const MAX_PENDING_MESSAGES: usize = 1_024;
49
50/// Default number of pending messages retained by the delivery queue.
51pub fn default_max_pending() -> usize {
52    super::manifest::integer_default(&text::DEFINITION.settings, "max_pending") as usize
53}
54/// Default delivery for user messages submitted during an active turn.
55pub fn default_delivery() -> ActiveMessageDelivery {
56    super::manifest::string_default(&text::DEFINITION.settings, "delivery")
57        .parse()
58        .expect("valid embedded delivery")
59}
60
61super::manifest::middleware_manifest! {
62/// Configuration and presentation metadata for message delivery.
63    "messages", text::DEFINITION, required: true, capability: None, settings: &text::DEFINITION.settings
64}
65
66const EDIT_COMMAND: &str = "edit";
67const STALE_EDIT: &str = "message is no longer queued";
68const INVALID_EDIT: &str = "message edit requires non-empty text";
69
70/// Prepares every conversation message and owns its durable delivery lifecycle.
71pub struct Messages {
72    max_pending: usize,
73    delivery: ActiveMessageDelivery,
74}
75
76impl Default for Messages {
77    fn default() -> Self {
78        Self {
79            max_pending: default_max_pending(),
80            delivery: default_delivery(),
81        }
82    }
83}
84
85impl Messages {
86    /// Creates message delivery with a bounded queue and active-turn default.
87    /// # Errors
88    ///
89    /// Returns an error if configuration is invalid or a required resource cannot be initialized.
90    pub fn new(max_pending: usize, delivery: ActiveMessageDelivery) -> Result<Self> {
91        if max_pending == 0 || max_pending > MAX_PENDING_MESSAGES {
92            return Err(Error::Config(format!(
93                "message queue limit must be between 1 and {MAX_PENDING_MESSAGES}"
94            )));
95        }
96        Ok(Self {
97            max_pending,
98            delivery,
99        })
100    }
101
102    fn remove_widget(&self, id: &str) -> FrontendEvent {
103        FrontendEvent::RemoveWidget {
104            capability: self.name().into(),
105            id: id.into(),
106        }
107    }
108
109    fn queued_widget(&self, id: &str, message: &MessageEvent) -> FrontendEvent {
110        let action = matches!(message.author, MessageAuthor::User).then(|| Op::CapabilityCommand {
111            capability: self.name().into(),
112            command: EDIT_COMMAND.into(),
113            arguments: id.into(),
114            input: Some(message.text.clone()),
115            target: None,
116        });
117        FrontendEvent::Widget {
118            capability: self.name().into(),
119            item: FrontendWidget {
120                id: id.into(),
121                slot: FrontendSlot::TranscriptTail,
122                text: message.text.clone(),
123                tone: FrontendTone::Neutral,
124                symbol: Some(match message.delivery {
125                    MessageDelivery::Turn => FrontendSymbol::Chat,
126                    MessageDelivery::Steer => FrontendSymbol::Custom("steer".into()),
127                    MessageDelivery::Queue => FrontendSymbol::Custom("queue".into()),
128                }),
129                icon_only: false,
130                progress: None,
131                content: None,
132                action,
133            },
134        }
135    }
136
137    fn prepare(
138        &self,
139        context: &MessageRouteContext<'_>,
140    ) -> std::result::Result<QueuedMessageBoundary, String> {
141        let Some(turn_id) = context.active_turn_id else {
142            return if context.message.target_turn_id.is_some() {
143                Err("message targeted a stale turn".into())
144            } else {
145                Ok(QueuedMessageBoundary::Turn)
146            };
147        };
148        if let Some(target) = &context.message.target_turn_id
149            && target != turn_id
150        {
151            return Err("message targeted a stale turn".into());
152        }
153        match context.message.requested_delivery.unwrap_or(self.delivery) {
154            ActiveMessageDelivery::Steer => Ok(QueuedMessageBoundary::Steer {
155                turn_id: turn_id.into(),
156            }),
157            ActiveMessageDelivery::Queue => Ok(QueuedMessageBoundary::Queue),
158        }
159    }
160}
161
162impl Middleware for Messages {
163    fn name(&self) -> &'static str {
164        MANIFEST.id
165    }
166
167    fn frontend(&self) -> FrontendContribution {
168        FrontendContribution {
169            capability: self.name().into(),
170            commands: vec![FrontendCommand {
171                name: voice::transcript::COMMAND.into(),
172                arguments: String::new(),
173                description: "Open the voice transcript".into(),
174                requires_idle: false,
175            }],
176            ..FrontendContribution::default()
177        }
178    }
179
180    fn render(&self, event: &EventMsg, _session_id: &str) -> Option<FrontendBlock> {
181        let EventMsg::Message(message) = event else {
182            return None;
183        };
184        let MessageAuthor::Source {
185            message_id,
186            source,
187            handle,
188            symbol,
189            ..
190        } = &message.author
191        else {
192            return None;
193        };
194        if matches!(source, crate::protocol::MessageSource::Session { .. }) {
195            return None;
196        }
197        Some(FrontendBlock {
198            id: Some(format!(
199                "message_received:{}:{session_id}:{message_id}",
200                source.id().len(),
201                session_id = source.id()
202            )),
203            title: format!("Message received from @{handle}"),
204            text: message.text.clone(),
205            symbol: Some(symbol.clone().unwrap_or(FrontendSymbol::Chat)),
206            ..Default::default()
207        })
208    }
209
210    fn command<'a>(
211        &'a self,
212        context: MiddlewareCommandContext<'a>,
213    ) -> BoxFuture<'a, Result<MiddlewareCommandOutput>> {
214        Box::pin(async move {
215            if context.command != voice::transcript::COMMAND {
216                return Err(Error::Unknown(format!(
217                    "messages command `{}`",
218                    context.command
219                )));
220            }
221            Ok(MiddlewareCommandOutput::events(vec![
222                voice::transcript::read_preview(
223                    context.checkpoints.as_ref(),
224                    context.session_id,
225                    context.arguments,
226                )
227                .await?,
228            ]))
229        })
230    }
231
232    fn handles_messages(&self) -> bool {
233        true
234    }
235
236    fn route_message(&self, context: &mut MessageRouteContext<'_>) -> Result<SubmissionResult> {
237        let boundary = match self.prepare(context) {
238            Ok(boundary) => boundary,
239            Err(message) => return Ok(SubmissionResult::Rejected(message)),
240        };
241        if context.queued_messages.count() >= self.max_pending {
242            return Ok(SubmissionResult::Rejected("message queue is full".into()));
243        }
244        let event = MessageEvent {
245            author: context.message.author.clone(),
246            delivery: boundary.delivery(),
247            text: context.message.text.clone(),
248            attachments: context.message.attachments.clone(),
249            reply: context.message.reply.clone(),
250            message_target: None,
251        };
252        if !context.queued_messages.enqueue(
253            context.submission_id,
254            boundary.clone(),
255            event.clone(),
256        )? {
257            return Ok(SubmissionResult::Rejected(
258                "message could not be queued".into(),
259            ));
260        }
261        if !matches!(boundary, QueuedMessageBoundary::Turn) {
262            context.events.push(EventMsg::Frontend(
263                self.queued_widget(context.submission_id, &event),
264            ));
265        }
266        Ok(SubmissionResult::Accepted {
267            input_changed: matches!(boundary, QueuedMessageBoundary::Steer { .. }),
268        })
269    }
270
271    fn message_boundary_events(&self, submission_id: &str) -> Vec<EventMsg> {
272        vec![EventMsg::Frontend(self.remove_widget(submission_id))]
273    }
274
275    fn active_command<'a>(
276        &'a self,
277        context: &'a mut ActiveCommandContext<'_>,
278    ) -> BoxFuture<'a, Result<Option<SubmissionResult>>> {
279        Box::pin(async move {
280            if context.command == voice::transcript::COMMAND {
281                let result = voice::transcript::read_preview(
282                    context.checkpoints,
283                    context.session_id,
284                    context.arguments,
285                )
286                .await;
287                return Ok(Some(match result {
288                    Ok(event) => {
289                        context.events.push(EventMsg::Frontend(event));
290                        SubmissionResult::Handled
291                    }
292                    Err(error) => SubmissionResult::Rejected(error.to_string()),
293                }));
294            }
295            if context.command != EDIT_COMMAND {
296                return Ok(None);
297            }
298            let Some(input) = context.input.filter(|input| !input.trim().is_empty()) else {
299                return Ok(Some(SubmissionResult::Rejected(INVALID_EDIT.into())));
300            };
301            if input.len() > MAX_CAPABILITY_INPUT_BYTES {
302                return Ok(Some(SubmissionResult::Rejected(
303                    "message exceeds editable size limit".into(),
304                )));
305            }
306            let Some(queued) = context.queued_messages.find(context.arguments) else {
307                return Ok(Some(SubmissionResult::Rejected(STALE_EDIT.into())));
308            };
309            let mut event = queued.event();
310            if !matches!(event.author, MessageAuthor::User) {
311                return Ok(Some(SubmissionResult::Rejected(
312                    "peer messages cannot be edited".into(),
313                )));
314            }
315            event.text = input.into();
316            let input_changed = event.delivery == MessageDelivery::Steer;
317            if !context.queued_messages.replace(
318                context.arguments,
319                context.submission_id,
320                event.clone(),
321            )? {
322                return Ok(Some(SubmissionResult::Rejected(STALE_EDIT.into())));
323            }
324            context
325                .events
326                .push(EventMsg::Frontend(self.remove_widget(context.arguments)));
327            context.events.push(EventMsg::Frontend(
328                self.queued_widget(context.submission_id, &event),
329            ));
330            Ok(Some(SubmissionResult::Accepted { input_changed }))
331        })
332    }
333
334    fn session_start<'a>(
335        &'a self,
336        context: &'a mut SessionStartContext<'_>,
337    ) -> BoxFuture<'a, Result<()>> {
338        Box::pin(async move {
339            if context.source() == SessionStartSource::Compact {
340                return Ok(());
341            }
342            for queued in context.queued_messages().views() {
343                (context.runtime.frontend)(self.queued_widget(queued.id(), &queued.event()))?;
344            }
345            if let Some(widget) = voice::transcript::restore_widget(
346                context.runtime.checkpoints.as_ref(),
347                &context.runtime.session_id,
348            )
349            .await?
350            {
351                (context.runtime.frontend)(widget)?;
352            }
353            Ok(())
354        })
355    }
356}
357
358#[cfg(test)]
359mod tests {
360    use crate::protocol::FrontendBlockRole;
361    #[test]
362    fn declared_deliveries_parse_at_the_protocol_boundary() {
363        for choice in
364            super::super::manifest::static_choices(&super::text::DEFINITION.settings, "delivery")
365                .iter()
366        {
367            assert_eq!(
368                choice
369                    .value
370                    .parse::<super::ActiveMessageDelivery>()
371                    .expect("declared delivery")
372                    .id(),
373                choice.value
374            );
375        }
376        assert!("other".parse::<super::ActiveMessageDelivery>().is_err());
377    }
378
379    use std::collections::BTreeMap;
380    use std::sync::Arc;
381
382    use super::*;
383    use crate::backend::checkpoint::QueuedMessage;
384    use crate::middleware::{
385        ActiveCommandContext, MessageQueue, MessageRouteContext, MiddlewareStack,
386    };
387    use crate::protocol::{MessageReply, MessageSubmission, MessageTarget, SessionFileReference};
388
389    #[test]
390    fn manifest_advertises_delivery_symbols() {
391        assert!(
392            MANIFEST
393                .feature(&[])
394                .settings
395                .iter()
396                .all(|setting| !setting.composer)
397        );
398        assert_eq!(
399            super::super::manifest::static_choices(&text::DEFINITION.settings, "delivery")
400                .iter()
401                .map(|choice| (choice.value.as_str(), choice.symbol.as_deref()))
402                .collect::<Vec<_>>(),
403            [("steer", Some("steer")), ("queue", Some("queue"))]
404        );
405    }
406
407    #[test]
408    fn queued_widgets_name_their_delivery() {
409        let messages = Messages::default();
410        let symbol = |delivery| {
411            let FrontendEvent::Widget { item, .. } = messages.queued_widget(
412                "message-1",
413                &MessageEvent {
414                    author: MessageAuthor::User,
415                    delivery,
416                    text: "hello".into(),
417                    attachments: Vec::new(),
418                    reply: None,
419                    message_target: None,
420                },
421            ) else {
422                panic!("queued message widget");
423            };
424            item.symbol
425        };
426
427        assert_eq!(
428            (
429                symbol(MessageDelivery::Steer),
430                symbol(MessageDelivery::Queue),
431            ),
432            (
433                Some(FrontendSymbol::Custom("steer".into())),
434                Some(FrontendSymbol::Custom("queue".into())),
435            )
436        );
437    }
438
439    fn user(delivery: Option<ActiveMessageDelivery>) -> MessageSubmission {
440        MessageSubmission {
441            author: MessageAuthor::User,
442            text: "hello".into(),
443            attachments: Vec::<SessionFileReference>::new(),
444            reply: None,
445            requested_delivery: delivery,
446            target_turn_id: None,
447        }
448    }
449
450    fn route(
451        stack: &MiddlewareStack,
452        queued: &mut Vec<QueuedMessage>,
453        message: &MessageSubmission,
454        active_turn_id: Option<&str>,
455    ) -> SubmissionResult {
456        stack
457            .route_message(&mut MessageRouteContext {
458                submission_id: "message-1",
459                message,
460                active_turn_id,
461                queued_messages: MessageQueue::new(queued),
462                events: &mut Vec::new(),
463            })
464            .expect("route message")
465    }
466
467    #[test]
468    fn active_user_uses_the_configured_queue_boundary() {
469        let stack = MiddlewareStack::new(vec![Arc::new(
470            Messages::new(4, ActiveMessageDelivery::Queue).expect("messages"),
471        )])
472        .expect("stack");
473        let mut queued = Vec::new();
474
475        let result = route(&stack, &mut queued, &user(None), Some("turn-1"));
476
477        assert_eq!(
478            result,
479            SubmissionResult::Accepted {
480                input_changed: false
481            }
482        );
483        assert!(
484            !stack
485                .messages_ready(&queued, "turn-1")
486                .expect("message input readiness")
487        );
488        assert_eq!(
489            stack
490                .next_turn(&mut queued)
491                .expect("next turn")
492                .expect("queued message")
493                .event
494                .message()
495                .map(|message| message.delivery),
496            Some(MessageDelivery::Queue)
497        );
498    }
499
500    #[test]
501    fn immediate_turn_does_not_publish_a_queued_widget() {
502        let stack = MiddlewareStack::new(vec![Arc::new(Messages::default())]).expect("stack");
503        let mut queued = Vec::new();
504        let mut events = Vec::new();
505
506        let result = stack
507            .route_message(&mut MessageRouteContext {
508                submission_id: "message-1",
509                message: &user(None),
510                active_turn_id: None,
511                queued_messages: MessageQueue::new(&mut queued),
512                events: &mut events,
513            })
514            .expect("route message");
515
516        assert_eq!(
517            result,
518            SubmissionResult::Accepted {
519                input_changed: false
520            }
521        );
522        assert!(events.is_empty());
523        assert_eq!(queued.len(), 1);
524    }
525
526    #[tokio::test]
527    async fn queued_message_edit_preserves_its_reply_snapshot() {
528        let stack = MiddlewareStack::new(vec![Arc::new(Messages::default())]).expect("stack");
529        let mut queued = Vec::new();
530        let mut message = user(Some(ActiveMessageDelivery::Queue));
531        message.reply = Some(MessageReply {
532            target: MessageTarget {
533                checkpoint_sequence: 5,
534                batch_item_count: 2,
535            },
536            text: "Earlier".into(),
537        });
538        route(&stack, &mut queued, &message, Some("turn-1"));
539        let mut events = Vec::new();
540        let metadata = BTreeMap::new();
541        let directory = tempfile::tempdir().expect("checkpoint directory");
542        let checkpoints = crate::backend::checkpoint::sqlite::SqliteCheckpoint::new(
543            directory.path().join("checkpoints.sqlite3"),
544        )
545        .expect("checkpoint store");
546
547        stack
548            .active_command(
549                MANIFEST.id,
550                &mut ActiveCommandContext {
551                    checkpoints: &checkpoints,
552                    submission_id: "message-2",
553                    session_id: "session-1",
554                    metadata: &metadata,
555                    active_turn_id: "turn-1",
556                    command: EDIT_COMMAND,
557                    arguments: "message-1",
558                    input: Some("Updated"),
559                    target: None,
560                    queued_messages: MessageQueue::new(&mut queued),
561                    events: &mut events,
562                },
563            )
564            .await
565            .expect("edit queued message")
566            .expect("message command");
567
568        let edited = queued[0].event();
569        assert_eq!(
570            (edited.text.as_str(), edited.reply),
571            ("Updated", message.reply)
572        );
573    }
574
575    #[test]
576    fn explicitly_steered_source_remains_non_authoritative_input() {
577        let stack = MiddlewareStack::new(vec![Arc::new(
578            Messages::new(4, ActiveMessageDelivery::Queue).expect("messages"),
579        )])
580        .expect("stack");
581        let peer = MessageSubmission {
582            author: MessageAuthor::Source {
583                message_id: "board-1".into(),
584                source: crate::protocol::MessageSource::Session {
585                    session_id: "peer-1".into(),
586                },
587                cause_id: None,
588                ancestry: Vec::new(),
589                handle: "worker".into(),
590                symbol: None,
591            },
592            text: "Review this.\n\nKeep the validation.\n".into(),
593            attachments: Vec::new(),
594            reply: None,
595            requested_delivery: Some(ActiveMessageDelivery::Steer),
596            target_turn_id: None,
597        };
598        let mut queued = Vec::new();
599
600        let result = route(&stack, &mut queued, &peer, Some("turn-1"));
601        let staged = stack
602            .stage_model_messages(&mut queued, "turn-1")
603            .expect("stage message");
604
605        assert_eq!(
606            result,
607            SubmissionResult::Accepted {
608                input_changed: true
609            }
610        );
611        assert!(staged[0].input.get("_mobius_internal").is_some());
612        assert_eq!(
613            staged[0].event.message().map(|message| message.delivery),
614            Some(MessageDelivery::Steer)
615        );
616        assert!(
617            Messages::default()
618                .render(&staged[0].event, "turn-1")
619                .is_none(),
620            "peer messages use their typed message event, not an activity block"
621        );
622        assert_eq!(staged[0].event.message().unwrap().text, peer.text);
623    }
624
625    #[test]
626    fn queued_external_event_waits_for_its_own_turn() {
627        let stack = MiddlewareStack::new(vec![Arc::new(Messages::default())]).expect("stack");
628        let source = MessageSubmission {
629            author: MessageAuthor::Source {
630                message_id: "report".into(),
631                source: crate::protocol::MessageSource::External {
632                    source_id: "routine".into(),
633                    event_id: "completed".into(),
634                },
635                cause_id: None,
636                ancestry: Vec::new(),
637                handle: "monitor".into(),
638                symbol: None,
639            },
640            text: "Monitoring completed.".into(),
641            attachments: Vec::new(),
642            reply: None,
643            requested_delivery: Some(ActiveMessageDelivery::Queue),
644            target_turn_id: None,
645        };
646        let mut queued = Vec::new();
647        assert_eq!(
648            route(&stack, &mut queued, &source, Some("user-turn")),
649            SubmissionResult::Accepted {
650                input_changed: false
651            }
652        );
653        assert!(
654            stack
655                .stage_model_messages(&mut queued, "user-turn")
656                .expect("stage")
657                .is_empty()
658        );
659        let report = stack.next_turn(&mut queued).expect("next").expect("report");
660        let block = Messages::default()
661            .render(&report.event, "report-turn")
662            .expect("external activity");
663        assert_eq!(block.role, FrontendBlockRole::Activity);
664        assert_eq!(block.title, "Message received from @monitor");
665        assert_eq!(block.text, source.text);
666        assert_eq!(block.symbol, Some(FrontendSymbol::Chat));
667        assert!(matches!(report.event, EventMsg::Message(event) if event.author == source.author));
668    }
669
670    #[test]
671    fn failed_turn_promotes_unstaged_steering_to_a_queued_turn() {
672        let stack = MiddlewareStack::new(vec![Arc::new(Messages::default())]).expect("stack");
673        let mut queued = Vec::new();
674        route(&stack, &mut queued, &user(None), Some("turn-1"));
675
676        stack
677            .finish_message_turn(
678                &mut queued,
679                "turn-1",
680                crate::backend::checkpoint::ExecutionOutcome::Failed,
681            )
682            .expect("promote failed turn");
683        let next = stack
684            .next_turn(&mut queued)
685            .expect("next turn")
686            .expect("promoted message");
687
688        assert_eq!(
689            next.event.message().map(|message| message.delivery),
690            Some(MessageDelivery::Queue)
691        );
692    }
693}