Skip to main content

meerkat_core/
session_store.rs

1//! SessionStore trait — canonical session persistence contract.
2//!
3//! This trait lives in `meerkat-core` so that custom storage implementations
4//! (Postgres, DynamoDB, etc.) can be written without depending on `meerkat-store`.
5//!
6//! # Snapshot = projection
7//!
8//! The `Session` row a `SessionStore` persists is a **projection of the
9//! canonical event log**. The event log (`EventStore`) is append-only at
10//! the trait level; the snapshot is a rebuildable materialization of
11//! replaying that log. Deleting a `.rkat/sessions/<id>/session.json` and
12//! replaying the event store produces an identical snapshot (the
13//! `CLAUDE.md` invariant).
14//!
15//! Wave-c C-H1 (F1 closure from the state-scope-audit) makes the
16//! append-only nature of that projection enforceable at the
17//! `SessionStore::save` boundary — see the trait docs on
18//! [`SessionStore`] and the [`append_only_save_guard`] helper.
19
20use async_trait::async_trait;
21use sha2::{Digest, Sha256};
22
23use crate::session::{SYSTEM_CONTEXT_SEPARATOR, SessionMeta};
24use crate::time_compat::SystemTime;
25use crate::types::{Message, SessionId, SystemMessage};
26use crate::{
27    Session, TranscriptHistoryState, TranscriptRewriteCommit, TranscriptRewriteSelection,
28    transcript_messages_digest,
29};
30
31/// Filter for listing sessions.
32#[derive(Debug, Clone, Default)]
33pub struct SessionFilter {
34    /// Only sessions created after this time.
35    pub created_after: Option<SystemTime>,
36    /// Only sessions updated after this time.
37    pub updated_after: Option<SystemTime>,
38    /// Maximum number of results.
39    pub limit: Option<usize>,
40    /// Offset for pagination.
41    pub offset: Option<usize>,
42}
43
44/// Errors from session store operations.
45///
46/// Backend-specific details (rusqlite, filesystem, etc.) are erased to strings
47/// so that the trait contract carries no I/O dependencies.
48#[derive(Debug, thiserror::Error)]
49pub enum SessionStoreError {
50    #[error("IO error: {0}")]
51    Io(#[from] std::io::Error),
52
53    #[error("Serialization error: {0}")]
54    Serialization(String),
55
56    #[error("Session not found: {0}")]
57    NotFound(SessionId),
58
59    #[error("Session corrupted: {0}")]
60    Corrupted(SessionId),
61
62    #[error(
63        "session {id} save rejected: new message count {new_len} is shorter than previously \
64         persisted {prev_len} without transcript-continuity proof"
65    )]
66    MonotonicityViolation {
67        id: SessionId,
68        prev_len: usize,
69        new_len: usize,
70    },
71
72    #[error(
73        "session {id} save rejected: incoming transcript is not a continuation of persisted revision {previous_revision}"
74    )]
75    TranscriptContinuityViolation {
76        id: SessionId,
77        previous_revision: String,
78        incoming_revision: String,
79        reason: String,
80    },
81
82    #[error(
83        "session {id} rewrite rejected: previous transcript revision {actual} did not match commit parent {expected}"
84    )]
85    TranscriptRevisionConflict {
86        id: SessionId,
87        expected: String,
88        actual: String,
89    },
90
91    #[error("session {id} rewrite rejected: {reason}")]
92    InvalidTranscriptRewrite { id: SessionId, reason: String },
93
94    #[error("Internal error: {0}")]
95    Internal(String),
96}
97
98/// Stable compare token for a full persisted session projection row.
99pub fn session_projection_cas_token(session: &Session) -> Result<String, SessionStoreError> {
100    let bytes = serde_json::to_vec(session).map_err(|err| {
101        SessionStoreError::Serialization(format!(
102            "failed to serialize session projection CAS token: {err}"
103        ))
104    })?;
105    Ok(format!("row-sha256:{:x}", Sha256::digest(bytes)))
106}
107
108/// Shared append-only guard for `SessionStore::save` implementations.
109///
110/// Backends call this at the top of their `save` method with the new
111/// session and the previously persisted row (or `None` if no prior row
112/// exists). Returns
113/// [`SessionStoreError::MonotonicityViolation`] when the new row's
114/// message count is strictly smaller than the previously persisted one
115/// without a transcript graph edge that proves a core-owned mutation.
116///
117/// The guard also rejects equal/longer saves whose retained prefix no longer
118/// matches the persisted transcript. A plain save may append or update
119/// metadata; same-session replacement must go through
120/// [`transcript_rewrite_save_guard`].
121pub fn append_only_save_guard(
122    incoming: &Session,
123    previous: Option<&Session>,
124) -> Result<(), SessionStoreError> {
125    incoming
126        .validate_transcript_history_state()
127        .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
128            id: incoming.id().clone(),
129            reason: format!("incoming transcript history state is malformed: {err}"),
130        })?;
131    let incoming_revision =
132        transcript_messages_digest(incoming.messages()).map_err(SessionStoreError::from)?;
133    let incoming_state = incoming.transcript_history_state().map_err(|err| {
134        SessionStoreError::InvalidTranscriptRewrite {
135            id: incoming.id().clone(),
136            reason: format!("incoming transcript history state is malformed: {err}"),
137        }
138    })?;
139    if let Some(state) = incoming_state.as_ref()
140        && state.head != incoming_revision
141    {
142        return Err(SessionStoreError::InvalidTranscriptRewrite {
143            id: incoming.id().clone(),
144            reason: format!(
145                "incoming transcript graph head {} does not match current message digest {incoming_revision}",
146                state.head
147            ),
148        });
149    }
150
151    let Some(previous) = previous else {
152        if incoming_state.is_some() {
153            return Err(SessionStoreError::InvalidTranscriptRewrite {
154                id: incoming.id().clone(),
155                reason: "incoming first save would seed transcript history state outside the rewrite/audit path"
156                    .to_string(),
157            });
158        }
159        validate_plain_save_transcript_history_preservation(
160            incoming,
161            None,
162            None,
163            incoming_state.as_ref(),
164        )?;
165        return Ok(());
166    };
167    let previous_state = previous.transcript_history_state().map_err(|err| {
168        SessionStoreError::InvalidTranscriptRewrite {
169            id: incoming.id().clone(),
170            reason: format!("previous transcript history state is malformed: {err}"),
171        }
172    })?;
173    let previous_had_history = previous_state.is_some();
174    let incoming_has_history = incoming_state.is_some();
175    if previous_had_history && !incoming_has_history {
176        return Err(SessionStoreError::InvalidTranscriptRewrite {
177            id: incoming.id().clone(),
178            reason: "incoming save would erase retained transcript history state".to_string(),
179        });
180    }
181    let previous_revision =
182        transcript_messages_digest(previous.messages()).map_err(SessionStoreError::from)?;
183    if previous_revision == incoming_revision {
184        validate_plain_save_transcript_history_preservation(
185            incoming,
186            Some(previous),
187            previous_state.as_ref(),
188            incoming_state.as_ref(),
189        )?;
190        return Ok(());
191    }
192
193    let prev_len = previous.messages().len();
194    let new_len = incoming.messages().len();
195    if new_len >= prev_len {
196        let incoming_prefix_revision = transcript_messages_digest(&incoming.messages()[..prev_len])
197            .map_err(SessionStoreError::from)?;
198        if incoming_prefix_revision == previous_revision {
199            validate_plain_save_transcript_history_preservation(
200                incoming,
201                Some(previous),
202                previous_state.as_ref(),
203                incoming_state.as_ref(),
204            )?;
205            return Ok(());
206        }
207    }
208    if incoming_preserves_conversation_tail_with_system_context_append(incoming, previous)? {
209        validate_plain_save_transcript_history_preservation(
210            incoming,
211            Some(previous),
212            previous_state.as_ref(),
213            incoming_state.as_ref(),
214        )?;
215        return Ok(());
216    }
217    if incoming_preserves_prefix_after_transient_notice_cleanup(incoming, previous)? {
218        validate_plain_save_transcript_history_preservation(
219            incoming,
220            Some(previous),
221            previous_state.as_ref(),
222            incoming_state.as_ref(),
223        )?;
224        return Ok(());
225    }
226    if new_len < prev_len {
227        return Err(SessionStoreError::MonotonicityViolation {
228            id: incoming.id().clone(),
229            prev_len,
230            new_len,
231        });
232    }
233
234    Err(SessionStoreError::TranscriptContinuityViolation {
235        id: incoming.id().clone(),
236        previous_revision,
237        incoming_revision,
238        reason: "incoming transcript neither preserves the persisted prefix nor records a graph edge from the persisted head".to_string(),
239    })
240}
241
242fn validate_plain_save_transcript_history_preservation(
243    incoming: &Session,
244    previous: Option<&Session>,
245    previous_state: Option<&TranscriptHistoryState>,
246    incoming_state: Option<&TranscriptHistoryState>,
247) -> Result<(), SessionStoreError> {
248    let Some(previous) = previous else {
249        if incoming_state.is_some() {
250            return Err(SessionStoreError::InvalidTranscriptRewrite {
251                id: incoming.id().clone(),
252                reason: "incoming first save would seed transcript history state outside the rewrite/audit path"
253                    .to_string(),
254            });
255        }
256        return Ok(());
257    };
258    if previous_state.is_none() && incoming_state.is_some() {
259        return Err(SessionStoreError::InvalidTranscriptRewrite {
260            id: incoming.id().clone(),
261            reason: "incoming append-only save would seed transcript history state outside the rewrite/audit path"
262                .to_string(),
263        });
264    }
265    let Some(previous_state) = previous_state else {
266        return Ok(());
267    };
268    let Some(incoming_state) = incoming_state else {
269        return Err(SessionStoreError::InvalidTranscriptRewrite {
270            id: incoming.id().clone(),
271            reason: "incoming append-only save would erase retained transcript history state"
272                .to_string(),
273        });
274    };
275    let previous_commits = previous_state.commits.as_slice();
276    let incoming_commits = incoming_state.commits.as_slice();
277    if incoming_commits != previous_commits {
278        return Err(SessionStoreError::InvalidTranscriptRewrite {
279            id: incoming.id().clone(),
280            reason: "incoming append-only save would change retained transcript rewrite commits"
281                .to_string(),
282        });
283    }
284    let retained_revisions_preserved =
285        transcript_revision_bodies_preserved(previous_state, incoming_state)?;
286    if retained_revisions_preserved
287        && incoming_state.revisions.len() == previous_state.revisions.len()
288        && incoming_state.head == previous_state.head
289    {
290        return Ok(());
291    }
292    if incoming_state.revisions.len() != previous_state.revisions.len() + 1
293        || !retained_revisions_preserved
294    {
295        return Err(SessionStoreError::InvalidTranscriptRewrite {
296            id: incoming.id().clone(),
297            reason: "incoming append-only save would change retained transcript revision graph"
298                .to_string(),
299        });
300    }
301    let incoming_revision =
302        transcript_messages_digest(incoming.messages()).map_err(SessionStoreError::from)?;
303    let previous_revision =
304        transcript_messages_digest(previous.messages()).map_err(SessionStoreError::from)?;
305    if previous_state.head != previous_revision {
306        return Err(SessionStoreError::InvalidTranscriptRewrite {
307            id: incoming.id().clone(),
308            reason: "previous transcript history head does not match persisted message digest"
309                .to_string(),
310        });
311    }
312    let added = &incoming_state.revisions[previous_state.revisions.len()];
313    if incoming_state.head != incoming_revision
314        || added.revision != incoming_revision
315        || added.parent_revision.as_deref() != Some(previous_state.head.as_str())
316        || transcript_messages_digest(&added.messages).map_err(SessionStoreError::from)?
317            != incoming_revision
318    {
319        return Err(SessionStoreError::InvalidTranscriptRewrite {
320            id: incoming.id().clone(),
321            reason: "incoming append-only save would add a transcript revision body that is not the current append"
322                .to_string(),
323        });
324    }
325    Ok(())
326}
327
328fn transcript_revision_bodies_preserved(
329    previous_state: &TranscriptHistoryState,
330    incoming_state: &TranscriptHistoryState,
331) -> Result<bool, SessionStoreError> {
332    if incoming_state.revisions.len() < previous_state.revisions.len() {
333        return Ok(false);
334    }
335    previous_state
336        .revisions
337        .iter()
338        .zip(incoming_state.revisions.iter())
339        .map(|(previous, incoming)| {
340            Ok(previous.revision == incoming.revision
341                && previous.parent_revision == incoming.parent_revision
342                && previous.created_at == incoming.created_at
343                && transcript_messages_digest(&previous.messages)
344                    .map_err(SessionStoreError::from)?
345                    == transcript_messages_digest(&incoming.messages)
346                        .map_err(SessionStoreError::from)?)
347        })
348        .try_fold(true, |acc, preserved| {
349            preserved.map(|preserved| acc && preserved)
350        })
351}
352
353fn validate_rewrite_save_retains_previous_commits(
354    incoming: &Session,
355    previous: &Session,
356    incoming_state: &TranscriptHistoryState,
357) -> Result<(), SessionStoreError> {
358    let previous_state = previous.transcript_history_state().map_err(|err| {
359        SessionStoreError::InvalidTranscriptRewrite {
360            id: incoming.id().clone(),
361            reason: format!("previous transcript history state is malformed: {err}"),
362        }
363    })?;
364    let Some(previous_state) = previous_state.as_ref() else {
365        return Ok(());
366    };
367    if incoming_state.commits.len() < previous_state.commits.len()
368        || incoming_state.commits[..previous_state.commits.len()] != previous_state.commits
369    {
370        return Err(SessionStoreError::InvalidTranscriptRewrite {
371            id: incoming.id().clone(),
372            reason: "incoming rewrite save would drop retained transcript rewrite commits"
373                .to_string(),
374        });
375    }
376    Ok(())
377}
378
379/// Validate that an authoritative projection write still targets the row that
380/// the caller proved continuity against.
381pub fn authoritative_projection_current_revision_guard(
382    incoming: &Session,
383    previous: Option<&Session>,
384    expected_current_revision: Option<&str>,
385) -> Result<(), SessionStoreError> {
386    let previous_token = previous.map(session_projection_cas_token).transpose()?;
387    if previous_token.as_deref() == expected_current_revision {
388        return Ok(());
389    }
390    let incoming_revision =
391        transcript_messages_digest(incoming.messages()).map_err(SessionStoreError::from)?;
392    Err(SessionStoreError::TranscriptContinuityViolation {
393        id: incoming.id().clone(),
394        previous_revision: previous_token.unwrap_or_else(|| "<missing>".to_string()),
395        incoming_revision,
396        reason: format!(
397            "authoritative projection expected persisted projection token {}, but current row has diverged",
398            expected_current_revision.unwrap_or("<missing>")
399        ),
400    })
401}
402
403fn incoming_preserves_conversation_tail_with_system_context_append(
404    incoming: &Session,
405    previous: &Session,
406) -> Result<bool, SessionStoreError> {
407    messages_preserve_conversation_tail_with_system_context_append(
408        incoming.messages(),
409        previous.messages(),
410    )
411}
412
413fn messages_preserve_conversation_tail_with_system_context_append(
414    incoming: &[Message],
415    previous: &[Message],
416) -> Result<bool, SessionStoreError> {
417    let (previous_system, previous_tail) = split_single_leading_system(previous);
418    let (incoming_system, incoming_tail) = split_single_leading_system(incoming);
419    let Some(incoming_system) = incoming_system else {
420        return Ok(false);
421    };
422    if !system_context_is_append(previous_system, incoming_system)? {
423        return Ok(false);
424    }
425    if incoming_tail.len() < previous_tail.len() {
426        return Ok(false);
427    }
428    let previous_tail_revision =
429        transcript_messages_digest(previous_tail).map_err(SessionStoreError::from)?;
430    let incoming_tail_prefix_revision =
431        transcript_messages_digest(&incoming_tail[..previous_tail.len()])
432            .map_err(SessionStoreError::from)?;
433    Ok(previous_tail_revision == incoming_tail_prefix_revision)
434}
435
436fn split_single_leading_system(messages: &[Message]) -> (Option<&SystemMessage>, &[Message]) {
437    match messages.first() {
438        Some(Message::System(system)) => (Some(system), &messages[1..]),
439        _ => (None, messages),
440    }
441}
442
443/// Decide whether `incoming` is a continuation of `previous` produced by a
444/// runtime system-context append.
445///
446/// The structural part — identical content, or `incoming = previous +
447/// separator + suffix` — is a transcript-continuity proof (content equality of
448/// the retained prefix), not classification. The SEMANTIC append-admission
449/// verdict ("is this incoming persisted prompt an admissible
450/// runtime-context-append continuation of the persisted one") is owned by the
451/// canonical [`SessionDocumentMachine`] — the same machine the staging path
452/// already drives for the four-way append disposition — not a handwritten shell
453/// reducer. This function extracts only the pure structural observations plus
454/// the typed [`SystemPromptMutationKind`] runtime-context-append marker, drives
455/// the machine's `ResolveSystemContextPersistAppendAdmission` input, and mirrors
456/// the emitted verdict (`Admit` -> `true`, `Reject` -> `false`). It fails closed
457/// if the machine refuses or emits no verdict.
458fn system_context_is_append(
459    previous: Option<&SystemMessage>,
460    incoming: &SystemMessage,
461) -> Result<bool, SessionStoreError> {
462    // Pure structural observations the shell computes; NO semantic decision.
463    let has_previous = previous.is_some();
464    let content_identical = previous.is_some_and(|previous| incoming.content == previous.content);
465    let content_extends_previous =
466        previous.is_some_and(|previous| incoming.content.starts_with(&previous.content));
467    let appended_starts_with_separator = previous.is_some_and(|previous| {
468        incoming
469            .content
470            .get(previous.content.len()..)
471            .is_some_and(|appended| appended.starts_with(SYSTEM_CONTEXT_SEPARATOR))
472    });
473    let incoming_is_runtime_context_append = incoming.mutation_kind.is_runtime_context_append();
474
475    let mut authority = crate::session_document::SessionDocumentMachineAuthority::new();
476    let effects = authority
477        .resolve_system_context_persist_append_admission(
478            has_previous,
479            content_identical,
480            content_extends_previous,
481            appended_starts_with_separator,
482            incoming_is_runtime_context_append,
483        )
484        .map_err(|err| {
485            SessionStoreError::Internal(format!(
486                "session document authority refused persist-time system-context append admission: {err}"
487            ))
488        })?;
489    effects
490        .into_iter()
491        .find_map(|effect| match effect {
492            crate::session_document::SessionDocumentEffect::SystemContextPersistAppendAdmissionResolved {
493                admission,
494            } => Some(matches!(
495                admission,
496                crate::session_document::SystemContextPersistAppendAdmission::Admit
497            )),
498            _ => None,
499        })
500        .ok_or_else(|| {
501            SessionStoreError::Internal(
502                "session document authority emitted no persist-time system-context append admission verdict".to_string(),
503            )
504        })
505}
506
507fn incoming_preserves_prefix_after_transient_notice_cleanup(
508    incoming: &Session,
509    previous: &Session,
510) -> Result<bool, SessionStoreError> {
511    let previous_without_transient = previous
512        .messages()
513        .iter()
514        .filter(|message| !is_transient_system_notice(message))
515        .cloned()
516        .collect::<Vec<_>>();
517    if previous_without_transient.len() == previous.messages().len()
518        || incoming.messages().len() < previous_without_transient.len()
519    {
520        return Ok(false);
521    }
522    let previous_revision =
523        transcript_messages_digest(&previous_without_transient).map_err(SessionStoreError::from)?;
524    let incoming_prefix_revision =
525        transcript_messages_digest(&incoming.messages()[..previous_without_transient.len()])
526            .map_err(SessionStoreError::from)?;
527    Ok(previous_revision == incoming_prefix_revision)
528}
529
530fn is_transient_system_notice(message: &Message) -> bool {
531    let Message::SystemNotice(notice) = message else {
532        return false;
533    };
534    notice.kind == crate::types::SystemNoticeKind::McpPending
535        && notice.blocks.iter().all(|block| {
536            matches!(
537                block,
538                crate::types::SystemNoticeBlock::Mcp {
539                    persisted: false,
540                    ..
541                }
542            )
543        })
544}
545
546/// Validate a runtime run-boundary snapshot.
547///
548/// Runtime turns normally append to the transcript, but core-owned turn
549/// mechanics such as compaction can also produce an audited internal rewrite.
550/// Runtime stores use this guard inside their atomic boundary commit: plain
551/// replacement is rejected, while an incoming snapshot carrying a typed rewrite
552/// commit from the currently persisted head is accepted through the same
553/// rewrite validator as [`SessionStore::save_transcript_rewrite`].
554pub fn run_boundary_snapshot_save_guard(
555    incoming: &Session,
556    previous: Option<&Session>,
557) -> Result<(), SessionStoreError> {
558    match append_only_save_guard(incoming, previous) {
559        Ok(()) => Ok(()),
560        Err(append_error) => {
561            if run_boundary_commitless_history_projection_save_guard(incoming, previous)? {
562                return Ok(());
563            }
564            let Some(previous) = previous else {
565                // First runtime-boundary commit for a session this authority
566                // has never snapshotted: adoption of a resumed/imported
567                // session. A typed rewrite graph carried in is audited by its
568                // own commits — validate every one against its retained
569                // bodies and require the graph head to match the incoming
570                // digest. Plain `SessionStore::save` keeps rejecting such
571                // seeds (the trait-level append-only contract); adoption is a
572                // runtime-authority decision, not an ordinary row write.
573                let incoming_revision = transcript_messages_digest(incoming.messages())
574                    .map_err(SessionStoreError::from)?;
575                if let Some(state) = incoming.transcript_history_state().map_err(|err| {
576                    SessionStoreError::InvalidTranscriptRewrite {
577                        id: incoming.id().clone(),
578                        reason: format!("incoming transcript history state is malformed: {err}"),
579                    }
580                })? && !state.commits.is_empty()
581                    && state.head == incoming_revision
582                {
583                    for commit in &state.commits {
584                        validate_transcript_rewrite_commit_bodies(incoming, commit, &state)?;
585                    }
586                    return Ok(());
587                }
588                return Err(append_error);
589            };
590            let incoming_revision =
591                transcript_messages_digest(incoming.messages()).map_err(SessionStoreError::from)?;
592            let Some(state) = incoming.transcript_history_state().map_err(|err| {
593                SessionStoreError::InvalidTranscriptRewrite {
594                    id: incoming.id().clone(),
595                    reason: format!("incoming transcript history state is malformed: {err}"),
596                }
597            })?
598            else {
599                return Err(append_error);
600            };
601            // append_only_save_guard's digest validation of the incoming
602            // history state was discarded with its error above; a
603            // digest-inconsistent witness body must not be able to prove a
604            // fork as a plain append on this branch either.
605            incoming
606                .validate_transcript_history_state()
607                .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
608                    id: incoming.id().clone(),
609                    reason: format!("incoming transcript history state is malformed: {err}"),
610                })?;
611            validate_rewrite_save_retains_previous_commits(incoming, previous, &state)?;
612            let commits = find_transcript_rewrite_commit_chain_extending_session(
613                &state,
614                previous,
615                &incoming_revision,
616            )?;
617            if commits.is_none()
618                && run_boundary_context_summary_tail_projection_save_guard(
619                    incoming, previous, &state,
620                )?
621            {
622                return Ok(());
623            }
624            let Some(commits) = commits else {
625                return Err(append_error);
626            };
627            let Some(commit) = commits.first() else {
628                if state.commits.is_empty() {
629                    return Err(append_error);
630                }
631                // Empty chain: the persisted row is already at (or past) the
632                // last rewrite and the incoming head extends it by plain
633                // appends. Unlike the non-empty chain below, no bridge guard
634                // runs here, so re-check the graph-head/message-digest
635                // agreement explicitly before accepting.
636                if state.head != incoming_revision {
637                    return Err(SessionStoreError::InvalidTranscriptRewrite {
638                        id: incoming.id().clone(),
639                        reason: format!(
640                            "incoming transcript graph head {} does not match current message digest {incoming_revision}",
641                            state.head
642                        ),
643                    });
644                }
645                for commit in &state.commits {
646                    validate_transcript_rewrite_commit_bodies(incoming, commit, &state)?;
647                }
648                return Ok(());
649            };
650            transcript_rewrite_bridge_save_guard(incoming, commit, &state, &incoming_revision)?;
651            // Validate every retained commit's recorded bodies, not only the
652            // walked chain: a plain-continuation proof can legitimately end
653            // the walk before trailing rebookkept commits, and those must
654            // stay digest-consistent to ride along (mirrors the empty-chain
655            // arm above).
656            for commit in &state.commits {
657                validate_transcript_rewrite_commit_bodies(incoming, commit, &state)?;
658            }
659            Ok(())
660        }
661    }
662}
663
664fn run_boundary_commitless_history_projection_save_guard(
665    incoming: &Session,
666    previous: Option<&Session>,
667) -> Result<bool, SessionStoreError> {
668    let Some(state) = incoming.transcript_history_state().map_err(|err| {
669        SessionStoreError::InvalidTranscriptRewrite {
670            id: incoming.id().clone(),
671            reason: format!("incoming transcript history state is malformed: {err}"),
672        }
673    })?
674    else {
675        return Ok(false);
676    };
677    if !state.commits.is_empty() {
678        return Ok(false);
679    }
680
681    let incoming_revision =
682        transcript_messages_digest(incoming.messages()).map_err(SessionStoreError::from)?;
683    if state.head != incoming_revision
684        || !state
685            .revisions
686            .iter()
687            .any(|body| body.revision == incoming_revision)
688    {
689        return Ok(false);
690    }
691
692    let mut projection_without_history = incoming.clone();
693    projection_without_history.clear_transcript_history_state();
694    if append_only_save_guard(&projection_without_history, previous).is_err() {
695        return Ok(false);
696    }
697
698    let Some(previous) = previous else {
699        return Ok(state.commits.is_empty());
700    };
701    if previous
702        .transcript_history_state()
703        .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
704            id: incoming.id().clone(),
705            reason: format!("previous transcript history state is malformed: {err}"),
706        })?
707        .is_some()
708    {
709        return Ok(false);
710    }
711
712    let previous_revision =
713        transcript_messages_digest(previous.messages()).map_err(SessionStoreError::from)?;
714    Ok(incoming_revision == previous_revision
715        || transcript_history_revision_extends(&state, &incoming_revision, &previous_revision))
716}
717
718fn run_boundary_context_summary_tail_projection_save_guard(
719    incoming: &Session,
720    previous: &Session,
721    state: &TranscriptHistoryState,
722) -> Result<bool, SessionStoreError> {
723    if state.commits.is_empty() {
724        return Ok(false);
725    }
726    incoming
727        .validate_transcript_history_state()
728        .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
729            id: incoming.id().clone(),
730            reason: format!("incoming transcript history state is malformed: {err}"),
731        })?;
732
733    let (incoming_system, incoming_tail) = match incoming.messages().split_first() {
734        Some((Message::System(system), tail)) => (Some(system), tail),
735        _ => (None, incoming.messages()),
736    };
737    let (previous_system, previous_tail) = match previous.messages().split_first() {
738        Some((Message::System(system), tail)) => (Some(system), tail),
739        _ => (None, previous.messages()),
740    };
741    if incoming_system.is_some() != previous_system.is_some()
742        || incoming_tail.len() <= previous_tail.len()
743    {
744        return Ok(false);
745    }
746    let Some(Message::User(summary)) = incoming_tail.first() else {
747        return Ok(false);
748    };
749    // Typed marker, not content classification: the runtime compaction producer
750    // stamps the rebuilt-transcript boundary message with the
751    // `CompactionSummary` transcript role. The save-guard admits the divergent
752    // rewrite parent only when that typed fact is present.
753    if !summary.transcript_role.is_compaction_summary() {
754        return Ok(false);
755    }
756
757    let retained_end = 1 + previous_tail.len();
758    let retained = &incoming_tail[1..retained_end];
759    let retained_revision =
760        transcript_messages_digest(retained).map_err(SessionStoreError::from)?;
761    let previous_revision =
762        transcript_messages_digest(previous_tail).map_err(SessionStoreError::from)?;
763    if retained_revision != previous_revision {
764        return Ok(false);
765    }
766
767    for commit in &state.commits {
768        validate_transcript_rewrite_commit_bodies(incoming, commit, state)?;
769    }
770    Ok(true)
771}
772
773/// Find the rewrite commit that authorizes replacing `previous_revision`,
774/// allowing the incoming head to extend the rewrite via normal append bodies.
775pub fn find_transcript_rewrite_commit_extending<'a>(
776    state: &'a TranscriptHistoryState,
777    previous_revision: &str,
778    incoming_revision: &str,
779) -> Option<&'a TranscriptRewriteCommit> {
780    find_transcript_rewrite_commit_chain_extending(state, previous_revision, incoming_revision)
781        .and_then(|commits| commits.into_iter().next())
782}
783
784/// Find the contiguous rewrite commits that connect `previous_revision` to the
785/// incoming head, allowing normal append bodies after the final rewrite.
786pub fn find_transcript_rewrite_commit_chain_extending<'a>(
787    state: &'a TranscriptHistoryState,
788    previous_revision: &str,
789    incoming_revision: &str,
790) -> Option<Vec<&'a TranscriptRewriteCommit>> {
791    let mut chain = Vec::new();
792    let mut cursor = previous_revision;
793    let mut visited = std::collections::BTreeSet::new();
794    loop {
795        if incoming_revision == cursor {
796            return Some(chain);
797        }
798        if !visited.insert(cursor.to_string()) {
799            return None;
800        }
801        let commit = state.commits.iter().find(|commit| {
802            (commit.parent_revision == cursor
803                || transcript_history_revision_extends(state, &commit.parent_revision, cursor))
804                && transcript_history_revision_extends(state, incoming_revision, &commit.revision)
805        });
806        let Some(commit) = commit else {
807            return transcript_history_revision_extends(state, incoming_revision, cursor)
808                .then_some(chain);
809        };
810        cursor = &commit.revision;
811        chain.push(commit);
812    }
813}
814
815/// Find a rewrite chain whose first parent may be an append-only continuation
816/// of a previously persisted snapshot.
817///
818/// Runtime-backed sessions can append messages in the runtime store before a
819/// core-owned compaction rewrite is checkpointed to the compatibility
820/// `SessionStore`. In that case the first rewrite commit's parent revision is
821/// not equal to the persisted row's digest, but its retained parent body proves
822/// a normal append path from that persisted row.
823pub fn find_transcript_rewrite_commit_chain_extending_session<'a>(
824    state: &'a TranscriptHistoryState,
825    previous: &Session,
826    incoming_revision: &str,
827) -> Result<Option<Vec<&'a TranscriptRewriteCommit>>, SessionStoreError> {
828    let previous_revision =
829        transcript_messages_digest(previous.messages()).map_err(SessionStoreError::from)?;
830    let mut chain = Vec::new();
831    let mut cursor = previous_revision.as_str();
832    let mut visited = std::collections::BTreeSet::new();
833    loop {
834        if incoming_revision == cursor {
835            return Ok(Some(chain));
836        }
837        if !visited.insert(cursor.to_string()) {
838            return Ok(None);
839        }
840
841        let Some(cursor_messages) = transcript_history_messages_for_revision(
842            state,
843            cursor,
844            &previous_revision,
845            previous.messages(),
846        ) else {
847            return Ok(None);
848        };
849
850        // Exact graph edges are authoritative: a commit recorded directly
851        // against this cursor advances the walk (and keeps that commit on
852        // the audited persistence chain). A commit whose revision is the
853        // cursor itself, or any revision this walk already visited, cannot
854        // make progress and is never selected.
855        let mut selected = None;
856        for commit in &state.commits {
857            if commit.revision == cursor || visited.contains(&commit.revision) {
858                continue;
859            }
860            if !transcript_history_revision_extends(state, incoming_revision, &commit.revision) {
861                continue;
862            }
863            if commit.parent_revision == cursor {
864                selected = Some(commit);
865                break;
866            }
867        }
868
869        // With no exact edge, a plain append continuation from this cursor
870        // completes the proof: the incoming transcript preserves the
871        // cursor's content and no further rewrite edge is needed. Proving
872        // this BEFORE the equivalence-based selection below is load-bearing:
873        // once the graph retains SEVERAL chained system-prompt-refresh
874        // commits, the refresh equivalence makes every retained refresh
875        // commit's parent body "extend" the cursor, so selection would walk
876        // an OLDER refresh commit forward onto the revision the cursor
877        // already reached and abort as a cycle — rejecting a valid append
878        // (chained resume refreshes with no turn in between, the idle mob
879        // member roster-drift shape).
880        if selected.is_none() {
881            if revision_body_preserves_append_continuation_prefix(
882                state,
883                incoming_revision,
884                cursor_messages,
885                cursor,
886                false,
887            )? {
888                return Ok(Some(chain));
889            }
890            // Only when neither an exact edge nor a plain continuation
891            // exists, fall back to the system-refresh equivalence: a refresh
892            // commit recorded against a rebookkept parent (the resume-time
893            // shape) whose parent body still extends the cursor.
894            for commit in &state.commits {
895                if commit.revision == cursor || visited.contains(&commit.revision) {
896                    continue;
897                }
898                if !transcript_history_revision_extends(state, incoming_revision, &commit.revision)
899                {
900                    continue;
901                }
902                if revision_body_preserves_append_continuation_prefix(
903                    state,
904                    &commit.parent_revision,
905                    cursor_messages,
906                    cursor,
907                    true,
908                )? {
909                    selected = Some(commit);
910                    break;
911                }
912            }
913        }
914
915        let Some(commit) = selected else {
916            return Ok(None);
917        };
918        cursor = &commit.revision;
919        chain.push(commit);
920    }
921}
922
923fn transcript_history_messages_for_revision<'a>(
924    state: &'a TranscriptHistoryState,
925    revision: &str,
926    previous_revision: &str,
927    previous_messages: &'a [Message],
928) -> Option<&'a [Message]> {
929    if revision == previous_revision {
930        return Some(previous_messages);
931    }
932    state
933        .revisions
934        .iter()
935        .find(|body| body.revision == revision)
936        .map(|body| body.messages.as_slice())
937}
938
939fn revision_body_preserves_append_continuation_prefix(
940    state: &TranscriptHistoryState,
941    revision: &str,
942    ancestor_messages: &[Message],
943    ancestor_revision: &str,
944    allow_leading_system_refresh: bool,
945) -> Result<bool, SessionStoreError> {
946    if revision == ancestor_revision {
947        return Ok(true);
948    }
949    let Some(body) = state
950        .revisions
951        .iter()
952        .find(|body| body.revision == revision)
953    else {
954        return Ok(false);
955    };
956    if body.messages.len() >= ancestor_messages.len() {
957        let prefix_revision = transcript_messages_digest(&body.messages[..ancestor_messages.len()])
958            .map_err(SessionStoreError::from)?;
959        if prefix_revision == ancestor_revision {
960            return Ok(true);
961        }
962    }
963    if messages_preserve_conversation_tail_with_system_context_append(
964        &body.messages,
965        ancestor_messages,
966    )? {
967        return Ok(true);
968    }
969    // The untyped leading-System-refresh equivalence bridges bookkeeping
970    // divergence between a persisted row and a rewrite commit's recorded
971    // PARENT body only. It must not prove the final plain-append
972    // continuation: that would admit an unaudited System replacement (a
973    // recorded refresh body with no typed commit) as an ordinary append.
974    Ok(allow_leading_system_refresh
975        && messages_preserve_tail_after_leading_system_refresh(&body.messages, ancestor_messages)?)
976}
977
978fn messages_preserve_tail_after_leading_system_refresh(
979    incoming: &[Message],
980    previous: &[Message],
981) -> Result<bool, SessionStoreError> {
982    let (Some(Message::System(_)), Some(Message::System(_))) = (incoming.first(), previous.first())
983    else {
984        return Ok(false);
985    };
986    if incoming.len() < previous.len() {
987        return Ok(false);
988    }
989    let previous_tail_len = previous.len().saturating_sub(1);
990    if previous_tail_len == 0 {
991        return Ok(true);
992    }
993    let previous_tail_revision =
994        transcript_messages_digest(&previous[1..]).map_err(SessionStoreError::from)?;
995    let incoming_tail = &incoming[1..];
996    if incoming_tail.len() < previous_tail_len {
997        return Ok(false);
998    }
999    let incoming_tail_prefix_revision =
1000        transcript_messages_digest(&incoming_tail[..previous_tail_len])
1001            .map_err(SessionStoreError::from)?;
1002    Ok(incoming_tail_prefix_revision == previous_tail_revision)
1003}
1004
1005fn transcript_history_revision_extends(
1006    state: &TranscriptHistoryState,
1007    descendant: &str,
1008    ancestor: &str,
1009) -> bool {
1010    if descendant == ancestor {
1011        return true;
1012    }
1013    let mut cursor = descendant;
1014    // Parent pointers are metadata, not digest-covered: bound the walk so a
1015    // crafted cyclic revision-parent chain fails closed instead of hanging.
1016    let mut visited = std::collections::BTreeSet::new();
1017    while let Some(body) = state.revisions.iter().find(|body| body.revision == cursor) {
1018        if !visited.insert(body.revision.clone()) {
1019            return false;
1020        }
1021        let Some(parent) = body.parent_revision.as_deref() else {
1022            return false;
1023        };
1024        if parent == ancestor {
1025            return true;
1026        }
1027        cursor = parent;
1028    }
1029    false
1030}
1031
1032fn transcript_rewrite_bridge_save_guard(
1033    incoming: &Session,
1034    commit: &TranscriptRewriteCommit,
1035    incoming_state: &TranscriptHistoryState,
1036    incoming_message_digest: &str,
1037) -> Result<(), SessionStoreError> {
1038    validate_transcript_rewrite_commit_bodies(incoming, commit, incoming_state)?;
1039    if incoming_state.head != incoming_message_digest {
1040        return Err(SessionStoreError::InvalidTranscriptRewrite {
1041            id: incoming.id().clone(),
1042            reason: format!(
1043                "incoming transcript graph head {} does not match current message digest {incoming_message_digest}",
1044                incoming_state.head
1045            ),
1046        });
1047    }
1048    if !transcript_history_revision_extends(
1049        incoming_state,
1050        incoming_message_digest,
1051        &commit.revision,
1052    ) {
1053        return Err(SessionStoreError::InvalidTranscriptRewrite {
1054            id: incoming.id().clone(),
1055            reason: format!(
1056                "incoming transcript head {incoming_message_digest} does not extend rewrite revision {}",
1057                commit.revision
1058            ),
1059        });
1060    }
1061    Ok(())
1062}
1063
1064/// Validate that a same-session shrink/replace save is backed by a typed
1065/// transcript rewrite commit.
1066pub fn transcript_rewrite_save_guard(
1067    incoming: &Session,
1068    previous: Option<&Session>,
1069    commit: &TranscriptRewriteCommit,
1070) -> Result<(), SessionStoreError> {
1071    incoming
1072        .validate_transcript_history_state()
1073        .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1074            id: incoming.id().clone(),
1075            reason: format!("incoming transcript history state is malformed: {err}"),
1076        })?;
1077    let Some(previous) = previous else {
1078        return Err(SessionStoreError::InvalidTranscriptRewrite {
1079            id: incoming.id().clone(),
1080            reason: "rewrite target has no previously persisted session".to_string(),
1081        });
1082    };
1083    if incoming.id() != previous.id() {
1084        return Err(SessionStoreError::InvalidTranscriptRewrite {
1085            id: incoming.id().clone(),
1086            reason: format!(
1087                "incoming session id {} differs from previous session id {}",
1088                incoming.id(),
1089                previous.id()
1090            ),
1091        });
1092    }
1093    let previous_revision = previous.transcript_revision().map_err(|err| {
1094        SessionStoreError::InvalidTranscriptRewrite {
1095            id: incoming.id().clone(),
1096            reason: format!("previous transcript revision is malformed: {err}"),
1097        }
1098    })?;
1099    if previous_revision != commit.parent_revision {
1100        return Err(SessionStoreError::TranscriptRevisionConflict {
1101            id: incoming.id().clone(),
1102            expected: commit.parent_revision.clone(),
1103            actual: previous_revision,
1104        });
1105    }
1106    let previous_message_digest =
1107        transcript_messages_digest(previous.messages()).map_err(|err| {
1108            SessionStoreError::InvalidTranscriptRewrite {
1109                id: incoming.id().clone(),
1110                reason: format!("previous current transcript is not digestible: {err}"),
1111            }
1112        })?;
1113    if previous_message_digest != commit.parent_revision {
1114        return Err(SessionStoreError::InvalidTranscriptRewrite {
1115            id: incoming.id().clone(),
1116            reason: format!(
1117                "previous current transcript digest {previous_message_digest} does not match commit parent {}",
1118                commit.parent_revision
1119            ),
1120        });
1121    }
1122    let incoming_revision = incoming.transcript_revision().map_err(|err| {
1123        SessionStoreError::InvalidTranscriptRewrite {
1124            id: incoming.id().clone(),
1125            reason: format!("incoming transcript revision is malformed: {err}"),
1126        }
1127    })?;
1128    if incoming_revision != commit.revision {
1129        return Err(SessionStoreError::InvalidTranscriptRewrite {
1130            id: incoming.id().clone(),
1131            reason: format!(
1132                "incoming transcript revision {incoming_revision} does not match commit revision {}",
1133                commit.revision
1134            ),
1135        });
1136    }
1137    let incoming_message_digest =
1138        transcript_messages_digest(incoming.messages()).map_err(|err| {
1139            SessionStoreError::InvalidTranscriptRewrite {
1140                id: incoming.id().clone(),
1141                reason: format!("incoming current transcript is not digestible: {err}"),
1142            }
1143        })?;
1144    if incoming_message_digest != commit.revision {
1145        return Err(SessionStoreError::InvalidTranscriptRewrite {
1146            id: incoming.id().clone(),
1147            reason: format!(
1148                "incoming current transcript digest {incoming_message_digest} does not match commit revision {}",
1149                commit.revision
1150            ),
1151        });
1152    }
1153    let Some(incoming_state) = incoming.transcript_history_state().map_err(|err| {
1154        SessionStoreError::InvalidTranscriptRewrite {
1155            id: incoming.id().clone(),
1156            reason: format!("incoming transcript history state is malformed: {err}"),
1157        }
1158    })?
1159    else {
1160        return Err(SessionStoreError::InvalidTranscriptRewrite {
1161            id: incoming.id().clone(),
1162            reason: "incoming rewrite did not persist a transcript revision graph".to_string(),
1163        });
1164    };
1165    validate_rewrite_save_retains_previous_commits(incoming, previous, &incoming_state)?;
1166    validate_transcript_rewrite_commit_bodies(incoming, commit, &incoming_state)
1167}
1168
1169fn validate_transcript_rewrite_commit_bodies(
1170    incoming: &Session,
1171    commit: &TranscriptRewriteCommit,
1172    incoming_state: &TranscriptHistoryState,
1173) -> Result<(), SessionStoreError> {
1174    if !incoming_state
1175        .commits
1176        .iter()
1177        .any(|persisted| persisted == commit)
1178    {
1179        return Err(SessionStoreError::InvalidTranscriptRewrite {
1180            id: incoming.id().clone(),
1181            reason: format!(
1182                "incoming rewrite did not persist the rewrite commit in the transcript graph (wanted {} -> {}, graph commits: {:?})",
1183                commit.parent_revision,
1184                commit.revision,
1185                incoming_state
1186                    .commits
1187                    .iter()
1188                    .map(|commit| (&commit.parent_revision, &commit.revision))
1189                    .collect::<Vec<_>>()
1190            ),
1191        });
1192    }
1193    let Some(parent_body) = incoming_state
1194        .revisions
1195        .iter()
1196        .find(|body| body.revision == commit.parent_revision)
1197    else {
1198        return Err(SessionStoreError::InvalidTranscriptRewrite {
1199            id: incoming.id().clone(),
1200            reason: format!(
1201                "incoming rewrite omitted parent revision body {}",
1202                commit.parent_revision
1203            ),
1204        });
1205    };
1206    let Some(revision_body) = incoming_state
1207        .revisions
1208        .iter()
1209        .find(|body| body.revision == commit.revision)
1210    else {
1211        return Err(SessionStoreError::InvalidTranscriptRewrite {
1212            id: incoming.id().clone(),
1213            reason: format!(
1214                "incoming rewrite omitted new revision body {}",
1215                commit.revision
1216            ),
1217        });
1218    };
1219    if parent_body.messages.len() != commit.messages_before
1220        || revision_body.messages.len() != commit.messages_after
1221    {
1222        return Err(SessionStoreError::InvalidTranscriptRewrite {
1223            id: incoming.id().clone(),
1224            reason: format!(
1225                "commit message counts {} -> {} do not match persisted rewrite {} -> {}",
1226                commit.messages_before,
1227                commit.messages_after,
1228                parent_body.messages.len(),
1229                revision_body.messages.len()
1230            ),
1231        });
1232    }
1233    let parent_body_revision =
1234        transcript_messages_digest(&parent_body.messages).map_err(|err| {
1235            SessionStoreError::InvalidTranscriptRewrite {
1236                id: incoming.id().clone(),
1237                reason: format!("parent revision body is not digestible: {err}"),
1238            }
1239        })?;
1240    if parent_body_revision != commit.parent_revision {
1241        return Err(SessionStoreError::InvalidTranscriptRewrite {
1242            id: incoming.id().clone(),
1243            reason: format!(
1244                "parent revision body digest {parent_body_revision} does not match commit parent {}",
1245                commit.parent_revision
1246            ),
1247        });
1248    }
1249    let (start, end) = match &commit.selection {
1250        TranscriptRewriteSelection::MessageRange { start, end } => (*start, *end),
1251    };
1252    if start > end || end > parent_body.messages.len() {
1253        return Err(SessionStoreError::InvalidTranscriptRewrite {
1254            id: incoming.id().clone(),
1255            reason: format!(
1256                "commit selection {start}..{end} is invalid for parent revision with {} messages",
1257                parent_body.messages.len()
1258            ),
1259        });
1260    }
1261    let original_span_digest = transcript_messages_digest(&parent_body.messages[start..end])
1262        .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1263            id: incoming.id().clone(),
1264            reason: format!("original span body is not digestible: {err}"),
1265        })?;
1266    if original_span_digest != commit.original_span_digest {
1267        return Err(SessionStoreError::InvalidTranscriptRewrite {
1268            id: incoming.id().clone(),
1269            reason: format!(
1270                "original span digest {original_span_digest} does not match commit digest {}",
1271                commit.original_span_digest
1272            ),
1273        });
1274    }
1275    let revision_body_digest =
1276        transcript_messages_digest(&revision_body.messages).map_err(|err| {
1277            SessionStoreError::InvalidTranscriptRewrite {
1278                id: incoming.id().clone(),
1279                reason: format!("new revision body is not digestible: {err}"),
1280            }
1281        })?;
1282    if revision_body_digest != commit.revision {
1283        return Err(SessionStoreError::InvalidTranscriptRewrite {
1284            id: incoming.id().clone(),
1285            reason: format!(
1286                "new revision body digest {revision_body_digest} does not match commit revision {}",
1287                commit.revision
1288            ),
1289        });
1290    }
1291    let removed_len = end - start;
1292    let retained_len = commit
1293        .messages_before
1294        .checked_sub(removed_len)
1295        .ok_or_else(|| SessionStoreError::InvalidTranscriptRewrite {
1296            id: incoming.id().clone(),
1297            reason: "commit removed more messages than it recorded before rewrite".to_string(),
1298        })?;
1299    let replacement_len = commit
1300        .messages_after
1301        .checked_sub(retained_len)
1302        .ok_or_else(|| SessionStoreError::InvalidTranscriptRewrite {
1303            id: incoming.id().clone(),
1304            reason: "commit message counts cannot describe a replacement span".to_string(),
1305        })?;
1306    let replacement_end = start.checked_add(replacement_len).ok_or_else(|| {
1307        SessionStoreError::InvalidTranscriptRewrite {
1308            id: incoming.id().clone(),
1309            reason: "replacement span end overflowed".to_string(),
1310        }
1311    })?;
1312    if replacement_end > revision_body.messages.len() {
1313        return Err(SessionStoreError::InvalidTranscriptRewrite {
1314            id: incoming.id().clone(),
1315            reason: format!(
1316                "replacement span {start}..{replacement_end} is invalid for revision with {} messages",
1317                revision_body.messages.len()
1318            ),
1319        });
1320    }
1321    let parent_prefix_digest =
1322        transcript_messages_digest(&parent_body.messages[..start]).map_err(|err| {
1323            SessionStoreError::InvalidTranscriptRewrite {
1324                id: incoming.id().clone(),
1325                reason: format!("parent prefix body is not digestible: {err}"),
1326            }
1327        })?;
1328    let revision_prefix_digest = transcript_messages_digest(&revision_body.messages[..start])
1329        .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1330            id: incoming.id().clone(),
1331            reason: format!("revision prefix body is not digestible: {err}"),
1332        })?;
1333    if parent_prefix_digest != revision_prefix_digest {
1334        return Err(SessionStoreError::InvalidTranscriptRewrite {
1335            id: incoming.id().clone(),
1336            reason: "rewrite revision changed messages before the selected span".to_string(),
1337        });
1338    }
1339    let parent_suffix_digest =
1340        transcript_messages_digest(&parent_body.messages[end..]).map_err(|err| {
1341            SessionStoreError::InvalidTranscriptRewrite {
1342                id: incoming.id().clone(),
1343                reason: format!("parent suffix body is not digestible: {err}"),
1344            }
1345        })?;
1346    let revision_suffix_digest =
1347        transcript_messages_digest(&revision_body.messages[replacement_end..]).map_err(|err| {
1348            SessionStoreError::InvalidTranscriptRewrite {
1349                id: incoming.id().clone(),
1350                reason: format!("revision suffix body is not digestible: {err}"),
1351            }
1352        })?;
1353    if parent_suffix_digest != revision_suffix_digest {
1354        return Err(SessionStoreError::InvalidTranscriptRewrite {
1355            id: incoming.id().clone(),
1356            reason: "rewrite revision changed messages after the selected span".to_string(),
1357        });
1358    }
1359    let replacement_digest = transcript_messages_digest(
1360        &revision_body.messages[start..replacement_end],
1361    )
1362    .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1363        id: incoming.id().clone(),
1364        reason: format!("replacement span body is not digestible: {err}"),
1365    })?;
1366    if replacement_digest != commit.replacement_digest {
1367        return Err(SessionStoreError::InvalidTranscriptRewrite {
1368            id: incoming.id().clone(),
1369            reason: format!(
1370                "replacement span digest {replacement_digest} does not match commit digest {}",
1371                commit.replacement_digest
1372            ),
1373        });
1374    }
1375    Ok(())
1376}
1377
1378impl From<serde_json::Error> for SessionStoreError {
1379    fn from(e: serde_json::Error) -> Self {
1380        Self::Serialization(e.to_string())
1381    }
1382}
1383
1384/// Abstraction over session storage backends.
1385///
1386/// All methods take `&self` — implementations must handle interior mutability.
1387/// Object-safe: consumed as `Arc<dyn SessionStore>` throughout the system.
1388///
1389/// # Append-only contract (F1 closure, wave-c C-H1)
1390///
1391/// The snapshot written by [`save`](Self::save) is a **projection of the
1392/// canonical event log** ([`crate::session_store`] doc: "snapshot =
1393/// projection"). Implementations that persist across calls MUST enforce
1394/// that the message vector stored for a given `SessionId` is monotonically
1395/// non-shrinking — a subsequent `save()` for the same id must not have a
1396/// smaller `messages().len()` than the previously persisted row.
1397///
1398/// Callers that need to produce a session with a shorter history must go
1399/// through [`Session::fork_at`], which rotates `SessionId` — a fork is a
1400/// new identity on a new event log, not a same-session truncation.
1401///
1402/// Backends are encouraged to assert this invariant in their `save`
1403/// implementation and return
1404/// [`SessionStoreError::MonotonicityViolation`] when a caller tries to
1405/// shrink a snapshot. The default implementations in `meerkat-store`
1406/// (`SqliteSessionStore`, `JsonlStore`, `MemoryStore`) all go through
1407/// the [`append_only_save_guard`] helper.
1408#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
1409#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
1410pub trait SessionStore: Send + Sync {
1411    /// Save a session (create or extend).
1412    ///
1413    /// Implementations MUST reject a save whose message history is
1414    /// shorter than the previously persisted row for the same `SessionId`
1415    /// — see the trait-level doc on the append-only contract.
1416    async fn save(&self, session: &Session) -> Result<(), SessionStoreError>;
1417
1418    /// Save a same-SessionId transcript rewrite.
1419    ///
1420    /// This is the only `SessionStore` path allowed to replace or shrink the
1421    /// current message projection. Implementations must validate `commit`
1422    /// against the previously persisted head before writing `session`.
1423    async fn save_transcript_rewrite(
1424        &self,
1425        session: &Session,
1426        commit: &TranscriptRewriteCommit,
1427    ) -> Result<(), SessionStoreError> {
1428        let _ = (session, commit);
1429        Err(SessionStoreError::Internal(
1430            "save_transcript_rewrite is not supported by this SessionStore".to_string(),
1431        ))
1432    }
1433
1434    /// Save a compatibility projection after a separate authority has already
1435    /// committed the session snapshot.
1436    ///
1437    /// This method is for runtime-backed services only: the runtime snapshot
1438    /// has already accepted the semantic mutation, and the `SessionStore` row is
1439    /// a rebuildable projection. Normal callers must use [`SessionStore::save`]
1440    /// or [`SessionStore::save_transcript_rewrite`] so the store boundary keeps
1441    /// enforcing append-only/CAS semantics.
1442    async fn save_authoritative_projection(
1443        &self,
1444        session: &Session,
1445    ) -> Result<(), SessionStoreError> {
1446        self.save(session).await
1447    }
1448
1449    /// Save an authoritative projection only if the persisted row is still the
1450    /// revision that the caller already validated.
1451    async fn save_authoritative_projection_if_current_revision(
1452        &self,
1453        session: &Session,
1454        expected_current_revision: Option<String>,
1455    ) -> Result<(), SessionStoreError> {
1456        let _ = (session, expected_current_revision);
1457        Err(SessionStoreError::Internal(
1458            "save_authoritative_projection_if_current_revision is not supported by this SessionStore"
1459                .to_string(),
1460        ))
1461    }
1462
1463    /// Load a session by ID.
1464    async fn load(&self, id: &SessionId) -> Result<Option<Session>, SessionStoreError>;
1465
1466    /// List sessions matching filter.
1467    async fn list(&self, filter: SessionFilter) -> Result<Vec<SessionMeta>, SessionStoreError>;
1468
1469    /// Delete a session.
1470    async fn delete(&self, id: &SessionId) -> Result<(), SessionStoreError>;
1471
1472    /// Delete a compatibility projection only if it is still the revision that
1473    /// the caller already validated as unsafe to expose.
1474    async fn delete_if_current_revision(
1475        &self,
1476        id: &SessionId,
1477        expected_current_revision: &str,
1478    ) -> Result<bool, SessionStoreError>;
1479
1480    /// Check if a session exists.
1481    async fn exists(&self, id: &SessionId) -> Result<bool, SessionStoreError> {
1482        Ok(self.load(id).await?.is_some())
1483    }
1484}
1485
1486#[cfg(test)]
1487mod tests {
1488    use super::*;
1489    use crate::types::{
1490        AssistantBlock, BlockAssistantMessage, StopReason, SystemMessage, SystemNoticeBlock,
1491        SystemNoticeKind, SystemNoticeMessage, UserMessage,
1492    };
1493
1494    /// A lagging persisted row must be walkable across chained refresh
1495    /// commits whose recorded parents were rebookkept (only the fuzzy
1496    /// refresh-equivalence edge can advance), without the walk re-selecting
1497    /// a refresh commit whose revision it already visited: at the later
1498    /// cursors, every OLDER refresh commit's parent body still "extends" the
1499    /// cursor under the refresh equivalence, and re-selecting one walks back
1500    /// onto visited territory and aborts as a cycle. Pins the
1501    /// visited-revision skip in both selection scans.
1502    #[test]
1503    #[allow(clippy::expect_used)]
1504    fn boundary_commit_walks_lagging_row_across_rebookkept_refresh_chain()
1505    -> Result<(), Box<dyn std::error::Error>> {
1506        let mut base = Session::new();
1507        base.push(Message::System(SystemMessage::new(
1508            "member prompt roster v1",
1509        )));
1510        base.push(Message::User(UserMessage::text(
1511            "the codeword is birch seventeen".to_string(),
1512        )));
1513        // The persisted row lags the whole rewrite graph (written before any
1514        // refresh boot, carrying no history state).
1515        let previous = base.clone();
1516        let v1 = base.transcript_revision()?;
1517
1518        // Three refresh boots chain commits onto the graph.
1519        let mut session = base;
1520        session.commit_transcript_rewrite(
1521            TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1522            vec![Message::System(SystemMessage::new(
1523                "member prompt roster v2",
1524            ))],
1525            crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
1526            Some("agent-factory/resume".to_string()),
1527            Some(v1),
1528        )?;
1529        let v2 = session.transcript_revision()?;
1530        session.commit_transcript_rewrite(
1531            TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1532            vec![Message::System(SystemMessage::new(
1533                "member prompt roster v3",
1534            ))],
1535            crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
1536            Some("agent-factory/resume".to_string()),
1537            Some(v2.clone()),
1538        )?;
1539        let v3 = session.transcript_revision()?;
1540        session.commit_transcript_rewrite(
1541            TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1542            vec![Message::System(SystemMessage::new(
1543                "member prompt roster v4",
1544            ))],
1545            crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
1546            Some("agent-factory/resume".to_string()),
1547            Some(v3.clone()),
1548        )?;
1549
1550        // Rebookkeep the recorded parents of the later refresh commits: each
1551        // now points at an equivalent parent body with re-stamped leading
1552        // System content (the re-created-authority shape), so no exact edge
1553        // exists from the walked cursors and only the refresh equivalence
1554        // can advance.
1555        let mut state = session
1556            .transcript_history_state()?
1557            .expect("chained refreshes retain history state");
1558        let rebookkeep = |state: &mut TranscriptHistoryState,
1559                          original_parent: &str,
1560                          stamp: &str|
1561         -> Result<String, Box<dyn std::error::Error>> {
1562            let body = state
1563                .revisions
1564                .iter()
1565                .find(|body| body.revision == original_parent)
1566                .expect("parent body retained")
1567                .clone();
1568            let mut messages = body.messages;
1569            messages[0] = Message::System(SystemMessage::new(stamp));
1570            let revision = transcript_messages_digest(&messages)?;
1571            // The rebookkept body chains off the revision it restamps, so
1572            // the graph stays a valid extension chain for the validator.
1573            state
1574                .revisions
1575                .push(crate::session::TranscriptRevisionBody {
1576                    revision: revision.clone(),
1577                    parent_revision: Some(original_parent.to_string()),
1578                    messages,
1579                    created_at: SystemTime::now(),
1580                });
1581            Ok(revision)
1582        };
1583        let v2_rebookkept = rebookkeep(&mut state, &v2, "member prompt roster v2 restamped")?;
1584        let v3_rebookkept = rebookkeep(&mut state, &v3, "member prompt roster v3 restamped")?;
1585        // Every refresh commit in this fixture rewrites message range 0..1,
1586        // so the rebookkept parent's recorded span is its leading System
1587        // message.
1588        let respan = |state: &TranscriptHistoryState,
1589                      parent: &str|
1590         -> Result<String, Box<dyn std::error::Error>> {
1591            let body = state
1592                .revisions
1593                .iter()
1594                .find(|body| body.revision == parent)
1595                .expect("rebookkept parent body retained");
1596            Ok(transcript_messages_digest(&body.messages[0..1])?)
1597        };
1598        state.commits[1].original_span_digest = respan(&state, &v2_rebookkept)?;
1599        state.commits[1].parent_revision = v2_rebookkept;
1600        state.commits[2].original_span_digest = respan(&state, &v3_rebookkept)?;
1601        state.commits[2].parent_revision = v3_rebookkept;
1602
1603        // The turn finally runs: two checkpointer-recorded appends.
1604        let mut incoming = session;
1605        incoming.set_metadata_unchecked_for_test(
1606            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1607            serde_json::to_value(&state)?,
1608        );
1609        let mut state = incoming
1610            .transcript_history_state()?
1611            .expect("history state survives rebookkeeping");
1612        for text in ["what was the codeword?", "birch seventeen"] {
1613            incoming.push(Message::User(UserMessage::text(text.to_string())));
1614            let appended_revision = incoming.transcript_revision()?;
1615            state
1616                .revisions
1617                .push(crate::session::TranscriptRevisionBody {
1618                    revision: appended_revision.clone(),
1619                    parent_revision: Some(state.head.clone()),
1620                    messages: incoming.messages().to_vec(),
1621                    created_at: SystemTime::now(),
1622                });
1623            state.head = appended_revision;
1624        }
1625        incoming.set_metadata_unchecked_for_test(
1626            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1627            serde_json::to_value(state)?,
1628        );
1629
1630        run_boundary_snapshot_save_guard(&incoming, Some(&previous))?;
1631        Ok(())
1632    }
1633
1634    /// Chained system-prompt-refresh commits with NO turn in between: the
1635    /// rewrite-chain walk must prove the plain append continuation from the
1636    /// persisted head instead of spuriously selecting an OLDER refresh
1637    /// commit (whose parent body also "extends" the head under the
1638    /// system-refresh equivalence), walking back onto its own cursor, and
1639    /// aborting as a cycle. Field regression (mobkit 0.7.23): idle mob
1640    /// members whose prompts carry drifting rosters get one refresh rewrite
1641    /// per boot; after two turn-less boots the next turn's run-boundary
1642    /// commit was rejected with "incoming append-only save would change
1643    /// retained transcript revision graph", permanently refusing resume.
1644    #[test]
1645    #[allow(clippy::expect_used)]
1646    fn boundary_commit_accepts_append_after_chained_promptless_system_refreshes()
1647    -> Result<(), Box<dyn std::error::Error>> {
1648        let mut base = Session::new();
1649        base.push(Message::System(SystemMessage::new(
1650            "member prompt roster v1",
1651        )));
1652        base.push(Message::User(UserMessage::text(
1653            "the codeword is birch seventeen".to_string(),
1654        )));
1655        let v1 = base.transcript_revision()?;
1656
1657        // Boot 1: resume refreshes the system prompt; the host dies before
1658        // any turn runs.
1659        let mut refreshed_once = base.clone();
1660        refreshed_once.commit_transcript_rewrite(
1661            TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1662            vec![Message::System(SystemMessage::new(
1663                "member prompt roster v2",
1664            ))],
1665            crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
1666            Some("agent-factory/resume".to_string()),
1667            Some(v1),
1668        )?;
1669        let v2 = refreshed_once.transcript_revision()?;
1670
1671        // Boot 2: another turn-less refresh chains onto the graph.
1672        let mut previous = refreshed_once.clone();
1673        previous.commit_transcript_rewrite(
1674            TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1675            vec![Message::System(SystemMessage::new(
1676                "member prompt roster v3",
1677            ))],
1678            crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
1679            Some("agent-factory/resume".to_string()),
1680            Some(v2),
1681        )?;
1682
1683        // Boot 3: the first turn finally runs. The intra-turn checkpointer
1684        // records one revision body per save, so the boundary commit's
1685        // incoming state carries MORE than one appended revision — the plain
1686        // +1 append validation cannot accept it and continuity must be
1687        // proven by the rewrite-chain walk.
1688        let mut incoming = previous.clone();
1689        let mut state = incoming
1690            .transcript_history_state()?
1691            .expect("chained refreshes retain history state");
1692        for text in ["what was the codeword?", "birch seventeen"] {
1693            incoming.push(Message::User(UserMessage::text(text.to_string())));
1694            let appended_revision = incoming.transcript_revision()?;
1695            state
1696                .revisions
1697                .push(crate::session::TranscriptRevisionBody {
1698                    revision: appended_revision.clone(),
1699                    parent_revision: Some(state.head.clone()),
1700                    messages: incoming.messages().to_vec(),
1701                    created_at: SystemTime::now(),
1702                });
1703            state.head = appended_revision;
1704        }
1705        incoming.set_metadata_unchecked_for_test(
1706            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1707            serde_json::to_value(state)?,
1708        );
1709
1710        run_boundary_snapshot_save_guard(&incoming, Some(&previous))?;
1711        Ok(())
1712    }
1713
1714    /// FOLD C: the canonical SessionDocumentMachine — not a handwritten shell
1715    /// boolean reducer — owns the live-vs-durable session-document authority
1716    /// verdict, the precedence (archived > uncommitted transcript > runtime
1717    /// system-context > stored transcript-revision), and the typed reason. This
1718    /// drives the classifier directly and asserts every authority/reason outcome
1719    /// and the precedence ordering.
1720    #[test]
1721    #[allow(clippy::expect_used)]
1722    fn classify_live_session_authority_is_decided_by_machine() {
1723        use crate::session_document::{
1724            LiveSessionAuthorityKind, LiveSessionAuthorityReason, SessionDocumentEffect,
1725            SessionDocumentMachineAuthority,
1726        };
1727
1728        fn classify(
1729            stored_transcript_diverged: bool,
1730            live_has_uncommitted_transcript: bool,
1731            runtime_system_context_diverged: bool,
1732            stored_is_archived: bool,
1733        ) -> (LiveSessionAuthorityKind, LiveSessionAuthorityReason) {
1734            let mut authority = SessionDocumentMachineAuthority::new();
1735            let effects = authority
1736                .classify_live_session_authority(
1737                    stored_transcript_diverged,
1738                    live_has_uncommitted_transcript,
1739                    runtime_system_context_diverged,
1740                    stored_is_archived,
1741                )
1742                .expect("classifier must resolve a verdict");
1743            effects
1744                .iter()
1745                .find_map(|effect| match effect {
1746                    SessionDocumentEffect::LiveSessionAuthorityClassified { authority, reason } => {
1747                        Some((*authority, *reason))
1748                    }
1749                    _ => None,
1750                })
1751                .expect("classifier must emit a verdict")
1752        }
1753
1754        // All four false -> LiveAuthoritative.
1755        let (kind, _) = classify(false, false, false, false);
1756        assert_eq!(kind, LiveSessionAuthorityKind::LiveAuthoritative);
1757
1758        // Each divergence (in isolation) -> DurableAuthoritative with its reason.
1759        assert_eq!(
1760            classify(true, false, false, false),
1761            (
1762                LiveSessionAuthorityKind::DurableAuthoritative,
1763                LiveSessionAuthorityReason::StoredTranscriptRevisionDiverged
1764            ),
1765        );
1766        assert_eq!(
1767            classify(false, true, false, false),
1768            (
1769                LiveSessionAuthorityKind::DurableAuthoritative,
1770                LiveSessionAuthorityReason::LiveUncommittedTranscript
1771            ),
1772        );
1773        assert_eq!(
1774            classify(false, false, true, false),
1775            (
1776                LiveSessionAuthorityKind::DurableAuthoritative,
1777                LiveSessionAuthorityReason::RuntimeSystemContextDiverged
1778            ),
1779        );
1780        assert_eq!(
1781            classify(false, false, false, true),
1782            (
1783                LiveSessionAuthorityKind::DurableAuthoritative,
1784                LiveSessionAuthorityReason::StoredArchived
1785            ),
1786        );
1787
1788        // Precedence: archived > uncommitted > system-context > revision.
1789        // When ALL four diverge, archived wins.
1790        assert_eq!(
1791            classify(true, true, true, true),
1792            (
1793                LiveSessionAuthorityKind::DurableAuthoritative,
1794                LiveSessionAuthorityReason::StoredArchived
1795            ),
1796        );
1797        // Not archived, but uncommitted + system-context + revision -> uncommitted.
1798        assert_eq!(
1799            classify(true, true, true, false),
1800            (
1801                LiveSessionAuthorityKind::DurableAuthoritative,
1802                LiveSessionAuthorityReason::LiveUncommittedTranscript
1803            ),
1804        );
1805        // Not archived, not uncommitted, but system-context + revision -> system-context.
1806        assert_eq!(
1807            classify(true, false, true, false),
1808            (
1809                LiveSessionAuthorityKind::DurableAuthoritative,
1810                LiveSessionAuthorityReason::RuntimeSystemContextDiverged
1811            ),
1812        );
1813    }
1814
1815    #[test]
1816    fn append_only_guard_rejects_leading_system_message_replacement() {
1817        let mut previous = Session::new();
1818        previous.push(Message::System(SystemMessage::new("original system")));
1819        previous.push(Message::User(UserMessage::text("hello".to_string())));
1820
1821        let mut incoming = previous.clone();
1822        let rewrite_result = incoming.replace_messages_internal(
1823            vec![
1824                Message::System(SystemMessage::new("rewritten system")),
1825                Message::User(UserMessage::text("hello".to_string())),
1826            ],
1827            crate::TranscriptRewriteReason::new("unit-test"),
1828        );
1829        assert!(
1830            rewrite_result.is_ok(),
1831            "typed rewrite should be constructible: {rewrite_result:?}"
1832        );
1833
1834        assert!(matches!(
1835            append_only_save_guard(&incoming, Some(&previous)),
1836            Err(SessionStoreError::TranscriptContinuityViolation { .. })
1837        ));
1838    }
1839
1840    #[test]
1841    fn append_only_guard_accepts_runtime_system_context_append()
1842    -> Result<(), Box<dyn std::error::Error>> {
1843        let mut previous = Session::new();
1844        previous.push(Message::System(SystemMessage::new("base system")));
1845        previous.push(Message::User(UserMessage::text("hello".to_string())));
1846
1847        let mut incoming = previous.clone();
1848        // The typed runtime-context-append producer stamps the system message's
1849        // mutation_kind so the save-guard admits the divergence from a typed
1850        // field, not the rendered `[Runtime System Context]` label.
1851        incoming.set_system_prompt_with_source(
1852            format!(
1853                "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nsource: unit-test\n\nextra context"
1854            ),
1855            crate::session_durable_config_authority::SessionSystemPromptSource::RuntimeContextAppend,
1856        )?;
1857
1858        assert!(append_only_save_guard(&incoming, Some(&previous)).is_ok());
1859        Ok(())
1860    }
1861
1862    #[test]
1863    fn append_only_guard_rejects_append_shaped_prompt_without_runtime_context_marker() {
1864        let mut previous = Session::new();
1865        previous.push(Message::System(SystemMessage::new("base system")));
1866        previous.push(Message::User(UserMessage::text("hello".to_string())));
1867
1868        // Same rendered shape as a runtime context append, but produced via a
1869        // direct mutation (mutation_kind != RuntimeContextAppend). The typed
1870        // gate must reject it — content prefix alone is not authority.
1871        let mut incoming = previous.clone();
1872        incoming.set_system_prompt(format!(
1873            "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nsource: forged\n\nextra context"
1874        ));
1875
1876        assert!(matches!(
1877            append_only_save_guard(&incoming, Some(&previous)),
1878            Err(SessionStoreError::TranscriptContinuityViolation { .. })
1879        ));
1880    }
1881
1882    #[test]
1883    fn append_only_guard_accepts_system_timestamp_refresh_without_content_change() {
1884        let mut previous = Session::new();
1885        previous.push(Message::System(SystemMessage::new("base system")));
1886
1887        let mut incoming = previous.clone();
1888        incoming.set_system_prompt("base system".to_string());
1889
1890        assert!(append_only_save_guard(&incoming, Some(&previous)).is_ok());
1891    }
1892
1893    #[test]
1894    fn run_boundary_guard_accepts_compaction_after_uncheckpointed_runtime_append()
1895    -> Result<(), Box<dyn std::error::Error>> {
1896        let mut previous = Session::new();
1897        previous.push(Message::System(SystemMessage::new("base system")));
1898        previous.push(Message::User(UserMessage::text("turn one".to_string())));
1899        previous.push(Message::BlockAssistant(BlockAssistantMessage {
1900            blocks: vec![AssistantBlock::Text {
1901                text: "answer one".to_string(),
1902                meta: None,
1903            }],
1904            stop_reason: StopReason::EndTurn,
1905            identity: crate::types::TranscriptMessageIdentity::default(),
1906            created_at: crate::types::message_timestamp_now(),
1907        }));
1908
1909        let mut parent = previous.clone();
1910        parent.set_system_prompt("refreshed runtime system projection".to_string());
1911        parent.push(Message::User(UserMessage::text(
1912            "runtime-only turn".to_string(),
1913        )));
1914        parent.push(Message::BlockAssistant(BlockAssistantMessage {
1915            blocks: vec![AssistantBlock::Text {
1916                text: "runtime-only answer".to_string(),
1917                meta: None,
1918            }],
1919            stop_reason: StopReason::EndTurn,
1920            identity: crate::types::TranscriptMessageIdentity::default(),
1921            created_at: crate::types::message_timestamp_now(),
1922        }));
1923        let parent_revision = parent.transcript_revision()?;
1924
1925        let mut incoming = parent.clone();
1926        let mut replacement = vec![
1927            parent.messages()[0].clone(),
1928            Message::User(UserMessage::compaction_summary(
1929                "[Context compacted] summary".to_string(),
1930            )),
1931        ];
1932        replacement.extend_from_slice(&parent.messages()[1..]);
1933        incoming.commit_transcript_rewrite(
1934            TranscriptRewriteSelection::MessageRange {
1935                start: 0,
1936                end: parent.messages().len(),
1937            },
1938            replacement,
1939            crate::TranscriptRewriteReason::new("compaction"),
1940            Some("meerkat-core".to_string()),
1941            Some(parent_revision),
1942        )?;
1943
1944        assert!(matches!(
1945            append_only_save_guard(&incoming, Some(&previous)),
1946            Err(SessionStoreError::TranscriptContinuityViolation { .. })
1947        ));
1948        assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
1949        Ok(())
1950    }
1951
1952    #[test]
1953    fn run_boundary_guard_accepts_compaction_with_retained_tail_window()
1954    -> Result<(), Box<dyn std::error::Error>> {
1955        let mut previous = Session::new();
1956        previous.push(Message::System(SystemMessage::new("base system")));
1957        previous.push(Message::User(UserMessage::text("turn one".to_string())));
1958        previous.push(Message::BlockAssistant(BlockAssistantMessage::new(
1959            vec![crate::types::AssistantBlock::Text {
1960                text: "answer one".to_string(),
1961                meta: None,
1962            }],
1963            StopReason::EndTurn,
1964        )));
1965
1966        let mut parent = previous.clone();
1967        parent.set_system_prompt("refreshed runtime system projection".to_string());
1968        parent.push(Message::SystemNotice(SystemNoticeMessage::new(
1969            SystemNoticeKind::Comms,
1970            "peer response queued",
1971        )));
1972        let parent_revision = parent.transcript_revision()?;
1973
1974        let mut incoming = parent.clone();
1975        let mut replacement = vec![
1976            parent.messages()[0].clone(),
1977            Message::User(UserMessage::compaction_summary(
1978                "[Context compacted] summary".to_string(),
1979            )),
1980        ];
1981        replacement.extend_from_slice(&parent.messages()[1..]);
1982        incoming.commit_transcript_rewrite(
1983            TranscriptRewriteSelection::MessageRange {
1984                start: 0,
1985                end: parent.messages().len(),
1986            },
1987            replacement,
1988            crate::TranscriptRewriteReason::new("compaction"),
1989            Some("meerkat-core".to_string()),
1990            Some(parent_revision),
1991        )?;
1992
1993        assert!(matches!(
1994            append_only_save_guard(&incoming, Some(&previous)),
1995            Err(SessionStoreError::TranscriptContinuityViolation { .. })
1996        ));
1997        assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
1998        Ok(())
1999    }
2000
2001    #[test]
2002    fn run_boundary_guard_rejects_commitless_history_parent_edge()
2003    -> Result<(), Box<dyn std::error::Error>> {
2004        let mut previous = Session::new();
2005        previous.push(Message::System(SystemMessage::new("base system")));
2006        previous.push(Message::User(UserMessage::text("turn one".to_string())));
2007        let previous_revision = previous.transcript_revision()?;
2008
2009        let mut incoming = previous.clone();
2010        incoming.set_system_prompt("forged replacement system".to_string());
2011        let incoming_revision = incoming.transcript_revision()?;
2012        let history = TranscriptHistoryState {
2013            head: incoming_revision.clone(),
2014            commits: Vec::new(),
2015            revisions: vec![
2016                crate::TranscriptRevisionBody {
2017                    revision: previous_revision,
2018                    parent_revision: None,
2019                    messages: previous.messages().to_vec(),
2020                    created_at: previous.updated_at(),
2021                },
2022                crate::TranscriptRevisionBody {
2023                    revision: incoming_revision,
2024                    parent_revision: Some(previous.transcript_revision()?),
2025                    messages: incoming.messages().to_vec(),
2026                    created_at: incoming.updated_at(),
2027                },
2028            ],
2029        };
2030        incoming.set_metadata_unchecked_for_test(
2031            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2032            serde_json::to_value(history)?,
2033        );
2034
2035        assert!(matches!(
2036            append_only_save_guard(&incoming, Some(&previous)),
2037            Err(SessionStoreError::TranscriptContinuityViolation { .. })
2038        ));
2039        assert!(matches!(
2040            run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2041            Err(SessionStoreError::TranscriptContinuityViolation { .. }
2042                | SessionStoreError::MonotonicityViolation { .. })
2043        ));
2044        Ok(())
2045    }
2046
2047    #[test]
2048    fn append_only_guard_rejects_history_head_that_does_not_match_current_messages()
2049    -> Result<(), Box<dyn std::error::Error>> {
2050        let mut previous = Session::new();
2051        previous.push(Message::User(UserMessage::text("persisted".to_string())));
2052
2053        let mut incoming = previous.clone();
2054        incoming.push(Message::User(UserMessage::text("append".to_string())));
2055        let poisoned_messages = vec![Message::User(UserMessage::text(
2056            "unrelated poisoned history".to_string(),
2057        ))];
2058        let poisoned_revision = transcript_messages_digest(&poisoned_messages)?;
2059        incoming.set_metadata_unchecked_for_test(
2060            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2061            serde_json::to_value(TranscriptHistoryState {
2062                head: poisoned_revision.clone(),
2063                commits: Vec::new(),
2064                revisions: vec![crate::TranscriptRevisionBody {
2065                    revision: poisoned_revision,
2066                    parent_revision: None,
2067                    messages: poisoned_messages,
2068                    created_at: incoming.updated_at(),
2069                }],
2070            })?,
2071        );
2072
2073        assert!(matches!(
2074            append_only_save_guard(&incoming, Some(&previous)),
2075            Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2076        ));
2077        assert!(matches!(
2078            append_only_save_guard(&incoming, None),
2079            Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2080        ));
2081        assert!(matches!(
2082            run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2083            Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2084        ));
2085        Ok(())
2086    }
2087
2088    #[test]
2089    fn append_only_guard_rejects_new_rewrite_commits_on_plain_append()
2090    -> Result<(), Box<dyn std::error::Error>> {
2091        let mut previous = Session::new();
2092        previous.push(Message::User(UserMessage::text("persisted".to_string())));
2093        let previous_revision = previous.transcript_revision()?;
2094
2095        let mut incoming = previous.clone();
2096        let appended = Message::BlockAssistant(BlockAssistantMessage {
2097            blocks: vec![AssistantBlock::Text {
2098                text: "plain append".to_string(),
2099                meta: None,
2100            }],
2101            stop_reason: StopReason::EndTurn,
2102            identity: crate::types::TranscriptMessageIdentity::default(),
2103            created_at: crate::types::message_timestamp_now(),
2104        });
2105        incoming.commit_transcript_rewrite(
2106            TranscriptRewriteSelection::MessageRange { start: 1, end: 1 },
2107            vec![appended],
2108            crate::TranscriptRewriteReason::new("forged-append"),
2109            Some("unit-test".to_string()),
2110            Some(previous_revision),
2111        )?;
2112
2113        assert!(matches!(
2114            append_only_save_guard(&incoming, Some(&previous)),
2115            Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2116        ));
2117        Ok(())
2118    }
2119
2120    #[test]
2121    fn append_only_guard_rejects_first_save_with_rewrite_commits()
2122    -> Result<(), Box<dyn std::error::Error>> {
2123        let mut incoming = Session::new();
2124        incoming.push(Message::User(UserMessage::text("seed".to_string())));
2125        let parent_messages = incoming.messages().to_vec();
2126        let parent_updated_at = incoming.updated_at();
2127        let parent_revision = incoming.transcript_revision()?;
2128        let commit = incoming.commit_transcript_rewrite(
2129            TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2130            vec![Message::User(UserMessage::text(
2131                "compacted seed".to_string(),
2132            ))],
2133            crate::TranscriptRewriteReason::new("compaction"),
2134            Some("meerkat-core".to_string()),
2135            Some(parent_revision),
2136        )?;
2137        let incoming_revision = incoming.transcript_revision()?;
2138        let commit_parent_revision = commit.parent_revision.clone();
2139        incoming.set_metadata_unchecked_for_test(
2140            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2141            serde_json::to_value(TranscriptHistoryState {
2142                head: incoming_revision.clone(),
2143                commits: vec![commit],
2144                revisions: vec![
2145                    crate::TranscriptRevisionBody {
2146                        revision: commit_parent_revision.clone(),
2147                        parent_revision: None,
2148                        messages: parent_messages,
2149                        created_at: parent_updated_at,
2150                    },
2151                    crate::TranscriptRevisionBody {
2152                        revision: incoming_revision,
2153                        parent_revision: Some(commit_parent_revision),
2154                        messages: incoming.messages().to_vec(),
2155                        created_at: incoming.updated_at(),
2156                    },
2157                ],
2158            })?,
2159        );
2160
2161        assert!(matches!(
2162            append_only_save_guard(&incoming, None),
2163            Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2164        ));
2165        Ok(())
2166    }
2167
2168    #[test]
2169    fn transcript_rewrite_guard_rejects_poisoned_history_graph()
2170    -> Result<(), Box<dyn std::error::Error>> {
2171        let mut previous = Session::new();
2172        previous.push(Message::User(UserMessage::text("persisted".to_string())));
2173        let parent_revision = previous.transcript_revision()?;
2174
2175        let mut first = previous.clone();
2176        let first_commit = first.commit_transcript_rewrite(
2177            TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2178            vec![Message::User(UserMessage::text(
2179                "compacted persisted".to_string(),
2180            ))],
2181            crate::TranscriptRewriteReason::new("compaction"),
2182            Some("unit-test".to_string()),
2183            Some(parent_revision),
2184        )?;
2185        let first_snapshot = first.clone();
2186
2187        first.commit_transcript_rewrite(
2188            TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2189            vec![Message::User(UserMessage::text(
2190                "uncommitted poisoned fork".to_string(),
2191            ))],
2192            crate::TranscriptRewriteReason::new("poison"),
2193            Some("unit-test".to_string()),
2194            Some(first_commit.revision.clone()),
2195        )?;
2196        let mut poisoned_state = first
2197            .transcript_history_state()?
2198            .ok_or_else(|| "second rewrite should retain history state".to_string())?;
2199        poisoned_state.head = first_commit.revision.clone();
2200
2201        let mut poisoned = first_snapshot;
2202        poisoned.set_metadata_unchecked_for_test(
2203            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2204            serde_json::to_value(poisoned_state)?,
2205        );
2206
2207        assert!(matches!(
2208            transcript_rewrite_save_guard(&poisoned, Some(&previous), &first_commit),
2209            Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
2210                if reason.contains("incoming transcript history state is malformed")
2211        ));
2212        Ok(())
2213    }
2214
2215    #[test]
2216    fn authoritative_projection_guard_rejects_changed_persisted_revision()
2217    -> Result<(), Box<dyn std::error::Error>> {
2218        let mut previous = Session::new();
2219        previous.push(Message::User(UserMessage::text("persisted A".to_string())));
2220        let expected_revision = previous.transcript_revision()?;
2221
2222        let mut current = previous.clone();
2223        current.push(Message::BlockAssistant(BlockAssistantMessage {
2224            blocks: vec![AssistantBlock::Text {
2225                text: "persisted B".to_string(),
2226                meta: None,
2227            }],
2228            stop_reason: StopReason::EndTurn,
2229            identity: crate::types::TranscriptMessageIdentity::default(),
2230            created_at: crate::types::message_timestamp_now(),
2231        }));
2232        let mut incoming = previous.clone();
2233        incoming.push(Message::User(UserMessage::text(
2234            "incoming from A".to_string(),
2235        )));
2236
2237        assert!(matches!(
2238            authoritative_projection_current_revision_guard(
2239                &incoming,
2240                Some(&current),
2241                Some(&expected_revision)
2242            ),
2243            Err(SessionStoreError::TranscriptContinuityViolation { .. })
2244        ));
2245        Ok(())
2246    }
2247
2248    #[test]
2249    fn append_only_guard_rejects_rewrite_commits_on_first_save()
2250    -> Result<(), Box<dyn std::error::Error>> {
2251        let mut incoming = Session::new();
2252        incoming.push(Message::User(UserMessage::text("persisted".to_string())));
2253        let parent_revision = incoming.transcript_revision()?;
2254        incoming.commit_transcript_rewrite(
2255            TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2256            vec![Message::User(UserMessage::text("rewritten".to_string()))],
2257            crate::TranscriptRewriteReason::new("forged-first-save"),
2258            Some("unit-test".to_string()),
2259            Some(parent_revision),
2260        )?;
2261
2262        assert!(matches!(
2263            append_only_save_guard(&incoming, None),
2264            Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2265        ));
2266        Ok(())
2267    }
2268
2269    #[test]
2270    fn append_only_guard_rejects_commitless_history_on_first_save()
2271    -> Result<(), Box<dyn std::error::Error>> {
2272        let mut incoming = Session::new();
2273        incoming.push(Message::User(UserMessage::text("persisted".to_string())));
2274        let incoming_revision = incoming.transcript_revision()?;
2275        incoming.set_metadata_unchecked_for_test(
2276            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2277            serde_json::to_value(TranscriptHistoryState {
2278                head: incoming_revision.clone(),
2279                commits: Vec::new(),
2280                revisions: vec![crate::TranscriptRevisionBody {
2281                    revision: incoming_revision,
2282                    parent_revision: None,
2283                    messages: incoming.messages().to_vec(),
2284                    created_at: incoming.updated_at(),
2285                }],
2286            })?,
2287        );
2288
2289        assert!(matches!(
2290            append_only_save_guard(&incoming, None),
2291            Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
2292                if reason.contains("first save would seed transcript history state")
2293        ));
2294        Ok(())
2295    }
2296
2297    #[test]
2298    fn run_boundary_guard_adopts_commit_carrying_history_on_first_commit()
2299    -> Result<(), Box<dyn std::error::Error>> {
2300        // A runtime authority adopting a session it never snapshotted
2301        // (resume/import over fresh runtime state) may receive a session that
2302        // already carries a typed rewrite graph — e.g. a resume-time
2303        // base-prompt refresh. The commits are the audit: the run-boundary
2304        // guard accepts the validated graph, while the plain trait-level
2305        // `SessionStore::save` contract keeps rejecting first-save seeds.
2306        let mut incoming = Session::new();
2307        incoming.set_system_prompt("old base".to_string());
2308        incoming.push(Message::User(UserMessage::text("hello".to_string())));
2309        incoming.commit_transcript_rewrite(
2310            crate::TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2311            vec![Message::System(SystemMessage::with_mutation_kind(
2312                "new base".to_string(),
2313                crate::types::SystemPromptMutationKind::ExplicitBuild,
2314            ))],
2315            crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
2316            None,
2317            None,
2318        )?;
2319
2320        assert!(run_boundary_snapshot_save_guard(&incoming, None).is_ok());
2321        assert!(matches!(
2322            append_only_save_guard(&incoming, None),
2323            Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2324        ));
2325        Ok(())
2326    }
2327
2328    #[test]
2329    fn run_boundary_guard_rejects_untyped_leading_system_refresh_after_head_rewrite()
2330    -> Result<(), Box<dyn std::error::Error>> {
2331        // A same-length rewrite commit sitting exactly at the persisted head
2332        // (the resume-refresh shape) must not widen acceptance to UNTYPED
2333        // leading-System replacements: a plain set_system_prompt records a
2334        // refresh body via the head refresh but carries no commit, and the
2335        // chain walker's fallback must not admit it as an ordinary append
2336        // continuation via the leading-system-refresh equivalence.
2337        let mut previous = Session::new();
2338        previous.set_system_prompt("original base".to_string());
2339        previous.push(Message::User(UserMessage::text("hello".to_string())));
2340        previous.commit_transcript_rewrite(
2341            crate::TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2342            vec![Message::System(SystemMessage::with_mutation_kind(
2343                "refreshed base".to_string(),
2344                crate::types::SystemPromptMutationKind::ExplicitBuild,
2345            ))],
2346            crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
2347            None,
2348            None,
2349        )?;
2350
2351        let mut incoming = previous.clone();
2352        incoming.set_system_prompt("untyped hijack".to_string());
2353
2354        assert!(matches!(
2355            run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2356            Err(SessionStoreError::TranscriptContinuityViolation { .. })
2357        ));
2358        Ok(())
2359    }
2360
2361    #[test]
2362    fn run_boundary_guard_accepts_plain_append_after_head_rewrite()
2363    -> Result<(), Box<dyn std::error::Error>> {
2364        // The empty-chain acceptance the cycle-skip exists for: the persisted
2365        // row is already AT the rewrite revision and the incoming snapshot
2366        // extends it by ordinary appends.
2367        let mut previous = Session::new();
2368        previous.set_system_prompt("original base".to_string());
2369        previous.push(Message::User(UserMessage::text("hello".to_string())));
2370        previous.commit_transcript_rewrite(
2371            crate::TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2372            vec![Message::System(SystemMessage::with_mutation_kind(
2373                "refreshed base".to_string(),
2374                crate::types::SystemPromptMutationKind::ExplicitBuild,
2375            ))],
2376            crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
2377            None,
2378            None,
2379        )?;
2380
2381        let mut incoming = previous.clone();
2382        incoming.push(Message::User(UserMessage::text(
2383            "post-rewrite turn".to_string(),
2384        )));
2385
2386        assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2387        Ok(())
2388    }
2389
2390    #[test]
2391    fn append_only_guard_rejects_commitless_history_seed_on_plain_append()
2392    -> Result<(), Box<dyn std::error::Error>> {
2393        let mut previous = Session::new();
2394        previous.push(Message::User(UserMessage::text("persisted".to_string())));
2395        let previous_revision = previous.transcript_revision()?;
2396
2397        let mut incoming = previous.clone();
2398        incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2399            blocks: vec![AssistantBlock::Text {
2400                text: "plain append".to_string(),
2401                meta: None,
2402            }],
2403            stop_reason: StopReason::EndTurn,
2404            identity: crate::types::TranscriptMessageIdentity::default(),
2405            created_at: crate::types::message_timestamp_now(),
2406        }));
2407        let incoming_revision = incoming.transcript_revision()?;
2408        incoming.set_metadata_unchecked_for_test(
2409            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2410            serde_json::to_value(TranscriptHistoryState {
2411                head: incoming_revision.clone(),
2412                commits: Vec::new(),
2413                revisions: vec![
2414                    crate::TranscriptRevisionBody {
2415                        revision: previous_revision,
2416                        parent_revision: None,
2417                        messages: previous.messages().to_vec(),
2418                        created_at: previous.updated_at(),
2419                    },
2420                    crate::TranscriptRevisionBody {
2421                        revision: incoming_revision,
2422                        parent_revision: Some(previous.transcript_revision()?),
2423                        messages: incoming.messages().to_vec(),
2424                        created_at: incoming.updated_at(),
2425                    },
2426                ],
2427            })?,
2428        );
2429
2430        assert!(matches!(
2431            append_only_save_guard(&incoming, Some(&previous)),
2432            Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
2433                if reason.contains("append-only save would seed transcript history state")
2434        ));
2435        Ok(())
2436    }
2437
2438    #[test]
2439    fn run_boundary_guard_accepts_commitless_history_seed_on_plain_append()
2440    -> Result<(), Box<dyn std::error::Error>> {
2441        let mut previous = Session::new();
2442        previous.push(Message::User(UserMessage::text("persisted".to_string())));
2443        let previous_revision = previous.transcript_revision()?;
2444
2445        let mut incoming = previous.clone();
2446        incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2447            blocks: vec![AssistantBlock::Text {
2448                text: "plain append".to_string(),
2449                meta: None,
2450            }],
2451            stop_reason: StopReason::EndTurn,
2452            identity: crate::types::TranscriptMessageIdentity::default(),
2453            created_at: crate::types::message_timestamp_now(),
2454        }));
2455        let incoming_revision = incoming.transcript_revision()?;
2456        incoming.set_metadata_unchecked_for_test(
2457            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2458            serde_json::to_value(TranscriptHistoryState {
2459                head: incoming_revision.clone(),
2460                commits: Vec::new(),
2461                revisions: vec![
2462                    crate::TranscriptRevisionBody {
2463                        revision: previous_revision.clone(),
2464                        parent_revision: None,
2465                        messages: previous.messages().to_vec(),
2466                        created_at: previous.updated_at(),
2467                    },
2468                    crate::TranscriptRevisionBody {
2469                        revision: incoming_revision,
2470                        parent_revision: Some(previous_revision),
2471                        messages: incoming.messages().to_vec(),
2472                        created_at: incoming.updated_at(),
2473                    },
2474                ],
2475            })?,
2476        );
2477
2478        assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2479        assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2480        Ok(())
2481    }
2482
2483    #[test]
2484    fn run_boundary_guard_accepts_retained_history_seed_on_plain_append()
2485    -> Result<(), Box<dyn std::error::Error>> {
2486        let mut original = Session::new();
2487        original.push(Message::User(UserMessage::text("verbose seed".to_string())));
2488        let original_revision = original.transcript_revision()?;
2489
2490        let mut previous = original.clone();
2491        previous.commit_transcript_rewrite(
2492            TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2493            vec![Message::User(UserMessage::text(
2494                "compacted seed".to_string(),
2495            ))],
2496            crate::TranscriptRewriteReason::new("compaction"),
2497            Some("meerkat-core".to_string()),
2498            Some(original_revision),
2499        )?;
2500        let previous_with_history = previous.clone();
2501        previous.clear_transcript_history_state();
2502
2503        let mut incoming = previous_with_history;
2504        incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2505            blocks: vec![AssistantBlock::Text {
2506                text: "plain append after retained history".to_string(),
2507                meta: None,
2508            }],
2509            stop_reason: StopReason::EndTurn,
2510            identity: crate::types::TranscriptMessageIdentity::default(),
2511            created_at: crate::types::message_timestamp_now(),
2512        }));
2513
2514        assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2515        assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2516        Ok(())
2517    }
2518
2519    #[test]
2520    fn run_boundary_guard_accepts_commitless_history_seed_on_first_snapshot()
2521    -> Result<(), Box<dyn std::error::Error>> {
2522        let mut incoming = Session::new();
2523        incoming.push(Message::User(UserMessage::text("persisted".to_string())));
2524        let incoming_revision = incoming.transcript_revision()?;
2525        incoming.set_metadata_unchecked_for_test(
2526            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2527            serde_json::to_value(TranscriptHistoryState {
2528                head: incoming_revision.clone(),
2529                commits: Vec::new(),
2530                revisions: vec![crate::TranscriptRevisionBody {
2531                    revision: incoming_revision,
2532                    parent_revision: None,
2533                    messages: incoming.messages().to_vec(),
2534                    created_at: incoming.updated_at(),
2535                }],
2536            })?,
2537        );
2538
2539        assert!(append_only_save_guard(&incoming, None).is_err());
2540        assert!(run_boundary_snapshot_save_guard(&incoming, None).is_ok());
2541        Ok(())
2542    }
2543
2544    #[test]
2545    fn run_boundary_guard_accepts_commitless_history_seed_on_initial_multi_revision_snapshot()
2546    -> Result<(), Box<dyn std::error::Error>> {
2547        let mut base = Session::new();
2548        base.push(Message::User(UserMessage::text("first".to_string())));
2549        let base_revision = base.transcript_revision()?;
2550
2551        let mut incoming = base.clone();
2552        incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2553            blocks: vec![AssistantBlock::Text {
2554                text: "second".to_string(),
2555                meta: None,
2556            }],
2557            stop_reason: StopReason::EndTurn,
2558            identity: crate::types::TranscriptMessageIdentity::default(),
2559            created_at: crate::types::message_timestamp_now(),
2560        }));
2561        let incoming_revision = incoming.transcript_revision()?;
2562        incoming.set_metadata_unchecked_for_test(
2563            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2564            serde_json::to_value(TranscriptHistoryState {
2565                head: incoming_revision.clone(),
2566                commits: Vec::new(),
2567                revisions: vec![
2568                    crate::TranscriptRevisionBody {
2569                        revision: base_revision.clone(),
2570                        parent_revision: None,
2571                        messages: base.messages().to_vec(),
2572                        created_at: base.updated_at(),
2573                    },
2574                    crate::TranscriptRevisionBody {
2575                        revision: incoming_revision,
2576                        parent_revision: Some(base_revision),
2577                        messages: incoming.messages().to_vec(),
2578                        created_at: incoming.updated_at(),
2579                    },
2580                ],
2581            })?,
2582        );
2583
2584        assert!(append_only_save_guard(&incoming, None).is_err());
2585        assert!(run_boundary_snapshot_save_guard(&incoming, None).is_ok());
2586        Ok(())
2587    }
2588
2589    #[test]
2590    fn append_only_guard_rejects_new_rewrite_commits_on_system_context_append()
2591    -> Result<(), Box<dyn std::error::Error>> {
2592        let mut previous = Session::new();
2593        previous.push(Message::System(SystemMessage::new("base system")));
2594        previous.push(Message::User(UserMessage::text("persisted".to_string())));
2595        let mut incoming = previous.clone();
2596        incoming.set_system_prompt_with_source(
2597            format!(
2598                "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nsource: unit-test\n\nextra context"
2599            ),
2600            crate::session_durable_config_authority::SessionSystemPromptSource::RuntimeContextAppend,
2601        )?;
2602        incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2603            blocks: vec![AssistantBlock::Text {
2604                text: "plain append".to_string(),
2605                meta: None,
2606            }],
2607            stop_reason: StopReason::EndTurn,
2608            identity: crate::types::TranscriptMessageIdentity::default(),
2609            created_at: crate::types::message_timestamp_now(),
2610        }));
2611        let incoming_revision = incoming.transcript_revision()?;
2612        incoming.set_metadata_unchecked_for_test(
2613            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2614            serde_json::to_value(TranscriptHistoryState {
2615                head: incoming_revision.clone(),
2616                commits: vec![TranscriptRewriteCommit {
2617                    parent_revision: previous.transcript_revision()?,
2618                    revision: incoming_revision.clone(),
2619                    selection: TranscriptRewriteSelection::MessageRange { start: 0, end: 0 },
2620                    original_span_digest: transcript_messages_digest(&[])?,
2621                    replacement_digest: transcript_messages_digest(&[])?,
2622                    messages_before: previous.messages().len(),
2623                    messages_after: incoming.messages().len(),
2624                    reason: crate::TranscriptRewriteReason::new("forged"),
2625                    actor: Some("unit-test".to_string()),
2626                    committed_at: incoming.updated_at(),
2627                }],
2628                revisions: vec![crate::TranscriptRevisionBody {
2629                    revision: incoming_revision,
2630                    parent_revision: None,
2631                    messages: incoming.messages().to_vec(),
2632                    created_at: incoming.updated_at(),
2633                }],
2634            })?,
2635        );
2636
2637        assert!(matches!(
2638            append_only_save_guard(&incoming, Some(&previous)),
2639            Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2640        ));
2641        Ok(())
2642    }
2643
2644    #[test]
2645    fn append_only_guard_rejects_new_rewrite_commits_on_transient_notice_cleanup()
2646    -> Result<(), Box<dyn std::error::Error>> {
2647        let mut previous = Session::new();
2648        previous.push(Message::SystemNotice(SystemNoticeMessage::new(
2649            SystemNoticeKind::Comms,
2650            "transient peer delivery notice",
2651        )));
2652        previous.push(Message::User(UserMessage::text("persisted".to_string())));
2653
2654        let mut incoming = Session::new();
2655        incoming.push(Message::User(UserMessage::text("persisted".to_string())));
2656        incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2657            blocks: vec![AssistantBlock::Text {
2658                text: "plain append after notice cleanup".to_string(),
2659                meta: None,
2660            }],
2661            stop_reason: StopReason::EndTurn,
2662            identity: crate::types::TranscriptMessageIdentity::default(),
2663            created_at: crate::types::message_timestamp_now(),
2664        }));
2665        let incoming_revision = incoming.transcript_revision()?;
2666        incoming.set_metadata_unchecked_for_test(
2667            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2668            serde_json::to_value(TranscriptHistoryState {
2669                head: incoming_revision.clone(),
2670                commits: vec![TranscriptRewriteCommit {
2671                    parent_revision: previous.transcript_revision()?,
2672                    revision: incoming_revision.clone(),
2673                    selection: TranscriptRewriteSelection::MessageRange { start: 0, end: 0 },
2674                    original_span_digest: transcript_messages_digest(&[])?,
2675                    replacement_digest: transcript_messages_digest(&[])?,
2676                    messages_before: previous.messages().len(),
2677                    messages_after: incoming.messages().len(),
2678                    reason: crate::TranscriptRewriteReason::new("forged"),
2679                    actor: Some("unit-test".to_string()),
2680                    committed_at: incoming.updated_at(),
2681                }],
2682                revisions: vec![crate::TranscriptRevisionBody {
2683                    revision: incoming_revision,
2684                    parent_revision: None,
2685                    messages: incoming.messages().to_vec(),
2686                    created_at: incoming.updated_at(),
2687                }],
2688            })?,
2689        );
2690
2691        assert!(matches!(
2692            append_only_save_guard(&incoming, Some(&previous)),
2693            Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2694        ));
2695        Ok(())
2696    }
2697
2698    #[test]
2699    fn run_boundary_guard_accepts_generated_context_summary_before_retained_tail()
2700    -> Result<(), Box<dyn std::error::Error>> {
2701        let mut previous = Session::new();
2702        previous.push(Message::System(SystemMessage::new(
2703            "runtime system before context refresh",
2704        )));
2705        previous.push(Message::User(UserMessage::text(
2706            "Turn 1 request".to_string(),
2707        )));
2708        previous.push(Message::BlockAssistant(BlockAssistantMessage {
2709            blocks: vec![AssistantBlock::Text {
2710                text: "Turn 1 answer".to_string(),
2711                meta: None,
2712            }],
2713            stop_reason: StopReason::EndTurn,
2714            identity: crate::types::TranscriptMessageIdentity::default(),
2715            created_at: crate::types::message_timestamp_now(),
2716        }));
2717
2718        let mut incoming = Session::with_id(previous.id().clone());
2719        incoming.push(Message::System(SystemMessage::new(
2720            "runtime system after context refresh",
2721        )));
2722        incoming.push(Message::User(UserMessage::text(
2723            "Verbose context that will be compacted".to_string(),
2724        )));
2725        for message in previous.messages()[1..].iter().cloned() {
2726            incoming.push(message);
2727        }
2728        incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2729            blocks: vec![AssistantBlock::Text {
2730                text: "Turn 2 generated answer".to_string(),
2731                meta: None,
2732            }],
2733            stop_reason: StopReason::EndTurn,
2734            identity: crate::types::TranscriptMessageIdentity::default(),
2735            created_at: crate::types::message_timestamp_now(),
2736        }));
2737        let parent_revision = incoming.transcript_revision()?;
2738        incoming.commit_transcript_rewrite(
2739            TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2740            vec![Message::User(UserMessage::compaction_summary(
2741                "[Context compacted] Earlier runtime context".to_string(),
2742            ))],
2743            crate::TranscriptRewriteReason::new("compaction"),
2744            Some("meerkat-core".to_string()),
2745            Some(parent_revision),
2746        )?;
2747
2748        assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2749        assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2750        Ok(())
2751    }
2752
2753    #[test]
2754    fn run_boundary_guard_rejects_context_summary_tail_without_compaction_summary_marker()
2755    -> Result<(), Box<dyn std::error::Error>> {
2756        let mut previous = Session::new();
2757        previous.push(Message::System(SystemMessage::new(
2758            "runtime system before context refresh",
2759        )));
2760        previous.push(Message::User(UserMessage::text(
2761            "Turn 1 request".to_string(),
2762        )));
2763        previous.push(Message::BlockAssistant(BlockAssistantMessage {
2764            blocks: vec![AssistantBlock::Text {
2765                text: "Turn 1 answer".to_string(),
2766                meta: None,
2767            }],
2768            stop_reason: StopReason::EndTurn,
2769            identity: crate::types::TranscriptMessageIdentity::default(),
2770            created_at: crate::types::message_timestamp_now(),
2771        }));
2772
2773        let mut incoming = Session::with_id(previous.id().clone());
2774        incoming.push(Message::System(SystemMessage::new(
2775            "runtime system after context refresh",
2776        )));
2777        incoming.push(Message::User(UserMessage::text(
2778            "Verbose context that will be compacted".to_string(),
2779        )));
2780        for message in previous.messages()[1..].iter().cloned() {
2781            incoming.push(message);
2782        }
2783        incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2784            blocks: vec![AssistantBlock::Text {
2785                text: "Turn 2 generated answer".to_string(),
2786                meta: None,
2787            }],
2788            stop_reason: StopReason::EndTurn,
2789            identity: crate::types::TranscriptMessageIdentity::default(),
2790            created_at: crate::types::message_timestamp_now(),
2791        }));
2792        let parent_revision = incoming.transcript_revision()?;
2793        // Same rendered shape (content begins with `[Context compacted]`) but the
2794        // summary message uses the ordinary conversational role. The typed gate
2795        // must reject it: rendered content alone is not authority.
2796        incoming.commit_transcript_rewrite(
2797            TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2798            vec![Message::User(UserMessage::text(
2799                "[Context compacted] Earlier runtime context".to_string(),
2800            ))],
2801            crate::TranscriptRewriteReason::new("compaction"),
2802            Some("meerkat-core".to_string()),
2803            Some(parent_revision),
2804        )?;
2805
2806        assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2807        assert!(matches!(
2808            run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2809            Err(SessionStoreError::TranscriptContinuityViolation { .. }
2810                | SessionStoreError::MonotonicityViolation { .. })
2811        ));
2812        Ok(())
2813    }
2814
2815    /// Ask 1 save-guard invariant: the injected-context transcript role must
2816    /// NOT satisfy the transcript-continuity save-guard. Only the
2817    /// runtime-minted `CompactionSummary` role admits a divergent rewrite
2818    /// parent (`is_compaction_summary()` stays `CompactionSummary`-only).
2819    #[test]
2820    fn run_boundary_guard_rejects_context_summary_tail_with_injected_context_marker()
2821    -> Result<(), Box<dyn std::error::Error>> {
2822        let mut previous = Session::new();
2823        previous.push(Message::System(SystemMessage::new(
2824            "runtime system before context refresh",
2825        )));
2826        previous.push(Message::User(UserMessage::text(
2827            "Turn 1 request".to_string(),
2828        )));
2829        previous.push(Message::BlockAssistant(BlockAssistantMessage {
2830            blocks: vec![AssistantBlock::Text {
2831                text: "Turn 1 answer".to_string(),
2832                meta: None,
2833            }],
2834            stop_reason: StopReason::EndTurn,
2835            identity: crate::types::TranscriptMessageIdentity::default(),
2836            created_at: crate::types::message_timestamp_now(),
2837        }));
2838
2839        let mut incoming = Session::with_id(previous.id().clone());
2840        incoming.push(Message::System(SystemMessage::new(
2841            "runtime system after context refresh",
2842        )));
2843        incoming.push(Message::User(UserMessage::text(
2844            "Verbose context that will be compacted".to_string(),
2845        )));
2846        for message in previous.messages()[1..].iter().cloned() {
2847            incoming.push(message);
2848        }
2849        incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2850            blocks: vec![AssistantBlock::Text {
2851                text: "Turn 2 generated answer".to_string(),
2852                meta: None,
2853            }],
2854            stop_reason: StopReason::EndTurn,
2855            identity: crate::types::TranscriptMessageIdentity::default(),
2856            created_at: crate::types::message_timestamp_now(),
2857        }));
2858        let parent_revision = incoming.transcript_revision()?;
2859        // Same rendered shape, but the boundary message carries the typed
2860        // injected-context role instead of the compaction-summary role. The
2861        // guard reads the typed marker: injected context is host-attached
2862        // ambient content, not a runtime compaction boundary.
2863        incoming.commit_transcript_rewrite(
2864            TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2865            vec![Message::User(UserMessage::injected_context(
2866                "[Context compacted] Earlier runtime context".to_string(),
2867            ))],
2868            crate::TranscriptRewriteReason::new("compaction"),
2869            Some("meerkat-core".to_string()),
2870            Some(parent_revision),
2871        )?;
2872
2873        assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2874        assert!(matches!(
2875            run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2876            Err(SessionStoreError::TranscriptContinuityViolation { .. }
2877                | SessionStoreError::MonotonicityViolation { .. })
2878        ));
2879        Ok(())
2880    }
2881
2882    #[test]
2883    fn run_boundary_guard_rejects_runtime_parent_with_inserted_message_before_tail()
2884    -> Result<(), Box<dyn std::error::Error>> {
2885        let mut previous = Session::new();
2886        previous.push(Message::System(SystemMessage::new("base system")));
2887        previous.push(Message::User(UserMessage::text("turn one".to_string())));
2888        previous.push(Message::BlockAssistant(BlockAssistantMessage {
2889            blocks: vec![AssistantBlock::Text {
2890                text: "answer one".to_string(),
2891                meta: None,
2892            }],
2893            stop_reason: StopReason::EndTurn,
2894            identity: crate::types::TranscriptMessageIdentity::default(),
2895            created_at: crate::types::message_timestamp_now(),
2896        }));
2897
2898        let parent_messages = vec![
2899            Message::System(SystemMessage::new("refreshed runtime system projection")),
2900            Message::User(UserMessage::text(
2901                "injected before retained tail".to_string(),
2902            )),
2903            previous.messages()[1].clone(),
2904            previous.messages()[2].clone(),
2905        ];
2906        let parent_revision = transcript_messages_digest(&parent_messages)?;
2907        let mut parent = previous.clone();
2908        parent.apply_transcript_history_state(TranscriptHistoryState {
2909            head: parent_revision.clone(),
2910            commits: Vec::new(),
2911            revisions: vec![crate::TranscriptRevisionBody {
2912                revision: parent_revision,
2913                parent_revision: None,
2914                messages: parent_messages,
2915                created_at: parent.updated_at(),
2916            }],
2917        })?;
2918        let parent_revision = parent.transcript_revision()?;
2919
2920        let mut incoming = parent.clone();
2921        incoming.commit_transcript_rewrite(
2922            TranscriptRewriteSelection::MessageRange {
2923                start: 0,
2924                end: parent.messages().len(),
2925            },
2926            vec![Message::User(UserMessage::text(
2927                "[Context compacted] summary".to_string(),
2928            ))],
2929            crate::TranscriptRewriteReason::new("compaction"),
2930            Some("meerkat-core".to_string()),
2931            Some(parent_revision),
2932        )?;
2933
2934        assert!(matches!(
2935            run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2936            Err(SessionStoreError::TranscriptContinuityViolation { .. }
2937                | SessionStoreError::MonotonicityViolation { .. })
2938        ));
2939        Ok(())
2940    }
2941
2942    #[test]
2943    fn run_boundary_guard_rejects_forged_parent_edge_before_real_rewrite_commit()
2944    -> Result<(), Box<dyn std::error::Error>> {
2945        let mut previous = Session::new();
2946        previous.push(Message::System(SystemMessage::new("base system")));
2947        previous.push(Message::User(UserMessage::text("turn one".to_string())));
2948        previous.push(Message::BlockAssistant(BlockAssistantMessage {
2949            blocks: vec![AssistantBlock::Text {
2950                text: "answer one".to_string(),
2951                meta: None,
2952            }],
2953            stop_reason: StopReason::EndTurn,
2954            identity: crate::types::TranscriptMessageIdentity::default(),
2955            created_at: crate::types::message_timestamp_now(),
2956        }));
2957        let previous_revision = previous.transcript_revision()?;
2958
2959        let forged_parent_messages = vec![
2960            Message::System(SystemMessage::new("refreshed runtime system projection")),
2961            Message::User(UserMessage::text(
2962                "forged insertion before retained tail".to_string(),
2963            )),
2964            previous.messages()[1].clone(),
2965            previous.messages()[2].clone(),
2966        ];
2967        let forged_parent_revision = transcript_messages_digest(&forged_parent_messages)?;
2968        let mut forged_parent = previous.clone();
2969        forged_parent.apply_transcript_history_state(TranscriptHistoryState {
2970            head: forged_parent_revision.clone(),
2971            commits: Vec::new(),
2972            revisions: vec![
2973                crate::TranscriptRevisionBody {
2974                    revision: previous_revision.clone(),
2975                    parent_revision: None,
2976                    messages: previous.messages().to_vec(),
2977                    created_at: previous.updated_at(),
2978                },
2979                crate::TranscriptRevisionBody {
2980                    revision: forged_parent_revision.clone(),
2981                    parent_revision: Some(previous_revision),
2982                    messages: forged_parent_messages,
2983                    created_at: forged_parent.updated_at(),
2984                },
2985            ],
2986        })?;
2987
2988        let mut incoming = forged_parent.clone();
2989        incoming.commit_transcript_rewrite(
2990            TranscriptRewriteSelection::MessageRange {
2991                start: 0,
2992                end: forged_parent.messages().len(),
2993            },
2994            vec![Message::User(UserMessage::text(
2995                "[Context compacted] forged branch".to_string(),
2996            ))],
2997            crate::TranscriptRewriteReason::new("compaction"),
2998            Some("meerkat-core".to_string()),
2999            Some(forged_parent_revision),
3000        )?;
3001
3002        assert!(matches!(
3003            run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
3004            Err(SessionStoreError::TranscriptContinuityViolation { .. }
3005                | SessionStoreError::MonotonicityViolation { .. })
3006        ));
3007        Ok(())
3008    }
3009
3010    #[test]
3011    fn append_only_guard_rejects_transient_mcp_pending_notice_cleanup_with_unaudited_commit()
3012    -> Result<(), crate::TranscriptEditError> {
3013        let mut previous = Session::new();
3014        previous.push(Message::User(UserMessage::text("hello".to_string())));
3015        previous.push(Message::SystemNotice(SystemNoticeMessage {
3016            kind: SystemNoticeKind::McpPending,
3017            body: Some("connecting".to_string()),
3018            blocks: vec![SystemNoticeBlock::Mcp {
3019                server_id: None,
3020                operation: None,
3021                phase: None,
3022                persisted: false,
3023                detail: Some("connecting".to_string()),
3024                pending_sources: vec!["test-server".to_string()],
3025            }],
3026            created_at: crate::types::message_timestamp_now(),
3027        }));
3028        previous.push(Message::BlockAssistant(BlockAssistantMessage::new(
3029            vec![crate::types::AssistantBlock::Text {
3030                text: "answer".to_string(),
3031                meta: None,
3032            }],
3033            StopReason::EndTurn,
3034        )));
3035
3036        let mut incoming = previous.clone();
3037        incoming.replace_messages_internal(
3038            previous
3039                .messages()
3040                .iter()
3041                .filter(|message| !matches!(message, Message::SystemNotice(_)))
3042                .cloned()
3043                .collect(),
3044            crate::TranscriptRewriteReason::new("unit-test"),
3045        )?;
3046        incoming.push(Message::User(UserMessage::text("again".to_string())));
3047
3048        assert!(matches!(
3049            append_only_save_guard(&incoming, Some(&previous)),
3050            Err(SessionStoreError::InvalidTranscriptRewrite { .. })
3051        ));
3052        Ok::<(), crate::TranscriptEditError>(())
3053    }
3054
3055    #[test]
3056    fn rewrite_chain_finder_crosses_normal_append_between_rewrites()
3057    -> Result<(), Box<dyn std::error::Error>> {
3058        let mut session = Session::new();
3059        session.push(Message::User(UserMessage::text("first".to_string())));
3060        session.push(Message::BlockAssistant(BlockAssistantMessage {
3061            blocks: vec![AssistantBlock::Text {
3062                text: "verbose first answer".to_string(),
3063                meta: None,
3064            }],
3065            stop_reason: StopReason::EndTurn,
3066            identity: crate::types::TranscriptMessageIdentity::default(),
3067            created_at: crate::types::message_timestamp_now(),
3068        }));
3069
3070        let original = session.transcript_revision()?;
3071        let first = session.commit_transcript_rewrite(
3072            TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
3073            vec![Message::BlockAssistant(BlockAssistantMessage {
3074                blocks: vec![AssistantBlock::Text {
3075                    text: "compact first answer".to_string(),
3076                    meta: None,
3077                }],
3078                stop_reason: StopReason::EndTurn,
3079                identity: crate::types::TranscriptMessageIdentity::default(),
3080                created_at: crate::types::message_timestamp_now(),
3081            })],
3082            crate::TranscriptRewriteReason::new("compaction"),
3083            Some("unit-test".to_string()),
3084            Some(original.clone()),
3085        )?;
3086
3087        session.push(Message::User(UserMessage::text("second".to_string())));
3088        session.push(Message::BlockAssistant(BlockAssistantMessage {
3089            blocks: vec![AssistantBlock::Text {
3090                text: "verbose second answer".to_string(),
3091                meta: None,
3092            }],
3093            stop_reason: StopReason::EndTurn,
3094            identity: crate::types::TranscriptMessageIdentity::default(),
3095            created_at: crate::types::message_timestamp_now(),
3096        }));
3097        let bridge = session.transcript_revision()?;
3098        assert_ne!(bridge, first.revision);
3099
3100        let second = session.commit_transcript_rewrite(
3101            TranscriptRewriteSelection::MessageRange { start: 3, end: 4 },
3102            vec![Message::BlockAssistant(BlockAssistantMessage {
3103                blocks: vec![AssistantBlock::Text {
3104                    text: "compact second answer".to_string(),
3105                    meta: None,
3106                }],
3107                stop_reason: StopReason::EndTurn,
3108                identity: crate::types::TranscriptMessageIdentity::default(),
3109                created_at: crate::types::message_timestamp_now(),
3110            })],
3111            crate::TranscriptRewriteReason::new("compaction"),
3112            Some("unit-test".to_string()),
3113            Some(bridge),
3114        )?;
3115        let state = session
3116            .transcript_history_state()?
3117            .ok_or_else(|| std::io::Error::other("missing transcript history state"))?;
3118
3119        let chain =
3120            find_transcript_rewrite_commit_chain_extending(&state, &original, &second.revision)
3121                .ok_or_else(|| {
3122                    std::io::Error::other(
3123                        "rewrite chain should extend through normal append bridge",
3124                    )
3125                })?;
3126        assert_eq!(chain.len(), 2);
3127        assert_eq!(chain[0].revision, first.revision);
3128        assert_eq!(chain[1].revision, second.revision);
3129        Ok(())
3130    }
3131
3132    #[test]
3133    fn run_boundary_guard_rejects_dropped_retained_rewrite_commits()
3134    -> Result<(), Box<dyn std::error::Error>> {
3135        let mut base = Session::new();
3136        base.push(Message::User(UserMessage::text("turn one".to_string())));
3137        base.push(Message::BlockAssistant(BlockAssistantMessage {
3138            blocks: vec![AssistantBlock::Text {
3139                text: "verbose answer".to_string(),
3140                meta: None,
3141            }],
3142            stop_reason: StopReason::EndTurn,
3143            identity: crate::types::TranscriptMessageIdentity::default(),
3144            created_at: crate::types::message_timestamp_now(),
3145        }));
3146        let base_revision = base.transcript_revision()?;
3147
3148        let mut previous = base.clone();
3149        let _retained_commit = previous.commit_transcript_rewrite(
3150            TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
3151            vec![Message::BlockAssistant(BlockAssistantMessage {
3152                blocks: vec![AssistantBlock::Text {
3153                    text: "first compact answer".to_string(),
3154                    meta: None,
3155                }],
3156                stop_reason: StopReason::EndTurn,
3157                identity: crate::types::TranscriptMessageIdentity::default(),
3158                created_at: crate::types::message_timestamp_now(),
3159            })],
3160            crate::TranscriptRewriteReason::new("compaction"),
3161            Some("unit-test".to_string()),
3162            Some(base_revision),
3163        )?;
3164        let previous_revision = previous.transcript_revision()?;
3165
3166        let mut incoming = previous.clone();
3167        let new_commit = incoming.commit_transcript_rewrite(
3168            TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
3169            vec![Message::BlockAssistant(BlockAssistantMessage {
3170                blocks: vec![AssistantBlock::Text {
3171                    text: "second compact answer".to_string(),
3172                    meta: None,
3173                }],
3174                stop_reason: StopReason::EndTurn,
3175                identity: crate::types::TranscriptMessageIdentity::default(),
3176                created_at: crate::types::message_timestamp_now(),
3177            })],
3178            crate::TranscriptRewriteReason::new("compaction"),
3179            Some("unit-test".to_string()),
3180            Some(previous_revision),
3181        )?;
3182        let mut state = incoming
3183            .transcript_history_state()?
3184            .ok_or_else(|| std::io::Error::other("incoming rewrite should retain history"))?;
3185        state.commits = vec![new_commit];
3186        incoming.set_metadata_unchecked_for_test(
3187            crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
3188            serde_json::to_value(state)?,
3189        );
3190
3191        assert!(matches!(
3192            run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
3193            Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
3194                if reason.contains("drop retained transcript rewrite commits")
3195        ));
3196        Ok(())
3197    }
3198
3199    // ------------------------------------------------------------------
3200    // FOLD 2: the persist-time system-context append-admission decision routes
3201    // through SessionDocumentMachine ResolveSystemContextPersistAppendAdmission
3202    // (the SAME machine the staging path drives). These tests pin that the
3203    // persist-time verdict matches a direct machine call for every shape, and
3204    // that the four admission cases behave exactly as the retired shell reducer.
3205    // ------------------------------------------------------------------
3206
3207    fn runtime_append_system(content: &str) -> SystemMessage {
3208        let mut system = SystemMessage::new(content);
3209        system.mutation_kind = crate::types::SystemPromptMutationKind::RuntimeContextAppend;
3210        system
3211    }
3212
3213    /// Direct machine call mirroring the persist-time observation extraction —
3214    /// the persist-time path MUST agree with this for every input shape.
3215    #[allow(clippy::expect_used)]
3216    fn machine_persist_append_admits(
3217        previous: Option<&SystemMessage>,
3218        incoming: &SystemMessage,
3219    ) -> bool {
3220        let has_previous = previous.is_some();
3221        let content_identical =
3222            previous.is_some_and(|previous| incoming.content == previous.content);
3223        let content_extends_previous =
3224            previous.is_some_and(|previous| incoming.content.starts_with(&previous.content));
3225        let appended_starts_with_separator = previous.is_some_and(|previous| {
3226            incoming
3227                .content
3228                .get(previous.content.len()..)
3229                .is_some_and(|appended| appended.starts_with(SYSTEM_CONTEXT_SEPARATOR))
3230        });
3231        let incoming_is_runtime_context_append = incoming.mutation_kind.is_runtime_context_append();
3232        let mut authority = crate::session_document::SessionDocumentMachineAuthority::new();
3233        let effects = authority
3234            .resolve_system_context_persist_append_admission(
3235                has_previous,
3236                content_identical,
3237                content_extends_previous,
3238                appended_starts_with_separator,
3239                incoming_is_runtime_context_append,
3240            )
3241            .expect("machine resolves persist-append admission");
3242        effects.into_iter().any(|effect| {
3243            matches!(
3244                effect,
3245                crate::session_document::SessionDocumentEffect::SystemContextPersistAppendAdmissionResolved {
3246                    admission: crate::session_document::SystemContextPersistAppendAdmission::Admit,
3247                }
3248            )
3249        })
3250    }
3251
3252    #[allow(clippy::expect_used)]
3253    fn assert_persist_append_matches_machine(
3254        previous: Option<&SystemMessage>,
3255        incoming: &SystemMessage,
3256        expected: bool,
3257    ) {
3258        let verdict =
3259            system_context_is_append(previous, incoming).expect("persist-time admission resolves");
3260        assert_eq!(verdict, expected, "persist-time verdict mismatch");
3261        assert_eq!(
3262            verdict,
3263            machine_persist_append_admits(previous, incoming),
3264            "persist-time verdict diverges from direct machine call"
3265        );
3266    }
3267
3268    #[test]
3269    fn persist_append_identical_content_admits() {
3270        let previous = SystemMessage::new("base system");
3271        let incoming = SystemMessage::new("base system");
3272        assert_persist_append_matches_machine(Some(&previous), &incoming, true);
3273    }
3274
3275    #[test]
3276    fn persist_append_separator_append_with_marker_admits() {
3277        let previous = SystemMessage::new("base system");
3278        let incoming = runtime_append_system(&format!(
3279            "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nextra"
3280        ));
3281        assert_persist_append_matches_machine(Some(&previous), &incoming, true);
3282    }
3283
3284    #[test]
3285    fn persist_append_shaped_without_marker_rejects() {
3286        let previous = SystemMessage::new("base system");
3287        // Append-shaped content but no runtime-context-append provenance marker.
3288        let incoming = SystemMessage::new(format!(
3289            "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nextra"
3290        ));
3291        assert_persist_append_matches_machine(Some(&previous), &incoming, false);
3292    }
3293
3294    #[test]
3295    fn persist_append_divergent_content_rejects() {
3296        let previous = SystemMessage::new("base system");
3297        let incoming = runtime_append_system("totally different");
3298        assert_persist_append_matches_machine(Some(&previous), &incoming, false);
3299    }
3300
3301    #[test]
3302    fn persist_append_no_previous_admits_only_with_marker() {
3303        let with_marker = runtime_append_system("brand new context");
3304        assert_persist_append_matches_machine(None, &with_marker, true);
3305
3306        let without_marker = SystemMessage::new("brand new context");
3307        assert_persist_append_matches_machine(None, &without_marker, false);
3308    }
3309
3310    fn assistant_with_bookkeeping(
3311        text: &str,
3312        run_id: Option<crate::lifecycle::RunId>,
3313        created_at: crate::types::MessageTimestamp,
3314    ) -> Message {
3315        Message::BlockAssistant(BlockAssistantMessage {
3316            blocks: vec![AssistantBlock::Text {
3317                text: text.to_string(),
3318                meta: None,
3319            }],
3320            stop_reason: StopReason::EndTurn,
3321            identity: crate::types::TranscriptMessageIdentity {
3322                interaction_id: None,
3323                run_id,
3324            },
3325            created_at,
3326        })
3327    }
3328
3329    /// Cold-restart resume regression (Ask B): a re-created runtime authority
3330    /// re-stamps run identity and timestamps on the transcript copy it
3331    /// re-projects. The transcript revision is a content address, so a
3332    /// bookkeeping-only difference on the shared prefix must not fail
3333    /// continuity.
3334    #[test]
3335    fn append_only_guard_accepts_rebookkept_prefix_identity_and_timestamps() {
3336        let base_time = crate::types::message_timestamp_now();
3337        let mut previous = Session::new();
3338        previous.push(Message::System(SystemMessage::new("base system")));
3339        previous.push(Message::User(UserMessage::text("turn one".to_string())));
3340        previous.push(assistant_with_bookkeeping(
3341            "answer one",
3342            Some(crate::lifecycle::RunId::new()),
3343            base_time,
3344        ));
3345
3346        let mut incoming = previous.clone();
3347        let mut rebookkept = previous.messages().to_vec();
3348        for message in &mut rebookkept {
3349            match message {
3350                Message::User(user) => {
3351                    user.created_at = base_time + chrono::Duration::hours(1);
3352                }
3353                Message::BlockAssistant(assistant) => {
3354                    assistant.identity = crate::types::TranscriptMessageIdentity {
3355                        interaction_id: None,
3356                        run_id: Some(crate::lifecycle::RunId::new()),
3357                    };
3358                    assistant.created_at = base_time + chrono::Duration::hours(1);
3359                }
3360                _ => {}
3361            }
3362        }
3363        rebookkept.push(Message::User(UserMessage::text("turn two".to_string())));
3364        incoming.messages = std::sync::Arc::new(rebookkept);
3365
3366        assert!(
3367            append_only_save_guard(&incoming, Some(&previous)).is_ok(),
3368            "bookkeeping-only prefix divergence must not fail continuity"
3369        );
3370    }
3371
3372    #[test]
3373    fn append_only_guard_rejects_content_divergence_despite_matching_bookkeeping() {
3374        let base_time = crate::types::message_timestamp_now();
3375        let run_id = crate::lifecycle::RunId::new();
3376        let mut previous = Session::new();
3377        previous.push(Message::User(UserMessage::text("turn one".to_string())));
3378        previous.push(assistant_with_bookkeeping(
3379            "answer one",
3380            Some(run_id),
3381            base_time,
3382        ));
3383
3384        let mut incoming = previous.clone();
3385        let mut diverged = previous.messages().to_vec();
3386        if let Message::BlockAssistant(assistant) = &mut diverged[1] {
3387            assistant.blocks = vec![AssistantBlock::Text {
3388                text: "a different answer".to_string(),
3389                meta: None,
3390            }];
3391        }
3392        diverged.push(Message::User(UserMessage::text("turn two".to_string())));
3393        incoming.messages = std::sync::Arc::new(diverged);
3394
3395        assert!(matches!(
3396            append_only_save_guard(&incoming, Some(&previous)),
3397            Err(SessionStoreError::TranscriptContinuityViolation { .. })
3398        ));
3399    }
3400}