1use 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#[derive(Debug, Clone, Default)]
33pub struct SessionFilter {
34 pub created_after: Option<SystemTime>,
36 pub updated_after: Option<SystemTime>,
38 pub limit: Option<usize>,
40 pub offset: Option<usize>,
42}
43
44#[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
98pub 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
108pub 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
379pub 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
443fn system_context_is_append(
459 previous: Option<&SystemMessage>,
460 incoming: &SystemMessage,
461) -> Result<bool, SessionStoreError> {
462 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
546pub 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 return Err(append_error);
566 };
567 let incoming_revision =
568 transcript_messages_digest(incoming.messages()).map_err(SessionStoreError::from)?;
569 let Some(state) = incoming.transcript_history_state().map_err(|err| {
570 SessionStoreError::InvalidTranscriptRewrite {
571 id: incoming.id().clone(),
572 reason: format!("incoming transcript history state is malformed: {err}"),
573 }
574 })?
575 else {
576 return Err(append_error);
577 };
578 validate_rewrite_save_retains_previous_commits(incoming, previous, &state)?;
579 let commits = find_transcript_rewrite_commit_chain_extending_session(
580 &state,
581 previous,
582 &incoming_revision,
583 )?;
584 if commits.is_none()
585 && run_boundary_context_summary_tail_projection_save_guard(
586 incoming, previous, &state,
587 )?
588 {
589 return Ok(());
590 }
591 let Some(commits) = commits else {
592 return Err(append_error);
593 };
594 let Some(commit) = commits.first() else {
595 if state.commits.is_empty() {
596 return Err(append_error);
597 }
598 for commit in &state.commits {
599 validate_transcript_rewrite_commit_bodies(incoming, commit, &state)?;
600 }
601 return Ok(());
602 };
603 transcript_rewrite_bridge_save_guard(incoming, commit, &state, &incoming_revision)?;
604 for commit in commits.iter().skip(1) {
605 validate_transcript_rewrite_commit_bodies(incoming, commit, &state)?;
606 }
607 Ok(())
608 }
609 }
610}
611
612fn run_boundary_commitless_history_projection_save_guard(
613 incoming: &Session,
614 previous: Option<&Session>,
615) -> Result<bool, SessionStoreError> {
616 let Some(state) = incoming.transcript_history_state().map_err(|err| {
617 SessionStoreError::InvalidTranscriptRewrite {
618 id: incoming.id().clone(),
619 reason: format!("incoming transcript history state is malformed: {err}"),
620 }
621 })?
622 else {
623 return Ok(false);
624 };
625 if !state.commits.is_empty() {
626 return Ok(false);
627 }
628
629 let incoming_revision =
630 transcript_messages_digest(incoming.messages()).map_err(SessionStoreError::from)?;
631 if state.head != incoming_revision
632 || !state
633 .revisions
634 .iter()
635 .any(|body| body.revision == incoming_revision)
636 {
637 return Ok(false);
638 }
639
640 let mut projection_without_history = incoming.clone();
641 projection_without_history.clear_transcript_history_state();
642 if append_only_save_guard(&projection_without_history, previous).is_err() {
643 return Ok(false);
644 }
645
646 let Some(previous) = previous else {
647 return Ok(state.commits.is_empty());
648 };
649 if previous
650 .transcript_history_state()
651 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
652 id: incoming.id().clone(),
653 reason: format!("previous transcript history state is malformed: {err}"),
654 })?
655 .is_some()
656 {
657 return Ok(false);
658 }
659
660 let previous_revision =
661 transcript_messages_digest(previous.messages()).map_err(SessionStoreError::from)?;
662 Ok(incoming_revision == previous_revision
663 || transcript_history_revision_extends(&state, &incoming_revision, &previous_revision))
664}
665
666fn run_boundary_context_summary_tail_projection_save_guard(
667 incoming: &Session,
668 previous: &Session,
669 state: &TranscriptHistoryState,
670) -> Result<bool, SessionStoreError> {
671 if state.commits.is_empty() {
672 return Ok(false);
673 }
674 incoming
675 .validate_transcript_history_state()
676 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
677 id: incoming.id().clone(),
678 reason: format!("incoming transcript history state is malformed: {err}"),
679 })?;
680
681 let (incoming_system, incoming_tail) = match incoming.messages().split_first() {
682 Some((Message::System(system), tail)) => (Some(system), tail),
683 _ => (None, incoming.messages()),
684 };
685 let (previous_system, previous_tail) = match previous.messages().split_first() {
686 Some((Message::System(system), tail)) => (Some(system), tail),
687 _ => (None, previous.messages()),
688 };
689 if incoming_system.is_some() != previous_system.is_some()
690 || incoming_tail.len() <= previous_tail.len()
691 {
692 return Ok(false);
693 }
694 let Some(Message::User(summary)) = incoming_tail.first() else {
695 return Ok(false);
696 };
697 if !summary.transcript_role.is_compaction_summary() {
702 return Ok(false);
703 }
704
705 let retained_end = 1 + previous_tail.len();
706 let retained = &incoming_tail[1..retained_end];
707 let retained_revision =
708 transcript_messages_digest(retained).map_err(SessionStoreError::from)?;
709 let previous_revision =
710 transcript_messages_digest(previous_tail).map_err(SessionStoreError::from)?;
711 if retained_revision != previous_revision {
712 return Ok(false);
713 }
714
715 for commit in &state.commits {
716 validate_transcript_rewrite_commit_bodies(incoming, commit, state)?;
717 }
718 Ok(true)
719}
720
721pub fn find_transcript_rewrite_commit_extending<'a>(
724 state: &'a TranscriptHistoryState,
725 previous_revision: &str,
726 incoming_revision: &str,
727) -> Option<&'a TranscriptRewriteCommit> {
728 find_transcript_rewrite_commit_chain_extending(state, previous_revision, incoming_revision)
729 .and_then(|commits| commits.into_iter().next())
730}
731
732pub fn find_transcript_rewrite_commit_chain_extending<'a>(
735 state: &'a TranscriptHistoryState,
736 previous_revision: &str,
737 incoming_revision: &str,
738) -> Option<Vec<&'a TranscriptRewriteCommit>> {
739 let mut chain = Vec::new();
740 let mut cursor = previous_revision;
741 let mut visited = std::collections::BTreeSet::new();
742 loop {
743 if incoming_revision == cursor {
744 return Some(chain);
745 }
746 if !visited.insert(cursor.to_string()) {
747 return None;
748 }
749 let commit = state.commits.iter().find(|commit| {
750 (commit.parent_revision == cursor
751 || transcript_history_revision_extends(state, &commit.parent_revision, cursor))
752 && transcript_history_revision_extends(state, incoming_revision, &commit.revision)
753 });
754 let Some(commit) = commit else {
755 return transcript_history_revision_extends(state, incoming_revision, cursor)
756 .then_some(chain);
757 };
758 cursor = &commit.revision;
759 chain.push(commit);
760 }
761}
762
763pub fn find_transcript_rewrite_commit_chain_extending_session<'a>(
772 state: &'a TranscriptHistoryState,
773 previous: &Session,
774 incoming_revision: &str,
775) -> Result<Option<Vec<&'a TranscriptRewriteCommit>>, SessionStoreError> {
776 let previous_revision =
777 transcript_messages_digest(previous.messages()).map_err(SessionStoreError::from)?;
778 let mut chain = Vec::new();
779 let mut cursor = previous_revision.as_str();
780 let mut visited = std::collections::BTreeSet::new();
781 loop {
782 if incoming_revision == cursor {
783 return Ok(Some(chain));
784 }
785 if !visited.insert(cursor.to_string()) {
786 return Ok(None);
787 }
788
789 let Some(cursor_messages) = transcript_history_messages_for_revision(
790 state,
791 cursor,
792 &previous_revision,
793 previous.messages(),
794 ) else {
795 return Ok(None);
796 };
797 let mut selected = None;
798 for commit in &state.commits {
799 if !transcript_history_revision_extends(state, incoming_revision, &commit.revision) {
800 continue;
801 }
802 let parent_extends_cursor = commit.parent_revision == cursor
803 || revision_body_preserves_append_continuation_prefix(
804 state,
805 &commit.parent_revision,
806 cursor_messages,
807 cursor,
808 )?;
809 if parent_extends_cursor {
810 selected = Some(commit);
811 break;
812 }
813 }
814
815 let Some(commit) = selected else {
816 if revision_body_preserves_append_continuation_prefix(
817 state,
818 incoming_revision,
819 cursor_messages,
820 cursor,
821 )? {
822 return Ok(Some(chain));
823 }
824 return Ok(None);
825 };
826 cursor = &commit.revision;
827 chain.push(commit);
828 }
829}
830
831fn transcript_history_messages_for_revision<'a>(
832 state: &'a TranscriptHistoryState,
833 revision: &str,
834 previous_revision: &str,
835 previous_messages: &'a [Message],
836) -> Option<&'a [Message]> {
837 if revision == previous_revision {
838 return Some(previous_messages);
839 }
840 state
841 .revisions
842 .iter()
843 .find(|body| body.revision == revision)
844 .map(|body| body.messages.as_slice())
845}
846
847fn revision_body_preserves_append_continuation_prefix(
848 state: &TranscriptHistoryState,
849 revision: &str,
850 ancestor_messages: &[Message],
851 ancestor_revision: &str,
852) -> Result<bool, SessionStoreError> {
853 if revision == ancestor_revision {
854 return Ok(true);
855 }
856 let Some(body) = state
857 .revisions
858 .iter()
859 .find(|body| body.revision == revision)
860 else {
861 return Ok(false);
862 };
863 if body.messages.len() >= ancestor_messages.len() {
864 let prefix_revision = transcript_messages_digest(&body.messages[..ancestor_messages.len()])
865 .map_err(SessionStoreError::from)?;
866 if prefix_revision == ancestor_revision {
867 return Ok(true);
868 }
869 }
870 Ok(
871 messages_preserve_conversation_tail_with_system_context_append(
872 &body.messages,
873 ancestor_messages,
874 )? || messages_preserve_tail_after_leading_system_refresh(
875 &body.messages,
876 ancestor_messages,
877 )?,
878 )
879}
880
881fn messages_preserve_tail_after_leading_system_refresh(
882 incoming: &[Message],
883 previous: &[Message],
884) -> Result<bool, SessionStoreError> {
885 let (Some(Message::System(_)), Some(Message::System(_))) = (incoming.first(), previous.first())
886 else {
887 return Ok(false);
888 };
889 if incoming.len() < previous.len() {
890 return Ok(false);
891 }
892 let previous_tail_len = previous.len().saturating_sub(1);
893 if previous_tail_len == 0 {
894 return Ok(true);
895 }
896 let previous_tail_revision =
897 transcript_messages_digest(&previous[1..]).map_err(SessionStoreError::from)?;
898 let incoming_tail = &incoming[1..];
899 if incoming_tail.len() < previous_tail_len {
900 return Ok(false);
901 }
902 let incoming_tail_prefix_revision =
903 transcript_messages_digest(&incoming_tail[..previous_tail_len])
904 .map_err(SessionStoreError::from)?;
905 Ok(incoming_tail_prefix_revision == previous_tail_revision)
906}
907
908fn transcript_history_revision_extends(
909 state: &TranscriptHistoryState,
910 descendant: &str,
911 ancestor: &str,
912) -> bool {
913 if descendant == ancestor {
914 return true;
915 }
916 let mut cursor = descendant;
917 while let Some(body) = state.revisions.iter().find(|body| body.revision == cursor) {
918 let Some(parent) = body.parent_revision.as_deref() else {
919 return false;
920 };
921 if parent == ancestor {
922 return true;
923 }
924 cursor = parent;
925 }
926 false
927}
928
929fn transcript_rewrite_bridge_save_guard(
930 incoming: &Session,
931 commit: &TranscriptRewriteCommit,
932 incoming_state: &TranscriptHistoryState,
933 incoming_message_digest: &str,
934) -> Result<(), SessionStoreError> {
935 validate_transcript_rewrite_commit_bodies(incoming, commit, incoming_state)?;
936 if incoming_state.head != incoming_message_digest {
937 return Err(SessionStoreError::InvalidTranscriptRewrite {
938 id: incoming.id().clone(),
939 reason: format!(
940 "incoming transcript graph head {} does not match current message digest {incoming_message_digest}",
941 incoming_state.head
942 ),
943 });
944 }
945 if !transcript_history_revision_extends(
946 incoming_state,
947 incoming_message_digest,
948 &commit.revision,
949 ) {
950 return Err(SessionStoreError::InvalidTranscriptRewrite {
951 id: incoming.id().clone(),
952 reason: format!(
953 "incoming transcript head {incoming_message_digest} does not extend rewrite revision {}",
954 commit.revision
955 ),
956 });
957 }
958 Ok(())
959}
960
961pub fn transcript_rewrite_save_guard(
964 incoming: &Session,
965 previous: Option<&Session>,
966 commit: &TranscriptRewriteCommit,
967) -> Result<(), SessionStoreError> {
968 incoming
969 .validate_transcript_history_state()
970 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
971 id: incoming.id().clone(),
972 reason: format!("incoming transcript history state is malformed: {err}"),
973 })?;
974 let Some(previous) = previous else {
975 return Err(SessionStoreError::InvalidTranscriptRewrite {
976 id: incoming.id().clone(),
977 reason: "rewrite target has no previously persisted session".to_string(),
978 });
979 };
980 if incoming.id() != previous.id() {
981 return Err(SessionStoreError::InvalidTranscriptRewrite {
982 id: incoming.id().clone(),
983 reason: format!(
984 "incoming session id {} differs from previous session id {}",
985 incoming.id(),
986 previous.id()
987 ),
988 });
989 }
990 let previous_revision = previous.transcript_revision().map_err(|err| {
991 SessionStoreError::InvalidTranscriptRewrite {
992 id: incoming.id().clone(),
993 reason: format!("previous transcript revision is malformed: {err}"),
994 }
995 })?;
996 if previous_revision != commit.parent_revision {
997 return Err(SessionStoreError::TranscriptRevisionConflict {
998 id: incoming.id().clone(),
999 expected: commit.parent_revision.clone(),
1000 actual: previous_revision,
1001 });
1002 }
1003 let previous_message_digest =
1004 transcript_messages_digest(previous.messages()).map_err(|err| {
1005 SessionStoreError::InvalidTranscriptRewrite {
1006 id: incoming.id().clone(),
1007 reason: format!("previous current transcript is not digestible: {err}"),
1008 }
1009 })?;
1010 if previous_message_digest != commit.parent_revision {
1011 return Err(SessionStoreError::InvalidTranscriptRewrite {
1012 id: incoming.id().clone(),
1013 reason: format!(
1014 "previous current transcript digest {previous_message_digest} does not match commit parent {}",
1015 commit.parent_revision
1016 ),
1017 });
1018 }
1019 let incoming_revision = incoming.transcript_revision().map_err(|err| {
1020 SessionStoreError::InvalidTranscriptRewrite {
1021 id: incoming.id().clone(),
1022 reason: format!("incoming transcript revision is malformed: {err}"),
1023 }
1024 })?;
1025 if incoming_revision != commit.revision {
1026 return Err(SessionStoreError::InvalidTranscriptRewrite {
1027 id: incoming.id().clone(),
1028 reason: format!(
1029 "incoming transcript revision {incoming_revision} does not match commit revision {}",
1030 commit.revision
1031 ),
1032 });
1033 }
1034 let incoming_message_digest =
1035 transcript_messages_digest(incoming.messages()).map_err(|err| {
1036 SessionStoreError::InvalidTranscriptRewrite {
1037 id: incoming.id().clone(),
1038 reason: format!("incoming current transcript is not digestible: {err}"),
1039 }
1040 })?;
1041 if incoming_message_digest != commit.revision {
1042 return Err(SessionStoreError::InvalidTranscriptRewrite {
1043 id: incoming.id().clone(),
1044 reason: format!(
1045 "incoming current transcript digest {incoming_message_digest} does not match commit revision {}",
1046 commit.revision
1047 ),
1048 });
1049 }
1050 let Some(incoming_state) = incoming.transcript_history_state().map_err(|err| {
1051 SessionStoreError::InvalidTranscriptRewrite {
1052 id: incoming.id().clone(),
1053 reason: format!("incoming transcript history state is malformed: {err}"),
1054 }
1055 })?
1056 else {
1057 return Err(SessionStoreError::InvalidTranscriptRewrite {
1058 id: incoming.id().clone(),
1059 reason: "incoming rewrite did not persist a transcript revision graph".to_string(),
1060 });
1061 };
1062 validate_rewrite_save_retains_previous_commits(incoming, previous, &incoming_state)?;
1063 validate_transcript_rewrite_commit_bodies(incoming, commit, &incoming_state)
1064}
1065
1066fn validate_transcript_rewrite_commit_bodies(
1067 incoming: &Session,
1068 commit: &TranscriptRewriteCommit,
1069 incoming_state: &TranscriptHistoryState,
1070) -> Result<(), SessionStoreError> {
1071 if !incoming_state
1072 .commits
1073 .iter()
1074 .any(|persisted| persisted == commit)
1075 {
1076 return Err(SessionStoreError::InvalidTranscriptRewrite {
1077 id: incoming.id().clone(),
1078 reason: format!(
1079 "incoming rewrite did not persist the rewrite commit in the transcript graph (wanted {} -> {}, graph commits: {:?})",
1080 commit.parent_revision,
1081 commit.revision,
1082 incoming_state
1083 .commits
1084 .iter()
1085 .map(|commit| (&commit.parent_revision, &commit.revision))
1086 .collect::<Vec<_>>()
1087 ),
1088 });
1089 }
1090 let Some(parent_body) = incoming_state
1091 .revisions
1092 .iter()
1093 .find(|body| body.revision == commit.parent_revision)
1094 else {
1095 return Err(SessionStoreError::InvalidTranscriptRewrite {
1096 id: incoming.id().clone(),
1097 reason: format!(
1098 "incoming rewrite omitted parent revision body {}",
1099 commit.parent_revision
1100 ),
1101 });
1102 };
1103 let Some(revision_body) = incoming_state
1104 .revisions
1105 .iter()
1106 .find(|body| body.revision == commit.revision)
1107 else {
1108 return Err(SessionStoreError::InvalidTranscriptRewrite {
1109 id: incoming.id().clone(),
1110 reason: format!(
1111 "incoming rewrite omitted new revision body {}",
1112 commit.revision
1113 ),
1114 });
1115 };
1116 if parent_body.messages.len() != commit.messages_before
1117 || revision_body.messages.len() != commit.messages_after
1118 {
1119 return Err(SessionStoreError::InvalidTranscriptRewrite {
1120 id: incoming.id().clone(),
1121 reason: format!(
1122 "commit message counts {} -> {} do not match persisted rewrite {} -> {}",
1123 commit.messages_before,
1124 commit.messages_after,
1125 parent_body.messages.len(),
1126 revision_body.messages.len()
1127 ),
1128 });
1129 }
1130 let parent_body_revision =
1131 transcript_messages_digest(&parent_body.messages).map_err(|err| {
1132 SessionStoreError::InvalidTranscriptRewrite {
1133 id: incoming.id().clone(),
1134 reason: format!("parent revision body is not digestible: {err}"),
1135 }
1136 })?;
1137 if parent_body_revision != commit.parent_revision {
1138 return Err(SessionStoreError::InvalidTranscriptRewrite {
1139 id: incoming.id().clone(),
1140 reason: format!(
1141 "parent revision body digest {parent_body_revision} does not match commit parent {}",
1142 commit.parent_revision
1143 ),
1144 });
1145 }
1146 let (start, end) = match &commit.selection {
1147 TranscriptRewriteSelection::MessageRange { start, end } => (*start, *end),
1148 };
1149 if start > end || end > parent_body.messages.len() {
1150 return Err(SessionStoreError::InvalidTranscriptRewrite {
1151 id: incoming.id().clone(),
1152 reason: format!(
1153 "commit selection {start}..{end} is invalid for parent revision with {} messages",
1154 parent_body.messages.len()
1155 ),
1156 });
1157 }
1158 let original_span_digest = transcript_messages_digest(&parent_body.messages[start..end])
1159 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1160 id: incoming.id().clone(),
1161 reason: format!("original span body is not digestible: {err}"),
1162 })?;
1163 if original_span_digest != commit.original_span_digest {
1164 return Err(SessionStoreError::InvalidTranscriptRewrite {
1165 id: incoming.id().clone(),
1166 reason: format!(
1167 "original span digest {original_span_digest} does not match commit digest {}",
1168 commit.original_span_digest
1169 ),
1170 });
1171 }
1172 let revision_body_digest =
1173 transcript_messages_digest(&revision_body.messages).map_err(|err| {
1174 SessionStoreError::InvalidTranscriptRewrite {
1175 id: incoming.id().clone(),
1176 reason: format!("new revision body is not digestible: {err}"),
1177 }
1178 })?;
1179 if revision_body_digest != commit.revision {
1180 return Err(SessionStoreError::InvalidTranscriptRewrite {
1181 id: incoming.id().clone(),
1182 reason: format!(
1183 "new revision body digest {revision_body_digest} does not match commit revision {}",
1184 commit.revision
1185 ),
1186 });
1187 }
1188 let removed_len = end - start;
1189 let retained_len = commit
1190 .messages_before
1191 .checked_sub(removed_len)
1192 .ok_or_else(|| SessionStoreError::InvalidTranscriptRewrite {
1193 id: incoming.id().clone(),
1194 reason: "commit removed more messages than it recorded before rewrite".to_string(),
1195 })?;
1196 let replacement_len = commit
1197 .messages_after
1198 .checked_sub(retained_len)
1199 .ok_or_else(|| SessionStoreError::InvalidTranscriptRewrite {
1200 id: incoming.id().clone(),
1201 reason: "commit message counts cannot describe a replacement span".to_string(),
1202 })?;
1203 let replacement_end = start.checked_add(replacement_len).ok_or_else(|| {
1204 SessionStoreError::InvalidTranscriptRewrite {
1205 id: incoming.id().clone(),
1206 reason: "replacement span end overflowed".to_string(),
1207 }
1208 })?;
1209 if replacement_end > revision_body.messages.len() {
1210 return Err(SessionStoreError::InvalidTranscriptRewrite {
1211 id: incoming.id().clone(),
1212 reason: format!(
1213 "replacement span {start}..{replacement_end} is invalid for revision with {} messages",
1214 revision_body.messages.len()
1215 ),
1216 });
1217 }
1218 let parent_prefix_digest =
1219 transcript_messages_digest(&parent_body.messages[..start]).map_err(|err| {
1220 SessionStoreError::InvalidTranscriptRewrite {
1221 id: incoming.id().clone(),
1222 reason: format!("parent prefix body is not digestible: {err}"),
1223 }
1224 })?;
1225 let revision_prefix_digest = transcript_messages_digest(&revision_body.messages[..start])
1226 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1227 id: incoming.id().clone(),
1228 reason: format!("revision prefix body is not digestible: {err}"),
1229 })?;
1230 if parent_prefix_digest != revision_prefix_digest {
1231 return Err(SessionStoreError::InvalidTranscriptRewrite {
1232 id: incoming.id().clone(),
1233 reason: "rewrite revision changed messages before the selected span".to_string(),
1234 });
1235 }
1236 let parent_suffix_digest =
1237 transcript_messages_digest(&parent_body.messages[end..]).map_err(|err| {
1238 SessionStoreError::InvalidTranscriptRewrite {
1239 id: incoming.id().clone(),
1240 reason: format!("parent suffix body is not digestible: {err}"),
1241 }
1242 })?;
1243 let revision_suffix_digest =
1244 transcript_messages_digest(&revision_body.messages[replacement_end..]).map_err(|err| {
1245 SessionStoreError::InvalidTranscriptRewrite {
1246 id: incoming.id().clone(),
1247 reason: format!("revision suffix body is not digestible: {err}"),
1248 }
1249 })?;
1250 if parent_suffix_digest != revision_suffix_digest {
1251 return Err(SessionStoreError::InvalidTranscriptRewrite {
1252 id: incoming.id().clone(),
1253 reason: "rewrite revision changed messages after the selected span".to_string(),
1254 });
1255 }
1256 let replacement_digest = transcript_messages_digest(
1257 &revision_body.messages[start..replacement_end],
1258 )
1259 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1260 id: incoming.id().clone(),
1261 reason: format!("replacement span body is not digestible: {err}"),
1262 })?;
1263 if replacement_digest != commit.replacement_digest {
1264 return Err(SessionStoreError::InvalidTranscriptRewrite {
1265 id: incoming.id().clone(),
1266 reason: format!(
1267 "replacement span digest {replacement_digest} does not match commit digest {}",
1268 commit.replacement_digest
1269 ),
1270 });
1271 }
1272 Ok(())
1273}
1274
1275impl From<serde_json::Error> for SessionStoreError {
1276 fn from(e: serde_json::Error) -> Self {
1277 Self::Serialization(e.to_string())
1278 }
1279}
1280
1281#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
1306#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
1307pub trait SessionStore: Send + Sync {
1308 async fn save(&self, session: &Session) -> Result<(), SessionStoreError>;
1314
1315 async fn save_transcript_rewrite(
1321 &self,
1322 session: &Session,
1323 commit: &TranscriptRewriteCommit,
1324 ) -> Result<(), SessionStoreError> {
1325 let _ = (session, commit);
1326 Err(SessionStoreError::Internal(
1327 "save_transcript_rewrite is not supported by this SessionStore".to_string(),
1328 ))
1329 }
1330
1331 async fn save_authoritative_projection(
1340 &self,
1341 session: &Session,
1342 ) -> Result<(), SessionStoreError> {
1343 self.save(session).await
1344 }
1345
1346 async fn save_authoritative_projection_if_current_revision(
1349 &self,
1350 session: &Session,
1351 expected_current_revision: Option<String>,
1352 ) -> Result<(), SessionStoreError> {
1353 let _ = (session, expected_current_revision);
1354 Err(SessionStoreError::Internal(
1355 "save_authoritative_projection_if_current_revision is not supported by this SessionStore"
1356 .to_string(),
1357 ))
1358 }
1359
1360 async fn load(&self, id: &SessionId) -> Result<Option<Session>, SessionStoreError>;
1362
1363 async fn list(&self, filter: SessionFilter) -> Result<Vec<SessionMeta>, SessionStoreError>;
1365
1366 async fn delete(&self, id: &SessionId) -> Result<(), SessionStoreError>;
1368
1369 async fn delete_if_current_revision(
1372 &self,
1373 id: &SessionId,
1374 expected_current_revision: &str,
1375 ) -> Result<bool, SessionStoreError>;
1376
1377 async fn exists(&self, id: &SessionId) -> Result<bool, SessionStoreError> {
1379 Ok(self.load(id).await?.is_some())
1380 }
1381}
1382
1383#[cfg(test)]
1384mod tests {
1385 use super::*;
1386 use crate::types::{
1387 AssistantBlock, BlockAssistantMessage, StopReason, SystemMessage, SystemNoticeBlock,
1388 SystemNoticeKind, SystemNoticeMessage, UserMessage,
1389 };
1390
1391 #[test]
1398 #[allow(clippy::expect_used)]
1399 fn classify_live_session_authority_is_decided_by_machine() {
1400 use crate::session_document::{
1401 LiveSessionAuthorityKind, LiveSessionAuthorityReason, SessionDocumentEffect,
1402 SessionDocumentMachineAuthority,
1403 };
1404
1405 fn classify(
1406 stored_transcript_diverged: bool,
1407 live_has_uncommitted_transcript: bool,
1408 runtime_system_context_diverged: bool,
1409 stored_is_archived: bool,
1410 ) -> (LiveSessionAuthorityKind, LiveSessionAuthorityReason) {
1411 let mut authority = SessionDocumentMachineAuthority::new();
1412 let effects = authority
1413 .classify_live_session_authority(
1414 stored_transcript_diverged,
1415 live_has_uncommitted_transcript,
1416 runtime_system_context_diverged,
1417 stored_is_archived,
1418 )
1419 .expect("classifier must resolve a verdict");
1420 effects
1421 .iter()
1422 .find_map(|effect| match effect {
1423 SessionDocumentEffect::LiveSessionAuthorityClassified { authority, reason } => {
1424 Some((*authority, *reason))
1425 }
1426 _ => None,
1427 })
1428 .expect("classifier must emit a verdict")
1429 }
1430
1431 let (kind, _) = classify(false, false, false, false);
1433 assert_eq!(kind, LiveSessionAuthorityKind::LiveAuthoritative);
1434
1435 assert_eq!(
1437 classify(true, false, false, false),
1438 (
1439 LiveSessionAuthorityKind::DurableAuthoritative,
1440 LiveSessionAuthorityReason::StoredTranscriptRevisionDiverged
1441 ),
1442 );
1443 assert_eq!(
1444 classify(false, true, false, false),
1445 (
1446 LiveSessionAuthorityKind::DurableAuthoritative,
1447 LiveSessionAuthorityReason::LiveUncommittedTranscript
1448 ),
1449 );
1450 assert_eq!(
1451 classify(false, false, true, false),
1452 (
1453 LiveSessionAuthorityKind::DurableAuthoritative,
1454 LiveSessionAuthorityReason::RuntimeSystemContextDiverged
1455 ),
1456 );
1457 assert_eq!(
1458 classify(false, false, false, true),
1459 (
1460 LiveSessionAuthorityKind::DurableAuthoritative,
1461 LiveSessionAuthorityReason::StoredArchived
1462 ),
1463 );
1464
1465 assert_eq!(
1468 classify(true, true, true, true),
1469 (
1470 LiveSessionAuthorityKind::DurableAuthoritative,
1471 LiveSessionAuthorityReason::StoredArchived
1472 ),
1473 );
1474 assert_eq!(
1476 classify(true, true, true, false),
1477 (
1478 LiveSessionAuthorityKind::DurableAuthoritative,
1479 LiveSessionAuthorityReason::LiveUncommittedTranscript
1480 ),
1481 );
1482 assert_eq!(
1484 classify(true, false, true, false),
1485 (
1486 LiveSessionAuthorityKind::DurableAuthoritative,
1487 LiveSessionAuthorityReason::RuntimeSystemContextDiverged
1488 ),
1489 );
1490 }
1491
1492 #[test]
1493 fn append_only_guard_rejects_leading_system_message_replacement() {
1494 let mut previous = Session::new();
1495 previous.push(Message::System(SystemMessage::new("original system")));
1496 previous.push(Message::User(UserMessage::text("hello".to_string())));
1497
1498 let mut incoming = previous.clone();
1499 let rewrite_result = incoming.replace_messages_internal(
1500 vec![
1501 Message::System(SystemMessage::new("rewritten system")),
1502 Message::User(UserMessage::text("hello".to_string())),
1503 ],
1504 crate::TranscriptRewriteReason::new("unit-test"),
1505 );
1506 assert!(
1507 rewrite_result.is_ok(),
1508 "typed rewrite should be constructible: {rewrite_result:?}"
1509 );
1510
1511 assert!(matches!(
1512 append_only_save_guard(&incoming, Some(&previous)),
1513 Err(SessionStoreError::TranscriptContinuityViolation { .. })
1514 ));
1515 }
1516
1517 #[test]
1518 fn append_only_guard_accepts_runtime_system_context_append()
1519 -> Result<(), Box<dyn std::error::Error>> {
1520 let mut previous = Session::new();
1521 previous.push(Message::System(SystemMessage::new("base system")));
1522 previous.push(Message::User(UserMessage::text("hello".to_string())));
1523
1524 let mut incoming = previous.clone();
1525 incoming.set_system_prompt_with_source(
1529 format!(
1530 "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nsource: unit-test\n\nextra context"
1531 ),
1532 crate::session_durable_config_authority::SessionSystemPromptSource::RuntimeContextAppend,
1533 )?;
1534
1535 assert!(append_only_save_guard(&incoming, Some(&previous)).is_ok());
1536 Ok(())
1537 }
1538
1539 #[test]
1540 fn append_only_guard_rejects_append_shaped_prompt_without_runtime_context_marker() {
1541 let mut previous = Session::new();
1542 previous.push(Message::System(SystemMessage::new("base system")));
1543 previous.push(Message::User(UserMessage::text("hello".to_string())));
1544
1545 let mut incoming = previous.clone();
1549 incoming.set_system_prompt(format!(
1550 "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nsource: forged\n\nextra context"
1551 ));
1552
1553 assert!(matches!(
1554 append_only_save_guard(&incoming, Some(&previous)),
1555 Err(SessionStoreError::TranscriptContinuityViolation { .. })
1556 ));
1557 }
1558
1559 #[test]
1560 fn append_only_guard_accepts_system_timestamp_refresh_without_content_change() {
1561 let mut previous = Session::new();
1562 previous.push(Message::System(SystemMessage::new("base system")));
1563
1564 let mut incoming = previous.clone();
1565 incoming.set_system_prompt("base system".to_string());
1566
1567 assert!(append_only_save_guard(&incoming, Some(&previous)).is_ok());
1568 }
1569
1570 #[test]
1571 fn run_boundary_guard_accepts_compaction_after_uncheckpointed_runtime_append()
1572 -> Result<(), Box<dyn std::error::Error>> {
1573 let mut previous = Session::new();
1574 previous.push(Message::System(SystemMessage::new("base system")));
1575 previous.push(Message::User(UserMessage::text("turn one".to_string())));
1576 previous.push(Message::BlockAssistant(BlockAssistantMessage {
1577 blocks: vec![AssistantBlock::Text {
1578 text: "answer one".to_string(),
1579 meta: None,
1580 }],
1581 stop_reason: StopReason::EndTurn,
1582 identity: crate::types::TranscriptMessageIdentity::default(),
1583 created_at: crate::types::message_timestamp_now(),
1584 }));
1585
1586 let mut parent = previous.clone();
1587 parent.set_system_prompt("refreshed runtime system projection".to_string());
1588 parent.push(Message::User(UserMessage::text(
1589 "runtime-only turn".to_string(),
1590 )));
1591 parent.push(Message::BlockAssistant(BlockAssistantMessage {
1592 blocks: vec![AssistantBlock::Text {
1593 text: "runtime-only answer".to_string(),
1594 meta: None,
1595 }],
1596 stop_reason: StopReason::EndTurn,
1597 identity: crate::types::TranscriptMessageIdentity::default(),
1598 created_at: crate::types::message_timestamp_now(),
1599 }));
1600 let parent_revision = parent.transcript_revision()?;
1601
1602 let mut incoming = parent.clone();
1603 let mut replacement = vec![
1604 parent.messages()[0].clone(),
1605 Message::User(UserMessage::compaction_summary(
1606 "[Context compacted] summary".to_string(),
1607 )),
1608 ];
1609 replacement.extend_from_slice(&parent.messages()[1..]);
1610 incoming.commit_transcript_rewrite(
1611 TranscriptRewriteSelection::MessageRange {
1612 start: 0,
1613 end: parent.messages().len(),
1614 },
1615 replacement,
1616 crate::TranscriptRewriteReason::new("compaction"),
1617 Some("meerkat-core".to_string()),
1618 Some(parent_revision),
1619 )?;
1620
1621 assert!(matches!(
1622 append_only_save_guard(&incoming, Some(&previous)),
1623 Err(SessionStoreError::TranscriptContinuityViolation { .. })
1624 ));
1625 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
1626 Ok(())
1627 }
1628
1629 #[test]
1630 fn run_boundary_guard_accepts_compaction_with_retained_tail_window()
1631 -> Result<(), Box<dyn std::error::Error>> {
1632 let mut previous = Session::new();
1633 previous.push(Message::System(SystemMessage::new("base system")));
1634 previous.push(Message::User(UserMessage::text("turn one".to_string())));
1635 previous.push(Message::BlockAssistant(BlockAssistantMessage::new(
1636 vec![crate::types::AssistantBlock::Text {
1637 text: "answer one".to_string(),
1638 meta: None,
1639 }],
1640 StopReason::EndTurn,
1641 )));
1642
1643 let mut parent = previous.clone();
1644 parent.set_system_prompt("refreshed runtime system projection".to_string());
1645 parent.push(Message::SystemNotice(SystemNoticeMessage::new(
1646 SystemNoticeKind::Comms,
1647 "peer response queued",
1648 )));
1649 let parent_revision = parent.transcript_revision()?;
1650
1651 let mut incoming = parent.clone();
1652 let mut replacement = vec![
1653 parent.messages()[0].clone(),
1654 Message::User(UserMessage::compaction_summary(
1655 "[Context compacted] summary".to_string(),
1656 )),
1657 ];
1658 replacement.extend_from_slice(&parent.messages()[1..]);
1659 incoming.commit_transcript_rewrite(
1660 TranscriptRewriteSelection::MessageRange {
1661 start: 0,
1662 end: parent.messages().len(),
1663 },
1664 replacement,
1665 crate::TranscriptRewriteReason::new("compaction"),
1666 Some("meerkat-core".to_string()),
1667 Some(parent_revision),
1668 )?;
1669
1670 assert!(matches!(
1671 append_only_save_guard(&incoming, Some(&previous)),
1672 Err(SessionStoreError::TranscriptContinuityViolation { .. })
1673 ));
1674 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
1675 Ok(())
1676 }
1677
1678 #[test]
1679 fn run_boundary_guard_rejects_commitless_history_parent_edge()
1680 -> Result<(), Box<dyn std::error::Error>> {
1681 let mut previous = Session::new();
1682 previous.push(Message::System(SystemMessage::new("base system")));
1683 previous.push(Message::User(UserMessage::text("turn one".to_string())));
1684 let previous_revision = previous.transcript_revision()?;
1685
1686 let mut incoming = previous.clone();
1687 incoming.set_system_prompt("forged replacement system".to_string());
1688 let incoming_revision = incoming.transcript_revision()?;
1689 let history = TranscriptHistoryState {
1690 head: incoming_revision.clone(),
1691 commits: Vec::new(),
1692 revisions: vec![
1693 crate::TranscriptRevisionBody {
1694 revision: previous_revision,
1695 parent_revision: None,
1696 messages: previous.messages().to_vec(),
1697 created_at: previous.updated_at(),
1698 },
1699 crate::TranscriptRevisionBody {
1700 revision: incoming_revision,
1701 parent_revision: Some(previous.transcript_revision()?),
1702 messages: incoming.messages().to_vec(),
1703 created_at: incoming.updated_at(),
1704 },
1705 ],
1706 };
1707 incoming.set_metadata_unchecked_for_test(
1708 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1709 serde_json::to_value(history)?,
1710 );
1711
1712 assert!(matches!(
1713 append_only_save_guard(&incoming, Some(&previous)),
1714 Err(SessionStoreError::TranscriptContinuityViolation { .. })
1715 ));
1716 assert!(matches!(
1717 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
1718 Err(SessionStoreError::TranscriptContinuityViolation { .. }
1719 | SessionStoreError::MonotonicityViolation { .. })
1720 ));
1721 Ok(())
1722 }
1723
1724 #[test]
1725 fn append_only_guard_rejects_history_head_that_does_not_match_current_messages()
1726 -> Result<(), Box<dyn std::error::Error>> {
1727 let mut previous = Session::new();
1728 previous.push(Message::User(UserMessage::text("persisted".to_string())));
1729
1730 let mut incoming = previous.clone();
1731 incoming.push(Message::User(UserMessage::text("append".to_string())));
1732 let poisoned_messages = vec![Message::User(UserMessage::text(
1733 "unrelated poisoned history".to_string(),
1734 ))];
1735 let poisoned_revision = transcript_messages_digest(&poisoned_messages)?;
1736 incoming.set_metadata_unchecked_for_test(
1737 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1738 serde_json::to_value(TranscriptHistoryState {
1739 head: poisoned_revision.clone(),
1740 commits: Vec::new(),
1741 revisions: vec![crate::TranscriptRevisionBody {
1742 revision: poisoned_revision,
1743 parent_revision: None,
1744 messages: poisoned_messages,
1745 created_at: incoming.updated_at(),
1746 }],
1747 })?,
1748 );
1749
1750 assert!(matches!(
1751 append_only_save_guard(&incoming, Some(&previous)),
1752 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
1753 ));
1754 assert!(matches!(
1755 append_only_save_guard(&incoming, None),
1756 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
1757 ));
1758 assert!(matches!(
1759 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
1760 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
1761 ));
1762 Ok(())
1763 }
1764
1765 #[test]
1766 fn append_only_guard_rejects_new_rewrite_commits_on_plain_append()
1767 -> Result<(), Box<dyn std::error::Error>> {
1768 let mut previous = Session::new();
1769 previous.push(Message::User(UserMessage::text("persisted".to_string())));
1770 let previous_revision = previous.transcript_revision()?;
1771
1772 let mut incoming = previous.clone();
1773 let appended = Message::BlockAssistant(BlockAssistantMessage {
1774 blocks: vec![AssistantBlock::Text {
1775 text: "plain append".to_string(),
1776 meta: None,
1777 }],
1778 stop_reason: StopReason::EndTurn,
1779 identity: crate::types::TranscriptMessageIdentity::default(),
1780 created_at: crate::types::message_timestamp_now(),
1781 });
1782 incoming.commit_transcript_rewrite(
1783 TranscriptRewriteSelection::MessageRange { start: 1, end: 1 },
1784 vec![appended],
1785 crate::TranscriptRewriteReason::new("forged-append"),
1786 Some("unit-test".to_string()),
1787 Some(previous_revision),
1788 )?;
1789
1790 assert!(matches!(
1791 append_only_save_guard(&incoming, Some(&previous)),
1792 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
1793 ));
1794 Ok(())
1795 }
1796
1797 #[test]
1798 fn append_only_guard_rejects_first_save_with_rewrite_commits()
1799 -> Result<(), Box<dyn std::error::Error>> {
1800 let mut incoming = Session::new();
1801 incoming.push(Message::User(UserMessage::text("seed".to_string())));
1802 let parent_messages = incoming.messages().to_vec();
1803 let parent_updated_at = incoming.updated_at();
1804 let parent_revision = incoming.transcript_revision()?;
1805 let commit = incoming.commit_transcript_rewrite(
1806 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1807 vec![Message::User(UserMessage::text(
1808 "compacted seed".to_string(),
1809 ))],
1810 crate::TranscriptRewriteReason::new("compaction"),
1811 Some("meerkat-core".to_string()),
1812 Some(parent_revision),
1813 )?;
1814 let incoming_revision = incoming.transcript_revision()?;
1815 let commit_parent_revision = commit.parent_revision.clone();
1816 incoming.set_metadata_unchecked_for_test(
1817 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1818 serde_json::to_value(TranscriptHistoryState {
1819 head: incoming_revision.clone(),
1820 commits: vec![commit],
1821 revisions: vec![
1822 crate::TranscriptRevisionBody {
1823 revision: commit_parent_revision.clone(),
1824 parent_revision: None,
1825 messages: parent_messages,
1826 created_at: parent_updated_at,
1827 },
1828 crate::TranscriptRevisionBody {
1829 revision: incoming_revision,
1830 parent_revision: Some(commit_parent_revision),
1831 messages: incoming.messages().to_vec(),
1832 created_at: incoming.updated_at(),
1833 },
1834 ],
1835 })?,
1836 );
1837
1838 assert!(matches!(
1839 append_only_save_guard(&incoming, None),
1840 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
1841 ));
1842 Ok(())
1843 }
1844
1845 #[test]
1846 fn transcript_rewrite_guard_rejects_poisoned_history_graph()
1847 -> Result<(), Box<dyn std::error::Error>> {
1848 let mut previous = Session::new();
1849 previous.push(Message::User(UserMessage::text("persisted".to_string())));
1850 let parent_revision = previous.transcript_revision()?;
1851
1852 let mut first = previous.clone();
1853 let first_commit = first.commit_transcript_rewrite(
1854 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1855 vec![Message::User(UserMessage::text(
1856 "compacted persisted".to_string(),
1857 ))],
1858 crate::TranscriptRewriteReason::new("compaction"),
1859 Some("unit-test".to_string()),
1860 Some(parent_revision),
1861 )?;
1862 let first_snapshot = first.clone();
1863
1864 first.commit_transcript_rewrite(
1865 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1866 vec![Message::User(UserMessage::text(
1867 "uncommitted poisoned fork".to_string(),
1868 ))],
1869 crate::TranscriptRewriteReason::new("poison"),
1870 Some("unit-test".to_string()),
1871 Some(first_commit.revision.clone()),
1872 )?;
1873 let mut poisoned_state = first
1874 .transcript_history_state()?
1875 .ok_or_else(|| "second rewrite should retain history state".to_string())?;
1876 poisoned_state.head = first_commit.revision.clone();
1877
1878 let mut poisoned = first_snapshot;
1879 poisoned.set_metadata_unchecked_for_test(
1880 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1881 serde_json::to_value(poisoned_state)?,
1882 );
1883
1884 assert!(matches!(
1885 transcript_rewrite_save_guard(&poisoned, Some(&previous), &first_commit),
1886 Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
1887 if reason.contains("incoming transcript history state is malformed")
1888 ));
1889 Ok(())
1890 }
1891
1892 #[test]
1893 fn authoritative_projection_guard_rejects_changed_persisted_revision()
1894 -> Result<(), Box<dyn std::error::Error>> {
1895 let mut previous = Session::new();
1896 previous.push(Message::User(UserMessage::text("persisted A".to_string())));
1897 let expected_revision = previous.transcript_revision()?;
1898
1899 let mut current = previous.clone();
1900 current.push(Message::BlockAssistant(BlockAssistantMessage {
1901 blocks: vec![AssistantBlock::Text {
1902 text: "persisted B".to_string(),
1903 meta: None,
1904 }],
1905 stop_reason: StopReason::EndTurn,
1906 identity: crate::types::TranscriptMessageIdentity::default(),
1907 created_at: crate::types::message_timestamp_now(),
1908 }));
1909 let mut incoming = previous.clone();
1910 incoming.push(Message::User(UserMessage::text(
1911 "incoming from A".to_string(),
1912 )));
1913
1914 assert!(matches!(
1915 authoritative_projection_current_revision_guard(
1916 &incoming,
1917 Some(¤t),
1918 Some(&expected_revision)
1919 ),
1920 Err(SessionStoreError::TranscriptContinuityViolation { .. })
1921 ));
1922 Ok(())
1923 }
1924
1925 #[test]
1926 fn append_only_guard_rejects_rewrite_commits_on_first_save()
1927 -> Result<(), Box<dyn std::error::Error>> {
1928 let mut incoming = Session::new();
1929 incoming.push(Message::User(UserMessage::text("persisted".to_string())));
1930 let parent_revision = incoming.transcript_revision()?;
1931 incoming.commit_transcript_rewrite(
1932 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1933 vec![Message::User(UserMessage::text("rewritten".to_string()))],
1934 crate::TranscriptRewriteReason::new("forged-first-save"),
1935 Some("unit-test".to_string()),
1936 Some(parent_revision),
1937 )?;
1938
1939 assert!(matches!(
1940 append_only_save_guard(&incoming, None),
1941 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
1942 ));
1943 Ok(())
1944 }
1945
1946 #[test]
1947 fn append_only_guard_rejects_commitless_history_on_first_save()
1948 -> Result<(), Box<dyn std::error::Error>> {
1949 let mut incoming = Session::new();
1950 incoming.push(Message::User(UserMessage::text("persisted".to_string())));
1951 let incoming_revision = incoming.transcript_revision()?;
1952 incoming.set_metadata_unchecked_for_test(
1953 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1954 serde_json::to_value(TranscriptHistoryState {
1955 head: incoming_revision.clone(),
1956 commits: Vec::new(),
1957 revisions: vec![crate::TranscriptRevisionBody {
1958 revision: incoming_revision,
1959 parent_revision: None,
1960 messages: incoming.messages().to_vec(),
1961 created_at: incoming.updated_at(),
1962 }],
1963 })?,
1964 );
1965
1966 assert!(matches!(
1967 append_only_save_guard(&incoming, None),
1968 Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
1969 if reason.contains("first save would seed transcript history state")
1970 ));
1971 Ok(())
1972 }
1973
1974 #[test]
1975 fn append_only_guard_rejects_commitless_history_seed_on_plain_append()
1976 -> Result<(), Box<dyn std::error::Error>> {
1977 let mut previous = Session::new();
1978 previous.push(Message::User(UserMessage::text("persisted".to_string())));
1979 let previous_revision = previous.transcript_revision()?;
1980
1981 let mut incoming = previous.clone();
1982 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
1983 blocks: vec![AssistantBlock::Text {
1984 text: "plain append".to_string(),
1985 meta: None,
1986 }],
1987 stop_reason: StopReason::EndTurn,
1988 identity: crate::types::TranscriptMessageIdentity::default(),
1989 created_at: crate::types::message_timestamp_now(),
1990 }));
1991 let incoming_revision = incoming.transcript_revision()?;
1992 incoming.set_metadata_unchecked_for_test(
1993 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1994 serde_json::to_value(TranscriptHistoryState {
1995 head: incoming_revision.clone(),
1996 commits: Vec::new(),
1997 revisions: vec![
1998 crate::TranscriptRevisionBody {
1999 revision: previous_revision,
2000 parent_revision: None,
2001 messages: previous.messages().to_vec(),
2002 created_at: previous.updated_at(),
2003 },
2004 crate::TranscriptRevisionBody {
2005 revision: incoming_revision,
2006 parent_revision: Some(previous.transcript_revision()?),
2007 messages: incoming.messages().to_vec(),
2008 created_at: incoming.updated_at(),
2009 },
2010 ],
2011 })?,
2012 );
2013
2014 assert!(matches!(
2015 append_only_save_guard(&incoming, Some(&previous)),
2016 Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
2017 if reason.contains("append-only save would seed transcript history state")
2018 ));
2019 Ok(())
2020 }
2021
2022 #[test]
2023 fn run_boundary_guard_accepts_commitless_history_seed_on_plain_append()
2024 -> Result<(), Box<dyn std::error::Error>> {
2025 let mut previous = Session::new();
2026 previous.push(Message::User(UserMessage::text("persisted".to_string())));
2027 let previous_revision = previous.transcript_revision()?;
2028
2029 let mut incoming = previous.clone();
2030 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2031 blocks: vec![AssistantBlock::Text {
2032 text: "plain append".to_string(),
2033 meta: None,
2034 }],
2035 stop_reason: StopReason::EndTurn,
2036 identity: crate::types::TranscriptMessageIdentity::default(),
2037 created_at: crate::types::message_timestamp_now(),
2038 }));
2039 let incoming_revision = incoming.transcript_revision()?;
2040 incoming.set_metadata_unchecked_for_test(
2041 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2042 serde_json::to_value(TranscriptHistoryState {
2043 head: incoming_revision.clone(),
2044 commits: Vec::new(),
2045 revisions: vec![
2046 crate::TranscriptRevisionBody {
2047 revision: previous_revision.clone(),
2048 parent_revision: None,
2049 messages: previous.messages().to_vec(),
2050 created_at: previous.updated_at(),
2051 },
2052 crate::TranscriptRevisionBody {
2053 revision: incoming_revision,
2054 parent_revision: Some(previous_revision),
2055 messages: incoming.messages().to_vec(),
2056 created_at: incoming.updated_at(),
2057 },
2058 ],
2059 })?,
2060 );
2061
2062 assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2063 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2064 Ok(())
2065 }
2066
2067 #[test]
2068 fn run_boundary_guard_accepts_retained_history_seed_on_plain_append()
2069 -> Result<(), Box<dyn std::error::Error>> {
2070 let mut original = Session::new();
2071 original.push(Message::User(UserMessage::text("verbose seed".to_string())));
2072 let original_revision = original.transcript_revision()?;
2073
2074 let mut previous = original.clone();
2075 previous.commit_transcript_rewrite(
2076 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2077 vec![Message::User(UserMessage::text(
2078 "compacted seed".to_string(),
2079 ))],
2080 crate::TranscriptRewriteReason::new("compaction"),
2081 Some("meerkat-core".to_string()),
2082 Some(original_revision),
2083 )?;
2084 let previous_with_history = previous.clone();
2085 previous.clear_transcript_history_state();
2086
2087 let mut incoming = previous_with_history;
2088 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2089 blocks: vec![AssistantBlock::Text {
2090 text: "plain append after retained history".to_string(),
2091 meta: None,
2092 }],
2093 stop_reason: StopReason::EndTurn,
2094 identity: crate::types::TranscriptMessageIdentity::default(),
2095 created_at: crate::types::message_timestamp_now(),
2096 }));
2097
2098 assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2099 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2100 Ok(())
2101 }
2102
2103 #[test]
2104 fn run_boundary_guard_accepts_commitless_history_seed_on_first_snapshot()
2105 -> Result<(), Box<dyn std::error::Error>> {
2106 let mut incoming = Session::new();
2107 incoming.push(Message::User(UserMessage::text("persisted".to_string())));
2108 let incoming_revision = incoming.transcript_revision()?;
2109 incoming.set_metadata_unchecked_for_test(
2110 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2111 serde_json::to_value(TranscriptHistoryState {
2112 head: incoming_revision.clone(),
2113 commits: Vec::new(),
2114 revisions: vec![crate::TranscriptRevisionBody {
2115 revision: incoming_revision,
2116 parent_revision: None,
2117 messages: incoming.messages().to_vec(),
2118 created_at: incoming.updated_at(),
2119 }],
2120 })?,
2121 );
2122
2123 assert!(append_only_save_guard(&incoming, None).is_err());
2124 assert!(run_boundary_snapshot_save_guard(&incoming, None).is_ok());
2125 Ok(())
2126 }
2127
2128 #[test]
2129 fn run_boundary_guard_accepts_commitless_history_seed_on_initial_multi_revision_snapshot()
2130 -> Result<(), Box<dyn std::error::Error>> {
2131 let mut base = Session::new();
2132 base.push(Message::User(UserMessage::text("first".to_string())));
2133 let base_revision = base.transcript_revision()?;
2134
2135 let mut incoming = base.clone();
2136 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2137 blocks: vec![AssistantBlock::Text {
2138 text: "second".to_string(),
2139 meta: None,
2140 }],
2141 stop_reason: StopReason::EndTurn,
2142 identity: crate::types::TranscriptMessageIdentity::default(),
2143 created_at: crate::types::message_timestamp_now(),
2144 }));
2145 let incoming_revision = incoming.transcript_revision()?;
2146 incoming.set_metadata_unchecked_for_test(
2147 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2148 serde_json::to_value(TranscriptHistoryState {
2149 head: incoming_revision.clone(),
2150 commits: Vec::new(),
2151 revisions: vec![
2152 crate::TranscriptRevisionBody {
2153 revision: base_revision.clone(),
2154 parent_revision: None,
2155 messages: base.messages().to_vec(),
2156 created_at: base.updated_at(),
2157 },
2158 crate::TranscriptRevisionBody {
2159 revision: incoming_revision,
2160 parent_revision: Some(base_revision),
2161 messages: incoming.messages().to_vec(),
2162 created_at: incoming.updated_at(),
2163 },
2164 ],
2165 })?,
2166 );
2167
2168 assert!(append_only_save_guard(&incoming, None).is_err());
2169 assert!(run_boundary_snapshot_save_guard(&incoming, None).is_ok());
2170 Ok(())
2171 }
2172
2173 #[test]
2174 fn append_only_guard_rejects_new_rewrite_commits_on_system_context_append()
2175 -> Result<(), Box<dyn std::error::Error>> {
2176 let mut previous = Session::new();
2177 previous.push(Message::System(SystemMessage::new("base system")));
2178 previous.push(Message::User(UserMessage::text("persisted".to_string())));
2179 let mut incoming = previous.clone();
2180 incoming.set_system_prompt_with_source(
2181 format!(
2182 "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nsource: unit-test\n\nextra context"
2183 ),
2184 crate::session_durable_config_authority::SessionSystemPromptSource::RuntimeContextAppend,
2185 )?;
2186 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2187 blocks: vec![AssistantBlock::Text {
2188 text: "plain append".to_string(),
2189 meta: None,
2190 }],
2191 stop_reason: StopReason::EndTurn,
2192 identity: crate::types::TranscriptMessageIdentity::default(),
2193 created_at: crate::types::message_timestamp_now(),
2194 }));
2195 let incoming_revision = incoming.transcript_revision()?;
2196 incoming.set_metadata_unchecked_for_test(
2197 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2198 serde_json::to_value(TranscriptHistoryState {
2199 head: incoming_revision.clone(),
2200 commits: vec![TranscriptRewriteCommit {
2201 parent_revision: previous.transcript_revision()?,
2202 revision: incoming_revision.clone(),
2203 selection: TranscriptRewriteSelection::MessageRange { start: 0, end: 0 },
2204 original_span_digest: transcript_messages_digest(&[])?,
2205 replacement_digest: transcript_messages_digest(&[])?,
2206 messages_before: previous.messages().len(),
2207 messages_after: incoming.messages().len(),
2208 reason: crate::TranscriptRewriteReason::new("forged"),
2209 actor: Some("unit-test".to_string()),
2210 committed_at: incoming.updated_at(),
2211 }],
2212 revisions: vec![crate::TranscriptRevisionBody {
2213 revision: incoming_revision,
2214 parent_revision: None,
2215 messages: incoming.messages().to_vec(),
2216 created_at: incoming.updated_at(),
2217 }],
2218 })?,
2219 );
2220
2221 assert!(matches!(
2222 append_only_save_guard(&incoming, Some(&previous)),
2223 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2224 ));
2225 Ok(())
2226 }
2227
2228 #[test]
2229 fn append_only_guard_rejects_new_rewrite_commits_on_transient_notice_cleanup()
2230 -> Result<(), Box<dyn std::error::Error>> {
2231 let mut previous = Session::new();
2232 previous.push(Message::SystemNotice(SystemNoticeMessage::new(
2233 SystemNoticeKind::Comms,
2234 "transient peer delivery notice",
2235 )));
2236 previous.push(Message::User(UserMessage::text("persisted".to_string())));
2237
2238 let mut incoming = Session::new();
2239 incoming.push(Message::User(UserMessage::text("persisted".to_string())));
2240 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2241 blocks: vec![AssistantBlock::Text {
2242 text: "plain append after notice cleanup".to_string(),
2243 meta: None,
2244 }],
2245 stop_reason: StopReason::EndTurn,
2246 identity: crate::types::TranscriptMessageIdentity::default(),
2247 created_at: crate::types::message_timestamp_now(),
2248 }));
2249 let incoming_revision = incoming.transcript_revision()?;
2250 incoming.set_metadata_unchecked_for_test(
2251 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2252 serde_json::to_value(TranscriptHistoryState {
2253 head: incoming_revision.clone(),
2254 commits: vec![TranscriptRewriteCommit {
2255 parent_revision: previous.transcript_revision()?,
2256 revision: incoming_revision.clone(),
2257 selection: TranscriptRewriteSelection::MessageRange { start: 0, end: 0 },
2258 original_span_digest: transcript_messages_digest(&[])?,
2259 replacement_digest: transcript_messages_digest(&[])?,
2260 messages_before: previous.messages().len(),
2261 messages_after: incoming.messages().len(),
2262 reason: crate::TranscriptRewriteReason::new("forged"),
2263 actor: Some("unit-test".to_string()),
2264 committed_at: incoming.updated_at(),
2265 }],
2266 revisions: vec![crate::TranscriptRevisionBody {
2267 revision: incoming_revision,
2268 parent_revision: None,
2269 messages: incoming.messages().to_vec(),
2270 created_at: incoming.updated_at(),
2271 }],
2272 })?,
2273 );
2274
2275 assert!(matches!(
2276 append_only_save_guard(&incoming, Some(&previous)),
2277 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2278 ));
2279 Ok(())
2280 }
2281
2282 #[test]
2283 fn run_boundary_guard_accepts_generated_context_summary_before_retained_tail()
2284 -> Result<(), Box<dyn std::error::Error>> {
2285 let mut previous = Session::new();
2286 previous.push(Message::System(SystemMessage::new(
2287 "runtime system before context refresh",
2288 )));
2289 previous.push(Message::User(UserMessage::text(
2290 "Turn 1 request".to_string(),
2291 )));
2292 previous.push(Message::BlockAssistant(BlockAssistantMessage {
2293 blocks: vec![AssistantBlock::Text {
2294 text: "Turn 1 answer".to_string(),
2295 meta: None,
2296 }],
2297 stop_reason: StopReason::EndTurn,
2298 identity: crate::types::TranscriptMessageIdentity::default(),
2299 created_at: crate::types::message_timestamp_now(),
2300 }));
2301
2302 let mut incoming = Session::with_id(previous.id().clone());
2303 incoming.push(Message::System(SystemMessage::new(
2304 "runtime system after context refresh",
2305 )));
2306 incoming.push(Message::User(UserMessage::text(
2307 "Verbose context that will be compacted".to_string(),
2308 )));
2309 for message in previous.messages()[1..].iter().cloned() {
2310 incoming.push(message);
2311 }
2312 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2313 blocks: vec![AssistantBlock::Text {
2314 text: "Turn 2 generated answer".to_string(),
2315 meta: None,
2316 }],
2317 stop_reason: StopReason::EndTurn,
2318 identity: crate::types::TranscriptMessageIdentity::default(),
2319 created_at: crate::types::message_timestamp_now(),
2320 }));
2321 let parent_revision = incoming.transcript_revision()?;
2322 incoming.commit_transcript_rewrite(
2323 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2324 vec![Message::User(UserMessage::compaction_summary(
2325 "[Context compacted] Earlier runtime context".to_string(),
2326 ))],
2327 crate::TranscriptRewriteReason::new("compaction"),
2328 Some("meerkat-core".to_string()),
2329 Some(parent_revision),
2330 )?;
2331
2332 assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2333 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2334 Ok(())
2335 }
2336
2337 #[test]
2338 fn run_boundary_guard_rejects_context_summary_tail_without_compaction_summary_marker()
2339 -> Result<(), Box<dyn std::error::Error>> {
2340 let mut previous = Session::new();
2341 previous.push(Message::System(SystemMessage::new(
2342 "runtime system before context refresh",
2343 )));
2344 previous.push(Message::User(UserMessage::text(
2345 "Turn 1 request".to_string(),
2346 )));
2347 previous.push(Message::BlockAssistant(BlockAssistantMessage {
2348 blocks: vec![AssistantBlock::Text {
2349 text: "Turn 1 answer".to_string(),
2350 meta: None,
2351 }],
2352 stop_reason: StopReason::EndTurn,
2353 identity: crate::types::TranscriptMessageIdentity::default(),
2354 created_at: crate::types::message_timestamp_now(),
2355 }));
2356
2357 let mut incoming = Session::with_id(previous.id().clone());
2358 incoming.push(Message::System(SystemMessage::new(
2359 "runtime system after context refresh",
2360 )));
2361 incoming.push(Message::User(UserMessage::text(
2362 "Verbose context that will be compacted".to_string(),
2363 )));
2364 for message in previous.messages()[1..].iter().cloned() {
2365 incoming.push(message);
2366 }
2367 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2368 blocks: vec![AssistantBlock::Text {
2369 text: "Turn 2 generated answer".to_string(),
2370 meta: None,
2371 }],
2372 stop_reason: StopReason::EndTurn,
2373 identity: crate::types::TranscriptMessageIdentity::default(),
2374 created_at: crate::types::message_timestamp_now(),
2375 }));
2376 let parent_revision = incoming.transcript_revision()?;
2377 incoming.commit_transcript_rewrite(
2381 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2382 vec![Message::User(UserMessage::text(
2383 "[Context compacted] Earlier runtime context".to_string(),
2384 ))],
2385 crate::TranscriptRewriteReason::new("compaction"),
2386 Some("meerkat-core".to_string()),
2387 Some(parent_revision),
2388 )?;
2389
2390 assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2391 assert!(matches!(
2392 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2393 Err(SessionStoreError::TranscriptContinuityViolation { .. }
2394 | SessionStoreError::MonotonicityViolation { .. })
2395 ));
2396 Ok(())
2397 }
2398
2399 #[test]
2404 fn run_boundary_guard_rejects_context_summary_tail_with_injected_context_marker()
2405 -> Result<(), Box<dyn std::error::Error>> {
2406 let mut previous = Session::new();
2407 previous.push(Message::System(SystemMessage::new(
2408 "runtime system before context refresh",
2409 )));
2410 previous.push(Message::User(UserMessage::text(
2411 "Turn 1 request".to_string(),
2412 )));
2413 previous.push(Message::BlockAssistant(BlockAssistantMessage {
2414 blocks: vec![AssistantBlock::Text {
2415 text: "Turn 1 answer".to_string(),
2416 meta: None,
2417 }],
2418 stop_reason: StopReason::EndTurn,
2419 identity: crate::types::TranscriptMessageIdentity::default(),
2420 created_at: crate::types::message_timestamp_now(),
2421 }));
2422
2423 let mut incoming = Session::with_id(previous.id().clone());
2424 incoming.push(Message::System(SystemMessage::new(
2425 "runtime system after context refresh",
2426 )));
2427 incoming.push(Message::User(UserMessage::text(
2428 "Verbose context that will be compacted".to_string(),
2429 )));
2430 for message in previous.messages()[1..].iter().cloned() {
2431 incoming.push(message);
2432 }
2433 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2434 blocks: vec![AssistantBlock::Text {
2435 text: "Turn 2 generated answer".to_string(),
2436 meta: None,
2437 }],
2438 stop_reason: StopReason::EndTurn,
2439 identity: crate::types::TranscriptMessageIdentity::default(),
2440 created_at: crate::types::message_timestamp_now(),
2441 }));
2442 let parent_revision = incoming.transcript_revision()?;
2443 incoming.commit_transcript_rewrite(
2448 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2449 vec![Message::User(UserMessage::injected_context(
2450 "[Context compacted] Earlier runtime context".to_string(),
2451 ))],
2452 crate::TranscriptRewriteReason::new("compaction"),
2453 Some("meerkat-core".to_string()),
2454 Some(parent_revision),
2455 )?;
2456
2457 assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2458 assert!(matches!(
2459 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2460 Err(SessionStoreError::TranscriptContinuityViolation { .. }
2461 | SessionStoreError::MonotonicityViolation { .. })
2462 ));
2463 Ok(())
2464 }
2465
2466 #[test]
2467 fn run_boundary_guard_rejects_runtime_parent_with_inserted_message_before_tail()
2468 -> Result<(), Box<dyn std::error::Error>> {
2469 let mut previous = Session::new();
2470 previous.push(Message::System(SystemMessage::new("base system")));
2471 previous.push(Message::User(UserMessage::text("turn one".to_string())));
2472 previous.push(Message::BlockAssistant(BlockAssistantMessage {
2473 blocks: vec![AssistantBlock::Text {
2474 text: "answer one".to_string(),
2475 meta: None,
2476 }],
2477 stop_reason: StopReason::EndTurn,
2478 identity: crate::types::TranscriptMessageIdentity::default(),
2479 created_at: crate::types::message_timestamp_now(),
2480 }));
2481
2482 let parent_messages = vec![
2483 Message::System(SystemMessage::new("refreshed runtime system projection")),
2484 Message::User(UserMessage::text(
2485 "injected before retained tail".to_string(),
2486 )),
2487 previous.messages()[1].clone(),
2488 previous.messages()[2].clone(),
2489 ];
2490 let parent_revision = transcript_messages_digest(&parent_messages)?;
2491 let mut parent = previous.clone();
2492 parent.apply_transcript_history_state(TranscriptHistoryState {
2493 head: parent_revision.clone(),
2494 commits: Vec::new(),
2495 revisions: vec![crate::TranscriptRevisionBody {
2496 revision: parent_revision,
2497 parent_revision: None,
2498 messages: parent_messages,
2499 created_at: parent.updated_at(),
2500 }],
2501 })?;
2502 let parent_revision = parent.transcript_revision()?;
2503
2504 let mut incoming = parent.clone();
2505 incoming.commit_transcript_rewrite(
2506 TranscriptRewriteSelection::MessageRange {
2507 start: 0,
2508 end: parent.messages().len(),
2509 },
2510 vec![Message::User(UserMessage::text(
2511 "[Context compacted] summary".to_string(),
2512 ))],
2513 crate::TranscriptRewriteReason::new("compaction"),
2514 Some("meerkat-core".to_string()),
2515 Some(parent_revision),
2516 )?;
2517
2518 assert!(matches!(
2519 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2520 Err(SessionStoreError::TranscriptContinuityViolation { .. }
2521 | SessionStoreError::MonotonicityViolation { .. })
2522 ));
2523 Ok(())
2524 }
2525
2526 #[test]
2527 fn run_boundary_guard_rejects_forged_parent_edge_before_real_rewrite_commit()
2528 -> Result<(), Box<dyn std::error::Error>> {
2529 let mut previous = Session::new();
2530 previous.push(Message::System(SystemMessage::new("base system")));
2531 previous.push(Message::User(UserMessage::text("turn one".to_string())));
2532 previous.push(Message::BlockAssistant(BlockAssistantMessage {
2533 blocks: vec![AssistantBlock::Text {
2534 text: "answer one".to_string(),
2535 meta: None,
2536 }],
2537 stop_reason: StopReason::EndTurn,
2538 identity: crate::types::TranscriptMessageIdentity::default(),
2539 created_at: crate::types::message_timestamp_now(),
2540 }));
2541 let previous_revision = previous.transcript_revision()?;
2542
2543 let forged_parent_messages = vec![
2544 Message::System(SystemMessage::new("refreshed runtime system projection")),
2545 Message::User(UserMessage::text(
2546 "forged insertion before retained tail".to_string(),
2547 )),
2548 previous.messages()[1].clone(),
2549 previous.messages()[2].clone(),
2550 ];
2551 let forged_parent_revision = transcript_messages_digest(&forged_parent_messages)?;
2552 let mut forged_parent = previous.clone();
2553 forged_parent.apply_transcript_history_state(TranscriptHistoryState {
2554 head: forged_parent_revision.clone(),
2555 commits: Vec::new(),
2556 revisions: vec![
2557 crate::TranscriptRevisionBody {
2558 revision: previous_revision.clone(),
2559 parent_revision: None,
2560 messages: previous.messages().to_vec(),
2561 created_at: previous.updated_at(),
2562 },
2563 crate::TranscriptRevisionBody {
2564 revision: forged_parent_revision.clone(),
2565 parent_revision: Some(previous_revision),
2566 messages: forged_parent_messages,
2567 created_at: forged_parent.updated_at(),
2568 },
2569 ],
2570 })?;
2571
2572 let mut incoming = forged_parent.clone();
2573 incoming.commit_transcript_rewrite(
2574 TranscriptRewriteSelection::MessageRange {
2575 start: 0,
2576 end: forged_parent.messages().len(),
2577 },
2578 vec![Message::User(UserMessage::text(
2579 "[Context compacted] forged branch".to_string(),
2580 ))],
2581 crate::TranscriptRewriteReason::new("compaction"),
2582 Some("meerkat-core".to_string()),
2583 Some(forged_parent_revision),
2584 )?;
2585
2586 assert!(matches!(
2587 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2588 Err(SessionStoreError::TranscriptContinuityViolation { .. }
2589 | SessionStoreError::MonotonicityViolation { .. })
2590 ));
2591 Ok(())
2592 }
2593
2594 #[test]
2595 fn append_only_guard_rejects_transient_mcp_pending_notice_cleanup_with_unaudited_commit()
2596 -> Result<(), crate::TranscriptEditError> {
2597 let mut previous = Session::new();
2598 previous.push(Message::User(UserMessage::text("hello".to_string())));
2599 previous.push(Message::SystemNotice(SystemNoticeMessage {
2600 kind: SystemNoticeKind::McpPending,
2601 body: Some("connecting".to_string()),
2602 blocks: vec![SystemNoticeBlock::Mcp {
2603 server_id: None,
2604 operation: None,
2605 phase: None,
2606 persisted: false,
2607 detail: Some("connecting".to_string()),
2608 pending_sources: vec!["test-server".to_string()],
2609 }],
2610 created_at: crate::types::message_timestamp_now(),
2611 }));
2612 previous.push(Message::BlockAssistant(BlockAssistantMessage::new(
2613 vec![crate::types::AssistantBlock::Text {
2614 text: "answer".to_string(),
2615 meta: None,
2616 }],
2617 StopReason::EndTurn,
2618 )));
2619
2620 let mut incoming = previous.clone();
2621 incoming.replace_messages_internal(
2622 previous
2623 .messages()
2624 .iter()
2625 .filter(|message| !matches!(message, Message::SystemNotice(_)))
2626 .cloned()
2627 .collect(),
2628 crate::TranscriptRewriteReason::new("unit-test"),
2629 )?;
2630 incoming.push(Message::User(UserMessage::text("again".to_string())));
2631
2632 assert!(matches!(
2633 append_only_save_guard(&incoming, Some(&previous)),
2634 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2635 ));
2636 Ok::<(), crate::TranscriptEditError>(())
2637 }
2638
2639 #[test]
2640 fn rewrite_chain_finder_crosses_normal_append_between_rewrites()
2641 -> Result<(), Box<dyn std::error::Error>> {
2642 let mut session = Session::new();
2643 session.push(Message::User(UserMessage::text("first".to_string())));
2644 session.push(Message::BlockAssistant(BlockAssistantMessage {
2645 blocks: vec![AssistantBlock::Text {
2646 text: "verbose first answer".to_string(),
2647 meta: None,
2648 }],
2649 stop_reason: StopReason::EndTurn,
2650 identity: crate::types::TranscriptMessageIdentity::default(),
2651 created_at: crate::types::message_timestamp_now(),
2652 }));
2653
2654 let original = session.transcript_revision()?;
2655 let first = session.commit_transcript_rewrite(
2656 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2657 vec![Message::BlockAssistant(BlockAssistantMessage {
2658 blocks: vec![AssistantBlock::Text {
2659 text: "compact first answer".to_string(),
2660 meta: None,
2661 }],
2662 stop_reason: StopReason::EndTurn,
2663 identity: crate::types::TranscriptMessageIdentity::default(),
2664 created_at: crate::types::message_timestamp_now(),
2665 })],
2666 crate::TranscriptRewriteReason::new("compaction"),
2667 Some("unit-test".to_string()),
2668 Some(original.clone()),
2669 )?;
2670
2671 session.push(Message::User(UserMessage::text("second".to_string())));
2672 session.push(Message::BlockAssistant(BlockAssistantMessage {
2673 blocks: vec![AssistantBlock::Text {
2674 text: "verbose second answer".to_string(),
2675 meta: None,
2676 }],
2677 stop_reason: StopReason::EndTurn,
2678 identity: crate::types::TranscriptMessageIdentity::default(),
2679 created_at: crate::types::message_timestamp_now(),
2680 }));
2681 let bridge = session.transcript_revision()?;
2682 assert_ne!(bridge, first.revision);
2683
2684 let second = session.commit_transcript_rewrite(
2685 TranscriptRewriteSelection::MessageRange { start: 3, end: 4 },
2686 vec![Message::BlockAssistant(BlockAssistantMessage {
2687 blocks: vec![AssistantBlock::Text {
2688 text: "compact second answer".to_string(),
2689 meta: None,
2690 }],
2691 stop_reason: StopReason::EndTurn,
2692 identity: crate::types::TranscriptMessageIdentity::default(),
2693 created_at: crate::types::message_timestamp_now(),
2694 })],
2695 crate::TranscriptRewriteReason::new("compaction"),
2696 Some("unit-test".to_string()),
2697 Some(bridge),
2698 )?;
2699 let state = session
2700 .transcript_history_state()?
2701 .ok_or_else(|| std::io::Error::other("missing transcript history state"))?;
2702
2703 let chain =
2704 find_transcript_rewrite_commit_chain_extending(&state, &original, &second.revision)
2705 .ok_or_else(|| {
2706 std::io::Error::other(
2707 "rewrite chain should extend through normal append bridge",
2708 )
2709 })?;
2710 assert_eq!(chain.len(), 2);
2711 assert_eq!(chain[0].revision, first.revision);
2712 assert_eq!(chain[1].revision, second.revision);
2713 Ok(())
2714 }
2715
2716 #[test]
2717 fn run_boundary_guard_rejects_dropped_retained_rewrite_commits()
2718 -> Result<(), Box<dyn std::error::Error>> {
2719 let mut base = Session::new();
2720 base.push(Message::User(UserMessage::text("turn one".to_string())));
2721 base.push(Message::BlockAssistant(BlockAssistantMessage {
2722 blocks: vec![AssistantBlock::Text {
2723 text: "verbose answer".to_string(),
2724 meta: None,
2725 }],
2726 stop_reason: StopReason::EndTurn,
2727 identity: crate::types::TranscriptMessageIdentity::default(),
2728 created_at: crate::types::message_timestamp_now(),
2729 }));
2730 let base_revision = base.transcript_revision()?;
2731
2732 let mut previous = base.clone();
2733 let _retained_commit = previous.commit_transcript_rewrite(
2734 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2735 vec![Message::BlockAssistant(BlockAssistantMessage {
2736 blocks: vec![AssistantBlock::Text {
2737 text: "first compact answer".to_string(),
2738 meta: None,
2739 }],
2740 stop_reason: StopReason::EndTurn,
2741 identity: crate::types::TranscriptMessageIdentity::default(),
2742 created_at: crate::types::message_timestamp_now(),
2743 })],
2744 crate::TranscriptRewriteReason::new("compaction"),
2745 Some("unit-test".to_string()),
2746 Some(base_revision),
2747 )?;
2748 let previous_revision = previous.transcript_revision()?;
2749
2750 let mut incoming = previous.clone();
2751 let new_commit = incoming.commit_transcript_rewrite(
2752 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2753 vec![Message::BlockAssistant(BlockAssistantMessage {
2754 blocks: vec![AssistantBlock::Text {
2755 text: "second compact answer".to_string(),
2756 meta: None,
2757 }],
2758 stop_reason: StopReason::EndTurn,
2759 identity: crate::types::TranscriptMessageIdentity::default(),
2760 created_at: crate::types::message_timestamp_now(),
2761 })],
2762 crate::TranscriptRewriteReason::new("compaction"),
2763 Some("unit-test".to_string()),
2764 Some(previous_revision),
2765 )?;
2766 let mut state = incoming
2767 .transcript_history_state()?
2768 .ok_or_else(|| std::io::Error::other("incoming rewrite should retain history"))?;
2769 state.commits = vec![new_commit];
2770 incoming.set_metadata_unchecked_for_test(
2771 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2772 serde_json::to_value(state)?,
2773 );
2774
2775 assert!(matches!(
2776 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2777 Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
2778 if reason.contains("drop retained transcript rewrite commits")
2779 ));
2780 Ok(())
2781 }
2782
2783 fn runtime_append_system(content: &str) -> SystemMessage {
2792 let mut system = SystemMessage::new(content);
2793 system.mutation_kind = crate::types::SystemPromptMutationKind::RuntimeContextAppend;
2794 system
2795 }
2796
2797 #[allow(clippy::expect_used)]
2800 fn machine_persist_append_admits(
2801 previous: Option<&SystemMessage>,
2802 incoming: &SystemMessage,
2803 ) -> bool {
2804 let has_previous = previous.is_some();
2805 let content_identical =
2806 previous.is_some_and(|previous| incoming.content == previous.content);
2807 let content_extends_previous =
2808 previous.is_some_and(|previous| incoming.content.starts_with(&previous.content));
2809 let appended_starts_with_separator = previous.is_some_and(|previous| {
2810 incoming
2811 .content
2812 .get(previous.content.len()..)
2813 .is_some_and(|appended| appended.starts_with(SYSTEM_CONTEXT_SEPARATOR))
2814 });
2815 let incoming_is_runtime_context_append = incoming.mutation_kind.is_runtime_context_append();
2816 let mut authority = crate::session_document::SessionDocumentMachineAuthority::new();
2817 let effects = authority
2818 .resolve_system_context_persist_append_admission(
2819 has_previous,
2820 content_identical,
2821 content_extends_previous,
2822 appended_starts_with_separator,
2823 incoming_is_runtime_context_append,
2824 )
2825 .expect("machine resolves persist-append admission");
2826 effects.into_iter().any(|effect| {
2827 matches!(
2828 effect,
2829 crate::session_document::SessionDocumentEffect::SystemContextPersistAppendAdmissionResolved {
2830 admission: crate::session_document::SystemContextPersistAppendAdmission::Admit,
2831 }
2832 )
2833 })
2834 }
2835
2836 #[allow(clippy::expect_used)]
2837 fn assert_persist_append_matches_machine(
2838 previous: Option<&SystemMessage>,
2839 incoming: &SystemMessage,
2840 expected: bool,
2841 ) {
2842 let verdict =
2843 system_context_is_append(previous, incoming).expect("persist-time admission resolves");
2844 assert_eq!(verdict, expected, "persist-time verdict mismatch");
2845 assert_eq!(
2846 verdict,
2847 machine_persist_append_admits(previous, incoming),
2848 "persist-time verdict diverges from direct machine call"
2849 );
2850 }
2851
2852 #[test]
2853 fn persist_append_identical_content_admits() {
2854 let previous = SystemMessage::new("base system");
2855 let incoming = SystemMessage::new("base system");
2856 assert_persist_append_matches_machine(Some(&previous), &incoming, true);
2857 }
2858
2859 #[test]
2860 fn persist_append_separator_append_with_marker_admits() {
2861 let previous = SystemMessage::new("base system");
2862 let incoming = runtime_append_system(&format!(
2863 "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nextra"
2864 ));
2865 assert_persist_append_matches_machine(Some(&previous), &incoming, true);
2866 }
2867
2868 #[test]
2869 fn persist_append_shaped_without_marker_rejects() {
2870 let previous = SystemMessage::new("base system");
2871 let incoming = SystemMessage::new(format!(
2873 "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nextra"
2874 ));
2875 assert_persist_append_matches_machine(Some(&previous), &incoming, false);
2876 }
2877
2878 #[test]
2879 fn persist_append_divergent_content_rejects() {
2880 let previous = SystemMessage::new("base system");
2881 let incoming = runtime_append_system("totally different");
2882 assert_persist_append_matches_machine(Some(&previous), &incoming, false);
2883 }
2884
2885 #[test]
2886 fn persist_append_no_previous_admits_only_with_marker() {
2887 let with_marker = runtime_append_system("brand new context");
2888 assert_persist_append_matches_machine(None, &with_marker, true);
2889
2890 let without_marker = SystemMessage::new("brand new context");
2891 assert_persist_append_matches_machine(None, &without_marker, false);
2892 }
2893
2894 fn assistant_with_bookkeeping(
2895 text: &str,
2896 run_id: Option<crate::lifecycle::RunId>,
2897 created_at: crate::types::MessageTimestamp,
2898 ) -> Message {
2899 Message::BlockAssistant(BlockAssistantMessage {
2900 blocks: vec![AssistantBlock::Text {
2901 text: text.to_string(),
2902 meta: None,
2903 }],
2904 stop_reason: StopReason::EndTurn,
2905 identity: crate::types::TranscriptMessageIdentity {
2906 interaction_id: None,
2907 run_id,
2908 },
2909 created_at,
2910 })
2911 }
2912
2913 #[test]
2919 fn append_only_guard_accepts_rebookkept_prefix_identity_and_timestamps() {
2920 let base_time = crate::types::message_timestamp_now();
2921 let mut previous = Session::new();
2922 previous.push(Message::System(SystemMessage::new("base system")));
2923 previous.push(Message::User(UserMessage::text("turn one".to_string())));
2924 previous.push(assistant_with_bookkeeping(
2925 "answer one",
2926 Some(crate::lifecycle::RunId::new()),
2927 base_time,
2928 ));
2929
2930 let mut incoming = previous.clone();
2931 let mut rebookkept = previous.messages().to_vec();
2932 for message in &mut rebookkept {
2933 match message {
2934 Message::User(user) => {
2935 user.created_at = base_time + chrono::Duration::hours(1);
2936 }
2937 Message::BlockAssistant(assistant) => {
2938 assistant.identity = crate::types::TranscriptMessageIdentity {
2939 interaction_id: None,
2940 run_id: Some(crate::lifecycle::RunId::new()),
2941 };
2942 assistant.created_at = base_time + chrono::Duration::hours(1);
2943 }
2944 _ => {}
2945 }
2946 }
2947 rebookkept.push(Message::User(UserMessage::text("turn two".to_string())));
2948 incoming.messages = std::sync::Arc::new(rebookkept);
2949
2950 assert!(
2951 append_only_save_guard(&incoming, Some(&previous)).is_ok(),
2952 "bookkeeping-only prefix divergence must not fail continuity"
2953 );
2954 }
2955
2956 #[test]
2957 fn append_only_guard_rejects_content_divergence_despite_matching_bookkeeping() {
2958 let base_time = crate::types::message_timestamp_now();
2959 let run_id = crate::lifecycle::RunId::new();
2960 let mut previous = Session::new();
2961 previous.push(Message::User(UserMessage::text("turn one".to_string())));
2962 previous.push(assistant_with_bookkeeping(
2963 "answer one",
2964 Some(run_id),
2965 base_time,
2966 ));
2967
2968 let mut incoming = previous.clone();
2969 let mut diverged = previous.messages().to_vec();
2970 if let Message::BlockAssistant(assistant) = &mut diverged[1] {
2971 assistant.blocks = vec![AssistantBlock::Text {
2972 text: "a different answer".to_string(),
2973 meta: None,
2974 }];
2975 }
2976 diverged.push(Message::User(UserMessage::text("turn two".to_string())));
2977 incoming.messages = std::sync::Arc::new(diverged);
2978
2979 assert!(matches!(
2980 append_only_save_guard(&incoming, Some(&previous)),
2981 Err(SessionStoreError::TranscriptContinuityViolation { .. })
2982 ));
2983 }
2984}