Skip to main content

mobius/middleware/
compaction.rs

1//! Context compaction policy and provider routing.
2
3use std::collections::BTreeSet;
4use std::sync::Arc;
5
6use super::Middleware;
7use super::ModelContext;
8use super::approximate_item_tokens;
9use super::attachments::is_attachment_materialization;
10use super::manifest::{MiddlewareManifest, MiddlewareSettingManifest};
11use super::scratchpad::is_projection_item;
12use serde_json::Value;
13use uuid::Uuid;
14
15use crate::BoxFuture;
16use crate::Error;
17use crate::Result;
18use crate::backend::checkpoint::ContextRewriteReason;
19use crate::backend::model::CompactOutput;
20use crate::backend::model::CompactRequest;
21use crate::backend::model::ModelRequest;
22use crate::backend::model::PromptCacheIdentity;
23use crate::backend::model::ToolDefinition;
24use crate::backend::model::ToolLoad;
25use crate::backend::model::internal_user_message;
26use crate::backend::model::prompt_cache_key;
27use crate::backend::model::reset_prompt_cache_breakpoint;
28use crate::backend::model::user_message;
29use crate::protocol::CONTEXT_COMPACTED_MARKER;
30use crate::protocol::EventMsg;
31use crate::protocol::FrontendBlock;
32use crate::protocol::FrontendTone;
33use crate::protocol::MESSAGE_METADATA_FIELD;
34use crate::protocol::internal_message_kind;
35use crate::protocol::is_internal_message;
36use crate::protocol::tool_complete_boundaries;
37
38mod text {
39    include!(concat!(
40        env!("OUT_DIR"),
41        "/src_middleware_compaction_text.rs"
42    ));
43}
44
45const KEEP_RECENT_TOKENS: usize = 20_000;
46const NATIVE_RETAINED_TOKENS: usize = 64_000;
47const MAX_SUMMARY_TOOL_RESULT_CHARS: usize = 2_000;
48const COMPACTION_RESERVE_TOKENS: i64 = 16_384;
49const _: () = {
50    assert!(text::DEFAULTS_COMPACTION_TOKENS >= 1);
51    assert!(text::SETTING_AT_TOKENS_STEP > 0);
52};
53/// Default compaction trigger for middleware instances without an override.
54pub const DEFAULT_COMPACTION_TOKENS: i64 = text::DEFAULTS_COMPACTION_TOKENS;
55const SETTINGS: &[MiddlewareSettingManifest] = &[MiddlewareSettingManifest::Integer {
56    id: "at_tokens",
57    label: text::SETTING_AT_TOKENS_LABEL,
58    description: text::SETTING_AT_TOKENS_DESCRIPTION,
59    min: 1,
60    max: None,
61    step: text::SETTING_AT_TOKENS_STEP,
62    default: DEFAULT_COMPACTION_TOKENS,
63}];
64
65/// Configuration and presentation metadata for compaction.
66pub const MANIFEST: MiddlewareManifest = MiddlewareManifest {
67    id: "compaction",
68    label: text::MANIFEST_LABEL,
69    description: text::MANIFEST_DESCRIPTION,
70    required: false,
71    default_enabled: true,
72    settings: SETTINGS,
73};
74
75/// Compacts visible context after a configurable token threshold.
76pub struct Compaction {
77    at_tokens: i64,
78}
79
80impl Default for Compaction {
81    fn default() -> Self {
82        Self {
83            at_tokens: DEFAULT_COMPACTION_TOKENS,
84        }
85    }
86}
87
88impl Compaction {
89    /// Creates a threshold-based compaction policy.
90    pub fn new(at_tokens: i64) -> Result<Self> {
91        if at_tokens <= 0 {
92            return Err(Error::Config(
93                "compaction threshold must be positive".into(),
94            ));
95        }
96        Ok(Self { at_tokens })
97    }
98
99    /// Returns the effective trigger after reserving response space.
100    #[must_use]
101    pub fn trigger_tokens(&self, context_window: i64) -> i64 {
102        self.at_tokens
103            .min(context_window.saturating_sub(COMPACTION_RESERVE_TOKENS))
104            .max(1)
105    }
106}
107
108impl Middleware for Compaction {
109    fn name(&self) -> &'static str {
110        MANIFEST.id
111    }
112
113    fn render(&self, event: &EventMsg, _session_id: &str) -> Option<FrontendBlock> {
114        matches!(event, EventMsg::ContextCompacted).then(|| FrontendBlock {
115            id: None,
116            group: None,
117            update: crate::protocol::FrontendBlockUpdate::Replace,
118            state: crate::protocol::FrontendBlockState::Complete,
119            role: crate::protocol::FrontendBlockRole::Notice,
120            title: text::RENDER_CONTEXT_COMPACTED.into(),
121            text: String::new(),
122            symbol: None,
123            files: Vec::new(),
124            format: crate::protocol::FrontendBlockFormat::PlainText,
125            tone: FrontendTone::Neutral,
126        })
127    }
128
129    fn pre_model<'a>(&'a self, context: &'a mut ModelContext<'_>) -> BoxFuture<'a, Result<()>> {
130        Box::pin(async move {
131            let estimated = context.estimated_input_tokens();
132            let observed = if contains_compaction(context.input()) {
133                estimated
134            } else {
135                context
136                    .last_usage
137                    .map_or(0, |usage| usage.input_tokens)
138                    .max(estimated)
139            };
140            if observed < self.trigger_tokens(context.context_window) || context.input().is_empty()
141            {
142                return Ok(());
143            }
144            context.pre_compact().await?;
145            if context.turn_stopped() {
146                return Ok(());
147            }
148            let catalog_revision = context.tools.revision()?;
149            let tool_load = retained_tool_load(
150                context.input(),
151                catalog_revision,
152                &context.tools.deferred_definitions(),
153            )?;
154            let output = if context.model.compaction_endpoint(context.provider)? {
155                let tools = context
156                    .tools
157                    .direct_definitions()
158                    .iter()
159                    .filter(|tool| context.available_tools.contains(&tool.name))
160                    .cloned()
161                    .collect::<Vec<_>>();
162                let deferred_tools = context
163                    .tools
164                    .deferred_definitions()
165                    .iter()
166                    .filter(|tool| context.available_tools.contains(&tool.name))
167                    .cloned()
168                    .collect::<Vec<_>>();
169                let cache_key = prompt_cache_key(context.session_id);
170                context
171                    .model
172                    .compact(
173                        context.provider,
174                        CompactRequest {
175                            session_id: context.session_id,
176                            prompt_cache: Some(PromptCacheIdentity {
177                                key: &cache_key,
178                                context_epoch: *context.context_epoch,
179                            }),
180                            instructions: context.instructions,
181                            input: context.input(),
182                            catalog_revision,
183                            tools: &tools,
184                            deferred_tools: &deferred_tools,
185                        },
186                    )
187                    .await?
188            } else {
189                summarize(context).await?
190            };
191            if output.output.is_empty() {
192                return Err(Error::Provider(
193                    "compaction returned an empty context".into(),
194                ));
195            }
196            let latest_turn_input = latest_turn_input(context.input());
197            let active_message_metadata = latest_turn_input
198                .as_ref()
199                .and_then(|active| active.item.get(MESSAGE_METADATA_FIELD))
200                .cloned();
201            let mut compacted = retain_native_context(context.input(), output.output);
202            compacted.retain(|item| !is_projection_item(item));
203            restore_input_private_fields(&mut compacted, latest_turn_input);
204            validate_active_message_metadata(&compacted, active_message_metadata.as_ref())?;
205            reset_prompt_cache_breakpoint(&mut compacted);
206            if let Some(tool_load) = tool_load {
207                compacted.push(tool_load);
208            }
209            context.rewrite_input(ContextRewriteReason::Compaction, compacted)?;
210            *context.compaction_count = context
211                .compaction_count
212                .checked_add(1)
213                .ok_or_else(|| Error::Checkpoint("compaction count overflow".into()))?;
214            context.record_transcript_item(internal_user_message(CONTEXT_COMPACTED_MARKER, ""));
215            context.usage.push(output.usage);
216            context.events.push(EventMsg::ContextCompacted);
217            context.post_compact().await?;
218            Ok(())
219        })
220    }
221}
222
223fn retained_tool_load(
224    input: &[Value],
225    catalog_revision: &str,
226    deferred_tools: &[ToolDefinition],
227) -> Result<Option<Value>> {
228    let deferred = deferred_tools
229        .iter()
230        .map(|tool| tool.name.as_str())
231        .collect::<BTreeSet<_>>();
232    let mut loaded = BTreeSet::new();
233    for item in input {
234        let Some(tool_load) = ToolLoad::from_input(item)? else {
235            continue;
236        };
237        if tool_load.catalog_revision == catalog_revision {
238            loaded.extend(
239                tool_load
240                    .tools
241                    .into_iter()
242                    .filter(|name| deferred.contains(name.as_str())),
243            );
244        }
245    }
246    Ok((!loaded.is_empty()).then(|| {
247        ToolLoad {
248            catalog_revision: catalog_revision.into(),
249            tools: loaded.into_iter().collect(),
250        }
251        .into_input()
252    }))
253}
254
255fn retain_native_context(input: &[Value], mut compacted: Vec<Value>) -> Vec<Value> {
256    if compacted.len() != 1
257        || compacted[0].get("type").and_then(Value::as_str) != Some("compaction")
258    {
259        return compacted;
260    }
261    let cut = recent_cut(input, NATIVE_RETAINED_TOKENS).unwrap_or(0);
262    let recent = &input[cut..];
263    let mut retained = Vec::new();
264    for (index, item) in recent.iter().enumerate() {
265        if (!is_internal_message(item) || item.get(MESSAGE_METADATA_FIELD).is_some())
266            && item.get("role").and_then(Value::as_str) == Some("user")
267        {
268            retained.push(item.clone());
269            if let Some(materialization) = recent.get(index + 1)
270                && is_attachment_materialization(materialization)
271            {
272                retained.push(materialization.clone());
273            }
274        }
275    }
276    retained.append(&mut compacted);
277    retained
278}
279
280struct LatestTurnInput<'a> {
281    item: &'a Value,
282    materialization: Option<&'a Value>,
283}
284
285fn latest_turn_input(input: &[Value]) -> Option<LatestTurnInput<'_>> {
286    let index = input
287        .iter()
288        .rposition(|item| item.get(MESSAGE_METADATA_FIELD).is_some())?;
289    Some(LatestTurnInput {
290        item: &input[index],
291        materialization: input
292            .get(index + 1)
293            .filter(|item| is_attachment_materialization(item)),
294    })
295}
296
297fn validate_active_message_metadata(input: &[Value], expected: Option<&Value>) -> Result<()> {
298    let Some(expected) = expected else {
299        return Ok(());
300    };
301    if latest_turn_input(input).and_then(|active| active.item.get(MESSAGE_METADATA_FIELD))
302        == Some(expected)
303    {
304        return Ok(());
305    }
306    Err(Error::Provider(
307        "compaction did not preserve active message metadata".into(),
308    ))
309}
310
311fn restore_input_private_fields(
312    compacted: &mut Vec<Value>,
313    latest_turn_input: Option<LatestTurnInput<'_>>,
314) {
315    let Some(LatestTurnInput {
316        item: input,
317        materialization,
318    }) = latest_turn_input
319    else {
320        return;
321    };
322    let Some(fields) = input.as_object() else {
323        return;
324    };
325    let private = fields
326        .iter()
327        .filter(|(name, _)| name.starts_with('_'))
328        .map(|(name, value)| (name.clone(), value.clone()))
329        .collect::<Vec<_>>();
330    if private.is_empty() && materialization.is_none() {
331        return;
332    }
333    let retained_index = compacted.iter().rposition(|item| {
334        item.get("role") == input.get("role") && item.get("content") == input.get("content")
335    });
336    let input_index = if let Some(index) = retained_index {
337        if let Some(fields) = compacted[index].as_object_mut() {
338            fields.extend(private);
339        }
340        index
341    } else {
342        compacted.push(input.clone());
343        compacted.len() - 1
344    };
345    restore_attachment_materialization(compacted, input_index, materialization);
346}
347
348fn restore_attachment_materialization(
349    compacted: &mut Vec<Value>,
350    user_index: usize,
351    materialization: Option<&Value>,
352) {
353    let Some(materialization) = materialization else {
354        return;
355    };
356    match compacted.get(user_index + 1) {
357        Some(retained) if retained == materialization => {}
358        Some(retained) if is_attachment_materialization(retained) => {
359            compacted[user_index + 1] = materialization.clone();
360        }
361        Some(_) | None => compacted.insert(user_index + 1, materialization.clone()),
362    }
363}
364
365async fn summarize(context: &ModelContext<'_>) -> Result<CompactOutput> {
366    let (prompt, recent) = prepare_summary(context.input())
367        .ok_or_else(|| Error::Provider("context has no safe history boundary to compact".into()))?;
368    let session_id = Uuid::new_v4().to_string();
369    let cache_key = prompt_cache_key(&session_id);
370    let input = [user_message(&prompt)];
371    let output = context
372        .model
373        .respond(
374            context.provider,
375            ModelRequest {
376                session_id: &session_id,
377                prompt_cache: Some(PromptCacheIdentity {
378                    key: &cache_key,
379                    context_epoch: *context.context_epoch,
380                }),
381                instructions: text::PROMPT_SUMMARY_SYSTEM,
382                input: &input,
383                catalog_revision: context.tools.revision()?,
384                tools: &[],
385                deferred_tools: &[],
386                allow_hosted_tools: false,
387                allow_continuation: false,
388            },
389            Arc::new(|_| Ok(())),
390        )
391        .await?;
392    let summary = output.text().trim();
393    if summary.is_empty() {
394        return Err(Error::Provider(
395            "model compaction returned no summary".into(),
396        ));
397    }
398    let mut compacted = Vec::with_capacity(recent.len() + 1);
399    compacted.push(internal_user_message(
400        "compaction",
401        &format!("<compacted_context>\n{summary}\n</compacted_context>"),
402    ));
403    for item in recent {
404        if ToolLoad::from_input(&item)?.is_none() {
405            compacted.push(item);
406        }
407    }
408    CompactOutput::from_output(compacted, output.usage().clone())
409}
410
411fn prepare_summary(input: &[Value]) -> Option<(String, Vec<Value>)> {
412    let cut = recent_cut(input, KEEP_RECENT_TOKENS)?;
413    let prompt = summary_prompt(&input[..cut])?;
414    Some((prompt, input[cut..].to_vec()))
415}
416
417fn recent_cut(input: &[Value], keep_tokens: usize) -> Option<usize> {
418    let mut accumulated = 0;
419    let mut desired = None;
420    for index in (0..input.len()).rev() {
421        accumulated += approximate_item_tokens(&input[index]);
422        if accumulated >= keep_tokens {
423            desired = Some(index);
424            break;
425        }
426    }
427    let desired = desired?;
428    let safe = safe_boundaries(input);
429    safe.iter()
430        .rev()
431        .copied()
432        .find(|&index| index > 0 && index <= desired)
433        .or_else(|| {
434            safe.iter()
435                .copied()
436                .find(|&index| index > desired && index < input.len())
437        })
438}
439
440fn safe_boundaries(input: &[Value]) -> Vec<usize> {
441    tool_complete_boundaries(input)
442        .into_iter()
443        .filter(|&boundary| boundary == input.len() || safe_start(&input[boundary]))
444        .collect()
445}
446
447fn safe_start(item: &Value) -> bool {
448    if is_attachment_materialization(item) {
449        return false;
450    }
451    match item.get("type").and_then(Value::as_str) {
452        Some("function_call") => true,
453        Some("message") | None => matches!(
454            item.get("role").and_then(Value::as_str),
455            Some("user" | "assistant")
456        ),
457        Some(_) => false,
458    }
459}
460
461fn summary_prompt(history: &[Value]) -> Option<String> {
462    let mut conversation = Vec::new();
463    let mut previous_summary = None;
464    for item in history {
465        if let Some(summary) = compacted_summary(item) {
466            previous_summary = Some(summary);
467        } else if let Some(serialized) = serialize_item(item) {
468            conversation.push(serialized);
469        }
470    }
471    if conversation.is_empty() {
472        return None;
473    }
474    let mut prompt = format!(
475        "<conversation>\n{}\n</conversation>\n",
476        conversation.join("\n\n")
477    );
478    if let Some(summary) = previous_summary {
479        prompt.push_str(&format!(
480            "\n<previous_summary>\n{summary}\n</previous_summary>\n"
481        ));
482    }
483    prompt.push_str(&format!("\n{}", text::PROMPT_SUMMARY_TASK));
484    Some(prompt)
485}
486
487fn serialize_item(item: &Value) -> Option<String> {
488    match item.get("type").and_then(Value::as_str) {
489        Some("function_call") => Some(format!(
490            "[Assistant tool call]: {}({})",
491            item.get("name").and_then(Value::as_str).unwrap_or("tool"),
492            value_text(item.get("arguments"))
493        )),
494        Some("function_call_output") => Some(format!(
495            "[Tool result]: {}",
496            truncate_chars(
497                &value_text(item.get("output")),
498                MAX_SUMMARY_TOOL_RESULT_CHARS
499            )
500        )),
501        Some("reasoning") => {
502            let text = content_text(item.get("summary"));
503            (!text.is_empty()).then(|| format!("[Assistant reasoning]: {text}"))
504        }
505        Some("message") | None => {
506            let role = item.get("role").and_then(Value::as_str)?;
507            let text = content_text(item.get("content"));
508            (!text.is_empty()).then(|| {
509                let label = if role == "assistant" {
510                    "Assistant"
511                } else {
512                    "User"
513                };
514                format!("[{label}]: {text}")
515            })
516        }
517        Some(_) => None,
518    }
519}
520
521fn compacted_summary(item: &Value) -> Option<String> {
522    if internal_message_kind(item) != Some("compaction") {
523        return None;
524    }
525    let text = content_text(item.get("content"));
526    text.strip_prefix("<compacted_context>")?
527        .strip_suffix("</compacted_context>")
528        .map(|summary| summary.trim().to_string())
529}
530
531fn contains_compaction(input: &[Value]) -> bool {
532    input.iter().any(|item| {
533        item.get("type").and_then(Value::as_str) == Some("compaction")
534            || compacted_summary(item).is_some()
535    })
536}
537
538fn content_text(value: Option<&Value>) -> String {
539    match value {
540        Some(Value::String(text)) => text.clone(),
541        Some(Value::Array(parts)) => parts
542            .iter()
543            .filter_map(|part| {
544                part.get("text")
545                    .or_else(|| part.get("content"))
546                    .and_then(Value::as_str)
547            })
548            .collect::<Vec<_>>()
549            .join("\n"),
550        Some(value) => value.to_string(),
551        None => String::new(),
552    }
553}
554
555fn value_text(value: Option<&Value>) -> String {
556    match value {
557        Some(Value::String(text)) => text.clone(),
558        Some(value) => value.to_string(),
559        None => String::new(),
560    }
561}
562
563fn truncate_chars(text: &str, limit: usize) -> String {
564    text.char_indices()
565        .nth(limit)
566        .map_or_else(|| text.to_string(), |(end, _)| format!("{}…", &text[..end]))
567}
568
569#[cfg(test)]
570mod tests {
571    use super::*;
572    use crate::backend::model::tool_output;
573
574    #[test]
575    fn recent_cut_keeps_parallel_calls_with_their_outputs() {
576        let input = vec![
577            user_message("old"),
578            serde_json::json!({
579                "type": "function_call",
580                "call_id": "a",
581                "name": "read",
582                "arguments": "{}"
583            }),
584            serde_json::json!({
585                "type": "function_call",
586                "call_id": "b",
587                "name": "read",
588                "arguments": "{}"
589            }),
590            tool_output("a", &"x".repeat(200), false),
591            tool_output("b", "done", false),
592        ];
593
594        assert_eq!(recent_cut(&input, 10), Some(1));
595    }
596
597    #[test]
598    fn trigger_reserves_space_from_the_live_context_window() {
599        let compaction = Compaction::default();
600
601        assert_eq!(compaction.trigger_tokens(128_000), 111_616);
602        assert_eq!(compaction.trigger_tokens(8_000), 1);
603        assert_eq!(
604            Compaction::new(4_000)
605                .expect("custom threshold")
606                .trigger_tokens(128_000),
607            4_000
608        );
609    }
610
611    #[test]
612    fn compaction_restores_private_fields_on_the_retained_user() {
613        let user = serde_json::json!({
614            "role": "user",
615            "content": [{"type": "input_text", "text": "inspect"}],
616            "_middleware_state": {"id": "state"}
617        });
618        let mut compacted = vec![
619            serde_json::json!({
620                "type": "message",
621                "id": "message-1",
622                "role": "user",
623                "status": "completed",
624                "content": [{"type": "input_text", "text": "inspect"}]
625            }),
626            serde_json::json!({"type": "compaction", "encrypted_content": "opaque"}),
627        ];
628
629        restore_input_private_fields(
630            &mut compacted,
631            Some(LatestTurnInput {
632                item: &user,
633                materialization: None,
634            }),
635        );
636
637        assert_eq!(compacted.len(), 2);
638        assert_eq!(compacted[0]["id"], "message-1");
639        assert_eq!(compacted[0]["status"], "completed");
640        assert_eq!(
641            compacted[0]["_middleware_state"],
642            serde_json::json!({"id": "state"})
643        );
644    }
645
646    #[test]
647    fn v2_compaction_retains_the_user_before_the_marker_without_its_tool_tail() {
648        let input = vec![
649            serde_json::json!({
650                "role": "developer",
651                "content": [{"type": "input_text", "text": "stale instructions"}]
652            }),
653            serde_json::json!({
654                "role": "user",
655                "content": [{"type": "input_text", "text": "inspect"}],
656                "_middleware_state": {"id": "state"}
657            }),
658            serde_json::json!({
659                "type": "function_call",
660                "call_id": "call-1",
661                "name": "read",
662                "arguments": "{}"
663            }),
664            tool_output("call-1", "large result", false),
665        ];
666        let compacted = vec![serde_json::json!({
667            "type": "compaction",
668            "encrypted_content": "opaque"
669        })];
670        let mut compacted = retain_native_context(&input, compacted);
671        restore_input_private_fields(&mut compacted, latest_turn_input(&input));
672
673        assert_eq!(compacted.len(), 2);
674        assert_eq!(compacted[0], input[1]);
675        assert_eq!(compacted[1]["type"], "compaction");
676    }
677
678    #[test]
679    fn compaction_restores_an_omitted_attachment_materialization_with_its_user() {
680        let user = crate::backend::model::message_input(&crate::protocol::MessageEvent {
681            author: crate::protocol::MessageAuthor::User,
682            delivery: crate::protocol::MessageDelivery::Turn,
683            text: "inspect".into(),
684            attachments: vec![crate::protocol::SessionFileReference {
685                id: "upload-1".into(),
686                name: "photo.png".into(),
687                size: 1,
688                media_type: "image/png".into(),
689            }],
690            message_target: None,
691        })
692        .expect("message input");
693        let materialization = internal_user_message(
694            crate::protocol::ATTACHMENT_CONTEXT_MARKER,
695            "attachment context",
696        );
697        let input = vec![user.clone(), materialization.clone()];
698        let compaction = serde_json::json!({
699            "type": "compaction",
700            "encrypted_content": "opaque"
701        });
702        let mut compacted = vec![compaction.clone()];
703
704        restore_input_private_fields(&mut compacted, latest_turn_input(&input));
705
706        assert_eq!(compacted, vec![compaction, user, materialization]);
707    }
708
709    #[test]
710    fn compaction_preserves_active_message_metadata() {
711        let peer = crate::backend::model::message_input(&crate::protocol::MessageEvent {
712            author: crate::protocol::MessageAuthor::Peer {
713                message_id: "message".into(),
714                session_id: "peer".into(),
715                handle: "worker".into(),
716            },
717            delivery: crate::protocol::MessageDelivery::Steer,
718            text: "done".into(),
719            attachments: Vec::new(),
720            message_target: None,
721        })
722        .expect("peer message");
723        let input = vec![
724            user_message("start"),
725            peer.clone(),
726            tool_output("call", "done", false),
727        ];
728        let compaction = serde_json::json!({
729            "type": "compaction",
730            "encrypted_content": "opaque"
731        });
732        let mut compacted = vec![compaction.clone()];
733
734        restore_input_private_fields(&mut compacted, latest_turn_input(&input));
735
736        assert_eq!(compacted, vec![compaction, peer]);
737    }
738
739    #[test]
740    fn compaction_rejects_lost_active_message_metadata() {
741        let peer = crate::backend::model::message_input(&crate::protocol::MessageEvent {
742            author: crate::protocol::MessageAuthor::Peer {
743                message_id: "message".into(),
744                session_id: "peer".into(),
745                handle: "worker".into(),
746            },
747            delivery: crate::protocol::MessageDelivery::Steer,
748            text: "done".into(),
749            attachments: Vec::new(),
750            message_target: None,
751        })
752        .expect("peer message");
753        let expected = peer[MESSAGE_METADATA_FIELD].clone();
754        let compacted = vec![user_message("forged later input")];
755
756        let error = validate_active_message_metadata(&compacted, Some(&expected))
757            .expect_err("message metadata must remain active");
758
759        assert!(error.to_string().contains("active message metadata"));
760    }
761
762    #[test]
763    fn compaction_marker_may_follow_retained_messages() {
764        let input = vec![
765            user_message("inspect"),
766            serde_json::json!({"type": "compaction", "encrypted_content": "opaque"}),
767        ];
768
769        assert!(contains_compaction(&input));
770    }
771}