use std::collections::VecDeque;
use sui_sdk_types::Object;
use sui_sdk_types::framework::EventBatch;
use sui_sdk_types::framework::EventCommitment;
use sui_sdk_types::framework::EventStreamHead;
use sui_sdk_types::framework::apply_stream_updates;
use super::envelope::AuthenticatedEvent;
use crate::light_client::error::LightClientError;
use crate::light_client::error::MmrMismatch;
#[derive(Debug, Default)]
pub(super) struct StreamState {
pub buffer: VecDeque<AuthenticatedEvent>,
pub local_head: EventStreamHead,
pub confirmed_through: u64,
pub events_scanned_through: u64,
}
impl StreamState {
pub(super) fn new(initial_head: EventStreamHead, confirmed_through: u64) -> Self {
Self {
buffer: VecDeque::new(),
local_head: initial_head,
confirmed_through,
events_scanned_through: confirmed_through,
}
}
}
pub(super) fn buffer_response_batch(
state: &mut StreamState,
events: Vec<AuthenticatedEvent>,
watermark_hi: Option<u64>,
) {
let last_event_cp = events.last().map(|e| e.checkpoint);
for event in events {
state.buffer.push_back(event);
}
if let Some(hi) = watermark_hi {
state.events_scanned_through = state.events_scanned_through.max(hi);
}
if let Some(cp) = last_event_cp {
state.events_scanned_through = state.events_scanned_through.max(cp);
}
}
pub(super) fn fold_and_reconcile(
state: &mut StreamState,
settlements: &[(u64, u64)],
chain_head: EventStreamHead,
reconcile_checkpoint: u64,
) -> Result<Vec<AuthenticatedEvent>, LightClientError> {
let last_settlement =
settlements
.last()
.copied()
.ok_or(LightClientError::UnexpectedObjectShape {
reason: "fold_and_reconcile called with no settlements; \
reconciliation needs at least one settlement to anchor the chain head",
})?;
if last_settlement.0 != reconcile_checkpoint {
return Err(LightClientError::UnexpectedObjectShape {
reason: "settlement boundaries disagree with reconciliation checkpoint",
});
}
let fold_count = state
.buffer
.iter()
.take_while(|e| e.checkpoint <= reconcile_checkpoint)
.count();
let batches = bucket_events_by_settlement(
state.buffer.iter().take(fold_count),
settlements,
state.confirmed_through,
)?;
let new_head =
apply_stream_updates(state.local_head.clone(), &batches).map_err(LightClientError::from)?;
if new_head != chain_head {
return Err(LightClientError::MmrMismatch(Box::new(MmrMismatch {
checkpoint: reconcile_checkpoint,
expected: chain_head,
actual: new_head,
})));
}
state.local_head = new_head;
state.confirmed_through = reconcile_checkpoint;
let released: Vec<AuthenticatedEvent> = state.buffer.drain(..fold_count).collect();
Ok(released)
}
fn bucket_events_by_settlement<'a, I>(
events: I,
settlements: &[(u64, u64)],
floor_checkpoint: u64,
) -> Result<Vec<EventBatch>, LightClientError>
where
I: IntoIterator<Item = &'a AuthenticatedEvent>,
{
let mut batches: Vec<EventBatch> = Vec::new();
let mut current_key: Option<(u64, u64)> = None;
let mut settlement_idx: usize = 0;
for event in events {
if event.checkpoint <= floor_checkpoint {
continue;
}
while settlement_idx < settlements.len()
&& (settlements[settlement_idx].0 < event.checkpoint
|| (settlements[settlement_idx].0 == event.checkpoint
&& settlements[settlement_idx].1 < event.transaction_index))
{
settlement_idx += 1;
}
if settlement_idx >= settlements.len() || settlements[settlement_idx].0 != event.checkpoint
{
return Err(LightClientError::UnexpectedObjectShape {
reason: "no settlement transaction covering buffered event",
});
}
let settlement_key = settlements[settlement_idx];
let commitment = EventCommitment {
checkpoint_seq: event.checkpoint,
transaction_idx: event.transaction_index,
event_idx: event.event_index as u64,
digest: event.event.digest(),
};
match current_key {
Some(key) if key == settlement_key => {
batches
.last_mut()
.expect("current_key set ⇒ at least one batch exists")
.commitments
.push(commitment);
}
_ => {
batches.push(EventBatch {
checkpoint_seq: event.checkpoint,
commitments: vec![commitment],
});
current_key = Some(settlement_key);
}
}
}
Ok(batches)
}
pub(super) fn extract_event_stream_head(
object: &Object,
) -> Result<EventStreamHead, LightClientError> {
const HEAD_OFFSET: usize = 64;
let move_struct = object
.as_struct()
.ok_or(LightClientError::UnexpectedObjectShape {
reason: "expected a Move struct (dynamic field), got a package",
})?;
let contents = move_struct.contents();
let head_bytes =
contents
.get(HEAD_OFFSET..)
.ok_or(LightClientError::UnexpectedObjectShape {
reason: "dynamic-field contents too short for EventStreamHead value",
})?;
bcs::from_bytes(head_bytes).map_err(LightClientError::from)
}
#[cfg(test)]
mod tests {
use super::*;
use sui_sdk_types::Address;
use sui_sdk_types::Digest;
use sui_sdk_types::Event;
use sui_sdk_types::Identifier;
use sui_sdk_types::StructTag;
use sui_sdk_types::U256;
fn sample_event(checkpoint: u64, tx_idx: u64, event_idx: u32) -> AuthenticatedEvent {
AuthenticatedEvent {
checkpoint,
transaction_index: tx_idx,
event_index: event_idx,
transaction_digest: Digest::new([0xaa; 32]),
event: Event {
package_id: Address::TWO,
module: Identifier::from_static("m"),
sender: Address::TWO,
type_: StructTag::new(
Address::TWO,
Identifier::from_static("m"),
Identifier::from_static("E"),
vec![],
),
contents: vec![checkpoint as u8, tx_idx as u8, event_idx as u8],
},
}
}
fn expected_chain_head(
events: &[AuthenticatedEvent],
settlements: &[(u64, u64)],
) -> EventStreamHead {
let batches = bucket_events_by_settlement(events.iter(), settlements, 0).unwrap();
apply_stream_updates(EventStreamHead::default(), &batches).unwrap()
}
fn u256_from_decimal(s: &str) -> U256 {
s.parse().expect("decimal U256 literal must parse")
}
#[test]
fn buffer_response_batch_defers_folding() {
let mut state = StreamState::new(EventStreamHead::default(), 0);
let events = vec![sample_event(7, 0, 0), sample_event(7, 0, 1)];
buffer_response_batch(&mut state, events, Some(7));
assert_eq!(state.local_head, EventStreamHead::default());
assert_eq!(state.buffer.len(), 2);
assert_eq!(state.events_scanned_through, 7);
}
#[test]
fn buffer_response_batch_advances_scan_floor_from_events() {
let mut state = StreamState::new(EventStreamHead::default(), 0);
buffer_response_batch(&mut state, vec![sample_event(11, 0, 0)], None);
assert_eq!(state.events_scanned_through, 11);
}
#[test]
fn buffer_response_batch_does_not_regress_scan_floor() {
let mut state = StreamState::new(EventStreamHead::default(), 0);
state.events_scanned_through = 20;
buffer_response_batch(&mut state, vec![sample_event(5, 0, 0)], Some(5));
assert_eq!(state.events_scanned_through, 20);
}
#[test]
fn fold_and_reconcile_single_settlement_per_checkpoint() {
let events = vec![
sample_event(7, 0, 0),
sample_event(7, 0, 1),
sample_event(7, 1, 0),
];
let settlements = vec![(7u64, 2u64)];
let chain_head = expected_chain_head(&events, &settlements);
let mut state = StreamState::new(EventStreamHead::default(), 0);
buffer_response_batch(&mut state, events.clone(), Some(7));
let released = fold_and_reconcile(&mut state, &settlements, chain_head.clone(), 7).unwrap();
assert_eq!(released, events);
assert_eq!(state.local_head, chain_head);
assert_eq!(state.confirmed_through, 7);
assert!(state.buffer.is_empty());
}
#[test]
fn fold_and_reconcile_multiple_settlements_per_checkpoint() {
let events = vec![
sample_event(10, 0, 0),
sample_event(10, 1, 0),
sample_event(10, 3, 0),
];
let settlements = vec![(10u64, 2u64), (10u64, 4u64)];
let chain_head = expected_chain_head(&events, &settlements);
let mut state = StreamState::new(EventStreamHead::default(), 0);
buffer_response_batch(&mut state, events.clone(), Some(10));
let single_batch = apply_stream_updates(
EventStreamHead::default(),
&[EventBatch {
checkpoint_seq: 10,
commitments: events
.iter()
.map(|e| EventCommitment {
checkpoint_seq: e.checkpoint,
transaction_idx: e.transaction_index,
event_idx: e.event_index as u64,
digest: e.event.digest(),
})
.collect(),
}],
)
.unwrap();
assert_ne!(single_batch, chain_head);
let released =
fold_and_reconcile(&mut state, &settlements, chain_head.clone(), 10).unwrap();
assert_eq!(released, events);
assert_eq!(state.local_head, chain_head);
assert_eq!(state.confirmed_through, 10);
}
#[test]
fn fold_and_reconcile_keeps_unfolded_tail_buffered() {
let events = vec![
sample_event(5, 0, 0),
sample_event(5, 0, 1),
sample_event(9, 0, 0), ];
let settlements = vec![(5u64, 0u64)];
let cp5_events: Vec<_> = events.iter().take(2).cloned().collect();
let chain_head_at_5 = expected_chain_head(&cp5_events, &settlements);
let mut state = StreamState::new(EventStreamHead::default(), 0);
buffer_response_batch(&mut state, events.clone(), Some(9));
let released =
fold_and_reconcile(&mut state, &settlements, chain_head_at_5.clone(), 5).unwrap();
assert_eq!(released, cp5_events);
assert_eq!(state.local_head, chain_head_at_5);
assert_eq!(state.confirmed_through, 5);
assert_eq!(state.buffer.len(), 1);
assert_eq!(state.buffer.front().unwrap().checkpoint, 9);
}
#[test]
fn fold_and_reconcile_rejects_divergence_without_mutating_observable_state() {
let events = vec![sample_event(3, 0, 0)];
let settlements = vec![(3u64, 0u64)];
let bogus = EventStreamHead {
mmr: vec![u256_from_decimal("999")],
checkpoint_seq: 3,
num_events: 1,
};
let mut state = StreamState::new(EventStreamHead::default(), 0);
buffer_response_batch(&mut state, events.clone(), Some(3));
let pre_local_head = state.local_head.clone();
let pre_buffer_len = state.buffer.len();
let pre_confirmed = state.confirmed_through;
let err = fold_and_reconcile(&mut state, &settlements, bogus.clone(), 3).unwrap_err();
assert!(matches!(err, LightClientError::MmrMismatch(_)));
assert_eq!(state.local_head, pre_local_head);
assert_eq!(state.buffer.len(), pre_buffer_len);
assert_eq!(state.confirmed_through, pre_confirmed);
}
#[test]
fn fold_and_reconcile_rejects_empty_settlements() {
let mut state = StreamState::new(EventStreamHead::default(), 0);
let err = fold_and_reconcile(&mut state, &[], EventStreamHead::default(), 0).unwrap_err();
assert!(matches!(
err,
LightClientError::UnexpectedObjectShape { .. }
));
}
#[test]
fn fold_and_reconcile_rejects_misaligned_reconcile_checkpoint() {
let settlements = vec![(5u64, 0u64)];
let mut state = StreamState::new(EventStreamHead::default(), 0);
let err = fold_and_reconcile(&mut state, &settlements, EventStreamHead::default(), 7)
.unwrap_err();
assert!(matches!(
err,
LightClientError::UnexpectedObjectShape { .. }
));
}
#[test]
fn fold_and_reconcile_rejects_event_without_matching_settlement() {
let events = vec![sample_event(5, 0, 0)];
let settlements = vec![(7u64, 0u64)];
let mut state = StreamState::new(EventStreamHead::default(), 0);
buffer_response_batch(&mut state, events, Some(7));
let err = fold_and_reconcile(&mut state, &settlements, EventStreamHead::default(), 7)
.unwrap_err();
assert!(matches!(
err,
LightClientError::UnexpectedObjectShape { .. }
));
}
#[test]
fn extract_event_stream_head_skips_uid_and_key_prefix() {
use sui_sdk_types::Object;
use sui_sdk_types::ObjectData;
use sui_sdk_types::Owner;
use sui_sdk_types::Version;
use sui_sdk_types::framework::derive_event_stream_head_object_id;
let head = EventStreamHead {
mmr: vec![U256::ZERO, U256::ONE],
checkpoint_seq: 17,
num_events: 5,
};
let stream_id = Address::TWO;
let mut contents = Vec::new();
let uid = derive_event_stream_head_object_id(stream_id);
contents.extend_from_slice(uid.as_bytes());
contents.extend_from_slice(stream_id.as_bytes());
contents.extend(bcs::to_bytes(&head).unwrap());
let move_struct = sui_sdk_types::MoveStruct::new(
StructTag::new(
Address::TWO,
Identifier::from_static("dynamic_field"),
Identifier::from_static("Field"),
vec![],
),
true,
Version::from(1u64),
contents,
)
.expect("contents are at least 32 bytes");
let object = Object::new(
ObjectData::Struct(move_struct),
Owner::Address(Address::ZERO),
Digest::new([0; 32]),
0,
);
let recovered = extract_event_stream_head(&object).unwrap();
assert_eq!(recovered, head);
}
#[test]
fn extract_event_stream_head_rejects_non_struct() {
use sui_sdk_types::MovePackage;
use sui_sdk_types::Object;
use sui_sdk_types::ObjectData;
use sui_sdk_types::Owner;
use sui_sdk_types::Version;
let pkg = MovePackage {
id: Address::TWO,
version: Version::from(1u64),
modules: Default::default(),
type_origin_table: Vec::new(),
linkage_table: Default::default(),
};
let object = Object::new(
ObjectData::Package(pkg),
Owner::Immutable,
Digest::new([0; 32]),
0,
);
let err = extract_event_stream_head(&object).unwrap_err();
assert!(matches!(
err,
LightClientError::UnexpectedObjectShape { .. }
));
}
}