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 let incoming_revision = transcript_messages_digest(incoming.messages())
574 .map_err(SessionStoreError::from)?;
575 if let Some(state) = incoming.transcript_history_state().map_err(|err| {
576 SessionStoreError::InvalidTranscriptRewrite {
577 id: incoming.id().clone(),
578 reason: format!("incoming transcript history state is malformed: {err}"),
579 }
580 })? && !state.commits.is_empty()
581 && state.head == incoming_revision
582 {
583 for commit in &state.commits {
584 validate_transcript_rewrite_commit_bodies(incoming, commit, &state)?;
585 }
586 return Ok(());
587 }
588 return Err(append_error);
589 };
590 let incoming_revision =
591 transcript_messages_digest(incoming.messages()).map_err(SessionStoreError::from)?;
592 let Some(state) = incoming.transcript_history_state().map_err(|err| {
593 SessionStoreError::InvalidTranscriptRewrite {
594 id: incoming.id().clone(),
595 reason: format!("incoming transcript history state is malformed: {err}"),
596 }
597 })?
598 else {
599 return Err(append_error);
600 };
601 incoming
606 .validate_transcript_history_state()
607 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
608 id: incoming.id().clone(),
609 reason: format!("incoming transcript history state is malformed: {err}"),
610 })?;
611 validate_rewrite_save_retains_previous_commits(incoming, previous, &state)?;
612 let commits = find_transcript_rewrite_commit_chain_extending_session(
613 &state,
614 previous,
615 &incoming_revision,
616 )?;
617 if commits.is_none()
618 && run_boundary_context_summary_tail_projection_save_guard(
619 incoming, previous, &state,
620 )?
621 {
622 return Ok(());
623 }
624 let Some(commits) = commits else {
625 return Err(append_error);
626 };
627 let Some(commit) = commits.first() else {
628 if state.commits.is_empty() {
629 return Err(append_error);
630 }
631 if state.head != incoming_revision {
637 return Err(SessionStoreError::InvalidTranscriptRewrite {
638 id: incoming.id().clone(),
639 reason: format!(
640 "incoming transcript graph head {} does not match current message digest {incoming_revision}",
641 state.head
642 ),
643 });
644 }
645 for commit in &state.commits {
646 validate_transcript_rewrite_commit_bodies(incoming, commit, &state)?;
647 }
648 return Ok(());
649 };
650 transcript_rewrite_bridge_save_guard(incoming, commit, &state, &incoming_revision)?;
651 for commit in &state.commits {
657 validate_transcript_rewrite_commit_bodies(incoming, commit, &state)?;
658 }
659 Ok(())
660 }
661 }
662}
663
664fn run_boundary_commitless_history_projection_save_guard(
665 incoming: &Session,
666 previous: Option<&Session>,
667) -> Result<bool, SessionStoreError> {
668 let Some(state) = incoming.transcript_history_state().map_err(|err| {
669 SessionStoreError::InvalidTranscriptRewrite {
670 id: incoming.id().clone(),
671 reason: format!("incoming transcript history state is malformed: {err}"),
672 }
673 })?
674 else {
675 return Ok(false);
676 };
677 if !state.commits.is_empty() {
678 return Ok(false);
679 }
680
681 let incoming_revision =
682 transcript_messages_digest(incoming.messages()).map_err(SessionStoreError::from)?;
683 if state.head != incoming_revision
684 || !state
685 .revisions
686 .iter()
687 .any(|body| body.revision == incoming_revision)
688 {
689 return Ok(false);
690 }
691
692 let mut projection_without_history = incoming.clone();
693 projection_without_history.clear_transcript_history_state();
694 if append_only_save_guard(&projection_without_history, previous).is_err() {
695 return Ok(false);
696 }
697
698 let Some(previous) = previous else {
699 return Ok(state.commits.is_empty());
700 };
701 if previous
702 .transcript_history_state()
703 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
704 id: incoming.id().clone(),
705 reason: format!("previous transcript history state is malformed: {err}"),
706 })?
707 .is_some()
708 {
709 return Ok(false);
710 }
711
712 let previous_revision =
713 transcript_messages_digest(previous.messages()).map_err(SessionStoreError::from)?;
714 Ok(incoming_revision == previous_revision
715 || transcript_history_revision_extends(&state, &incoming_revision, &previous_revision))
716}
717
718fn run_boundary_context_summary_tail_projection_save_guard(
719 incoming: &Session,
720 previous: &Session,
721 state: &TranscriptHistoryState,
722) -> Result<bool, SessionStoreError> {
723 if state.commits.is_empty() {
724 return Ok(false);
725 }
726 incoming
727 .validate_transcript_history_state()
728 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
729 id: incoming.id().clone(),
730 reason: format!("incoming transcript history state is malformed: {err}"),
731 })?;
732
733 let (incoming_system, incoming_tail) = match incoming.messages().split_first() {
734 Some((Message::System(system), tail)) => (Some(system), tail),
735 _ => (None, incoming.messages()),
736 };
737 let (previous_system, previous_tail) = match previous.messages().split_first() {
738 Some((Message::System(system), tail)) => (Some(system), tail),
739 _ => (None, previous.messages()),
740 };
741 if incoming_system.is_some() != previous_system.is_some()
742 || incoming_tail.len() <= previous_tail.len()
743 {
744 return Ok(false);
745 }
746 let Some(Message::User(summary)) = incoming_tail.first() else {
747 return Ok(false);
748 };
749 if !summary.transcript_role.is_compaction_summary() {
754 return Ok(false);
755 }
756
757 let retained_end = 1 + previous_tail.len();
758 let retained = &incoming_tail[1..retained_end];
759 let retained_revision =
760 transcript_messages_digest(retained).map_err(SessionStoreError::from)?;
761 let previous_revision =
762 transcript_messages_digest(previous_tail).map_err(SessionStoreError::from)?;
763 if retained_revision != previous_revision {
764 return Ok(false);
765 }
766
767 for commit in &state.commits {
768 validate_transcript_rewrite_commit_bodies(incoming, commit, state)?;
769 }
770 Ok(true)
771}
772
773pub fn find_transcript_rewrite_commit_extending<'a>(
776 state: &'a TranscriptHistoryState,
777 previous_revision: &str,
778 incoming_revision: &str,
779) -> Option<&'a TranscriptRewriteCommit> {
780 find_transcript_rewrite_commit_chain_extending(state, previous_revision, incoming_revision)
781 .and_then(|commits| commits.into_iter().next())
782}
783
784pub fn find_transcript_rewrite_commit_chain_extending<'a>(
787 state: &'a TranscriptHistoryState,
788 previous_revision: &str,
789 incoming_revision: &str,
790) -> Option<Vec<&'a TranscriptRewriteCommit>> {
791 let mut chain = Vec::new();
792 let mut cursor = previous_revision;
793 let mut visited = std::collections::BTreeSet::new();
794 loop {
795 if incoming_revision == cursor {
796 return Some(chain);
797 }
798 if !visited.insert(cursor.to_string()) {
799 return None;
800 }
801 let commit = state.commits.iter().find(|commit| {
802 (commit.parent_revision == cursor
803 || transcript_history_revision_extends(state, &commit.parent_revision, cursor))
804 && transcript_history_revision_extends(state, incoming_revision, &commit.revision)
805 });
806 let Some(commit) = commit else {
807 return transcript_history_revision_extends(state, incoming_revision, cursor)
808 .then_some(chain);
809 };
810 cursor = &commit.revision;
811 chain.push(commit);
812 }
813}
814
815pub fn find_transcript_rewrite_commit_chain_extending_session<'a>(
824 state: &'a TranscriptHistoryState,
825 previous: &Session,
826 incoming_revision: &str,
827) -> Result<Option<Vec<&'a TranscriptRewriteCommit>>, SessionStoreError> {
828 let previous_revision =
829 transcript_messages_digest(previous.messages()).map_err(SessionStoreError::from)?;
830 let mut chain = Vec::new();
831 let mut cursor = previous_revision.as_str();
832 let mut visited = std::collections::BTreeSet::new();
833 loop {
834 if incoming_revision == cursor {
835 return Ok(Some(chain));
836 }
837 if !visited.insert(cursor.to_string()) {
838 return Ok(None);
839 }
840
841 let Some(cursor_messages) = transcript_history_messages_for_revision(
842 state,
843 cursor,
844 &previous_revision,
845 previous.messages(),
846 ) else {
847 return Ok(None);
848 };
849
850 let mut selected = None;
856 for commit in &state.commits {
857 if commit.revision == cursor || visited.contains(&commit.revision) {
858 continue;
859 }
860 if !transcript_history_revision_extends(state, incoming_revision, &commit.revision) {
861 continue;
862 }
863 if commit.parent_revision == cursor {
864 selected = Some(commit);
865 break;
866 }
867 }
868
869 if selected.is_none() {
881 if revision_body_preserves_append_continuation_prefix(
882 state,
883 incoming_revision,
884 cursor_messages,
885 cursor,
886 false,
887 )? {
888 return Ok(Some(chain));
889 }
890 for commit in &state.commits {
895 if commit.revision == cursor || visited.contains(&commit.revision) {
896 continue;
897 }
898 if !transcript_history_revision_extends(state, incoming_revision, &commit.revision)
899 {
900 continue;
901 }
902 if revision_body_preserves_append_continuation_prefix(
903 state,
904 &commit.parent_revision,
905 cursor_messages,
906 cursor,
907 true,
908 )? {
909 selected = Some(commit);
910 break;
911 }
912 }
913 }
914
915 let Some(commit) = selected else {
916 return Ok(None);
917 };
918 cursor = &commit.revision;
919 chain.push(commit);
920 }
921}
922
923fn transcript_history_messages_for_revision<'a>(
924 state: &'a TranscriptHistoryState,
925 revision: &str,
926 previous_revision: &str,
927 previous_messages: &'a [Message],
928) -> Option<&'a [Message]> {
929 if revision == previous_revision {
930 return Some(previous_messages);
931 }
932 state
933 .revisions
934 .iter()
935 .find(|body| body.revision == revision)
936 .map(|body| body.messages.as_slice())
937}
938
939fn revision_body_preserves_append_continuation_prefix(
940 state: &TranscriptHistoryState,
941 revision: &str,
942 ancestor_messages: &[Message],
943 ancestor_revision: &str,
944 allow_leading_system_refresh: bool,
945) -> Result<bool, SessionStoreError> {
946 if revision == ancestor_revision {
947 return Ok(true);
948 }
949 let Some(body) = state
950 .revisions
951 .iter()
952 .find(|body| body.revision == revision)
953 else {
954 return Ok(false);
955 };
956 if body.messages.len() >= ancestor_messages.len() {
957 let prefix_revision = transcript_messages_digest(&body.messages[..ancestor_messages.len()])
958 .map_err(SessionStoreError::from)?;
959 if prefix_revision == ancestor_revision {
960 return Ok(true);
961 }
962 }
963 if messages_preserve_conversation_tail_with_system_context_append(
964 &body.messages,
965 ancestor_messages,
966 )? {
967 return Ok(true);
968 }
969 Ok(allow_leading_system_refresh
975 && messages_preserve_tail_after_leading_system_refresh(&body.messages, ancestor_messages)?)
976}
977
978fn messages_preserve_tail_after_leading_system_refresh(
979 incoming: &[Message],
980 previous: &[Message],
981) -> Result<bool, SessionStoreError> {
982 let (Some(Message::System(_)), Some(Message::System(_))) = (incoming.first(), previous.first())
983 else {
984 return Ok(false);
985 };
986 if incoming.len() < previous.len() {
987 return Ok(false);
988 }
989 let previous_tail_len = previous.len().saturating_sub(1);
990 if previous_tail_len == 0 {
991 return Ok(true);
992 }
993 let previous_tail_revision =
994 transcript_messages_digest(&previous[1..]).map_err(SessionStoreError::from)?;
995 let incoming_tail = &incoming[1..];
996 if incoming_tail.len() < previous_tail_len {
997 return Ok(false);
998 }
999 let incoming_tail_prefix_revision =
1000 transcript_messages_digest(&incoming_tail[..previous_tail_len])
1001 .map_err(SessionStoreError::from)?;
1002 Ok(incoming_tail_prefix_revision == previous_tail_revision)
1003}
1004
1005fn transcript_history_revision_extends(
1006 state: &TranscriptHistoryState,
1007 descendant: &str,
1008 ancestor: &str,
1009) -> bool {
1010 if descendant == ancestor {
1011 return true;
1012 }
1013 let mut cursor = descendant;
1014 let mut visited = std::collections::BTreeSet::new();
1017 while let Some(body) = state.revisions.iter().find(|body| body.revision == cursor) {
1018 if !visited.insert(body.revision.clone()) {
1019 return false;
1020 }
1021 let Some(parent) = body.parent_revision.as_deref() else {
1022 return false;
1023 };
1024 if parent == ancestor {
1025 return true;
1026 }
1027 cursor = parent;
1028 }
1029 false
1030}
1031
1032fn transcript_rewrite_bridge_save_guard(
1033 incoming: &Session,
1034 commit: &TranscriptRewriteCommit,
1035 incoming_state: &TranscriptHistoryState,
1036 incoming_message_digest: &str,
1037) -> Result<(), SessionStoreError> {
1038 validate_transcript_rewrite_commit_bodies(incoming, commit, incoming_state)?;
1039 if incoming_state.head != incoming_message_digest {
1040 return Err(SessionStoreError::InvalidTranscriptRewrite {
1041 id: incoming.id().clone(),
1042 reason: format!(
1043 "incoming transcript graph head {} does not match current message digest {incoming_message_digest}",
1044 incoming_state.head
1045 ),
1046 });
1047 }
1048 if !transcript_history_revision_extends(
1049 incoming_state,
1050 incoming_message_digest,
1051 &commit.revision,
1052 ) {
1053 return Err(SessionStoreError::InvalidTranscriptRewrite {
1054 id: incoming.id().clone(),
1055 reason: format!(
1056 "incoming transcript head {incoming_message_digest} does not extend rewrite revision {}",
1057 commit.revision
1058 ),
1059 });
1060 }
1061 Ok(())
1062}
1063
1064pub fn transcript_rewrite_save_guard(
1067 incoming: &Session,
1068 previous: Option<&Session>,
1069 commit: &TranscriptRewriteCommit,
1070) -> Result<(), SessionStoreError> {
1071 incoming
1072 .validate_transcript_history_state()
1073 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1074 id: incoming.id().clone(),
1075 reason: format!("incoming transcript history state is malformed: {err}"),
1076 })?;
1077 let Some(previous) = previous else {
1078 return Err(SessionStoreError::InvalidTranscriptRewrite {
1079 id: incoming.id().clone(),
1080 reason: "rewrite target has no previously persisted session".to_string(),
1081 });
1082 };
1083 if incoming.id() != previous.id() {
1084 return Err(SessionStoreError::InvalidTranscriptRewrite {
1085 id: incoming.id().clone(),
1086 reason: format!(
1087 "incoming session id {} differs from previous session id {}",
1088 incoming.id(),
1089 previous.id()
1090 ),
1091 });
1092 }
1093 let previous_revision = previous.transcript_revision().map_err(|err| {
1094 SessionStoreError::InvalidTranscriptRewrite {
1095 id: incoming.id().clone(),
1096 reason: format!("previous transcript revision is malformed: {err}"),
1097 }
1098 })?;
1099 if previous_revision != commit.parent_revision {
1100 return Err(SessionStoreError::TranscriptRevisionConflict {
1101 id: incoming.id().clone(),
1102 expected: commit.parent_revision.clone(),
1103 actual: previous_revision,
1104 });
1105 }
1106 let previous_message_digest =
1107 transcript_messages_digest(previous.messages()).map_err(|err| {
1108 SessionStoreError::InvalidTranscriptRewrite {
1109 id: incoming.id().clone(),
1110 reason: format!("previous current transcript is not digestible: {err}"),
1111 }
1112 })?;
1113 if previous_message_digest != commit.parent_revision {
1114 return Err(SessionStoreError::InvalidTranscriptRewrite {
1115 id: incoming.id().clone(),
1116 reason: format!(
1117 "previous current transcript digest {previous_message_digest} does not match commit parent {}",
1118 commit.parent_revision
1119 ),
1120 });
1121 }
1122 let incoming_revision = incoming.transcript_revision().map_err(|err| {
1123 SessionStoreError::InvalidTranscriptRewrite {
1124 id: incoming.id().clone(),
1125 reason: format!("incoming transcript revision is malformed: {err}"),
1126 }
1127 })?;
1128 if incoming_revision != commit.revision {
1129 return Err(SessionStoreError::InvalidTranscriptRewrite {
1130 id: incoming.id().clone(),
1131 reason: format!(
1132 "incoming transcript revision {incoming_revision} does not match commit revision {}",
1133 commit.revision
1134 ),
1135 });
1136 }
1137 let incoming_message_digest =
1138 transcript_messages_digest(incoming.messages()).map_err(|err| {
1139 SessionStoreError::InvalidTranscriptRewrite {
1140 id: incoming.id().clone(),
1141 reason: format!("incoming current transcript is not digestible: {err}"),
1142 }
1143 })?;
1144 if incoming_message_digest != commit.revision {
1145 return Err(SessionStoreError::InvalidTranscriptRewrite {
1146 id: incoming.id().clone(),
1147 reason: format!(
1148 "incoming current transcript digest {incoming_message_digest} does not match commit revision {}",
1149 commit.revision
1150 ),
1151 });
1152 }
1153 let Some(incoming_state) = incoming.transcript_history_state().map_err(|err| {
1154 SessionStoreError::InvalidTranscriptRewrite {
1155 id: incoming.id().clone(),
1156 reason: format!("incoming transcript history state is malformed: {err}"),
1157 }
1158 })?
1159 else {
1160 return Err(SessionStoreError::InvalidTranscriptRewrite {
1161 id: incoming.id().clone(),
1162 reason: "incoming rewrite did not persist a transcript revision graph".to_string(),
1163 });
1164 };
1165 validate_rewrite_save_retains_previous_commits(incoming, previous, &incoming_state)?;
1166 validate_transcript_rewrite_commit_bodies(incoming, commit, &incoming_state)
1167}
1168
1169fn validate_transcript_rewrite_commit_bodies(
1170 incoming: &Session,
1171 commit: &TranscriptRewriteCommit,
1172 incoming_state: &TranscriptHistoryState,
1173) -> Result<(), SessionStoreError> {
1174 if !incoming_state
1175 .commits
1176 .iter()
1177 .any(|persisted| persisted == commit)
1178 {
1179 return Err(SessionStoreError::InvalidTranscriptRewrite {
1180 id: incoming.id().clone(),
1181 reason: format!(
1182 "incoming rewrite did not persist the rewrite commit in the transcript graph (wanted {} -> {}, graph commits: {:?})",
1183 commit.parent_revision,
1184 commit.revision,
1185 incoming_state
1186 .commits
1187 .iter()
1188 .map(|commit| (&commit.parent_revision, &commit.revision))
1189 .collect::<Vec<_>>()
1190 ),
1191 });
1192 }
1193 let Some(parent_body) = incoming_state
1194 .revisions
1195 .iter()
1196 .find(|body| body.revision == commit.parent_revision)
1197 else {
1198 return Err(SessionStoreError::InvalidTranscriptRewrite {
1199 id: incoming.id().clone(),
1200 reason: format!(
1201 "incoming rewrite omitted parent revision body {}",
1202 commit.parent_revision
1203 ),
1204 });
1205 };
1206 let Some(revision_body) = incoming_state
1207 .revisions
1208 .iter()
1209 .find(|body| body.revision == commit.revision)
1210 else {
1211 return Err(SessionStoreError::InvalidTranscriptRewrite {
1212 id: incoming.id().clone(),
1213 reason: format!(
1214 "incoming rewrite omitted new revision body {}",
1215 commit.revision
1216 ),
1217 });
1218 };
1219 if parent_body.messages.len() != commit.messages_before
1220 || revision_body.messages.len() != commit.messages_after
1221 {
1222 return Err(SessionStoreError::InvalidTranscriptRewrite {
1223 id: incoming.id().clone(),
1224 reason: format!(
1225 "commit message counts {} -> {} do not match persisted rewrite {} -> {}",
1226 commit.messages_before,
1227 commit.messages_after,
1228 parent_body.messages.len(),
1229 revision_body.messages.len()
1230 ),
1231 });
1232 }
1233 let parent_body_revision =
1234 transcript_messages_digest(&parent_body.messages).map_err(|err| {
1235 SessionStoreError::InvalidTranscriptRewrite {
1236 id: incoming.id().clone(),
1237 reason: format!("parent revision body is not digestible: {err}"),
1238 }
1239 })?;
1240 if parent_body_revision != commit.parent_revision {
1241 return Err(SessionStoreError::InvalidTranscriptRewrite {
1242 id: incoming.id().clone(),
1243 reason: format!(
1244 "parent revision body digest {parent_body_revision} does not match commit parent {}",
1245 commit.parent_revision
1246 ),
1247 });
1248 }
1249 let (start, end) = match &commit.selection {
1250 TranscriptRewriteSelection::MessageRange { start, end } => (*start, *end),
1251 };
1252 if start > end || end > parent_body.messages.len() {
1253 return Err(SessionStoreError::InvalidTranscriptRewrite {
1254 id: incoming.id().clone(),
1255 reason: format!(
1256 "commit selection {start}..{end} is invalid for parent revision with {} messages",
1257 parent_body.messages.len()
1258 ),
1259 });
1260 }
1261 let original_span_digest = transcript_messages_digest(&parent_body.messages[start..end])
1262 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1263 id: incoming.id().clone(),
1264 reason: format!("original span body is not digestible: {err}"),
1265 })?;
1266 if original_span_digest != commit.original_span_digest {
1267 return Err(SessionStoreError::InvalidTranscriptRewrite {
1268 id: incoming.id().clone(),
1269 reason: format!(
1270 "original span digest {original_span_digest} does not match commit digest {}",
1271 commit.original_span_digest
1272 ),
1273 });
1274 }
1275 let revision_body_digest =
1276 transcript_messages_digest(&revision_body.messages).map_err(|err| {
1277 SessionStoreError::InvalidTranscriptRewrite {
1278 id: incoming.id().clone(),
1279 reason: format!("new revision body is not digestible: {err}"),
1280 }
1281 })?;
1282 if revision_body_digest != commit.revision {
1283 return Err(SessionStoreError::InvalidTranscriptRewrite {
1284 id: incoming.id().clone(),
1285 reason: format!(
1286 "new revision body digest {revision_body_digest} does not match commit revision {}",
1287 commit.revision
1288 ),
1289 });
1290 }
1291 let removed_len = end - start;
1292 let retained_len = commit
1293 .messages_before
1294 .checked_sub(removed_len)
1295 .ok_or_else(|| SessionStoreError::InvalidTranscriptRewrite {
1296 id: incoming.id().clone(),
1297 reason: "commit removed more messages than it recorded before rewrite".to_string(),
1298 })?;
1299 let replacement_len = commit
1300 .messages_after
1301 .checked_sub(retained_len)
1302 .ok_or_else(|| SessionStoreError::InvalidTranscriptRewrite {
1303 id: incoming.id().clone(),
1304 reason: "commit message counts cannot describe a replacement span".to_string(),
1305 })?;
1306 let replacement_end = start.checked_add(replacement_len).ok_or_else(|| {
1307 SessionStoreError::InvalidTranscriptRewrite {
1308 id: incoming.id().clone(),
1309 reason: "replacement span end overflowed".to_string(),
1310 }
1311 })?;
1312 if replacement_end > revision_body.messages.len() {
1313 return Err(SessionStoreError::InvalidTranscriptRewrite {
1314 id: incoming.id().clone(),
1315 reason: format!(
1316 "replacement span {start}..{replacement_end} is invalid for revision with {} messages",
1317 revision_body.messages.len()
1318 ),
1319 });
1320 }
1321 let parent_prefix_digest =
1322 transcript_messages_digest(&parent_body.messages[..start]).map_err(|err| {
1323 SessionStoreError::InvalidTranscriptRewrite {
1324 id: incoming.id().clone(),
1325 reason: format!("parent prefix body is not digestible: {err}"),
1326 }
1327 })?;
1328 let revision_prefix_digest = transcript_messages_digest(&revision_body.messages[..start])
1329 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1330 id: incoming.id().clone(),
1331 reason: format!("revision prefix body is not digestible: {err}"),
1332 })?;
1333 if parent_prefix_digest != revision_prefix_digest {
1334 return Err(SessionStoreError::InvalidTranscriptRewrite {
1335 id: incoming.id().clone(),
1336 reason: "rewrite revision changed messages before the selected span".to_string(),
1337 });
1338 }
1339 let parent_suffix_digest =
1340 transcript_messages_digest(&parent_body.messages[end..]).map_err(|err| {
1341 SessionStoreError::InvalidTranscriptRewrite {
1342 id: incoming.id().clone(),
1343 reason: format!("parent suffix body is not digestible: {err}"),
1344 }
1345 })?;
1346 let revision_suffix_digest =
1347 transcript_messages_digest(&revision_body.messages[replacement_end..]).map_err(|err| {
1348 SessionStoreError::InvalidTranscriptRewrite {
1349 id: incoming.id().clone(),
1350 reason: format!("revision suffix body is not digestible: {err}"),
1351 }
1352 })?;
1353 if parent_suffix_digest != revision_suffix_digest {
1354 return Err(SessionStoreError::InvalidTranscriptRewrite {
1355 id: incoming.id().clone(),
1356 reason: "rewrite revision changed messages after the selected span".to_string(),
1357 });
1358 }
1359 let replacement_digest = transcript_messages_digest(
1360 &revision_body.messages[start..replacement_end],
1361 )
1362 .map_err(|err| SessionStoreError::InvalidTranscriptRewrite {
1363 id: incoming.id().clone(),
1364 reason: format!("replacement span body is not digestible: {err}"),
1365 })?;
1366 if replacement_digest != commit.replacement_digest {
1367 return Err(SessionStoreError::InvalidTranscriptRewrite {
1368 id: incoming.id().clone(),
1369 reason: format!(
1370 "replacement span digest {replacement_digest} does not match commit digest {}",
1371 commit.replacement_digest
1372 ),
1373 });
1374 }
1375 Ok(())
1376}
1377
1378impl From<serde_json::Error> for SessionStoreError {
1379 fn from(e: serde_json::Error) -> Self {
1380 Self::Serialization(e.to_string())
1381 }
1382}
1383
1384#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
1409#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
1410pub trait SessionStore: Send + Sync {
1411 async fn save(&self, session: &Session) -> Result<(), SessionStoreError>;
1417
1418 async fn save_transcript_rewrite(
1424 &self,
1425 session: &Session,
1426 commit: &TranscriptRewriteCommit,
1427 ) -> Result<(), SessionStoreError> {
1428 let _ = (session, commit);
1429 Err(SessionStoreError::Internal(
1430 "save_transcript_rewrite is not supported by this SessionStore".to_string(),
1431 ))
1432 }
1433
1434 async fn save_authoritative_projection(
1443 &self,
1444 session: &Session,
1445 ) -> Result<(), SessionStoreError> {
1446 self.save(session).await
1447 }
1448
1449 async fn save_authoritative_projection_if_current_revision(
1452 &self,
1453 session: &Session,
1454 expected_current_revision: Option<String>,
1455 ) -> Result<(), SessionStoreError> {
1456 let _ = (session, expected_current_revision);
1457 Err(SessionStoreError::Internal(
1458 "save_authoritative_projection_if_current_revision is not supported by this SessionStore"
1459 .to_string(),
1460 ))
1461 }
1462
1463 async fn load(&self, id: &SessionId) -> Result<Option<Session>, SessionStoreError>;
1465
1466 async fn list(&self, filter: SessionFilter) -> Result<Vec<SessionMeta>, SessionStoreError>;
1468
1469 async fn delete(&self, id: &SessionId) -> Result<(), SessionStoreError>;
1471
1472 async fn delete_if_current_revision(
1475 &self,
1476 id: &SessionId,
1477 expected_current_revision: &str,
1478 ) -> Result<bool, SessionStoreError>;
1479
1480 async fn exists(&self, id: &SessionId) -> Result<bool, SessionStoreError> {
1482 Ok(self.load(id).await?.is_some())
1483 }
1484}
1485
1486#[cfg(test)]
1487mod tests {
1488 use super::*;
1489 use crate::types::{
1490 AssistantBlock, BlockAssistantMessage, StopReason, SystemMessage, SystemNoticeBlock,
1491 SystemNoticeKind, SystemNoticeMessage, UserMessage,
1492 };
1493
1494 #[test]
1503 #[allow(clippy::expect_used)]
1504 fn boundary_commit_walks_lagging_row_across_rebookkept_refresh_chain()
1505 -> Result<(), Box<dyn std::error::Error>> {
1506 let mut base = Session::new();
1507 base.push(Message::System(SystemMessage::new(
1508 "member prompt roster v1",
1509 )));
1510 base.push(Message::User(UserMessage::text(
1511 "the codeword is birch seventeen".to_string(),
1512 )));
1513 let previous = base.clone();
1516 let v1 = base.transcript_revision()?;
1517
1518 let mut session = base;
1520 session.commit_transcript_rewrite(
1521 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1522 vec![Message::System(SystemMessage::new(
1523 "member prompt roster v2",
1524 ))],
1525 crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
1526 Some("agent-factory/resume".to_string()),
1527 Some(v1),
1528 )?;
1529 let v2 = session.transcript_revision()?;
1530 session.commit_transcript_rewrite(
1531 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1532 vec![Message::System(SystemMessage::new(
1533 "member prompt roster v3",
1534 ))],
1535 crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
1536 Some("agent-factory/resume".to_string()),
1537 Some(v2.clone()),
1538 )?;
1539 let v3 = session.transcript_revision()?;
1540 session.commit_transcript_rewrite(
1541 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1542 vec![Message::System(SystemMessage::new(
1543 "member prompt roster v4",
1544 ))],
1545 crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
1546 Some("agent-factory/resume".to_string()),
1547 Some(v3.clone()),
1548 )?;
1549
1550 let mut state = session
1556 .transcript_history_state()?
1557 .expect("chained refreshes retain history state");
1558 let rebookkeep = |state: &mut TranscriptHistoryState,
1559 original_parent: &str,
1560 stamp: &str|
1561 -> Result<String, Box<dyn std::error::Error>> {
1562 let body = state
1563 .revisions
1564 .iter()
1565 .find(|body| body.revision == original_parent)
1566 .expect("parent body retained")
1567 .clone();
1568 let mut messages = body.messages;
1569 messages[0] = Message::System(SystemMessage::new(stamp));
1570 let revision = transcript_messages_digest(&messages)?;
1571 state
1574 .revisions
1575 .push(crate::session::TranscriptRevisionBody {
1576 revision: revision.clone(),
1577 parent_revision: Some(original_parent.to_string()),
1578 messages,
1579 created_at: SystemTime::now(),
1580 });
1581 Ok(revision)
1582 };
1583 let v2_rebookkept = rebookkeep(&mut state, &v2, "member prompt roster v2 restamped")?;
1584 let v3_rebookkept = rebookkeep(&mut state, &v3, "member prompt roster v3 restamped")?;
1585 let respan = |state: &TranscriptHistoryState,
1589 parent: &str|
1590 -> Result<String, Box<dyn std::error::Error>> {
1591 let body = state
1592 .revisions
1593 .iter()
1594 .find(|body| body.revision == parent)
1595 .expect("rebookkept parent body retained");
1596 Ok(transcript_messages_digest(&body.messages[0..1])?)
1597 };
1598 state.commits[1].original_span_digest = respan(&state, &v2_rebookkept)?;
1599 state.commits[1].parent_revision = v2_rebookkept;
1600 state.commits[2].original_span_digest = respan(&state, &v3_rebookkept)?;
1601 state.commits[2].parent_revision = v3_rebookkept;
1602
1603 let mut incoming = session;
1605 incoming.set_metadata_unchecked_for_test(
1606 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1607 serde_json::to_value(&state)?,
1608 );
1609 let mut state = incoming
1610 .transcript_history_state()?
1611 .expect("history state survives rebookkeeping");
1612 for text in ["what was the codeword?", "birch seventeen"] {
1613 incoming.push(Message::User(UserMessage::text(text.to_string())));
1614 let appended_revision = incoming.transcript_revision()?;
1615 state
1616 .revisions
1617 .push(crate::session::TranscriptRevisionBody {
1618 revision: appended_revision.clone(),
1619 parent_revision: Some(state.head.clone()),
1620 messages: incoming.messages().to_vec(),
1621 created_at: SystemTime::now(),
1622 });
1623 state.head = appended_revision;
1624 }
1625 incoming.set_metadata_unchecked_for_test(
1626 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1627 serde_json::to_value(state)?,
1628 );
1629
1630 run_boundary_snapshot_save_guard(&incoming, Some(&previous))?;
1631 Ok(())
1632 }
1633
1634 #[test]
1645 #[allow(clippy::expect_used)]
1646 fn boundary_commit_accepts_append_after_chained_promptless_system_refreshes()
1647 -> Result<(), Box<dyn std::error::Error>> {
1648 let mut base = Session::new();
1649 base.push(Message::System(SystemMessage::new(
1650 "member prompt roster v1",
1651 )));
1652 base.push(Message::User(UserMessage::text(
1653 "the codeword is birch seventeen".to_string(),
1654 )));
1655 let v1 = base.transcript_revision()?;
1656
1657 let mut refreshed_once = base.clone();
1660 refreshed_once.commit_transcript_rewrite(
1661 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1662 vec![Message::System(SystemMessage::new(
1663 "member prompt roster v2",
1664 ))],
1665 crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
1666 Some("agent-factory/resume".to_string()),
1667 Some(v1),
1668 )?;
1669 let v2 = refreshed_once.transcript_revision()?;
1670
1671 let mut previous = refreshed_once.clone();
1673 previous.commit_transcript_rewrite(
1674 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
1675 vec![Message::System(SystemMessage::new(
1676 "member prompt roster v3",
1677 ))],
1678 crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
1679 Some("agent-factory/resume".to_string()),
1680 Some(v2),
1681 )?;
1682
1683 let mut incoming = previous.clone();
1689 let mut state = incoming
1690 .transcript_history_state()?
1691 .expect("chained refreshes retain history state");
1692 for text in ["what was the codeword?", "birch seventeen"] {
1693 incoming.push(Message::User(UserMessage::text(text.to_string())));
1694 let appended_revision = incoming.transcript_revision()?;
1695 state
1696 .revisions
1697 .push(crate::session::TranscriptRevisionBody {
1698 revision: appended_revision.clone(),
1699 parent_revision: Some(state.head.clone()),
1700 messages: incoming.messages().to_vec(),
1701 created_at: SystemTime::now(),
1702 });
1703 state.head = appended_revision;
1704 }
1705 incoming.set_metadata_unchecked_for_test(
1706 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
1707 serde_json::to_value(state)?,
1708 );
1709
1710 run_boundary_snapshot_save_guard(&incoming, Some(&previous))?;
1711 Ok(())
1712 }
1713
1714 #[test]
1721 #[allow(clippy::expect_used)]
1722 fn classify_live_session_authority_is_decided_by_machine() {
1723 use crate::session_document::{
1724 LiveSessionAuthorityKind, LiveSessionAuthorityReason, SessionDocumentEffect,
1725 SessionDocumentMachineAuthority,
1726 };
1727
1728 fn classify(
1729 stored_transcript_diverged: bool,
1730 live_has_uncommitted_transcript: bool,
1731 runtime_system_context_diverged: bool,
1732 stored_is_archived: bool,
1733 ) -> (LiveSessionAuthorityKind, LiveSessionAuthorityReason) {
1734 let mut authority = SessionDocumentMachineAuthority::new();
1735 let effects = authority
1736 .classify_live_session_authority(
1737 stored_transcript_diverged,
1738 live_has_uncommitted_transcript,
1739 runtime_system_context_diverged,
1740 stored_is_archived,
1741 )
1742 .expect("classifier must resolve a verdict");
1743 effects
1744 .iter()
1745 .find_map(|effect| match effect {
1746 SessionDocumentEffect::LiveSessionAuthorityClassified { authority, reason } => {
1747 Some((*authority, *reason))
1748 }
1749 _ => None,
1750 })
1751 .expect("classifier must emit a verdict")
1752 }
1753
1754 let (kind, _) = classify(false, false, false, false);
1756 assert_eq!(kind, LiveSessionAuthorityKind::LiveAuthoritative);
1757
1758 assert_eq!(
1760 classify(true, false, false, false),
1761 (
1762 LiveSessionAuthorityKind::DurableAuthoritative,
1763 LiveSessionAuthorityReason::StoredTranscriptRevisionDiverged
1764 ),
1765 );
1766 assert_eq!(
1767 classify(false, true, false, false),
1768 (
1769 LiveSessionAuthorityKind::DurableAuthoritative,
1770 LiveSessionAuthorityReason::LiveUncommittedTranscript
1771 ),
1772 );
1773 assert_eq!(
1774 classify(false, false, true, false),
1775 (
1776 LiveSessionAuthorityKind::DurableAuthoritative,
1777 LiveSessionAuthorityReason::RuntimeSystemContextDiverged
1778 ),
1779 );
1780 assert_eq!(
1781 classify(false, false, false, true),
1782 (
1783 LiveSessionAuthorityKind::DurableAuthoritative,
1784 LiveSessionAuthorityReason::StoredArchived
1785 ),
1786 );
1787
1788 assert_eq!(
1791 classify(true, true, true, true),
1792 (
1793 LiveSessionAuthorityKind::DurableAuthoritative,
1794 LiveSessionAuthorityReason::StoredArchived
1795 ),
1796 );
1797 assert_eq!(
1799 classify(true, true, true, false),
1800 (
1801 LiveSessionAuthorityKind::DurableAuthoritative,
1802 LiveSessionAuthorityReason::LiveUncommittedTranscript
1803 ),
1804 );
1805 assert_eq!(
1807 classify(true, false, true, false),
1808 (
1809 LiveSessionAuthorityKind::DurableAuthoritative,
1810 LiveSessionAuthorityReason::RuntimeSystemContextDiverged
1811 ),
1812 );
1813 }
1814
1815 #[test]
1816 fn append_only_guard_rejects_leading_system_message_replacement() {
1817 let mut previous = Session::new();
1818 previous.push(Message::System(SystemMessage::new("original system")));
1819 previous.push(Message::User(UserMessage::text("hello".to_string())));
1820
1821 let mut incoming = previous.clone();
1822 let rewrite_result = incoming.replace_messages_internal(
1823 vec![
1824 Message::System(SystemMessage::new("rewritten system")),
1825 Message::User(UserMessage::text("hello".to_string())),
1826 ],
1827 crate::TranscriptRewriteReason::new("unit-test"),
1828 );
1829 assert!(
1830 rewrite_result.is_ok(),
1831 "typed rewrite should be constructible: {rewrite_result:?}"
1832 );
1833
1834 assert!(matches!(
1835 append_only_save_guard(&incoming, Some(&previous)),
1836 Err(SessionStoreError::TranscriptContinuityViolation { .. })
1837 ));
1838 }
1839
1840 #[test]
1841 fn append_only_guard_accepts_runtime_system_context_append()
1842 -> Result<(), Box<dyn std::error::Error>> {
1843 let mut previous = Session::new();
1844 previous.push(Message::System(SystemMessage::new("base system")));
1845 previous.push(Message::User(UserMessage::text("hello".to_string())));
1846
1847 let mut incoming = previous.clone();
1848 incoming.set_system_prompt_with_source(
1852 format!(
1853 "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nsource: unit-test\n\nextra context"
1854 ),
1855 crate::session_durable_config_authority::SessionSystemPromptSource::RuntimeContextAppend,
1856 )?;
1857
1858 assert!(append_only_save_guard(&incoming, Some(&previous)).is_ok());
1859 Ok(())
1860 }
1861
1862 #[test]
1863 fn append_only_guard_rejects_append_shaped_prompt_without_runtime_context_marker() {
1864 let mut previous = Session::new();
1865 previous.push(Message::System(SystemMessage::new("base system")));
1866 previous.push(Message::User(UserMessage::text("hello".to_string())));
1867
1868 let mut incoming = previous.clone();
1872 incoming.set_system_prompt(format!(
1873 "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nsource: forged\n\nextra context"
1874 ));
1875
1876 assert!(matches!(
1877 append_only_save_guard(&incoming, Some(&previous)),
1878 Err(SessionStoreError::TranscriptContinuityViolation { .. })
1879 ));
1880 }
1881
1882 #[test]
1883 fn append_only_guard_accepts_system_timestamp_refresh_without_content_change() {
1884 let mut previous = Session::new();
1885 previous.push(Message::System(SystemMessage::new("base system")));
1886
1887 let mut incoming = previous.clone();
1888 incoming.set_system_prompt("base system".to_string());
1889
1890 assert!(append_only_save_guard(&incoming, Some(&previous)).is_ok());
1891 }
1892
1893 #[test]
1894 fn run_boundary_guard_accepts_compaction_after_uncheckpointed_runtime_append()
1895 -> Result<(), Box<dyn std::error::Error>> {
1896 let mut previous = Session::new();
1897 previous.push(Message::System(SystemMessage::new("base system")));
1898 previous.push(Message::User(UserMessage::text("turn one".to_string())));
1899 previous.push(Message::BlockAssistant(BlockAssistantMessage {
1900 blocks: vec![AssistantBlock::Text {
1901 text: "answer one".to_string(),
1902 meta: None,
1903 }],
1904 stop_reason: StopReason::EndTurn,
1905 identity: crate::types::TranscriptMessageIdentity::default(),
1906 created_at: crate::types::message_timestamp_now(),
1907 }));
1908
1909 let mut parent = previous.clone();
1910 parent.set_system_prompt("refreshed runtime system projection".to_string());
1911 parent.push(Message::User(UserMessage::text(
1912 "runtime-only turn".to_string(),
1913 )));
1914 parent.push(Message::BlockAssistant(BlockAssistantMessage {
1915 blocks: vec![AssistantBlock::Text {
1916 text: "runtime-only answer".to_string(),
1917 meta: None,
1918 }],
1919 stop_reason: StopReason::EndTurn,
1920 identity: crate::types::TranscriptMessageIdentity::default(),
1921 created_at: crate::types::message_timestamp_now(),
1922 }));
1923 let parent_revision = parent.transcript_revision()?;
1924
1925 let mut incoming = parent.clone();
1926 let mut replacement = vec![
1927 parent.messages()[0].clone(),
1928 Message::User(UserMessage::compaction_summary(
1929 "[Context compacted] summary".to_string(),
1930 )),
1931 ];
1932 replacement.extend_from_slice(&parent.messages()[1..]);
1933 incoming.commit_transcript_rewrite(
1934 TranscriptRewriteSelection::MessageRange {
1935 start: 0,
1936 end: parent.messages().len(),
1937 },
1938 replacement,
1939 crate::TranscriptRewriteReason::new("compaction"),
1940 Some("meerkat-core".to_string()),
1941 Some(parent_revision),
1942 )?;
1943
1944 assert!(matches!(
1945 append_only_save_guard(&incoming, Some(&previous)),
1946 Err(SessionStoreError::TranscriptContinuityViolation { .. })
1947 ));
1948 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
1949 Ok(())
1950 }
1951
1952 #[test]
1953 fn run_boundary_guard_accepts_compaction_with_retained_tail_window()
1954 -> Result<(), Box<dyn std::error::Error>> {
1955 let mut previous = Session::new();
1956 previous.push(Message::System(SystemMessage::new("base system")));
1957 previous.push(Message::User(UserMessage::text("turn one".to_string())));
1958 previous.push(Message::BlockAssistant(BlockAssistantMessage::new(
1959 vec![crate::types::AssistantBlock::Text {
1960 text: "answer one".to_string(),
1961 meta: None,
1962 }],
1963 StopReason::EndTurn,
1964 )));
1965
1966 let mut parent = previous.clone();
1967 parent.set_system_prompt("refreshed runtime system projection".to_string());
1968 parent.push(Message::SystemNotice(SystemNoticeMessage::new(
1969 SystemNoticeKind::Comms,
1970 "peer response queued",
1971 )));
1972 let parent_revision = parent.transcript_revision()?;
1973
1974 let mut incoming = parent.clone();
1975 let mut replacement = vec![
1976 parent.messages()[0].clone(),
1977 Message::User(UserMessage::compaction_summary(
1978 "[Context compacted] summary".to_string(),
1979 )),
1980 ];
1981 replacement.extend_from_slice(&parent.messages()[1..]);
1982 incoming.commit_transcript_rewrite(
1983 TranscriptRewriteSelection::MessageRange {
1984 start: 0,
1985 end: parent.messages().len(),
1986 },
1987 replacement,
1988 crate::TranscriptRewriteReason::new("compaction"),
1989 Some("meerkat-core".to_string()),
1990 Some(parent_revision),
1991 )?;
1992
1993 assert!(matches!(
1994 append_only_save_guard(&incoming, Some(&previous)),
1995 Err(SessionStoreError::TranscriptContinuityViolation { .. })
1996 ));
1997 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
1998 Ok(())
1999 }
2000
2001 #[test]
2002 fn run_boundary_guard_rejects_commitless_history_parent_edge()
2003 -> Result<(), Box<dyn std::error::Error>> {
2004 let mut previous = Session::new();
2005 previous.push(Message::System(SystemMessage::new("base system")));
2006 previous.push(Message::User(UserMessage::text("turn one".to_string())));
2007 let previous_revision = previous.transcript_revision()?;
2008
2009 let mut incoming = previous.clone();
2010 incoming.set_system_prompt("forged replacement system".to_string());
2011 let incoming_revision = incoming.transcript_revision()?;
2012 let history = TranscriptHistoryState {
2013 head: incoming_revision.clone(),
2014 commits: Vec::new(),
2015 revisions: vec![
2016 crate::TranscriptRevisionBody {
2017 revision: previous_revision,
2018 parent_revision: None,
2019 messages: previous.messages().to_vec(),
2020 created_at: previous.updated_at(),
2021 },
2022 crate::TranscriptRevisionBody {
2023 revision: incoming_revision,
2024 parent_revision: Some(previous.transcript_revision()?),
2025 messages: incoming.messages().to_vec(),
2026 created_at: incoming.updated_at(),
2027 },
2028 ],
2029 };
2030 incoming.set_metadata_unchecked_for_test(
2031 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2032 serde_json::to_value(history)?,
2033 );
2034
2035 assert!(matches!(
2036 append_only_save_guard(&incoming, Some(&previous)),
2037 Err(SessionStoreError::TranscriptContinuityViolation { .. })
2038 ));
2039 assert!(matches!(
2040 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2041 Err(SessionStoreError::TranscriptContinuityViolation { .. }
2042 | SessionStoreError::MonotonicityViolation { .. })
2043 ));
2044 Ok(())
2045 }
2046
2047 #[test]
2048 fn append_only_guard_rejects_history_head_that_does_not_match_current_messages()
2049 -> Result<(), Box<dyn std::error::Error>> {
2050 let mut previous = Session::new();
2051 previous.push(Message::User(UserMessage::text("persisted".to_string())));
2052
2053 let mut incoming = previous.clone();
2054 incoming.push(Message::User(UserMessage::text("append".to_string())));
2055 let poisoned_messages = vec![Message::User(UserMessage::text(
2056 "unrelated poisoned history".to_string(),
2057 ))];
2058 let poisoned_revision = transcript_messages_digest(&poisoned_messages)?;
2059 incoming.set_metadata_unchecked_for_test(
2060 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2061 serde_json::to_value(TranscriptHistoryState {
2062 head: poisoned_revision.clone(),
2063 commits: Vec::new(),
2064 revisions: vec![crate::TranscriptRevisionBody {
2065 revision: poisoned_revision,
2066 parent_revision: None,
2067 messages: poisoned_messages,
2068 created_at: incoming.updated_at(),
2069 }],
2070 })?,
2071 );
2072
2073 assert!(matches!(
2074 append_only_save_guard(&incoming, Some(&previous)),
2075 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2076 ));
2077 assert!(matches!(
2078 append_only_save_guard(&incoming, None),
2079 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2080 ));
2081 assert!(matches!(
2082 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2083 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2084 ));
2085 Ok(())
2086 }
2087
2088 #[test]
2089 fn append_only_guard_rejects_new_rewrite_commits_on_plain_append()
2090 -> Result<(), Box<dyn std::error::Error>> {
2091 let mut previous = Session::new();
2092 previous.push(Message::User(UserMessage::text("persisted".to_string())));
2093 let previous_revision = previous.transcript_revision()?;
2094
2095 let mut incoming = previous.clone();
2096 let appended = Message::BlockAssistant(BlockAssistantMessage {
2097 blocks: vec![AssistantBlock::Text {
2098 text: "plain append".to_string(),
2099 meta: None,
2100 }],
2101 stop_reason: StopReason::EndTurn,
2102 identity: crate::types::TranscriptMessageIdentity::default(),
2103 created_at: crate::types::message_timestamp_now(),
2104 });
2105 incoming.commit_transcript_rewrite(
2106 TranscriptRewriteSelection::MessageRange { start: 1, end: 1 },
2107 vec![appended],
2108 crate::TranscriptRewriteReason::new("forged-append"),
2109 Some("unit-test".to_string()),
2110 Some(previous_revision),
2111 )?;
2112
2113 assert!(matches!(
2114 append_only_save_guard(&incoming, Some(&previous)),
2115 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2116 ));
2117 Ok(())
2118 }
2119
2120 #[test]
2121 fn append_only_guard_rejects_first_save_with_rewrite_commits()
2122 -> Result<(), Box<dyn std::error::Error>> {
2123 let mut incoming = Session::new();
2124 incoming.push(Message::User(UserMessage::text("seed".to_string())));
2125 let parent_messages = incoming.messages().to_vec();
2126 let parent_updated_at = incoming.updated_at();
2127 let parent_revision = incoming.transcript_revision()?;
2128 let commit = incoming.commit_transcript_rewrite(
2129 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2130 vec![Message::User(UserMessage::text(
2131 "compacted seed".to_string(),
2132 ))],
2133 crate::TranscriptRewriteReason::new("compaction"),
2134 Some("meerkat-core".to_string()),
2135 Some(parent_revision),
2136 )?;
2137 let incoming_revision = incoming.transcript_revision()?;
2138 let commit_parent_revision = commit.parent_revision.clone();
2139 incoming.set_metadata_unchecked_for_test(
2140 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2141 serde_json::to_value(TranscriptHistoryState {
2142 head: incoming_revision.clone(),
2143 commits: vec![commit],
2144 revisions: vec![
2145 crate::TranscriptRevisionBody {
2146 revision: commit_parent_revision.clone(),
2147 parent_revision: None,
2148 messages: parent_messages,
2149 created_at: parent_updated_at,
2150 },
2151 crate::TranscriptRevisionBody {
2152 revision: incoming_revision,
2153 parent_revision: Some(commit_parent_revision),
2154 messages: incoming.messages().to_vec(),
2155 created_at: incoming.updated_at(),
2156 },
2157 ],
2158 })?,
2159 );
2160
2161 assert!(matches!(
2162 append_only_save_guard(&incoming, None),
2163 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2164 ));
2165 Ok(())
2166 }
2167
2168 #[test]
2169 fn transcript_rewrite_guard_rejects_poisoned_history_graph()
2170 -> Result<(), Box<dyn std::error::Error>> {
2171 let mut previous = Session::new();
2172 previous.push(Message::User(UserMessage::text("persisted".to_string())));
2173 let parent_revision = previous.transcript_revision()?;
2174
2175 let mut first = previous.clone();
2176 let first_commit = first.commit_transcript_rewrite(
2177 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2178 vec![Message::User(UserMessage::text(
2179 "compacted persisted".to_string(),
2180 ))],
2181 crate::TranscriptRewriteReason::new("compaction"),
2182 Some("unit-test".to_string()),
2183 Some(parent_revision),
2184 )?;
2185 let first_snapshot = first.clone();
2186
2187 first.commit_transcript_rewrite(
2188 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2189 vec![Message::User(UserMessage::text(
2190 "uncommitted poisoned fork".to_string(),
2191 ))],
2192 crate::TranscriptRewriteReason::new("poison"),
2193 Some("unit-test".to_string()),
2194 Some(first_commit.revision.clone()),
2195 )?;
2196 let mut poisoned_state = first
2197 .transcript_history_state()?
2198 .ok_or_else(|| "second rewrite should retain history state".to_string())?;
2199 poisoned_state.head = first_commit.revision.clone();
2200
2201 let mut poisoned = first_snapshot;
2202 poisoned.set_metadata_unchecked_for_test(
2203 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2204 serde_json::to_value(poisoned_state)?,
2205 );
2206
2207 assert!(matches!(
2208 transcript_rewrite_save_guard(&poisoned, Some(&previous), &first_commit),
2209 Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
2210 if reason.contains("incoming transcript history state is malformed")
2211 ));
2212 Ok(())
2213 }
2214
2215 #[test]
2216 fn authoritative_projection_guard_rejects_changed_persisted_revision()
2217 -> Result<(), Box<dyn std::error::Error>> {
2218 let mut previous = Session::new();
2219 previous.push(Message::User(UserMessage::text("persisted A".to_string())));
2220 let expected_revision = previous.transcript_revision()?;
2221
2222 let mut current = previous.clone();
2223 current.push(Message::BlockAssistant(BlockAssistantMessage {
2224 blocks: vec![AssistantBlock::Text {
2225 text: "persisted B".to_string(),
2226 meta: None,
2227 }],
2228 stop_reason: StopReason::EndTurn,
2229 identity: crate::types::TranscriptMessageIdentity::default(),
2230 created_at: crate::types::message_timestamp_now(),
2231 }));
2232 let mut incoming = previous.clone();
2233 incoming.push(Message::User(UserMessage::text(
2234 "incoming from A".to_string(),
2235 )));
2236
2237 assert!(matches!(
2238 authoritative_projection_current_revision_guard(
2239 &incoming,
2240 Some(¤t),
2241 Some(&expected_revision)
2242 ),
2243 Err(SessionStoreError::TranscriptContinuityViolation { .. })
2244 ));
2245 Ok(())
2246 }
2247
2248 #[test]
2249 fn append_only_guard_rejects_rewrite_commits_on_first_save()
2250 -> Result<(), Box<dyn std::error::Error>> {
2251 let mut incoming = Session::new();
2252 incoming.push(Message::User(UserMessage::text("persisted".to_string())));
2253 let parent_revision = incoming.transcript_revision()?;
2254 incoming.commit_transcript_rewrite(
2255 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2256 vec![Message::User(UserMessage::text("rewritten".to_string()))],
2257 crate::TranscriptRewriteReason::new("forged-first-save"),
2258 Some("unit-test".to_string()),
2259 Some(parent_revision),
2260 )?;
2261
2262 assert!(matches!(
2263 append_only_save_guard(&incoming, None),
2264 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2265 ));
2266 Ok(())
2267 }
2268
2269 #[test]
2270 fn append_only_guard_rejects_commitless_history_on_first_save()
2271 -> Result<(), Box<dyn std::error::Error>> {
2272 let mut incoming = Session::new();
2273 incoming.push(Message::User(UserMessage::text("persisted".to_string())));
2274 let incoming_revision = incoming.transcript_revision()?;
2275 incoming.set_metadata_unchecked_for_test(
2276 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2277 serde_json::to_value(TranscriptHistoryState {
2278 head: incoming_revision.clone(),
2279 commits: Vec::new(),
2280 revisions: vec![crate::TranscriptRevisionBody {
2281 revision: incoming_revision,
2282 parent_revision: None,
2283 messages: incoming.messages().to_vec(),
2284 created_at: incoming.updated_at(),
2285 }],
2286 })?,
2287 );
2288
2289 assert!(matches!(
2290 append_only_save_guard(&incoming, None),
2291 Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
2292 if reason.contains("first save would seed transcript history state")
2293 ));
2294 Ok(())
2295 }
2296
2297 #[test]
2298 fn run_boundary_guard_adopts_commit_carrying_history_on_first_commit()
2299 -> Result<(), Box<dyn std::error::Error>> {
2300 let mut incoming = Session::new();
2307 incoming.set_system_prompt("old base".to_string());
2308 incoming.push(Message::User(UserMessage::text("hello".to_string())));
2309 incoming.commit_transcript_rewrite(
2310 crate::TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2311 vec![Message::System(SystemMessage::with_mutation_kind(
2312 "new base".to_string(),
2313 crate::types::SystemPromptMutationKind::ExplicitBuild,
2314 ))],
2315 crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
2316 None,
2317 None,
2318 )?;
2319
2320 assert!(run_boundary_snapshot_save_guard(&incoming, None).is_ok());
2321 assert!(matches!(
2322 append_only_save_guard(&incoming, None),
2323 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2324 ));
2325 Ok(())
2326 }
2327
2328 #[test]
2329 fn run_boundary_guard_rejects_untyped_leading_system_refresh_after_head_rewrite()
2330 -> Result<(), Box<dyn std::error::Error>> {
2331 let mut previous = Session::new();
2338 previous.set_system_prompt("original base".to_string());
2339 previous.push(Message::User(UserMessage::text("hello".to_string())));
2340 previous.commit_transcript_rewrite(
2341 crate::TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2342 vec![Message::System(SystemMessage::with_mutation_kind(
2343 "refreshed base".to_string(),
2344 crate::types::SystemPromptMutationKind::ExplicitBuild,
2345 ))],
2346 crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
2347 None,
2348 None,
2349 )?;
2350
2351 let mut incoming = previous.clone();
2352 incoming.set_system_prompt("untyped hijack".to_string());
2353
2354 assert!(matches!(
2355 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2356 Err(SessionStoreError::TranscriptContinuityViolation { .. })
2357 ));
2358 Ok(())
2359 }
2360
2361 #[test]
2362 fn run_boundary_guard_accepts_plain_append_after_head_rewrite()
2363 -> Result<(), Box<dyn std::error::Error>> {
2364 let mut previous = Session::new();
2368 previous.set_system_prompt("original base".to_string());
2369 previous.push(Message::User(UserMessage::text("hello".to_string())));
2370 previous.commit_transcript_rewrite(
2371 crate::TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2372 vec![Message::System(SystemMessage::with_mutation_kind(
2373 "refreshed base".to_string(),
2374 crate::types::SystemPromptMutationKind::ExplicitBuild,
2375 ))],
2376 crate::TranscriptRewriteReason::new("resume-system-prompt-refresh"),
2377 None,
2378 None,
2379 )?;
2380
2381 let mut incoming = previous.clone();
2382 incoming.push(Message::User(UserMessage::text(
2383 "post-rewrite turn".to_string(),
2384 )));
2385
2386 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2387 Ok(())
2388 }
2389
2390 #[test]
2391 fn append_only_guard_rejects_commitless_history_seed_on_plain_append()
2392 -> Result<(), Box<dyn std::error::Error>> {
2393 let mut previous = Session::new();
2394 previous.push(Message::User(UserMessage::text("persisted".to_string())));
2395 let previous_revision = previous.transcript_revision()?;
2396
2397 let mut incoming = previous.clone();
2398 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2399 blocks: vec![AssistantBlock::Text {
2400 text: "plain append".to_string(),
2401 meta: None,
2402 }],
2403 stop_reason: StopReason::EndTurn,
2404 identity: crate::types::TranscriptMessageIdentity::default(),
2405 created_at: crate::types::message_timestamp_now(),
2406 }));
2407 let incoming_revision = incoming.transcript_revision()?;
2408 incoming.set_metadata_unchecked_for_test(
2409 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2410 serde_json::to_value(TranscriptHistoryState {
2411 head: incoming_revision.clone(),
2412 commits: Vec::new(),
2413 revisions: vec![
2414 crate::TranscriptRevisionBody {
2415 revision: previous_revision,
2416 parent_revision: None,
2417 messages: previous.messages().to_vec(),
2418 created_at: previous.updated_at(),
2419 },
2420 crate::TranscriptRevisionBody {
2421 revision: incoming_revision,
2422 parent_revision: Some(previous.transcript_revision()?),
2423 messages: incoming.messages().to_vec(),
2424 created_at: incoming.updated_at(),
2425 },
2426 ],
2427 })?,
2428 );
2429
2430 assert!(matches!(
2431 append_only_save_guard(&incoming, Some(&previous)),
2432 Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
2433 if reason.contains("append-only save would seed transcript history state")
2434 ));
2435 Ok(())
2436 }
2437
2438 #[test]
2439 fn run_boundary_guard_accepts_commitless_history_seed_on_plain_append()
2440 -> Result<(), Box<dyn std::error::Error>> {
2441 let mut previous = Session::new();
2442 previous.push(Message::User(UserMessage::text("persisted".to_string())));
2443 let previous_revision = previous.transcript_revision()?;
2444
2445 let mut incoming = previous.clone();
2446 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2447 blocks: vec![AssistantBlock::Text {
2448 text: "plain append".to_string(),
2449 meta: None,
2450 }],
2451 stop_reason: StopReason::EndTurn,
2452 identity: crate::types::TranscriptMessageIdentity::default(),
2453 created_at: crate::types::message_timestamp_now(),
2454 }));
2455 let incoming_revision = incoming.transcript_revision()?;
2456 incoming.set_metadata_unchecked_for_test(
2457 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2458 serde_json::to_value(TranscriptHistoryState {
2459 head: incoming_revision.clone(),
2460 commits: Vec::new(),
2461 revisions: vec![
2462 crate::TranscriptRevisionBody {
2463 revision: previous_revision.clone(),
2464 parent_revision: None,
2465 messages: previous.messages().to_vec(),
2466 created_at: previous.updated_at(),
2467 },
2468 crate::TranscriptRevisionBody {
2469 revision: incoming_revision,
2470 parent_revision: Some(previous_revision),
2471 messages: incoming.messages().to_vec(),
2472 created_at: incoming.updated_at(),
2473 },
2474 ],
2475 })?,
2476 );
2477
2478 assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2479 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2480 Ok(())
2481 }
2482
2483 #[test]
2484 fn run_boundary_guard_accepts_retained_history_seed_on_plain_append()
2485 -> Result<(), Box<dyn std::error::Error>> {
2486 let mut original = Session::new();
2487 original.push(Message::User(UserMessage::text("verbose seed".to_string())));
2488 let original_revision = original.transcript_revision()?;
2489
2490 let mut previous = original.clone();
2491 previous.commit_transcript_rewrite(
2492 TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
2493 vec![Message::User(UserMessage::text(
2494 "compacted seed".to_string(),
2495 ))],
2496 crate::TranscriptRewriteReason::new("compaction"),
2497 Some("meerkat-core".to_string()),
2498 Some(original_revision),
2499 )?;
2500 let previous_with_history = previous.clone();
2501 previous.clear_transcript_history_state();
2502
2503 let mut incoming = previous_with_history;
2504 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2505 blocks: vec![AssistantBlock::Text {
2506 text: "plain append after retained history".to_string(),
2507 meta: None,
2508 }],
2509 stop_reason: StopReason::EndTurn,
2510 identity: crate::types::TranscriptMessageIdentity::default(),
2511 created_at: crate::types::message_timestamp_now(),
2512 }));
2513
2514 assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2515 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2516 Ok(())
2517 }
2518
2519 #[test]
2520 fn run_boundary_guard_accepts_commitless_history_seed_on_first_snapshot()
2521 -> Result<(), Box<dyn std::error::Error>> {
2522 let mut incoming = Session::new();
2523 incoming.push(Message::User(UserMessage::text("persisted".to_string())));
2524 let incoming_revision = incoming.transcript_revision()?;
2525 incoming.set_metadata_unchecked_for_test(
2526 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2527 serde_json::to_value(TranscriptHistoryState {
2528 head: incoming_revision.clone(),
2529 commits: Vec::new(),
2530 revisions: vec![crate::TranscriptRevisionBody {
2531 revision: incoming_revision,
2532 parent_revision: None,
2533 messages: incoming.messages().to_vec(),
2534 created_at: incoming.updated_at(),
2535 }],
2536 })?,
2537 );
2538
2539 assert!(append_only_save_guard(&incoming, None).is_err());
2540 assert!(run_boundary_snapshot_save_guard(&incoming, None).is_ok());
2541 Ok(())
2542 }
2543
2544 #[test]
2545 fn run_boundary_guard_accepts_commitless_history_seed_on_initial_multi_revision_snapshot()
2546 -> Result<(), Box<dyn std::error::Error>> {
2547 let mut base = Session::new();
2548 base.push(Message::User(UserMessage::text("first".to_string())));
2549 let base_revision = base.transcript_revision()?;
2550
2551 let mut incoming = base.clone();
2552 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2553 blocks: vec![AssistantBlock::Text {
2554 text: "second".to_string(),
2555 meta: None,
2556 }],
2557 stop_reason: StopReason::EndTurn,
2558 identity: crate::types::TranscriptMessageIdentity::default(),
2559 created_at: crate::types::message_timestamp_now(),
2560 }));
2561 let incoming_revision = incoming.transcript_revision()?;
2562 incoming.set_metadata_unchecked_for_test(
2563 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2564 serde_json::to_value(TranscriptHistoryState {
2565 head: incoming_revision.clone(),
2566 commits: Vec::new(),
2567 revisions: vec![
2568 crate::TranscriptRevisionBody {
2569 revision: base_revision.clone(),
2570 parent_revision: None,
2571 messages: base.messages().to_vec(),
2572 created_at: base.updated_at(),
2573 },
2574 crate::TranscriptRevisionBody {
2575 revision: incoming_revision,
2576 parent_revision: Some(base_revision),
2577 messages: incoming.messages().to_vec(),
2578 created_at: incoming.updated_at(),
2579 },
2580 ],
2581 })?,
2582 );
2583
2584 assert!(append_only_save_guard(&incoming, None).is_err());
2585 assert!(run_boundary_snapshot_save_guard(&incoming, None).is_ok());
2586 Ok(())
2587 }
2588
2589 #[test]
2590 fn append_only_guard_rejects_new_rewrite_commits_on_system_context_append()
2591 -> Result<(), Box<dyn std::error::Error>> {
2592 let mut previous = Session::new();
2593 previous.push(Message::System(SystemMessage::new("base system")));
2594 previous.push(Message::User(UserMessage::text("persisted".to_string())));
2595 let mut incoming = previous.clone();
2596 incoming.set_system_prompt_with_source(
2597 format!(
2598 "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nsource: unit-test\n\nextra context"
2599 ),
2600 crate::session_durable_config_authority::SessionSystemPromptSource::RuntimeContextAppend,
2601 )?;
2602 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2603 blocks: vec![AssistantBlock::Text {
2604 text: "plain append".to_string(),
2605 meta: None,
2606 }],
2607 stop_reason: StopReason::EndTurn,
2608 identity: crate::types::TranscriptMessageIdentity::default(),
2609 created_at: crate::types::message_timestamp_now(),
2610 }));
2611 let incoming_revision = incoming.transcript_revision()?;
2612 incoming.set_metadata_unchecked_for_test(
2613 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2614 serde_json::to_value(TranscriptHistoryState {
2615 head: incoming_revision.clone(),
2616 commits: vec![TranscriptRewriteCommit {
2617 parent_revision: previous.transcript_revision()?,
2618 revision: incoming_revision.clone(),
2619 selection: TranscriptRewriteSelection::MessageRange { start: 0, end: 0 },
2620 original_span_digest: transcript_messages_digest(&[])?,
2621 replacement_digest: transcript_messages_digest(&[])?,
2622 messages_before: previous.messages().len(),
2623 messages_after: incoming.messages().len(),
2624 reason: crate::TranscriptRewriteReason::new("forged"),
2625 actor: Some("unit-test".to_string()),
2626 committed_at: incoming.updated_at(),
2627 }],
2628 revisions: vec![crate::TranscriptRevisionBody {
2629 revision: incoming_revision,
2630 parent_revision: None,
2631 messages: incoming.messages().to_vec(),
2632 created_at: incoming.updated_at(),
2633 }],
2634 })?,
2635 );
2636
2637 assert!(matches!(
2638 append_only_save_guard(&incoming, Some(&previous)),
2639 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2640 ));
2641 Ok(())
2642 }
2643
2644 #[test]
2645 fn append_only_guard_rejects_new_rewrite_commits_on_transient_notice_cleanup()
2646 -> Result<(), Box<dyn std::error::Error>> {
2647 let mut previous = Session::new();
2648 previous.push(Message::SystemNotice(SystemNoticeMessage::new(
2649 SystemNoticeKind::Comms,
2650 "transient peer delivery notice",
2651 )));
2652 previous.push(Message::User(UserMessage::text("persisted".to_string())));
2653
2654 let mut incoming = Session::new();
2655 incoming.push(Message::User(UserMessage::text("persisted".to_string())));
2656 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2657 blocks: vec![AssistantBlock::Text {
2658 text: "plain append after notice cleanup".to_string(),
2659 meta: None,
2660 }],
2661 stop_reason: StopReason::EndTurn,
2662 identity: crate::types::TranscriptMessageIdentity::default(),
2663 created_at: crate::types::message_timestamp_now(),
2664 }));
2665 let incoming_revision = incoming.transcript_revision()?;
2666 incoming.set_metadata_unchecked_for_test(
2667 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
2668 serde_json::to_value(TranscriptHistoryState {
2669 head: incoming_revision.clone(),
2670 commits: vec![TranscriptRewriteCommit {
2671 parent_revision: previous.transcript_revision()?,
2672 revision: incoming_revision.clone(),
2673 selection: TranscriptRewriteSelection::MessageRange { start: 0, end: 0 },
2674 original_span_digest: transcript_messages_digest(&[])?,
2675 replacement_digest: transcript_messages_digest(&[])?,
2676 messages_before: previous.messages().len(),
2677 messages_after: incoming.messages().len(),
2678 reason: crate::TranscriptRewriteReason::new("forged"),
2679 actor: Some("unit-test".to_string()),
2680 committed_at: incoming.updated_at(),
2681 }],
2682 revisions: vec![crate::TranscriptRevisionBody {
2683 revision: incoming_revision,
2684 parent_revision: None,
2685 messages: incoming.messages().to_vec(),
2686 created_at: incoming.updated_at(),
2687 }],
2688 })?,
2689 );
2690
2691 assert!(matches!(
2692 append_only_save_guard(&incoming, Some(&previous)),
2693 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
2694 ));
2695 Ok(())
2696 }
2697
2698 #[test]
2699 fn run_boundary_guard_accepts_generated_context_summary_before_retained_tail()
2700 -> Result<(), Box<dyn std::error::Error>> {
2701 let mut previous = Session::new();
2702 previous.push(Message::System(SystemMessage::new(
2703 "runtime system before context refresh",
2704 )));
2705 previous.push(Message::User(UserMessage::text(
2706 "Turn 1 request".to_string(),
2707 )));
2708 previous.push(Message::BlockAssistant(BlockAssistantMessage {
2709 blocks: vec![AssistantBlock::Text {
2710 text: "Turn 1 answer".to_string(),
2711 meta: None,
2712 }],
2713 stop_reason: StopReason::EndTurn,
2714 identity: crate::types::TranscriptMessageIdentity::default(),
2715 created_at: crate::types::message_timestamp_now(),
2716 }));
2717
2718 let mut incoming = Session::with_id(previous.id().clone());
2719 incoming.push(Message::System(SystemMessage::new(
2720 "runtime system after context refresh",
2721 )));
2722 incoming.push(Message::User(UserMessage::text(
2723 "Verbose context that will be compacted".to_string(),
2724 )));
2725 for message in previous.messages()[1..].iter().cloned() {
2726 incoming.push(message);
2727 }
2728 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2729 blocks: vec![AssistantBlock::Text {
2730 text: "Turn 2 generated answer".to_string(),
2731 meta: None,
2732 }],
2733 stop_reason: StopReason::EndTurn,
2734 identity: crate::types::TranscriptMessageIdentity::default(),
2735 created_at: crate::types::message_timestamp_now(),
2736 }));
2737 let parent_revision = incoming.transcript_revision()?;
2738 incoming.commit_transcript_rewrite(
2739 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2740 vec![Message::User(UserMessage::compaction_summary(
2741 "[Context compacted] Earlier runtime context".to_string(),
2742 ))],
2743 crate::TranscriptRewriteReason::new("compaction"),
2744 Some("meerkat-core".to_string()),
2745 Some(parent_revision),
2746 )?;
2747
2748 assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2749 assert!(run_boundary_snapshot_save_guard(&incoming, Some(&previous)).is_ok());
2750 Ok(())
2751 }
2752
2753 #[test]
2754 fn run_boundary_guard_rejects_context_summary_tail_without_compaction_summary_marker()
2755 -> Result<(), Box<dyn std::error::Error>> {
2756 let mut previous = Session::new();
2757 previous.push(Message::System(SystemMessage::new(
2758 "runtime system before context refresh",
2759 )));
2760 previous.push(Message::User(UserMessage::text(
2761 "Turn 1 request".to_string(),
2762 )));
2763 previous.push(Message::BlockAssistant(BlockAssistantMessage {
2764 blocks: vec![AssistantBlock::Text {
2765 text: "Turn 1 answer".to_string(),
2766 meta: None,
2767 }],
2768 stop_reason: StopReason::EndTurn,
2769 identity: crate::types::TranscriptMessageIdentity::default(),
2770 created_at: crate::types::message_timestamp_now(),
2771 }));
2772
2773 let mut incoming = Session::with_id(previous.id().clone());
2774 incoming.push(Message::System(SystemMessage::new(
2775 "runtime system after context refresh",
2776 )));
2777 incoming.push(Message::User(UserMessage::text(
2778 "Verbose context that will be compacted".to_string(),
2779 )));
2780 for message in previous.messages()[1..].iter().cloned() {
2781 incoming.push(message);
2782 }
2783 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2784 blocks: vec![AssistantBlock::Text {
2785 text: "Turn 2 generated answer".to_string(),
2786 meta: None,
2787 }],
2788 stop_reason: StopReason::EndTurn,
2789 identity: crate::types::TranscriptMessageIdentity::default(),
2790 created_at: crate::types::message_timestamp_now(),
2791 }));
2792 let parent_revision = incoming.transcript_revision()?;
2793 incoming.commit_transcript_rewrite(
2797 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2798 vec![Message::User(UserMessage::text(
2799 "[Context compacted] Earlier runtime context".to_string(),
2800 ))],
2801 crate::TranscriptRewriteReason::new("compaction"),
2802 Some("meerkat-core".to_string()),
2803 Some(parent_revision),
2804 )?;
2805
2806 assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2807 assert!(matches!(
2808 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2809 Err(SessionStoreError::TranscriptContinuityViolation { .. }
2810 | SessionStoreError::MonotonicityViolation { .. })
2811 ));
2812 Ok(())
2813 }
2814
2815 #[test]
2820 fn run_boundary_guard_rejects_context_summary_tail_with_injected_context_marker()
2821 -> Result<(), Box<dyn std::error::Error>> {
2822 let mut previous = Session::new();
2823 previous.push(Message::System(SystemMessage::new(
2824 "runtime system before context refresh",
2825 )));
2826 previous.push(Message::User(UserMessage::text(
2827 "Turn 1 request".to_string(),
2828 )));
2829 previous.push(Message::BlockAssistant(BlockAssistantMessage {
2830 blocks: vec![AssistantBlock::Text {
2831 text: "Turn 1 answer".to_string(),
2832 meta: None,
2833 }],
2834 stop_reason: StopReason::EndTurn,
2835 identity: crate::types::TranscriptMessageIdentity::default(),
2836 created_at: crate::types::message_timestamp_now(),
2837 }));
2838
2839 let mut incoming = Session::with_id(previous.id().clone());
2840 incoming.push(Message::System(SystemMessage::new(
2841 "runtime system after context refresh",
2842 )));
2843 incoming.push(Message::User(UserMessage::text(
2844 "Verbose context that will be compacted".to_string(),
2845 )));
2846 for message in previous.messages()[1..].iter().cloned() {
2847 incoming.push(message);
2848 }
2849 incoming.push(Message::BlockAssistant(BlockAssistantMessage {
2850 blocks: vec![AssistantBlock::Text {
2851 text: "Turn 2 generated answer".to_string(),
2852 meta: None,
2853 }],
2854 stop_reason: StopReason::EndTurn,
2855 identity: crate::types::TranscriptMessageIdentity::default(),
2856 created_at: crate::types::message_timestamp_now(),
2857 }));
2858 let parent_revision = incoming.transcript_revision()?;
2859 incoming.commit_transcript_rewrite(
2864 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
2865 vec![Message::User(UserMessage::injected_context(
2866 "[Context compacted] Earlier runtime context".to_string(),
2867 ))],
2868 crate::TranscriptRewriteReason::new("compaction"),
2869 Some("meerkat-core".to_string()),
2870 Some(parent_revision),
2871 )?;
2872
2873 assert!(append_only_save_guard(&incoming, Some(&previous)).is_err());
2874 assert!(matches!(
2875 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2876 Err(SessionStoreError::TranscriptContinuityViolation { .. }
2877 | SessionStoreError::MonotonicityViolation { .. })
2878 ));
2879 Ok(())
2880 }
2881
2882 #[test]
2883 fn run_boundary_guard_rejects_runtime_parent_with_inserted_message_before_tail()
2884 -> Result<(), Box<dyn std::error::Error>> {
2885 let mut previous = Session::new();
2886 previous.push(Message::System(SystemMessage::new("base system")));
2887 previous.push(Message::User(UserMessage::text("turn one".to_string())));
2888 previous.push(Message::BlockAssistant(BlockAssistantMessage {
2889 blocks: vec![AssistantBlock::Text {
2890 text: "answer one".to_string(),
2891 meta: None,
2892 }],
2893 stop_reason: StopReason::EndTurn,
2894 identity: crate::types::TranscriptMessageIdentity::default(),
2895 created_at: crate::types::message_timestamp_now(),
2896 }));
2897
2898 let parent_messages = vec![
2899 Message::System(SystemMessage::new("refreshed runtime system projection")),
2900 Message::User(UserMessage::text(
2901 "injected before retained tail".to_string(),
2902 )),
2903 previous.messages()[1].clone(),
2904 previous.messages()[2].clone(),
2905 ];
2906 let parent_revision = transcript_messages_digest(&parent_messages)?;
2907 let mut parent = previous.clone();
2908 parent.apply_transcript_history_state(TranscriptHistoryState {
2909 head: parent_revision.clone(),
2910 commits: Vec::new(),
2911 revisions: vec![crate::TranscriptRevisionBody {
2912 revision: parent_revision,
2913 parent_revision: None,
2914 messages: parent_messages,
2915 created_at: parent.updated_at(),
2916 }],
2917 })?;
2918 let parent_revision = parent.transcript_revision()?;
2919
2920 let mut incoming = parent.clone();
2921 incoming.commit_transcript_rewrite(
2922 TranscriptRewriteSelection::MessageRange {
2923 start: 0,
2924 end: parent.messages().len(),
2925 },
2926 vec![Message::User(UserMessage::text(
2927 "[Context compacted] summary".to_string(),
2928 ))],
2929 crate::TranscriptRewriteReason::new("compaction"),
2930 Some("meerkat-core".to_string()),
2931 Some(parent_revision),
2932 )?;
2933
2934 assert!(matches!(
2935 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
2936 Err(SessionStoreError::TranscriptContinuityViolation { .. }
2937 | SessionStoreError::MonotonicityViolation { .. })
2938 ));
2939 Ok(())
2940 }
2941
2942 #[test]
2943 fn run_boundary_guard_rejects_forged_parent_edge_before_real_rewrite_commit()
2944 -> Result<(), Box<dyn std::error::Error>> {
2945 let mut previous = Session::new();
2946 previous.push(Message::System(SystemMessage::new("base system")));
2947 previous.push(Message::User(UserMessage::text("turn one".to_string())));
2948 previous.push(Message::BlockAssistant(BlockAssistantMessage {
2949 blocks: vec![AssistantBlock::Text {
2950 text: "answer one".to_string(),
2951 meta: None,
2952 }],
2953 stop_reason: StopReason::EndTurn,
2954 identity: crate::types::TranscriptMessageIdentity::default(),
2955 created_at: crate::types::message_timestamp_now(),
2956 }));
2957 let previous_revision = previous.transcript_revision()?;
2958
2959 let forged_parent_messages = vec![
2960 Message::System(SystemMessage::new("refreshed runtime system projection")),
2961 Message::User(UserMessage::text(
2962 "forged insertion before retained tail".to_string(),
2963 )),
2964 previous.messages()[1].clone(),
2965 previous.messages()[2].clone(),
2966 ];
2967 let forged_parent_revision = transcript_messages_digest(&forged_parent_messages)?;
2968 let mut forged_parent = previous.clone();
2969 forged_parent.apply_transcript_history_state(TranscriptHistoryState {
2970 head: forged_parent_revision.clone(),
2971 commits: Vec::new(),
2972 revisions: vec![
2973 crate::TranscriptRevisionBody {
2974 revision: previous_revision.clone(),
2975 parent_revision: None,
2976 messages: previous.messages().to_vec(),
2977 created_at: previous.updated_at(),
2978 },
2979 crate::TranscriptRevisionBody {
2980 revision: forged_parent_revision.clone(),
2981 parent_revision: Some(previous_revision),
2982 messages: forged_parent_messages,
2983 created_at: forged_parent.updated_at(),
2984 },
2985 ],
2986 })?;
2987
2988 let mut incoming = forged_parent.clone();
2989 incoming.commit_transcript_rewrite(
2990 TranscriptRewriteSelection::MessageRange {
2991 start: 0,
2992 end: forged_parent.messages().len(),
2993 },
2994 vec![Message::User(UserMessage::text(
2995 "[Context compacted] forged branch".to_string(),
2996 ))],
2997 crate::TranscriptRewriteReason::new("compaction"),
2998 Some("meerkat-core".to_string()),
2999 Some(forged_parent_revision),
3000 )?;
3001
3002 assert!(matches!(
3003 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
3004 Err(SessionStoreError::TranscriptContinuityViolation { .. }
3005 | SessionStoreError::MonotonicityViolation { .. })
3006 ));
3007 Ok(())
3008 }
3009
3010 #[test]
3011 fn append_only_guard_rejects_transient_mcp_pending_notice_cleanup_with_unaudited_commit()
3012 -> Result<(), crate::TranscriptEditError> {
3013 let mut previous = Session::new();
3014 previous.push(Message::User(UserMessage::text("hello".to_string())));
3015 previous.push(Message::SystemNotice(SystemNoticeMessage {
3016 kind: SystemNoticeKind::McpPending,
3017 body: Some("connecting".to_string()),
3018 blocks: vec![SystemNoticeBlock::Mcp {
3019 server_id: None,
3020 operation: None,
3021 phase: None,
3022 persisted: false,
3023 detail: Some("connecting".to_string()),
3024 pending_sources: vec!["test-server".to_string()],
3025 }],
3026 created_at: crate::types::message_timestamp_now(),
3027 }));
3028 previous.push(Message::BlockAssistant(BlockAssistantMessage::new(
3029 vec![crate::types::AssistantBlock::Text {
3030 text: "answer".to_string(),
3031 meta: None,
3032 }],
3033 StopReason::EndTurn,
3034 )));
3035
3036 let mut incoming = previous.clone();
3037 incoming.replace_messages_internal(
3038 previous
3039 .messages()
3040 .iter()
3041 .filter(|message| !matches!(message, Message::SystemNotice(_)))
3042 .cloned()
3043 .collect(),
3044 crate::TranscriptRewriteReason::new("unit-test"),
3045 )?;
3046 incoming.push(Message::User(UserMessage::text("again".to_string())));
3047
3048 assert!(matches!(
3049 append_only_save_guard(&incoming, Some(&previous)),
3050 Err(SessionStoreError::InvalidTranscriptRewrite { .. })
3051 ));
3052 Ok::<(), crate::TranscriptEditError>(())
3053 }
3054
3055 #[test]
3056 fn rewrite_chain_finder_crosses_normal_append_between_rewrites()
3057 -> Result<(), Box<dyn std::error::Error>> {
3058 let mut session = Session::new();
3059 session.push(Message::User(UserMessage::text("first".to_string())));
3060 session.push(Message::BlockAssistant(BlockAssistantMessage {
3061 blocks: vec![AssistantBlock::Text {
3062 text: "verbose first answer".to_string(),
3063 meta: None,
3064 }],
3065 stop_reason: StopReason::EndTurn,
3066 identity: crate::types::TranscriptMessageIdentity::default(),
3067 created_at: crate::types::message_timestamp_now(),
3068 }));
3069
3070 let original = session.transcript_revision()?;
3071 let first = session.commit_transcript_rewrite(
3072 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
3073 vec![Message::BlockAssistant(BlockAssistantMessage {
3074 blocks: vec![AssistantBlock::Text {
3075 text: "compact first answer".to_string(),
3076 meta: None,
3077 }],
3078 stop_reason: StopReason::EndTurn,
3079 identity: crate::types::TranscriptMessageIdentity::default(),
3080 created_at: crate::types::message_timestamp_now(),
3081 })],
3082 crate::TranscriptRewriteReason::new("compaction"),
3083 Some("unit-test".to_string()),
3084 Some(original.clone()),
3085 )?;
3086
3087 session.push(Message::User(UserMessage::text("second".to_string())));
3088 session.push(Message::BlockAssistant(BlockAssistantMessage {
3089 blocks: vec![AssistantBlock::Text {
3090 text: "verbose second answer".to_string(),
3091 meta: None,
3092 }],
3093 stop_reason: StopReason::EndTurn,
3094 identity: crate::types::TranscriptMessageIdentity::default(),
3095 created_at: crate::types::message_timestamp_now(),
3096 }));
3097 let bridge = session.transcript_revision()?;
3098 assert_ne!(bridge, first.revision);
3099
3100 let second = session.commit_transcript_rewrite(
3101 TranscriptRewriteSelection::MessageRange { start: 3, end: 4 },
3102 vec![Message::BlockAssistant(BlockAssistantMessage {
3103 blocks: vec![AssistantBlock::Text {
3104 text: "compact second answer".to_string(),
3105 meta: None,
3106 }],
3107 stop_reason: StopReason::EndTurn,
3108 identity: crate::types::TranscriptMessageIdentity::default(),
3109 created_at: crate::types::message_timestamp_now(),
3110 })],
3111 crate::TranscriptRewriteReason::new("compaction"),
3112 Some("unit-test".to_string()),
3113 Some(bridge),
3114 )?;
3115 let state = session
3116 .transcript_history_state()?
3117 .ok_or_else(|| std::io::Error::other("missing transcript history state"))?;
3118
3119 let chain =
3120 find_transcript_rewrite_commit_chain_extending(&state, &original, &second.revision)
3121 .ok_or_else(|| {
3122 std::io::Error::other(
3123 "rewrite chain should extend through normal append bridge",
3124 )
3125 })?;
3126 assert_eq!(chain.len(), 2);
3127 assert_eq!(chain[0].revision, first.revision);
3128 assert_eq!(chain[1].revision, second.revision);
3129 Ok(())
3130 }
3131
3132 #[test]
3133 fn run_boundary_guard_rejects_dropped_retained_rewrite_commits()
3134 -> Result<(), Box<dyn std::error::Error>> {
3135 let mut base = Session::new();
3136 base.push(Message::User(UserMessage::text("turn one".to_string())));
3137 base.push(Message::BlockAssistant(BlockAssistantMessage {
3138 blocks: vec![AssistantBlock::Text {
3139 text: "verbose answer".to_string(),
3140 meta: None,
3141 }],
3142 stop_reason: StopReason::EndTurn,
3143 identity: crate::types::TranscriptMessageIdentity::default(),
3144 created_at: crate::types::message_timestamp_now(),
3145 }));
3146 let base_revision = base.transcript_revision()?;
3147
3148 let mut previous = base.clone();
3149 let _retained_commit = previous.commit_transcript_rewrite(
3150 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
3151 vec![Message::BlockAssistant(BlockAssistantMessage {
3152 blocks: vec![AssistantBlock::Text {
3153 text: "first compact answer".to_string(),
3154 meta: None,
3155 }],
3156 stop_reason: StopReason::EndTurn,
3157 identity: crate::types::TranscriptMessageIdentity::default(),
3158 created_at: crate::types::message_timestamp_now(),
3159 })],
3160 crate::TranscriptRewriteReason::new("compaction"),
3161 Some("unit-test".to_string()),
3162 Some(base_revision),
3163 )?;
3164 let previous_revision = previous.transcript_revision()?;
3165
3166 let mut incoming = previous.clone();
3167 let new_commit = incoming.commit_transcript_rewrite(
3168 TranscriptRewriteSelection::MessageRange { start: 1, end: 2 },
3169 vec![Message::BlockAssistant(BlockAssistantMessage {
3170 blocks: vec![AssistantBlock::Text {
3171 text: "second compact answer".to_string(),
3172 meta: None,
3173 }],
3174 stop_reason: StopReason::EndTurn,
3175 identity: crate::types::TranscriptMessageIdentity::default(),
3176 created_at: crate::types::message_timestamp_now(),
3177 })],
3178 crate::TranscriptRewriteReason::new("compaction"),
3179 Some("unit-test".to_string()),
3180 Some(previous_revision),
3181 )?;
3182 let mut state = incoming
3183 .transcript_history_state()?
3184 .ok_or_else(|| std::io::Error::other("incoming rewrite should retain history"))?;
3185 state.commits = vec![new_commit];
3186 incoming.set_metadata_unchecked_for_test(
3187 crate::session::SESSION_TRANSCRIPT_HISTORY_STATE_KEY,
3188 serde_json::to_value(state)?,
3189 );
3190
3191 assert!(matches!(
3192 run_boundary_snapshot_save_guard(&incoming, Some(&previous)),
3193 Err(SessionStoreError::InvalidTranscriptRewrite { reason, .. })
3194 if reason.contains("drop retained transcript rewrite commits")
3195 ));
3196 Ok(())
3197 }
3198
3199 fn runtime_append_system(content: &str) -> SystemMessage {
3208 let mut system = SystemMessage::new(content);
3209 system.mutation_kind = crate::types::SystemPromptMutationKind::RuntimeContextAppend;
3210 system
3211 }
3212
3213 #[allow(clippy::expect_used)]
3216 fn machine_persist_append_admits(
3217 previous: Option<&SystemMessage>,
3218 incoming: &SystemMessage,
3219 ) -> bool {
3220 let has_previous = previous.is_some();
3221 let content_identical =
3222 previous.is_some_and(|previous| incoming.content == previous.content);
3223 let content_extends_previous =
3224 previous.is_some_and(|previous| incoming.content.starts_with(&previous.content));
3225 let appended_starts_with_separator = previous.is_some_and(|previous| {
3226 incoming
3227 .content
3228 .get(previous.content.len()..)
3229 .is_some_and(|appended| appended.starts_with(SYSTEM_CONTEXT_SEPARATOR))
3230 });
3231 let incoming_is_runtime_context_append = incoming.mutation_kind.is_runtime_context_append();
3232 let mut authority = crate::session_document::SessionDocumentMachineAuthority::new();
3233 let effects = authority
3234 .resolve_system_context_persist_append_admission(
3235 has_previous,
3236 content_identical,
3237 content_extends_previous,
3238 appended_starts_with_separator,
3239 incoming_is_runtime_context_append,
3240 )
3241 .expect("machine resolves persist-append admission");
3242 effects.into_iter().any(|effect| {
3243 matches!(
3244 effect,
3245 crate::session_document::SessionDocumentEffect::SystemContextPersistAppendAdmissionResolved {
3246 admission: crate::session_document::SystemContextPersistAppendAdmission::Admit,
3247 }
3248 )
3249 })
3250 }
3251
3252 #[allow(clippy::expect_used)]
3253 fn assert_persist_append_matches_machine(
3254 previous: Option<&SystemMessage>,
3255 incoming: &SystemMessage,
3256 expected: bool,
3257 ) {
3258 let verdict =
3259 system_context_is_append(previous, incoming).expect("persist-time admission resolves");
3260 assert_eq!(verdict, expected, "persist-time verdict mismatch");
3261 assert_eq!(
3262 verdict,
3263 machine_persist_append_admits(previous, incoming),
3264 "persist-time verdict diverges from direct machine call"
3265 );
3266 }
3267
3268 #[test]
3269 fn persist_append_identical_content_admits() {
3270 let previous = SystemMessage::new("base system");
3271 let incoming = SystemMessage::new("base system");
3272 assert_persist_append_matches_machine(Some(&previous), &incoming, true);
3273 }
3274
3275 #[test]
3276 fn persist_append_separator_append_with_marker_admits() {
3277 let previous = SystemMessage::new("base system");
3278 let incoming = runtime_append_system(&format!(
3279 "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nextra"
3280 ));
3281 assert_persist_append_matches_machine(Some(&previous), &incoming, true);
3282 }
3283
3284 #[test]
3285 fn persist_append_shaped_without_marker_rejects() {
3286 let previous = SystemMessage::new("base system");
3287 let incoming = SystemMessage::new(format!(
3289 "base system{SYSTEM_CONTEXT_SEPARATOR}[Runtime System Context]\nextra"
3290 ));
3291 assert_persist_append_matches_machine(Some(&previous), &incoming, false);
3292 }
3293
3294 #[test]
3295 fn persist_append_divergent_content_rejects() {
3296 let previous = SystemMessage::new("base system");
3297 let incoming = runtime_append_system("totally different");
3298 assert_persist_append_matches_machine(Some(&previous), &incoming, false);
3299 }
3300
3301 #[test]
3302 fn persist_append_no_previous_admits_only_with_marker() {
3303 let with_marker = runtime_append_system("brand new context");
3304 assert_persist_append_matches_machine(None, &with_marker, true);
3305
3306 let without_marker = SystemMessage::new("brand new context");
3307 assert_persist_append_matches_machine(None, &without_marker, false);
3308 }
3309
3310 fn assistant_with_bookkeeping(
3311 text: &str,
3312 run_id: Option<crate::lifecycle::RunId>,
3313 created_at: crate::types::MessageTimestamp,
3314 ) -> Message {
3315 Message::BlockAssistant(BlockAssistantMessage {
3316 blocks: vec![AssistantBlock::Text {
3317 text: text.to_string(),
3318 meta: None,
3319 }],
3320 stop_reason: StopReason::EndTurn,
3321 identity: crate::types::TranscriptMessageIdentity {
3322 interaction_id: None,
3323 run_id,
3324 },
3325 created_at,
3326 })
3327 }
3328
3329 #[test]
3335 fn append_only_guard_accepts_rebookkept_prefix_identity_and_timestamps() {
3336 let base_time = crate::types::message_timestamp_now();
3337 let mut previous = Session::new();
3338 previous.push(Message::System(SystemMessage::new("base system")));
3339 previous.push(Message::User(UserMessage::text("turn one".to_string())));
3340 previous.push(assistant_with_bookkeeping(
3341 "answer one",
3342 Some(crate::lifecycle::RunId::new()),
3343 base_time,
3344 ));
3345
3346 let mut incoming = previous.clone();
3347 let mut rebookkept = previous.messages().to_vec();
3348 for message in &mut rebookkept {
3349 match message {
3350 Message::User(user) => {
3351 user.created_at = base_time + chrono::Duration::hours(1);
3352 }
3353 Message::BlockAssistant(assistant) => {
3354 assistant.identity = crate::types::TranscriptMessageIdentity {
3355 interaction_id: None,
3356 run_id: Some(crate::lifecycle::RunId::new()),
3357 };
3358 assistant.created_at = base_time + chrono::Duration::hours(1);
3359 }
3360 _ => {}
3361 }
3362 }
3363 rebookkept.push(Message::User(UserMessage::text("turn two".to_string())));
3364 incoming.messages = std::sync::Arc::new(rebookkept);
3365
3366 assert!(
3367 append_only_save_guard(&incoming, Some(&previous)).is_ok(),
3368 "bookkeeping-only prefix divergence must not fail continuity"
3369 );
3370 }
3371
3372 #[test]
3373 fn append_only_guard_rejects_content_divergence_despite_matching_bookkeeping() {
3374 let base_time = crate::types::message_timestamp_now();
3375 let run_id = crate::lifecycle::RunId::new();
3376 let mut previous = Session::new();
3377 previous.push(Message::User(UserMessage::text("turn one".to_string())));
3378 previous.push(assistant_with_bookkeeping(
3379 "answer one",
3380 Some(run_id),
3381 base_time,
3382 ));
3383
3384 let mut incoming = previous.clone();
3385 let mut diverged = previous.messages().to_vec();
3386 if let Message::BlockAssistant(assistant) = &mut diverged[1] {
3387 assistant.blocks = vec![AssistantBlock::Text {
3388 text: "a different answer".to_string(),
3389 meta: None,
3390 }];
3391 }
3392 diverged.push(Message::User(UserMessage::text("turn two".to_string())));
3393 incoming.messages = std::sync::Arc::new(diverged);
3394
3395 assert!(matches!(
3396 append_only_save_guard(&incoming, Some(&previous)),
3397 Err(SessionStoreError::TranscriptContinuityViolation { .. })
3398 ));
3399 }
3400}