diff --git a/crates/engine/src/streams.rs b/crates/engine/src/streams.rs index 0803a80f..40f3b2fc 100755 --- a/crates/engine/src/streams.rs +++ b/crates/engine/src/streams.rs @@ -111,13 +111,16 @@ pub async fn handle_get_shard_iterator( "SequenceNumber is required for AT_SEQUENCE_NUMBER iterator type".to_owned(), ) })?; - let n = raw.parse::().map_err(|_| { + let n = raw.parse::().map_err(|_| { DynamoDbError::ValidationException("Invalid SequenceNumber".to_owned()) })?; // n == 0: sequence 0 is the first possible record, so "at 0" // means "read from the beginning" — same as TRIM_HORIZON. if n > 0 { - format!("{:021}", n - 1) + // Pad to the backend's stored width so lexicographic order + // matches numeric order against stored sequence numbers. + let width = ctx.storage.sequence_number_width(); + format!("{:0>width$}", n - 1) } else { String::new() } @@ -262,3 +265,32 @@ fn storage_to_dynamo(e: StorageError) -> DynamoDbError { } } } + +#[cfg(test)] +mod tests { + /// AT_SEQUENCE_NUMBER converts to AFTER by subtracting 1 and padding to + /// the backend's stored width. Verify that an unpadded client input is + /// normalised correctly for both the 21-digit (postgres/sqlite/mongodb) + /// and 23-digit (cassandra) cases, and that the edge case n=0 maps to + /// TRIM_HORIZON (empty string). + #[test] + fn at_sequence_number_padding() { + let cases: &[(&str, usize, &str)] = &[ + ("5", 21, "000000000000000000004"), + ("000000000000000000005", 21, "000000000000000000004"), + ("5", 23, "00000000000000000000004"), + ("00000000000000000000005", 23, "00000000000000000000004"), + ("0", 21, ""), + ("1", 21, "000000000000000000000"), + ]; + for (input, width, expected) in cases { + let n: u128 = input.parse().unwrap(); + let result = if n > 0 { + format!("{:0>width$}", n - 1) + } else { + String::new() + }; + assert_eq!(&result, expected, "input={input} width={width}"); + } + } +} diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs index 185dcd21..607c1faf 100755 --- a/crates/storage/src/lib.rs +++ b/crates/storage/src/lib.rs @@ -75,6 +75,9 @@ use extenddb_core::types::{ use error::StorageError; +/// Default fixed width for stream record sequence numbers in characters (left-padded). +pub(crate) const DEFAULT_SEQ_NUMBER_WIDTH: usize = 21; + // Type aliases for complex return types used in trait methods. /// Result of an update/put/delete that may return old and/or new item images. pub type ItemPairResult = Result<(Option, Option), StorageError>; @@ -609,6 +612,14 @@ pub trait StreamEngine: Send + Sync { /// Generate the next sequence number for a shard. fn next_sequence_number(&self, shard_id: &str) -> BoxFuture<'_, Result>; + /// Width of sequence numbers stored by this backend, in decimal digits. + /// Client-supplied sequence numbers (e.g. from `GetShardIterator`) must be + /// zero-padded to this width before lexicographic comparison against stored + /// values. Defaults to 21 (PostgreSQL, SQLite, MongoDB). + fn sequence_number_width(&self) -> usize { + DEFAULT_SEQ_NUMBER_WIDTH + } + /// Validate that a shard exists for the given stream ARN. /// /// Returns `Ok(())` if the shard exists and belongs to the stream.