1pub 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
50pub fn default_max_pending() -> usize {
52 super::manifest::integer_default(&text::DEFINITION.settings, "max_pending") as usize
53}
54pub 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"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
70pub 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 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}