use std::collections::VecDeque;
use std::sync::{Arc, Mutex, PoisonError};
use sha2::{Digest, Sha256};
use crate::session_store::SessionMessageRowPrefixAccumulator;
use crate::types::Message;
const BOUNDARY_RING_CAPACITY: usize = 8;
#[derive(Debug, Clone)]
struct Midstate {
hasher: Sha256,
covered: usize,
}
impl Midstate {
fn finalize(&self) -> String {
let mut hasher = self.hasher.clone();
hasher.update(b"]");
let digest = hasher.finalize();
let mut out = String::with_capacity(digest.len() * 2 + 7);
out.push_str("sha256:");
const HEX: &[u8; 16] = b"0123456789abcdef";
for byte in digest {
out.push(HEX[(byte >> 4) as usize] as char);
out.push(HEX[(byte & 0x0f) as usize] as char);
}
out
}
fn absorb(&mut self, message: &Message) -> Result<(), serde_json::Error> {
if self.covered > 0 {
self.hasher.update(b",");
}
let canonical = super::canonicalize_message_for_digest(message);
let bytes = serde_json::to_vec(&canonical)?;
crate::digest_observability::record_content_digest_bytes(bytes.len() as u64);
self.hasher.update(bytes);
self.covered += 1;
Ok(())
}
}
#[derive(Debug, Clone, Default)]
struct AccumulatorState {
stream_a: Option<Midstate>,
boundaries: VecDeque<Midstate>,
exact_row_anchor: Option<SessionMessageRowPrefixAccumulator>,
exact_row_committed: Option<SessionMessageRowPrefixAccumulator>,
exact_row_current: Option<SessionMessageRowPrefixAccumulator>,
lazy_exact_row_current_count: Option<u64>,
epoch: u64,
parked: Option<Box<AccumulatorState>>,
}
#[derive(Debug, Default)]
pub(crate) struct TranscriptDigestAccumulator {
state: Mutex<Box<AccumulatorState>>,
}
impl Clone for TranscriptDigestAccumulator {
fn clone(&self) -> Self {
Self {
state: Mutex::new(self.locked().clone()),
}
}
}
thread_local! {
static CROSS_CHECK_ENABLED: std::cell::Cell<bool> = const { std::cell::Cell::new(true) };
}
pub(crate) fn take_verification_sample() -> bool {
if cfg!(test) {
CROSS_CHECK_ENABLED.with(std::cell::Cell::get)
} else {
false
}
}
impl TranscriptDigestAccumulator {
fn locked(&self) -> std::sync::MutexGuard<'_, Box<AccumulatorState>> {
self.state.lock().unwrap_or_else(PoisonError::into_inner)
}
pub(crate) fn epoch(&self) -> u64 {
self.locked().epoch
}
fn invalidate(&mut self) {
let state = self.state.get_mut().unwrap_or_else(PoisonError::into_inner);
state.stream_a = None;
state.boundaries.clear();
state.exact_row_anchor = None;
state.exact_row_committed = None;
state.exact_row_current = None;
state.lazy_exact_row_current_count = None;
state.parked = None;
state.epoch = state.epoch.saturating_add(1);
}
fn extend(&mut self, appended: &[Message]) {
let state = self.state.get_mut().unwrap_or_else(PoisonError::into_inner);
state.lazy_exact_row_current_count =
state
.lazy_exact_row_current_count
.and_then(|current_count| {
u64::try_from(appended.len())
.ok()
.and_then(|appended_count| current_count.checked_add(appended_count))
});
if let Some(current) = state.exact_row_current.take() {
let serialized = appended
.iter()
.map(serde_json::to_vec)
.collect::<Result<Vec<_>, _>>();
state.exact_row_current = serialized
.ok()
.and_then(|rows| current.extend_serialized_rows(&rows).ok());
}
if let Some(stream) = state.stream_a.as_mut() {
for message in appended {
if stream.absorb(message).is_err() {
state.stream_a = None;
state.boundaries.clear();
state.epoch = state.epoch.saturating_add(1);
return;
}
}
}
}
fn begin_in_place_scan(&mut self) {
let state = self.state.get_mut().unwrap_or_else(PoisonError::into_inner);
if state.stream_a.is_none()
&& state.boundaries.is_empty()
&& state.exact_row_anchor.is_none()
&& state.exact_row_committed.is_none()
&& state.exact_row_current.is_none()
&& state.lazy_exact_row_current_count.is_none()
{
return;
}
let parked = AccumulatorState {
stream_a: state.stream_a.take(),
boundaries: std::mem::take(&mut state.boundaries),
exact_row_anchor: state.exact_row_anchor.take(),
exact_row_committed: state.exact_row_committed.take(),
exact_row_current: state.exact_row_current.take(),
lazy_exact_row_current_count: state.lazy_exact_row_current_count.take(),
epoch: state.epoch,
parked: None,
};
state.parked = Some(Box::new(parked));
}
fn finish_in_place_scan(&mut self, lowest_mutated_index: Option<usize>) {
let state = self.state.get_mut().unwrap_or_else(PoisonError::into_inner);
let Some(parked) = state.parked.take() else {
if lowest_mutated_index.is_some() {
state.stream_a = None;
state.boundaries.clear();
state.exact_row_anchor = None;
state.exact_row_committed = None;
state.exact_row_current = None;
state.lazy_exact_row_current_count = None;
state.epoch = state.epoch.saturating_add(1);
}
return;
};
if let Some(lowest_mutated_index) = lowest_mutated_index {
state.exact_row_anchor = parked.exact_row_anchor.filter(|anchor| {
u64::try_from(lowest_mutated_index).is_ok_and(|index| index >= anchor.row_count())
});
state.exact_row_committed = parked.exact_row_committed.filter(|committed| {
u64::try_from(lowest_mutated_index)
.is_ok_and(|index| index >= committed.row_count())
});
state.epoch = state.epoch.saturating_add(1);
return;
}
state.stream_a = parked.stream_a;
state.boundaries = parked.boundaries;
state.exact_row_anchor = parked.exact_row_anchor;
state.exact_row_committed = parked.exact_row_committed;
state.exact_row_current = parked.exact_row_current;
state.lazy_exact_row_current_count = parked.lazy_exact_row_current_count;
}
fn rebuild_exact_row_current_from_retained_prefix(&mut self, messages: &[Message]) {
let state = self.state.get_mut().unwrap_or_else(PoisonError::into_inner);
let Some(base) = [
state.exact_row_anchor.as_ref(),
state.exact_row_committed.as_ref(),
]
.into_iter()
.flatten()
.filter(|prefix| prefix.row_count() <= messages.len() as u64)
.max_by_key(|prefix| prefix.row_count())
.cloned() else {
return;
};
let Ok(start) = usize::try_from(base.row_count()) else {
return;
};
let serialized = messages[start..]
.iter()
.map(serde_json::to_vec)
.collect::<Result<Vec<_>, _>>();
state.exact_row_current = serialized
.ok()
.and_then(|rows| base.extend_serialized_rows(&rows).ok());
}
fn install_exact_row_prefix(&self, prefix: SessionMessageRowPrefixAccumulator) {
let mut state = self.locked();
state.lazy_exact_row_current_count = None;
if state.exact_row_anchor.is_none() {
state.exact_row_anchor = Some(prefix.clone());
}
state.exact_row_committed = Some(prefix.clone());
state.exact_row_current = Some(prefix);
}
fn install_exact_row_lineage(
&self,
anchor: SessionMessageRowPrefixAccumulator,
current: SessionMessageRowPrefixAccumulator,
) {
let mut state = self.locked();
state.lazy_exact_row_current_count = None;
state.exact_row_anchor = Some(anchor);
if state.exact_row_committed.is_none() {
state.exact_row_committed = Some(current.clone());
}
state.exact_row_current = Some(current);
}
fn exact_row_prefix_at(&self, row_count: u64) -> Option<SessionMessageRowPrefixAccumulator> {
let state = self.locked();
state
.exact_row_current
.as_ref()
.filter(|prefix| prefix.row_count() == row_count)
.or_else(|| {
state
.exact_row_committed
.as_ref()
.filter(|prefix| prefix.row_count() == row_count)
})
.or_else(|| {
state
.exact_row_anchor
.as_ref()
.filter(|prefix| prefix.row_count() == row_count)
})
.cloned()
}
fn mark_lazy_exact_row_current(&self, row_count: u64) {
let mut state = self.locked();
if state.exact_row_current.is_none() {
state.lazy_exact_row_current_count = Some(row_count);
}
}
fn lazy_exact_row_current_matches(&self, row_count: u64) -> bool {
self.locked().lazy_exact_row_current_count == Some(row_count)
}
fn exact_row_lineage_extends(
&self,
anchor: &SessionMessageRowPrefixAccumulator,
current_count: u64,
) -> bool {
let state = self.locked();
(state.exact_row_anchor.as_ref() == Some(anchor)
|| state.exact_row_committed.as_ref() == Some(anchor))
&& state
.exact_row_current
.as_ref()
.is_some_and(|current| current.row_count() == current_count)
}
fn digest(&self, messages: &[Message]) -> Result<String, serde_json::Error> {
if let Some(witness) = self.witness(messages) {
let mut state = self.locked();
if let Some(stream) = state.stream_a.clone() {
record_boundary(&mut state, &stream);
}
return Ok(witness);
}
let mut state = self.locked();
let mut stream = Midstate {
hasher: Sha256::new(),
covered: 0,
};
stream.hasher.update(b"[");
crate::digest_observability::record_content_digest_computation();
for message in messages {
stream.absorb(message)?;
}
let digest = stream.finalize();
record_boundary(&mut state, &stream);
state.stream_a = Some(stream);
Ok(digest)
}
fn witness(&self, messages: &[Message]) -> Option<String> {
let state = self.locked();
let stream = state.stream_a.as_ref()?;
if stream.covered != messages.len() {
return None;
}
let digest = stream.finalize();
drop(state);
if take_verification_sample()
&& let Ok(recomputed) = super::transcript_messages_digest_uncounted(messages)
{
assert_eq!(
digest, recomputed,
"transcript digest accumulator served a stale witness: a message-mutation \
seam extended or replaced the transcript without invalidating the midstate"
);
}
Some(digest)
}
fn prefix_witness(&self, messages: &[Message], count: usize) -> Option<String> {
if count > messages.len() {
return None;
}
let state = self.locked();
let boundary = state
.boundaries
.iter()
.find(|midstate| midstate.covered == count)
.or_else(|| {
state
.stream_a
.as_ref()
.filter(|stream| stream.covered == count)
})?;
let digest = boundary.finalize();
drop(state);
if take_verification_sample()
&& let Ok(recomputed) = super::transcript_messages_digest_uncounted(&messages[..count])
{
assert_eq!(
digest, recomputed,
"transcript digest accumulator served a stale prefix witness: a \
message-mutation seam rewrote a retained prefix without invalidating the ring"
);
}
Some(digest)
}
}
fn record_boundary(state: &mut AccumulatorState, stream: &Midstate) {
if state
.boundaries
.iter()
.any(|midstate| midstate.covered == stream.covered)
{
return;
}
if state.boundaries.len() >= BOUNDARY_RING_CAPACITY {
state.boundaries.pop_front();
}
state.boundaries.push_back(stream.clone());
}
#[derive(Debug)]
pub(crate) struct TranscriptMessages {
messages: Arc<Vec<Message>>,
accumulator: Box<TranscriptDigestAccumulator>,
}
impl Default for TranscriptMessages {
fn default() -> Self {
let accumulator = TranscriptDigestAccumulator::default();
let empty = SessionMessageRowPrefixAccumulator::empty();
accumulator.install_exact_row_lineage(empty.clone(), empty);
Self {
messages: Arc::default(),
accumulator: Box::new(accumulator),
}
}
}
impl Clone for TranscriptMessages {
fn clone(&self) -> Self {
Self {
messages: Arc::clone(&self.messages),
accumulator: Box::new((*self.accumulator).clone()),
}
}
}
impl std::ops::Deref for TranscriptMessages {
type Target = Vec<Message>;
fn deref(&self) -> &Self::Target {
&self.messages
}
}
impl TranscriptMessages {
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn arc(&self) -> &Arc<Vec<Message>> {
&self.messages
}
pub(crate) fn from_vec(messages: Vec<Message>) -> Self {
Self {
messages: Arc::new(messages),
accumulator: Box::default(),
}
}
pub(crate) fn from_fresh_branch(messages: Vec<Message>) -> Self {
let transcript = Self::from_vec(messages);
if let Ok(prefix) = SessionMessageRowPrefixAccumulator::from_messages(&transcript.messages)
{
transcript
.accumulator
.install_exact_row_lineage(prefix.clone(), prefix);
}
transcript
}
pub(crate) fn install_exact_row_prefix(
&self,
prefix: SessionMessageRowPrefixAccumulator,
) -> bool {
if prefix.row_count() != self.messages.len() as u64 {
return false;
}
self.accumulator.install_exact_row_prefix(prefix);
true
}
pub(crate) fn install_exact_row_lineage(
&self,
anchor: SessionMessageRowPrefixAccumulator,
current: SessionMessageRowPrefixAccumulator,
) -> bool {
if anchor.row_count() > current.row_count()
|| current.row_count() != self.messages.len() as u64
{
return false;
}
self.accumulator.install_exact_row_lineage(anchor, current);
true
}
pub(crate) fn exact_row_prefix_at(
&self,
row_count: u64,
) -> Option<SessionMessageRowPrefixAccumulator> {
if let Some(prefix) = self.accumulator.exact_row_prefix_at(row_count) {
return Some(prefix);
}
if u64::try_from(self.messages.len()).ok() != Some(row_count)
|| !self.accumulator.lazy_exact_row_current_matches(row_count)
{
return None;
}
let prefix = SessionMessageRowPrefixAccumulator::from_messages(&self.messages).ok()?;
self.accumulator
.install_exact_row_lineage(prefix.clone(), prefix.clone());
Some(prefix)
}
pub(crate) fn mark_lazy_whole_blob_row_lineage(&self) {
if let Ok(row_count) = u64::try_from(self.messages.len()) {
self.accumulator.mark_lazy_exact_row_current(row_count);
}
}
pub(crate) fn exact_row_lineage_extends(
&self,
anchor: &SessionMessageRowPrefixAccumulator,
current_count: u64,
) -> bool {
self.accumulator
.exact_row_lineage_extends(anchor, current_count)
}
pub(crate) fn push(&mut self, message: Message) {
let Self {
messages,
accumulator,
} = self;
let inner = Arc::make_mut(messages);
inner.push(message);
let appended = &inner[inner.len() - 1..];
accumulator.extend(appended);
}
pub(crate) fn extend_batch(&mut self, appended: Vec<Message>) {
if appended.is_empty() {
return;
}
let Self {
messages,
accumulator,
} = self;
let inner = Arc::make_mut(messages);
let start = inner.len();
inner.extend(appended);
accumulator.extend(&inner[start..]);
}
pub(crate) fn replace(&mut self, messages: Vec<Message>) {
self.messages = Arc::new(messages);
self.accumulator.invalidate();
}
#[cfg(test)]
pub(crate) fn mutate_in_place(&mut self) -> &mut Vec<Message> {
self.accumulator.invalidate();
Arc::make_mut(&mut self.messages)
}
pub(crate) fn begin_in_place_scan(&mut self) -> &mut Vec<Message> {
self.accumulator.begin_in_place_scan();
Arc::make_mut(&mut self.messages)
}
pub(crate) fn finish_in_place_scan(&mut self, lowest_mutated_index: Option<usize>) {
self.accumulator.finish_in_place_scan(lowest_mutated_index);
if lowest_mutated_index.is_some() {
self.accumulator
.rebuild_exact_row_current_from_retained_prefix(&self.messages);
}
}
pub(crate) fn digest(&self) -> Result<String, serde_json::Error> {
self.accumulator.digest(&self.messages)
}
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn digest_witness(&self) -> Option<String> {
self.accumulator.witness(&self.messages)
}
pub(crate) fn prefix_digest_witness(&self, count: usize) -> Option<String> {
self.accumulator.prefix_witness(&self.messages, count)
}
pub(crate) fn mutation_epoch(&self) -> u64 {
self.accumulator.epoch()
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
use crate::types::{Message, UserMessage};
fn user(text: &str) -> Message {
Message::User(UserMessage::text(text))
}
fn transcript(count: usize) -> Vec<Message> {
(0..count).map(|i| user(&format!("m{i}"))).collect()
}
#[test]
fn canonicalize_messages_for_digest_is_element_wise() {
let messages = transcript(6);
let whole = super::super::canonicalize_messages_for_digest(&messages);
let per_message = messages
.iter()
.map(super::super::canonicalize_message_for_digest)
.collect::<Vec<_>>();
assert_eq!(whole, per_message);
let array_bytes = serde_json::to_vec(&whole).unwrap();
let mut streamed = Vec::from(b"[".as_slice());
for (index, message) in per_message.iter().enumerate() {
if index > 0 {
streamed.extend_from_slice(b",");
}
streamed.extend_from_slice(&serde_json::to_vec(message).unwrap());
}
streamed.extend_from_slice(b"]");
assert_eq!(array_bytes, streamed);
}
#[test]
fn seeded_digest_matches_full_recompute() {
for count in [0usize, 1, 2, 7, 40] {
let messages = TranscriptMessages::from_vec(transcript(count));
assert_eq!(
messages.digest().unwrap(),
super::super::transcript_messages_digest(&messages).unwrap(),
"count {count}"
);
}
}
#[test]
fn appended_digest_matches_full_recompute() {
let mut messages = TranscriptMessages::from_vec(transcript(3));
let _ = messages.digest().unwrap();
messages.push(user("appended"));
assert!(messages.digest_witness().is_some());
assert_eq!(
messages.digest().unwrap(),
super::super::transcript_messages_digest(&messages).unwrap()
);
messages.extend_batch(vec![user("a"), user("b")]);
assert!(messages.digest_witness().is_some());
assert_eq!(
messages.digest().unwrap(),
super::super::transcript_messages_digest(&messages).unwrap()
);
}
#[test]
fn unseeded_accumulator_has_no_witness() {
let messages = TranscriptMessages::from_vec(transcript(3));
assert!(messages.digest_witness().is_none());
assert!(messages.prefix_digest_witness(2).is_none());
}
#[test]
fn replacement_invalidates_the_witness() {
let mut messages = TranscriptMessages::from_vec(transcript(3));
let _ = messages.digest().unwrap();
let epoch = messages.mutation_epoch();
messages.replace(transcript(2));
assert!(messages.digest_witness().is_none());
assert!(messages.mutation_epoch() > epoch);
assert_eq!(
messages.digest().unwrap(),
super::super::transcript_messages_digest(&messages).unwrap()
);
}
#[test]
fn in_place_mutation_invalidates_the_witness() {
let mut messages = TranscriptMessages::from_vec(transcript(3));
let _ = messages.digest().unwrap();
messages.mutate_in_place()[0] = user("rewritten");
assert!(messages.digest_witness().is_none());
assert_eq!(
messages.digest().unwrap(),
super::super::transcript_messages_digest(&messages).unwrap()
);
}
#[test]
fn unmutated_in_place_scan_keeps_the_witness() {
let mut messages = TranscriptMessages::from_vec(transcript(3));
let seeded = messages.digest().unwrap();
let _buffer = messages.begin_in_place_scan();
assert!(
messages.digest_witness().is_none(),
"a parked accumulator must not serve a witness"
);
messages.finish_in_place_scan(None);
assert_eq!(messages.digest_witness(), Some(seeded));
}
#[test]
fn mutated_in_place_scan_drops_the_witness() {
let mut messages = TranscriptMessages::from_vec(transcript(3));
let _ = messages.digest().unwrap();
{
let buffer = messages.begin_in_place_scan();
buffer[1] = user("changed");
}
messages.finish_in_place_scan(Some(1));
assert!(messages.digest_witness().is_none());
assert_eq!(
messages.digest().unwrap(),
super::super::transcript_messages_digest(&messages).unwrap()
);
}
fn exact_row_prefix(
messages: &[Message],
) -> crate::session_store::SessionMessageRowPrefixAccumulator {
let rows = messages
.iter()
.map(serde_json::to_vec)
.collect::<Result<Vec<_>, _>>()
.unwrap();
crate::session_store::SessionMessageRowPrefixAccumulator::from_serialized_rows(&rows)
.unwrap()
}
#[test]
fn fresh_transcript_owns_exact_genesis_row_lineage() {
let mut messages = TranscriptMessages::default();
assert_eq!(
messages.exact_row_prefix_at(0),
Some(crate::session_store::SessionMessageRowPrefixAccumulator::empty())
);
messages.push(user("first"));
assert_eq!(
messages.exact_row_prefix_at(1),
Some(exact_row_prefix(&messages)),
"ordinary appends must extend fresh construction authority without a full rescan"
);
}
#[test]
fn committed_boundary_survives_appends_beside_graph_anchor_and_live_current() {
let mut messages = TranscriptMessages::default();
let graph_anchor = crate::session_store::SessionMessageRowPrefixAccumulator::empty();
messages.push(user("committed"));
let committed = exact_row_prefix(&messages);
assert!(messages.install_exact_row_prefix(committed.clone()));
messages.push(user("live tail"));
let current = exact_row_prefix(&messages);
assert_eq!(messages.exact_row_prefix_at(0), Some(graph_anchor));
assert_eq!(messages.exact_row_prefix_at(1), Some(committed));
assert_eq!(messages.exact_row_prefix_at(2), Some(current));
assert!(
messages.exact_row_lineage_extends(
&messages.exact_row_prefix_at(1).expect("committed boundary"),
2,
),
"ordinary successor preparation must accept the durable boundary beside the graph anchor"
);
}
#[test]
fn graph_replay_preflight_preserves_existing_committed_boundary() {
let mut messages = TranscriptMessages::default();
messages.push(user("committed predecessor"));
let committed = exact_row_prefix(&messages);
assert!(messages.install_exact_row_prefix(committed.clone()));
let rewritten = TranscriptMessages::from_vec(vec![user("rewritten endpoint")]);
let graph_anchor = exact_row_prefix(&rewritten);
let mut rewritten_with_tail = rewritten;
rewritten_with_tail.push(user("live tail"));
let live_current = exact_row_prefix(&rewritten_with_tail);
messages.push(user("live tail"));
assert!(messages.install_exact_row_lineage(graph_anchor.clone(), live_current.clone()));
assert_ne!(graph_anchor, committed);
assert_eq!(messages.exact_row_prefix_at(1), Some(committed));
assert_eq!(messages.exact_row_prefix_at(2), Some(live_current));
}
#[test]
fn suffix_only_in_place_scan_preserves_exact_durable_row_anchor() {
let mut messages = TranscriptMessages::from_vec(transcript(4));
let anchor = exact_row_prefix(&messages[..2]);
let current = exact_row_prefix(&messages);
assert!(messages.install_exact_row_lineage(anchor.clone(), current.clone()));
{
let buffer = messages.begin_in_place_scan();
buffer[2] = user("externalized suffix");
}
messages.finish_in_place_scan(Some(2));
let rebuilt_current = exact_row_prefix(&messages);
assert_eq!(messages.exact_row_prefix_at(2), Some(anchor));
assert_ne!(
rebuilt_current, current,
"the rebuilt current-row prefix must bind the changed suffix"
);
assert_eq!(
messages.exact_row_prefix_at(4),
Some(rebuilt_current),
"a suffix rewrite must rebuild the exact current-row prefix from the retained anchor"
);
assert!(
messages.digest_witness().is_none(),
"a suffix rewrite must invalidate the whole-transcript digest witness"
);
}
#[test]
fn in_place_scan_inside_exact_durable_prefix_drops_row_anchor() {
let mut messages = TranscriptMessages::from_vec(transcript(4));
let anchor = exact_row_prefix(&messages[..2]);
let current = exact_row_prefix(&messages);
assert!(messages.install_exact_row_lineage(anchor, current));
{
let buffer = messages.begin_in_place_scan();
buffer[1] = user("rewritten durable prefix");
}
messages.finish_in_place_scan(Some(1));
assert!(messages.exact_row_prefix_at(2).is_none());
assert!(messages.exact_row_prefix_at(4).is_none());
}
#[test]
fn abandoned_in_place_scan_fails_safe() {
let mut messages = TranscriptMessages::from_vec(transcript(3));
let _ = messages.digest().unwrap();
let buffer = messages.begin_in_place_scan();
buffer[2] = user("changed");
assert!(messages.digest_witness().is_none());
assert_eq!(
messages.digest().unwrap(),
super::super::transcript_messages_digest(&messages).unwrap()
);
}
#[test]
fn boundary_ring_answers_prefix_queries_after_appends() {
let mut messages = TranscriptMessages::from_vec(transcript(5));
let boundary = messages.digest().unwrap();
messages.extend_batch(vec![user("x"), user("y")]);
assert_eq!(messages.prefix_digest_witness(5), Some(boundary));
assert_eq!(
messages.prefix_digest_witness(5).unwrap(),
super::super::transcript_messages_digest(&messages[..5]).unwrap()
);
assert!(messages.prefix_digest_witness(4).is_none());
}
#[test]
fn boundary_ring_is_bounded_and_dropped_on_invalidation() {
let mut messages = TranscriptMessages::from_vec(Vec::new());
for _ in 0..(BOUNDARY_RING_CAPACITY + 4) {
messages.push(user("m"));
let _ = messages.digest().unwrap();
}
let retained = (0..=(BOUNDARY_RING_CAPACITY + 4))
.filter(|count| messages.prefix_digest_witness(*count).is_some())
.count();
assert!(
retained <= BOUNDARY_RING_CAPACITY,
"boundary ring grew past its bound: {retained}"
);
messages.mutate_in_place();
assert_eq!(
(0..=(BOUNDARY_RING_CAPACITY + 4))
.filter(|count| messages.prefix_digest_witness(*count).is_some())
.count(),
0
);
}
#[test]
fn randomized_mutation_sequences_match_full_recompute() {
let mut seed = 0x5eed_1234_u64;
let mut next = move || {
seed = seed
.wrapping_mul(6364136223846793005)
.wrapping_add(1442695040888963407);
(seed >> 33) as usize
};
let mut messages = TranscriptMessages::from_vec(transcript(4));
for step in 0..200 {
match next() % 6 {
0 => messages.push(user(&format!("p{step}"))),
1 => messages.extend_batch(vec![user(&format!("b{step}")), user("b2")]),
2 => messages.replace(transcript(next() % 9)),
3 => {
let buffer = messages.mutate_in_place();
if !buffer.is_empty() {
let index = next() % buffer.len();
buffer[index] = user(&format!("r{step}"));
}
}
4 => {
messages.begin_in_place_scan();
messages.finish_in_place_scan(None);
}
_ => {
let buffer = messages.begin_in_place_scan();
let mutated = if buffer.is_empty() {
None
} else {
let index = next() % buffer.len();
buffer[index] = user(&format!("s{step}"));
Some(index)
};
messages.finish_in_place_scan(mutated);
}
}
assert_eq!(
messages.digest().unwrap(),
super::super::transcript_messages_digest(&messages).unwrap(),
"step {step}"
);
let count = messages.len();
if count > 1 {
messages.push(user("tail"));
assert_eq!(
messages.prefix_digest_witness(count),
Some(super::super::transcript_messages_digest(&messages[..count]).unwrap()),
"prefix witness at step {step}"
);
}
}
}
}