diff --git a/Cargo.lock b/Cargo.lock index e700060547..02d7fa2075 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4447,6 +4447,7 @@ dependencies = [ "miden-node-utils", "miden-protocol", "prost", + "rand 0.10.2", "reqwest", "rusqlite", "tempfile", diff --git a/bin/note-transport/Cargo.toml b/bin/note-transport/Cargo.toml index 8b8e777ad5..b3a6a1b5e5 100644 --- a/bin/note-transport/Cargo.toml +++ b/bin/note-transport/Cargo.toml @@ -11,6 +11,10 @@ repository.workspace = true rust-version.workspace = true version.workspace = true +[package.metadata.cargo-shear] +# The build script links migrations through the generated migrator. +ignored-paths = ["src/db/migrations/*.rs"] + [lints] workspace = true @@ -30,6 +34,7 @@ miden-node-tracing = { workspace = true } miden-node-utils = { workspace = true } miden-protocol = { workspace = true } prost = { workspace = true } +rand = { workspace = true } rusqlite = { workspace = true } thiserror = { workspace = true } tokio = { features = ["macros", "net", "rt-multi-thread"], workspace = true } diff --git a/bin/note-transport/README.md b/bin/note-transport/README.md index d2a9ed18fa..0c99fded20 100644 --- a/bin/note-transport/README.md +++ b/bin/note-transport/README.md @@ -46,23 +46,37 @@ without chain lookup; an absent hint differs from block zero. A retry with the same note ID succeeds and keeps the first envelope, timestamp, and cursor. This also applies when the retry supplies a different block hint or storage is full. -`FetchNotes` accepts at most 128 tags and an exclusive cursor. Start with cursor zero. Use each response cursor for the -next request with the same set of tags. Results follow insertion order across all requested tags. Duplicate tags do not -duplicate results. Empty pages retain the request cursor. The response `has_more` field indicates that another page is -available. +`FetchNotes` accepts at most 128 tags and an exclusive cursor with a `fixed64` database nonce and a `fixed64` sequence. +Omit the cursor to start from the first retained note. Store the complete response cursor and use it for the next +request with the same set of tags. Results follow insertion order across all requested tags. Duplicate tags do not +duplicate results. Every successful response includes a cursor. An initial empty page returns the current nonce and +sequence zero. Later empty pages retain the request cursor. The response `has_more` field indicates that another page is +available. Nonce zero is valid and has no special meaning. -A cursor belongs to the requested set of tags. Reset the cursor to zero when you add or remove tags. Reordering tags or -changing duplicate tags does not change the set. A restart can return notes that you fetched before. Use note IDs to -remove duplicate results. +A cursor belongs to the requested set of tags. Clear the cursor when you add or remove tags. Reordering tags or changing +duplicate tags does not change the set. Restarting a fetch from the beginning can return notes that you fetched before. +Use note IDs to remove duplicate results. -For example, after you fetch tag A through cursor 100, reset the cursor to zero when you add tag B. If you reuse cursor -100, you skip retained notes for tag B with cursors at or below 100. +For example, after you fetch tag A through cursor 100, clear the cursor when you add tag B. If you reuse cursor 100, you +skip retained notes for tag B with cursors at or below 100. A page contains at most 500 notes and 3 MiB of canonically serialized header and detail bytes. The encoded response also fits the default 4 MiB gRPC client limit. The default per-note limit is 512,000 bytes. `--max-note-size` can change it up to the page limit. The required `--max-storage-bytes` limits retained header and detail bytes. It does not include -SQLite indexes, database metadata, or WAL disk usage. Cursors use the positive signed 64-bit range supported by SQLite. -Recipients must poll before notes expire. +SQLite indexes, database metadata, or WAL disk usage. Cursor sequences use the positive signed 64-bit range supported by +SQLite. Sequence zero starts a fetch. Nonces use the complete unsigned 64-bit range. Recipients must poll before notes +expire. + +The service generates a random nonce when it creates the database. Schema migration initializes the nonce for existing +databases. Ordinary service restarts, repeated migrations, and retention cleanup preserve the nonce. A cursor from a +different database generation returns `FAILED_PRECONDITION`. Clear the cursor and fetch again after this error. +Deduplicate results by note ID. The nonce detects a generation change; it does not recover lost notes or authenticate +cursors. + +The structured cursor is incompatible with the scalar cursor API. Coordinate client and server upgrades. Discard +persisted scalar cursors when upgrading clients. Run the database migration before starting the updated service. +Restoring an older backup also restores its nonce. That recovery procedure must rotate the nonce before the service +starts. This service does not provide a nonce rotation command. Malformed requests return `INVALID_ARGUMENT`. Note size and storage capacity limits return `RESOURCE_EXHAUSTED`. Storage failures return `INTERNAL` and are logged by the service. diff --git a/bin/note-transport/src/db/migrations/002_cursor_nonce.rs b/bin/note-transport/src/db/migrations/002_cursor_nonce.rs new file mode 100644 index 0000000000..092781c446 --- /dev/null +++ b/bin/note-transport/src/db/migrations/002_cursor_nonce.rs @@ -0,0 +1,23 @@ +use rusqlite::Transaction; + +/// Assigns a database nonce without changing notes or storage counters. +pub fn migrate(tx: &Transaction<'_>) -> anyhow::Result<()> { + tx.execute_batch( + "CREATE TABLE storage_metadata_with_nonce ( + singleton INTEGER PRIMARY KEY CHECK (singleton = 1), + retained_bytes INTEGER NOT NULL CHECK (retained_bytes >= 0), + nonce BLOB NOT NULL CHECK (length(nonce) = 8) + ) STRICT;", + )?; + let nonce = rand::random::().to_le_bytes(); + tx.execute( + "INSERT INTO storage_metadata_with_nonce (singleton, retained_bytes, nonce) + SELECT singleton, retained_bytes, ?1 FROM storage_metadata", + [nonce.as_slice()], + )?; + tx.execute_batch( + "DROP TABLE storage_metadata; + ALTER TABLE storage_metadata_with_nonce RENAME TO storage_metadata;", + )?; + Ok(()) +} diff --git a/bin/note-transport/src/db/mod.rs b/bin/note-transport/src/db/mod.rs index a85edc15ab..80ecd22fbc 100644 --- a/bin/note-transport/src/db/mod.rs +++ b/bin/note-transport/src/db/mod.rs @@ -52,12 +52,22 @@ pub enum StorageError { Capacity(String), #[error("invalid cursor")] InvalidCursor, + #[error("cursor belongs to another database generation; clear the cursor and retry")] + StaleCursor, #[error("{0}")] InvalidData(String), } +/// A position in one database generation. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct Cursor { + pub nonce: u64, + pub sequence: u64, +} + #[derive(Debug)] pub struct FetchPage { + pub cursor: Cursor, pub notes: Vec, pub has_more: bool, } @@ -134,7 +144,7 @@ pub async fn store_note( "accepting this note would exceed the {max_retained_bytes} byte limit" ))); } - queries::update_storage_metadata(tx, next_retained)?; + queries::update_retained_bytes(tx, next_retained)?; Ok(StoreResult::Inserted) }) .await @@ -145,15 +155,17 @@ pub async fn store_note( pub async fn fetch_notes( reader: &DbReader, tags: Vec, - cursor: u64, + cursor: Option, ) -> Result { - let cursor = i64::try_from(cursor).map_err(|_| StorageError::InvalidCursor)?; - if tags.is_empty() { - return Ok(FetchPage { notes: vec![], has_more: false }); - } + let sequence = cursor.map_or(0, |cursor| cursor.sequence); + let sequence = i64::try_from(sequence).map_err(|_| StorageError::InvalidCursor)?; reader .read("fetch_notes", move |tx| { - queries::fetch_notes(tx, tags.into_iter().map(NoteTag::new).collect(), cursor) + let nonce = queries::select_nonce(tx)?; + if cursor.is_some_and(|cursor| cursor.nonce != nonce) { + return Err(StorageError::StaleCursor); + } + queries::fetch_notes(tx, tags.into_iter().map(NoteTag::new).collect(), sequence, nonce) }) .await } diff --git a/bin/note-transport/src/db/queries/fetch_notes/mod.rs b/bin/note-transport/src/db/queries/fetch_notes/mod.rs index 3625e3ac5c..e6bab5579d 100644 --- a/bin/note-transport/src/db/queries/fetch_notes/mod.rs +++ b/bin/note-transport/src/db/queries/fetch_notes/mod.rs @@ -1,13 +1,21 @@ use miden_node_db::sqlite::{InList, ReadTx}; use miden_protocol::note::NoteTag; -use crate::db::{FETCH_NOTES_MAX_BYTES, FETCH_NOTES_MAX_ROWS, FetchPage, StorageError, StoredNote}; +use crate::db::{ + Cursor, + FETCH_NOTES_MAX_BYTES, + FETCH_NOTES_MAX_ROWS, + FetchPage, + StorageError, + StoredNote, +}; /// Returns a page in cursor order. Row and byte limits apply before note blobs are loaded. pub fn fetch_notes( tx: &ReadTx<'_>, tags: Vec, cursor: i64, + nonce: u64, ) -> Result { let tags = InList::from_values(tags); let rows = tx.query( @@ -35,5 +43,12 @@ pub fn fetch_notes( let candidate_count = rows.first().map_or(0, |(_, count)| *count); let notes: Vec<_> = rows.into_iter().map(|(note, _)| note).collect(); let has_more = candidate_count > i64::try_from(notes.len()).expect("page length fits i64"); - Ok(FetchPage { notes, has_more }) + let sequence = notes.last().map_or(cursor, |note| note.seq); + let sequence = u64::try_from(sequence) + .map_err(|_| StorageError::InvalidData("invalid stored note sequence".into()))?; + Ok(FetchPage { + notes, + has_more, + cursor: Cursor { nonce, sequence }, + }) } diff --git a/bin/note-transport/src/db/queries/mod.rs b/bin/note-transport/src/db/queries/mod.rs index 65b504b940..0cf0c16d57 100644 --- a/bin/note-transport/src/db/queries/mod.rs +++ b/bin/note-transport/src/db/queries/mod.rs @@ -6,12 +6,15 @@ pub use note_exists::note_exists; mod insert_note; pub use insert_note::insert_note; -mod update_storage_metadata; -pub use update_storage_metadata::update_storage_metadata; +mod update_retained_bytes; +pub use update_retained_bytes::update_retained_bytes; mod select_retained_bytes; pub use select_retained_bytes::select_retained_bytes; +mod select_nonce; +pub use select_nonce::select_nonce; + mod delete_notes_created_before; pub use delete_notes_created_before::delete_notes_created_before; diff --git a/bin/note-transport/src/db/queries/select_nonce/mod.rs b/bin/note-transport/src/db/queries/select_nonce/mod.rs new file mode 100644 index 0000000000..154833db6d --- /dev/null +++ b/bin/note-transport/src/db/queries/select_nonce/mod.rs @@ -0,0 +1,17 @@ +use miden_node_db::sqlite::ReadTx; + +use crate::db::StorageError; + +/// Reads the retained payload size. +pub fn select_nonce(tx: &ReadTx<'_>) -> Result { + let nonce = tx + .query(include_str!("select_nonce.sql"), &[], |row| row.get::>(0))? + .into_iter() + .next() + .ok_or_else(|| StorageError::InvalidData("storage metadata is missing".into()))?; + + nonce + .try_into() + .map_err(|_| StorageError::InvalidData("invalid database nonce length".into())) + .map(u64::from_le_bytes) +} diff --git a/bin/note-transport/src/db/queries/select_nonce/select_nonce.sql b/bin/note-transport/src/db/queries/select_nonce/select_nonce.sql new file mode 100644 index 0000000000..386920d072 --- /dev/null +++ b/bin/note-transport/src/db/queries/select_nonce/select_nonce.sql @@ -0,0 +1,3 @@ +SELECT nonce +FROM storage_metadata +WHERE singleton = 1; diff --git a/bin/note-transport/src/db/queries/update_storage_metadata/mod.rs b/bin/note-transport/src/db/queries/update_retained_bytes/mod.rs similarity index 69% rename from bin/note-transport/src/db/queries/update_storage_metadata/mod.rs rename to bin/note-transport/src/db/queries/update_retained_bytes/mod.rs index 5ed5efbc90..6d5ad201a1 100644 --- a/bin/note-transport/src/db/queries/update_storage_metadata/mod.rs +++ b/bin/note-transport/src/db/queries/update_retained_bytes/mod.rs @@ -2,7 +2,7 @@ use miden_node_db::DatabaseError; use miden_node_db::sqlite::WriteTx; /// Updates the retained payload size in the current transaction. -pub fn update_storage_metadata(tx: &WriteTx<'_>, retained_bytes: i64) -> Result<(), DatabaseError> { +pub fn update_retained_bytes(tx: &WriteTx<'_>, retained_bytes: i64) -> Result<(), DatabaseError> { tx.execute(include_str!("update_storage_metadata.sql"), &[&retained_bytes])?; Ok(()) } diff --git a/bin/note-transport/src/db/queries/update_storage_metadata/update_storage_metadata.sql b/bin/note-transport/src/db/queries/update_retained_bytes/update_storage_metadata.sql similarity index 100% rename from bin/note-transport/src/db/queries/update_storage_metadata/update_storage_metadata.sql rename to bin/note-transport/src/db/queries/update_retained_bytes/update_storage_metadata.sql diff --git a/bin/note-transport/src/db/tests.rs b/bin/note-transport/src/db/tests.rs index 438156b80c..e894e17daf 100644 --- a/bin/note-transport/src/db/tests.rs +++ b/bin/note-transport/src/db/tests.rs @@ -84,7 +84,7 @@ async fn retained_note_roundtrips_after_reopening() { drop((writer, reader)); let (_writer, reader) = load(&dir.path().join("notes.sqlite3")).unwrap(); - let page = fetch_notes(&reader, vec![u32::MAX], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![u32::MAX], None).await.unwrap(); assert_eq!(page.notes.len(), 1); let retained = &page.notes[0]; assert_eq!(retained.header, original.header); @@ -105,14 +105,14 @@ async fn retry_at_capacity_preserves_first_write() { StoreResult::Inserted ); let mut retry = original.clone(); - let created_at = fetch_notes(&reader, vec![42], 0).await.unwrap().notes[0].created_at; + let created_at = fetch_notes(&reader, vec![42], None).await.unwrap().notes[0].created_at; retry.after_block_num = Some(BlockNumber::from(99)); assert_eq!(store_note(&writer, retry, limit).await.unwrap(), StoreResult::AlreadyPresent); assert!(matches!( store_note(&writer, note(2, 42), limit).await, Err(StorageError::Capacity(_)) )); - let page = fetch_notes(&reader, vec![42], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![42], None).await.unwrap(); assert_eq!(page.notes.len(), 1); assert_eq!(page.notes[0].after_block_num, original.after_block_num); assert_eq!(page.notes[0].created_at, created_at); @@ -126,19 +126,19 @@ async fn multi_tag_fetch_has_stable_bounded_pages() { for seed in 1..=503 { store_note(&writer, note(seed, seed % 2), u64::MAX).await.unwrap(); } - let page = fetch_notes(&reader, vec![1, 0, 1], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![1, 0, 1], None).await.unwrap(); assert_eq!(page.notes.len(), FETCH_NOTES_MAX_ROWS); assert!(page.has_more); assert_eq!( page.notes.iter().map(|n| n.seq).collect::>(), (1..=500).collect::>() ); - let page = fetch_notes(&reader, vec![0, 1], 500).await.unwrap(); + let page = fetch_notes(&reader, vec![0, 1], Some(page.cursor)).await.unwrap(); assert_eq!(page.notes.len(), 3); assert!(!page.has_more); - assert!(fetch_notes(&reader, vec![], 0).await.unwrap().notes.is_empty()); + assert!(fetch_notes(&reader, vec![], None).await.unwrap().notes.is_empty()); assert!(matches!( - fetch_notes(&reader, vec![1], u64::MAX).await, + fetch_notes(&reader, vec![1], Some(Cursor { nonce: 0, sequence: u64::MAX })).await, Err(StorageError::InvalidCursor) )); } @@ -152,10 +152,10 @@ async fn fetch_byte_limit_and_oversize_rejection() { .await .unwrap(); } - let page = fetch_notes(&reader, vec![42], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![42], None).await.unwrap(); assert_eq!(page.notes.len(), 2); assert!(page.has_more); - assert_eq!(fetch_notes(&reader, vec![42], 2).await.unwrap().notes.len(), 2); + assert_eq!(fetch_notes(&reader, vec![42], Some(page.cursor)).await.unwrap().notes.len(), 2); let oversized = note_with_advice(5, 42, FETCH_NOTES_MAX_BYTES / 8); assert!(matches!( store_note(&writer, oversized, u64::MAX).await, @@ -182,7 +182,7 @@ async fn cursor_exhaustion_rejects_insert_without_losing_existing_notes() { store_note(&writer, note(1, 42), u64::MAX).await.unwrap(), StoreResult::AlreadyPresent ); - assert_eq!(fetch_notes(&reader, vec![42], 0).await.unwrap().notes.len(), 1); + assert_eq!(fetch_notes(&reader, vec![42], None).await.unwrap().notes.len(), 1); } #[tokio::test] @@ -203,7 +203,7 @@ async fn failed_insert_rolls_back_capacity_and_cursor() { .await .unwrap(); store_note(&writer, item, limit).await.unwrap(); - assert_eq!(fetch_notes(&reader, vec![42], 0).await.unwrap().notes[0].seq, 1); + assert_eq!(fetch_notes(&reader, vec![42], None).await.unwrap().notes[0].seq, 1); } #[tokio::test] @@ -214,7 +214,7 @@ async fn concurrent_writes_share_capacity() { let (first, second) = tokio::join!(store_note(&writer, item, limit), store_note(&writer, note(2, 42), limit),); assert_eq!(usize::from(first.is_ok()) + usize::from(second.is_ok()), 1); - assert_eq!(fetch_notes(&reader, vec![42], 0).await.unwrap().notes.len(), 1); + assert_eq!(fetch_notes(&reader, vec![42], None).await.unwrap().notes.len(), 1); } fn now_micros() -> i64 { @@ -263,10 +263,10 @@ async fn insertion_deletes_at_most_ten_oldest_notes_with_cursor_ties() { seed_expired(&writer, seed, timestamp, u64::MAX).await; } store_note(&writer, note(14, 42), u64::MAX).await.unwrap(); - let page = fetch_notes(&reader, vec![42], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![42], None).await.unwrap(); assert_eq!(page.notes.iter().map(|note| note.seq).collect::>(), vec![1, 3, 4, 14]); store_note(&writer, note(15, 42), u64::MAX).await.unwrap(); - let page = fetch_notes(&reader, vec![42], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![42], None).await.unwrap(); assert_eq!(page.notes.iter().map(|note| note.seq).collect::>(), vec![14, 15]); } @@ -286,7 +286,7 @@ async fn insertion_uses_cleanup_capacity_and_preserves_cursor_after_reopen() { store_note(&writer, note(5, 42), size * 2).await, Err(StorageError::Capacity(_)) )); - let page = fetch_notes(&reader, vec![42], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![42], None).await.unwrap(); assert_eq!(page.notes.iter().map(|note| note.seq).collect::>(), vec![3, 4]); } @@ -304,13 +304,13 @@ async fn insufficient_reclaimed_capacity_rolls_back_deletions_and_cursor() { store_note(&writer, oversized, limit).await, Err(StorageError::Capacity(_)) )); - let page = fetch_notes(&reader, vec![42], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![42], None).await.unwrap(); assert_eq!( page.notes.iter().map(|note| note.seq).collect::>(), (1..=12).collect::>() ); store_note(&writer, note(13, 42), limit).await.unwrap(); - let page = fetch_notes(&reader, vec![42], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![42], None).await.unwrap(); assert_eq!(page.notes.iter().map(|note| note.seq).collect::>(), vec![11, 12, 13]); } @@ -321,9 +321,9 @@ async fn duplicate_retry_and_reads_do_not_delete_expired_notes() { for seed in 2..=12 { seed_expired(&writer, seed, 0, u64::MAX).await; } - assert_eq!(fetch_notes(&reader, vec![42], 0).await.unwrap().notes.len(), 12); + assert_eq!(fetch_notes(&reader, vec![42], None).await.unwrap().notes.len(), 12); assert_eq!(store_note(&writer, original, 1).await.unwrap(), StoreResult::AlreadyPresent); - assert_eq!(fetch_notes(&reader, vec![42], 0).await.unwrap().notes.len(), 12); + assert_eq!(fetch_notes(&reader, vec![42], None).await.unwrap().notes.len(), 12); } #[tokio::test] @@ -336,7 +336,7 @@ async fn failed_cleanup_rolls_back_insert_and_deletions() { Ok::<_, miden_node_db::DatabaseError>(()) }).await.unwrap(); assert!(store_note(&writer, note(3, 42), u64::MAX).await.is_err()); - let page = fetch_notes(&reader, vec![42], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![42], None).await.unwrap(); assert_eq!(page.notes.iter().map(|note| note.seq).collect::>(), vec![1, 2]); writer .write("allow deletion", |tx| { @@ -346,7 +346,7 @@ async fn failed_cleanup_rolls_back_insert_and_deletions() { .await .unwrap(); store_note(&writer, note(3, 42), u64::MAX).await.unwrap(); - assert_eq!(fetch_notes(&reader, vec![42], 0).await.unwrap().notes[0].seq, 3); + assert_eq!(fetch_notes(&reader, vec![42], None).await.unwrap().notes[0].seq, 3); } #[tokio::test] @@ -359,12 +359,12 @@ async fn cleanup_preserves_notes_at_the_retention_boundary() { .write("check retention boundary", |tx| { let removed = queries::delete_notes_created_before(tx, 2, CLEANUP_MAX_NOTES)?; let retained = queries::select_retained_bytes(tx)?; - queries::update_storage_metadata(tx, retained - removed)?; + queries::update_retained_bytes(tx, retained - removed)?; Ok::<_, StorageError>(()) }) .await .unwrap(); - let page = fetch_notes(&reader, vec![42], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![42], None).await.unwrap(); assert_eq!(page.notes.iter().map(|note| note.seq).collect::>(), vec![2, 3]); } @@ -378,6 +378,30 @@ async fn insertion_uses_the_configured_retention_period() { super::store_note(&writer, note(3, 42), NonZeroU64::MAX, NonZeroU32::new(7).unwrap()) .await .unwrap(); - let page = fetch_notes(&reader, vec![42], 0).await.unwrap(); + let page = fetch_notes(&reader, vec![42], None).await.unwrap(); assert_eq!(page.notes.iter().map(|note| note.seq).collect::>(), vec![2, 3]); } + +#[tokio::test] +async fn cursor_survives_cleanup_of_all_notes() { + let (_dir, writer, reader) = database(); + seed_expired(&writer, 1, 0, u64::MAX).await; + let first = fetch_notes(&reader, vec![42], None).await.unwrap(); + writer + .write("remove expired notes", |tx| { + let removed = queries::delete_notes_created_before(tx, 1, CLEANUP_MAX_NOTES)?; + let retained_bytes = queries::select_retained_bytes(tx)?; + queries::update_retained_bytes(tx, retained_bytes - removed)?; + Ok::<_, StorageError>(()) + }) + .await + .unwrap(); + let empty = fetch_notes(&reader, vec![42], Some(first.cursor)).await.unwrap(); + assert!(empty.notes.is_empty()); + assert_eq!(empty.cursor, first.cursor); + store_note(&writer, note(2, 42), u64::MAX).await.unwrap(); + let next = fetch_notes(&reader, vec![42], Some(empty.cursor)).await.unwrap(); + assert_eq!(next.notes.len(), 1); + assert_eq!(next.cursor.sequence, 2); + assert_eq!(next.cursor.nonce, first.cursor.nonce); +} diff --git a/bin/note-transport/src/server/mod.rs b/bin/note-transport/src/server/mod.rs index c2f0352a6b..37c0149561 100644 --- a/bin/note-transport/src/server/mod.rs +++ b/bin/note-transport/src/server/mod.rs @@ -3,6 +3,7 @@ use std::num::{NonZeroU32, NonZeroU64, NonZeroUsize}; use miden_node_db::sqlite::{DbReader, DbWriter}; use miden_node_proto::errors::ConversionError; use miden_node_proto::generated::note_transport::{ + FetchNotesCursor, FetchNotesRequest, FetchNotesResponse, SendNoteRequest, @@ -198,7 +199,7 @@ impl FetchNotes for Server { if request.tags.len() > 128 { return Err(tonic::Status::invalid_argument("at most 128 tags are allowed")); } - if request.cursor > i64::MAX as u64 { + if request.cursor.is_some_and(|cursor| cursor.sequence > i64::MAX as u64) { return Err(tonic::Status::invalid_argument("invalid cursor")); } request.tags.sort_unstable(); @@ -217,15 +218,27 @@ impl FetchNotes for Server { _: &MetadataMap, _: &Extensions, ) -> tonic::Result { - let page = db::fetch_notes(&self.reader, request.tags, request.cursor) + let cursor = request.cursor.map(|cursor| db::Cursor { + nonce: cursor.nonce, + sequence: cursor.sequence, + }); + let page = db::fetch_notes(&self.reader, request.tags, cursor) .await .map_err(storage_status)?; - let mut cursor = request.cursor; + let mut cursor = FetchNotesCursor { + nonce: page.cursor.nonce, + sequence: cursor.map_or(0, |cursor| cursor.sequence), + }; let mut notes = Vec::with_capacity(page.notes.len()); let mut has_more = page.has_more; - // Reserve the fixed64 cursor and the boolean continuation field. - let mut response_bytes = 11; + // Reserve space for a nonzero sequence and a continuation flag. + let mut response_bytes = FetchNotesResponse { + notes: vec![], + cursor: Some(FetchNotesCursor { sequence: u64::MAX, ..cursor }), + has_more: true, + } + .encoded_len(); for note in page.notes { let next_cursor = u64::try_from(note.seq).map_err(|error| { error!(error, target: LOG_TARGET, "Invalid stored note cursor"); @@ -253,19 +266,22 @@ impl FetchNotes for Server { break; } response_bytes += field_bytes; - cursor = next_cursor; + cursor.sequence = next_cursor; notes.push(note); } info!(target: LOG_TARGET, "Notes fetched", - note_transport.returned = notes.len(), note_transport.cursor = cursor, + note_transport.returned = notes.len(), note_transport.cursor = cursor.sequence, note_transport.has_more = has_more); - Ok(FetchNotesResponse { notes, cursor, has_more }) + Ok(FetchNotesResponse { notes, cursor: Some(cursor), has_more }) } } fn storage_status(error: db::StorageError) -> tonic::Status { match error { db::StorageError::Capacity(message) => tonic::Status::resource_exhausted(message), + db::StorageError::StaleCursor => tonic::Status::failed_precondition( + "cursor belongs to another database generation; clear the cursor and retry", + ), db::StorageError::InvalidCursor => tonic::Status::invalid_argument("invalid cursor"), error => { error!(error, target: LOG_TARGET, "Note storage operation failed"); diff --git a/bin/note-transport/src/server/tests.rs b/bin/note-transport/src/server/tests.rs index 36b15a7680..7c6a5074ac 100644 --- a/bin/note-transport/src/server/tests.rs +++ b/bin/note-transport/src/server/tests.rs @@ -59,13 +59,13 @@ async fn send_fetch_preserves_optional_hint_presence() { } let page = FetchNotes::full( &server, - Request::new(FetchNotesRequest { tags: vec![8, 7, 7], cursor: 0 }), + Request::new(FetchNotesRequest { tags: vec![8, 7, 7], cursor: None }), ) .await .unwrap(); assert_eq!(page.notes, vec![first, second]); assert!(!page.has_more); - assert!(page.cursor > 0); + assert!(page.cursor.unwrap().sequence > 0); let empty = FetchNotes::full( &server, Request::new(FetchNotesRequest { tags: vec![7, 8], cursor: page.cursor }), @@ -85,7 +85,7 @@ async fn duplicate_send_preserves_first_note_and_cursor() { SendNote::full(&server, Request::new(SendNoteRequest { note: Some(first.clone()) })) .await .unwrap(); - let request = FetchNotesRequest { tags: vec![7], cursor: 0 }; + let request = FetchNotesRequest { tags: vec![7], cursor: None }; let before = FetchNotes::full(&server, Request::new(request.clone())).await.unwrap(); // Retry sending the same note with a different `after_block_num` and assert that it's not @@ -115,7 +115,10 @@ async fn invalid_cursor_has_generic_message() { let (_dir, server) = server(Config::default()); let error = FetchNotes::full( &server, - Request::new(FetchNotesRequest { tags: vec![7], cursor: u64::MAX }), + Request::new(FetchNotesRequest { + tags: vec![7], + cursor: Some(FetchNotesCursor { nonce: 0, sequence: u64::MAX }), + }), ) .await .unwrap_err(); @@ -178,8 +181,11 @@ async fn rejects_oversized_notes_and_invalid_fetches() { .unwrap_err(); assert_eq!(error.code(), tonic::Code::ResourceExhausted); for request in [ - FetchNotesRequest { tags: vec![7; 129], cursor: 0 }, - FetchNotesRequest { tags: vec![], cursor: u64::MAX }, + FetchNotesRequest { tags: vec![7; 129], cursor: None }, + FetchNotesRequest { + tags: vec![], + cursor: Some(FetchNotesCursor { nonce: 0, sequence: u64::MAX }), + }, ] { let error = FetchNotes::full(&server, Request::new(request)).await.unwrap_err(); assert_eq!(error.code(), tonic::Code::InvalidArgument); @@ -197,7 +203,7 @@ async fn capacity_failure_does_not_store_the_note() { .unwrap_err(); assert_eq!(error.code(), tonic::Code::ResourceExhausted); let page = - FetchNotes::full(&server, Request::new(FetchNotesRequest { tags: vec![7], cursor: 0 })) + FetchNotes::full(&server, Request::new(FetchNotesRequest { tags: vec![7], cursor: None })) .await .unwrap(); assert!(page.notes.is_empty()); @@ -262,7 +268,7 @@ async fn grpc_health_reflection_web_and_shutdown() { .await .unwrap(); let response = client - .fetch_notes(FetchNotesRequest { tags: vec![123], cursor: 0 }) + .fetch_notes(FetchNotesRequest { tags: vec![123], cursor: None }) .await .unwrap() .into_inner(); @@ -296,9 +302,12 @@ async fn grpc_health_reflection_web_and_shutdown() { let length = u32::from_be_bytes(frame[1..5].try_into().unwrap()) as usize; assert_eq!(SendNoteResponse::decode(&frame[5..5 + length]).unwrap(), SendNoteResponse {}); - let web = - grpc_web_request(address, "FetchNotes", FetchNotesRequest { tags: vec![123], cursor: 0 }) - .await; + let web = grpc_web_request( + address, + "FetchNotes", + FetchNotesRequest { tags: vec![123], cursor: None }, + ) + .await; assert_eq!(web.status(), reqwest::StatusCode::OK); assert_eq!(web.headers()["access-control-allow-origin"], "*"); let frame = web.bytes().await.unwrap(); @@ -357,7 +366,7 @@ async fn large_pages_fit_default_grpc_client_and_resume_without_gaps() { .await .unwrap(); let first = client - .fetch_notes(FetchNotesRequest { tags: vec![77], cursor: 0 }) + .fetch_notes(FetchNotesRequest { tags: vec![77], cursor: None }) .await .unwrap() .into_inner(); @@ -384,3 +393,157 @@ async fn large_pages_fit_default_grpc_client_and_resume_without_gaps() { shutdown.cancel(); task.await.unwrap().unwrap(); } + +#[tokio::test] +async fn rejects_cursor_from_another_database() { + let (_first_dir, first) = server(Config::default()); + let (_second_dir, second) = server(Config::default()); + let page = + FetchNotes::full(&first, Request::new(FetchNotesRequest { tags: vec![7], cursor: None })) + .await + .unwrap(); + for tags in [vec![7], vec![]] { + let error = FetchNotes::full( + &second, + Request::new(FetchNotesRequest { tags, cursor: page.cursor }), + ) + .await + .unwrap_err(); + assert_eq!(error.code(), tonic::Code::FailedPrecondition); + } +} + +#[tokio::test] +async fn cursor_survives_reopening_database() { + let (dir, server) = server(Config::default()); + SendNote::full(&server, Request::new(SendNoteRequest { note: Some(note(1, 7)) })) + .await + .unwrap(); + let first = + FetchNotes::full(&server, Request::new(FetchNotesRequest { tags: vec![7], cursor: None })) + .await + .unwrap(); + drop(server); + let path = dir.path().join("notes.sqlite3"); + db::migrate(&path).unwrap(); + let (writer, reader) = db::load(&path).unwrap(); + let server = Server::new(Config::default(), writer, reader).unwrap(); + let empty = FetchNotes::full( + &server, + Request::new(FetchNotesRequest { tags: vec![7], cursor: first.cursor }), + ) + .await + .unwrap(); + assert!(empty.notes.is_empty()); + assert_eq!(empty.cursor, first.cursor); + let next = note(2, 7); + SendNote::full(&server, Request::new(SendNoteRequest { note: Some(next.clone()) })) + .await + .unwrap(); + let page = FetchNotes::full( + &server, + Request::new(FetchNotesRequest { tags: vec![7], cursor: empty.cursor }), + ) + .await + .unwrap(); + assert_eq!(page.notes, vec![next]); +} + +#[tokio::test] +async fn stale_cursor_fails_before_and_after_sequence_catches_up() { + let (_first_dir, first) = server(Config::default()); + let (_second_dir, second) = server(Config::default()); + SendNote::full(&first, Request::new(SendNoteRequest { note: Some(note(1, 7)) })) + .await + .unwrap(); + let previous = + FetchNotes::full(&first, Request::new(FetchNotesRequest { tags: vec![7], cursor: None })) + .await + .unwrap(); + for count in 0..=2 { + if count > 0 { + SendNote::full(&second, Request::new(SendNoteRequest { note: Some(note(count, 7)) })) + .await + .unwrap(); + } + let error = FetchNotes::full( + &second, + Request::new(FetchNotesRequest { tags: vec![7], cursor: previous.cursor }), + ) + .await + .unwrap_err(); + assert_eq!(error.code(), tonic::Code::FailedPrecondition); + } + let restarted = + FetchNotes::full(&second, Request::new(FetchNotesRequest { tags: vec![7], cursor: None })) + .await + .unwrap(); + assert_eq!(restarted.notes, vec![note(1, 7), note(2, 7)]); + assert_eq!(restarted.cursor.unwrap().sequence, 2); +} + +#[tokio::test] +async fn empty_pages_and_nonce_extremes_roundtrip() { + let (_dir, server) = server(Config::default()); + for nonce in [0, u64::MAX] { + server + .writer + .write("set database nonce", move |tx| { + tx.execute( + "UPDATE storage_metadata SET nonce = ?1", + &[&nonce.to_le_bytes().to_vec()], + )?; + Ok::<_, db::StorageError>(()) + }) + .await + .unwrap(); + for tags in [vec![7], vec![]] { + let initial = FetchNotes::full( + &server, + Request::new(FetchNotesRequest { tags: tags.clone(), cursor: None }), + ) + .await + .unwrap(); + assert!(initial.notes.is_empty()); + assert!(!initial.has_more); + assert_eq!(initial.cursor, Some(FetchNotesCursor { nonce, sequence: 0 })); + let response = FetchNotesResponse::decode(initial.encode_to_vec().as_slice()).unwrap(); + let request = FetchNotesRequest { + tags: tags.clone(), + cursor: response.cursor, + }; + let request = FetchNotesRequest::decode(request.encode_to_vec().as_slice()).unwrap(); + let next = FetchNotes::full(&server, Request::new(request)).await.unwrap(); + assert_eq!(next.cursor, initial.cursor); + assert!(next.notes.is_empty()); + let error = FetchNotes::full( + &server, + Request::new(FetchNotesRequest { + tags, + cursor: Some(FetchNotesCursor { nonce: nonce ^ 1, sequence: 0 }), + }), + ) + .await + .unwrap_err(); + assert_eq!(error.code(), tonic::Code::FailedPrecondition); + } + } +} + +#[tokio::test] +async fn missing_metadata_fails_even_without_tags() { + let (_dir, server) = server(Config::default()); + server + .writer + .write("remove metadata", |tx| { + tx.execute("DELETE FROM storage_metadata", &[])?; + Ok::<_, db::StorageError>(()) + }) + .await + .unwrap(); + let error = + FetchNotes::full(&server, Request::new(FetchNotesRequest { tags: vec![], cursor: None })) + .await + .unwrap_err(); + assert_eq!(error.code(), tonic::Code::Internal); +} diff --git a/proto/proto/note_transport.proto b/proto/proto/note_transport.proto index bbe09768f5..d493a644c0 100644 --- a/proto/proto/note_transport.proto +++ b/proto/proto/note_transport.proto @@ -26,22 +26,39 @@ message TransportNote { optional blockchain.BlockNumber after_block_num = 3; } +// Identifies a position in a stable sequence of stored notes. +message FetchNotesCursor { + // Random nonce that identifies the sequence. Service restarts and note cleanup preserve it. + // Recreating the database assigns a new nonce and invalidates previous cursors. + // All unsigned 64-bit values are valid. + fixed64 nonce = 1; + // Exclusive lower bound. Zero starts from the first retained note. + fixed64 sequence = 2; +} + message FetchNotesRequest { - // At most 128 tags. Duplicate tags do not duplicate results. + // Tags to fetch notes for. Accepts up to 128 tags. + // Duplicate tags do not cause duplicate results. repeated fixed32 tags = 1; - // Exclusive lower bound. Zero starts from the first retained note. - // The cursor belongs to the requested set of tags. - // Reset it to zero when you add or remove tags. - // Reordering tags or changing duplicate tags does not change the set. - // A restart can return notes that you fetched before. - // Use note IDs to remove duplicate results. - fixed64 cursor = 2; + + // Continue from the cursor returned by the previous response. + // Omit the cursor to start from the first retained note. + // + // Reuse the cursor only for the same set of tags. Clear it when you add or remove a tag. + // Reordering tags or adding or removing duplicate tags does not require a reset. + // + // If the cursor nonce does not match the database nonce, the request returns + // FAILED_PRECONDITION. Clear the cursor and retry from the first retained note. + // This restart can return notes already fetched. Use note IDs to remove duplicates. + FetchNotesCursor cursor = 2; } message FetchNotesResponse { repeated TransportNote notes = 1; - // Cursor of the last returned note. An empty page keeps the request cursor. - fixed64 cursor = 2; + // Always present. Contains the sequence of the last returned note. + // An empty page keeps the validated request cursor. + // An initial empty page returns the current nonce and sequence zero. + FetchNotesCursor cursor = 2; // Indicates that another page is available for these tags. bool has_more = 3; }