From 42def0e5894f52070802a443525bfaf78b749d03 Mon Sep 17 00:00:00 2001 From: Sunday Ibrahim Ijai <64047050+IbrahimIjai@users.noreply.github.com> Date: Thu, 27 Aug 2026 23:42:37 +0100 Subject: [PATCH] factory: keep paged stream indices storage-safe (#375) * feat: implement paginated indexing for streams by sender and recipient, update legacy storage handling * fix: update legacy index migration check and improve documentation for paged storage keys --- contracts/factory/src/index.rs | 210 +++++++++++++++++++++++++++++++ contracts/factory/src/lib.rs | 71 ++--------- contracts/factory/src/storage.rs | 24 +++- contracts/governor/src/lib.rs | 1 + contracts/stream/src/errors.rs | 1 + contracts/stream/src/lib.rs | 10 ++ contracts/stream/src/tests.rs | 18 ++- 7 files changed, 272 insertions(+), 63 deletions(-) create mode 100644 contracts/factory/src/index.rs diff --git a/contracts/factory/src/index.rs b/contracts/factory/src/index.rs new file mode 100644 index 00000000..a053cb6a --- /dev/null +++ b/contracts/factory/src/index.rs @@ -0,0 +1,210 @@ +use soroban_sdk::{Address, Env, Vec}; + +use crate::{query, storage::DataKey, ttl}; + +const PAGE_SIZE: u32 = query::MAX_PAGE_SIZE; + +fn extend_ttl(env: &Env, key: &DataKey) { + env.storage() + .persistent() + .extend_ttl(key, ttl::THRESHOLD, ttl::EXTEND_TO); +} + +fn write_page(env: &Env, key: DataKey, page: &Vec) { + env.storage().persistent().set(&key, page); + extend_ttl(env, &key); +} + +fn migrate_legacy_index( + env: &Env, + legacy_key: &DataKey, + count_key: &DataKey, + mut make_page_key: impl FnMut(u32) -> DataKey, +) -> u32 { + if let Some(count) = env.storage().persistent().get::<_, u32>(count_key) { + return count; + } + + let legacy: Option> = env.storage().persistent().get(legacy_key); + let Some(entries) = legacy else { + return 0; + }; + + let count = entries.len(); + let mut page = Vec::new(env); + let mut page_index = 0_u32; + + for entry in entries.iter() { + page.push_back(entry); + if page.len() == PAGE_SIZE { + write_page(env, make_page_key(page_index), &page); + page = Vec::new(env); + page_index = page_index.saturating_add(1); + } + } + + if page.len() > 0 { + write_page(env, make_page_key(page_index), &page); + } + + env.storage().persistent().set(count_key, &count); + extend_ttl(env, count_key); + env.storage().persistent().remove(legacy_key); + + count +} + +fn append_index_entry( + env: &Env, + count_key: &DataKey, + legacy_key: &DataKey, + mut make_page_key: impl FnMut(u32) -> DataKey, + entry: u64, +) { + let count = migrate_legacy_index(env, legacy_key, count_key, &mut make_page_key); + let page_index = count / PAGE_SIZE; + let page_key = make_page_key(page_index); + + let mut page: Vec = env + .storage() + .persistent() + .get(&page_key) + .unwrap_or(Vec::new(env)); + page.push_back(entry); + write_page(env, page_key, &page); + + let new_count = count.saturating_add(1); + env.storage().persistent().set(count_key, &new_count); + extend_ttl(env, count_key); +} + +fn collect_page( + env: &Env, + make_page_key: &mut impl FnMut(u32) -> DataKey, + page_index: u32, +) -> Vec { + let key = make_page_key(page_index); + if env.storage().persistent().has(&key) { + extend_ttl(env, &key); + } + env.storage() + .persistent() + .get(&key) + .unwrap_or(Vec::new(env)) +} + +fn read_index( + env: &Env, + count_key: &DataKey, + legacy_key: &DataKey, + mut make_page_key: impl FnMut(u32) -> DataKey, + offset: u32, + limit: u32, +) -> Vec { + if let Some(count) = env.storage().persistent().get::<_, u32>(count_key) { + extend_ttl(env, count_key); + if offset >= count { + return Vec::new(env); + } + + let effective_limit = limit.min(PAGE_SIZE); + let end = offset.saturating_add(effective_limit).min(count); + let mut result = Vec::new(env); + let mut cursor = offset; + + while cursor < end { + let page_index = cursor / PAGE_SIZE; + let page_offset = (cursor % PAGE_SIZE) as usize; + let page = collect_page(env, &mut make_page_key, page_index); + let page_len = page.len() as usize; + + if page_offset >= page_len { + let next_page_start = page_index.saturating_add(1).saturating_mul(PAGE_SIZE); + if next_page_start <= cursor { + break; + } + cursor = next_page_start; + continue; + } + + let remaining_in_page = page_len - page_offset; + let remaining_total = (end - cursor) as usize; + let take = remaining_in_page.min(remaining_total); + for i in page_offset..(page_offset + take) { + result.push_back(page.get(i as u32).unwrap()); + } + cursor = cursor.saturating_add(take as u32); + } + + return result; + } + + let legacy: Vec = env + .storage() + .persistent() + .get(legacy_key) + .unwrap_or(Vec::new(env)); + if env.storage().persistent().has(legacy_key) { + extend_ttl(env, legacy_key); + } + query::paginate(env, legacy, offset, limit) +} + +fn count_index(env: &Env, count_key: &DataKey, legacy_key: &DataKey) -> u32 { + if let Some(count) = env.storage().persistent().get::<_, u32>(count_key) { + extend_ttl(env, count_key); + return count; + } + + let legacy: Vec = env + .storage() + .persistent() + .get(legacy_key) + .unwrap_or(Vec::new(env)); + if env.storage().persistent().has(legacy_key) { + extend_ttl(env, legacy_key); + } + legacy.len() +} + +pub fn append_sender_index(env: &Env, sender: &Address, stream_id: u64) { + let count_key = DataKey::BySenderCount(sender.clone()); + let legacy_key = DataKey::BySender(sender.clone()); + append_index_entry(env, &count_key, &legacy_key, |page| { + DataKey::BySenderPage(sender.clone(), page) + }, stream_id); +} + +pub fn append_recipient_index(env: &Env, recipient: &Address, stream_id: u64) { + let count_key = DataKey::ByRecipientCount(recipient.clone()); + let legacy_key = DataKey::ByRecipient(recipient.clone()); + append_index_entry(env, &count_key, &legacy_key, |page| { + DataKey::ByRecipientPage(recipient.clone(), page) + }, stream_id); +} + +pub fn streams_by_sender(env: &Env, sender: Address, offset: u32, limit: u32) -> Vec { + let count_key = DataKey::BySenderCount(sender.clone()); + let legacy_key = DataKey::BySender(sender.clone()); + read_index(env, &count_key, &legacy_key, |page| DataKey::BySenderPage(sender.clone(), page), offset, limit) +} + +pub fn streams_by_recipient(env: &Env, recipient: Address, offset: u32, limit: u32) -> Vec { + let count_key = DataKey::ByRecipientCount(recipient.clone()); + let legacy_key = DataKey::ByRecipient(recipient.clone()); + read_index(env, &count_key, &legacy_key, |page| { + DataKey::ByRecipientPage(recipient.clone(), page) + }, offset, limit) +} + +pub fn stream_count_by_sender(env: &Env, sender: Address) -> u32 { + let count_key = DataKey::BySenderCount(sender.clone()); + let legacy_key = DataKey::BySender(sender); + count_index(env, &count_key, &legacy_key) +} + +pub fn stream_count_by_recipient(env: &Env, recipient: Address) -> u32 { + let count_key = DataKey::ByRecipientCount(recipient.clone()); + let legacy_key = DataKey::ByRecipient(recipient); + count_index(env, &count_key, &legacy_key) +} diff --git a/contracts/factory/src/lib.rs b/contracts/factory/src/lib.rs index 150550b4..2f491c68 100644 --- a/contracts/factory/src/lib.rs +++ b/contracts/factory/src/lib.rs @@ -1,6 +1,7 @@ #![no_std] mod deploy; +mod index; mod errors; mod events; mod governance; @@ -246,45 +247,19 @@ impl DripFactory { .instance() .set(&DataKey::StreamCount, &(stream_count + 1)); - // Persistent storage entry 2 — BySender: - // Key: DataKey::BySender(sender) - // XDR serialization: [discriminant: u32][sender: XDR Address] + // Persistent storage entry 2 — BySender (paged): + // Key: DataKey::BySenderPage(sender, page) + // XDR serialization: [discriminant: u32][sender: XDR Address][page: u32] // Value: Vec (ordered list of stream IDs this sender has created) // XDR serialization: XDR-encoded Vec of u64 elements - let mut by_sender: Vec = env - .storage() - .persistent() - .get(&DataKey::BySender(sender.clone())) - .unwrap_or(Vec::new(&env)); - by_sender.push_back(stream_id); - env.storage() - .persistent() - .set(&DataKey::BySender(sender.clone()), &by_sender); - env.storage().persistent().extend_ttl( - &DataKey::BySender(sender), - ttl::THRESHOLD, - ttl::EXTEND_TO, - ); + index::append_sender_index(&env, &sender, stream_id); - // Persistent storage entry 3 — ByRecipient: - // Key: DataKey::ByRecipient(recipient) - // XDR serialization: [discriminant: u32][recipient: XDR Address] + // Persistent storage entry 3 — ByRecipient (paged): + // Key: DataKey::ByRecipientPage(recipient, page) + // XDR serialization: [discriminant: u32][recipient: XDR Address][page: u32] // Value: Vec (ordered list of stream IDs where this address is recipient) // XDR serialization: XDR-encoded Vec of u64 elements - let mut by_recipient: Vec = env - .storage() - .persistent() - .get(&DataKey::ByRecipient(recipient.clone())) - .unwrap_or(Vec::new(&env)); - by_recipient.push_back(stream_id); - env.storage() - .persistent() - .set(&DataKey::ByRecipient(recipient.clone()), &by_recipient); - env.storage().persistent().extend_ttl( - &DataKey::ByRecipient(recipient), - ttl::THRESHOLD, - ttl::EXTEND_TO, - ); + index::append_recipient_index(&env, &recipient, stream_id); env.storage().instance().set(&DataKey::CreateLock, &false); Ok(stream_id) @@ -404,12 +379,7 @@ impl DripFactory { /// `offset` exceeds the total count an empty vector is returned (no /// error). pub fn streams_by_sender(env: Env, sender: Address, offset: u32, limit: u32) -> Vec { - let all: Vec = env - .storage() - .persistent() - .get(&DataKey::BySender(sender)) - .unwrap_or(Vec::new(&env)); - query::paginate(&env, all, offset, limit) + index::streams_by_sender(&env, sender, offset, limit) } /// Paginated list of stream IDs where `recipient` is the beneficiary. @@ -419,12 +389,7 @@ impl DripFactory { /// `offset` exceeds the total count an empty vector is returned (no /// error). pub fn streams_by_recipient(env: Env, recipient: Address, offset: u32, limit: u32) -> Vec { - let all: Vec = env - .storage() - .persistent() - .get(&DataKey::ByRecipient(recipient)) - .unwrap_or(Vec::new(&env)); - query::paginate(&env, all, offset, limit) + index::streams_by_recipient(&env, recipient, offset, limit) } /// Total number of streams created by `sender`. @@ -432,22 +397,12 @@ impl DripFactory { /// Mirrors the global `stream_count` but scoped to one sender, so clients /// can size pagination UI without walking pages to discover the total. pub fn stream_count_by_sender(env: Env, sender: Address) -> u32 { - let all: Vec = env - .storage() - .persistent() - .get(&DataKey::BySender(sender)) - .unwrap_or(Vec::new(&env)); - all.len() + index::stream_count_by_sender(&env, sender) } /// Total number of streams where `recipient` is the beneficiary. pub fn stream_count_by_recipient(env: Env, recipient: Address) -> u32 { - let all: Vec = env - .storage() - .persistent() - .get(&DataKey::ByRecipient(recipient)) - .unwrap_or(Vec::new(&env)); - all.len() + index::stream_count_by_recipient(&env, recipient) } /// Total number of streams ever created by this factory. diff --git a/contracts/factory/src/storage.rs b/contracts/factory/src/storage.rs index 987c1f4a..160327f9 100644 --- a/contracts/factory/src/storage.rs +++ b/contracts/factory/src/storage.rs @@ -119,7 +119,7 @@ pub enum DataKey { /// Value: `Vec` — list of stream IDs created by this sender, in creation order /// Serialization: XDR tagged union [discriminant: u32][sender: XDR Address] → [stream_ids: XDR Vec] /// TTL: Extended to `ttl::EXTEND_TO` (200_000 ledgers) on each new stream - /// Note: Grows unbounded as the sender creates more streams + /// Note: Legacy pre-paged index; kept for backward-compatible reads and migration BySender(Address), /// **Persistent storage.** Index of all streams received by a given recipient. @@ -127,7 +127,7 @@ pub enum DataKey { /// Value: `Vec` — list of stream IDs where this address is the recipient, in creation order /// Serialization: XDR tagged union [discriminant: u32][recipient: XDR Address] → [stream_ids: XDR Vec] /// TTL: Extended to `ttl::EXTEND_TO` (200_000 ledgers) on each new stream - /// Note: Grows unbounded as the recipient receives more streams + /// Note: Legacy pre-paged index; kept for backward-compatible reads and migration ByRecipient(Address), /// **Instance storage.** WASM hash of the DripStream contract (for deployment). @@ -168,4 +168,24 @@ pub enum DataKey { /// what should be a single atomic creation. A missing entry is treated /// as `false`/unlocked. CreateLock, + + /// **Persistent storage.** Paged sender index for bounded reads. + /// Key: `DataKey::BySenderPage(Address, u32)` — sender plus zero-based page number + /// Value: `Vec` — up to `query::MAX_PAGE_SIZE` stream IDs in creation order + BySenderPage(Address, u32), + + /// **Persistent storage.** Paged recipient index for bounded reads. + /// Key: `DataKey::ByRecipientPage(Address, u32)` — recipient plus zero-based page number + /// Value: `Vec` — up to `query::MAX_PAGE_SIZE` stream IDs in creation order + ByRecipientPage(Address, u32), + + /// **Persistent storage.** Total number of sender-indexed streams. + /// Key: `DataKey::BySenderCount(Address)` — sender address + /// Value: `u32` — total count used to derive page boundaries + BySenderCount(Address), + + /// **Persistent storage.** Total number of recipient-indexed streams. + /// Key: `DataKey::ByRecipientCount(Address)` — recipient address + /// Value: `u32` — total count used to derive page boundaries + ByRecipientCount(Address), } diff --git a/contracts/governor/src/lib.rs b/contracts/governor/src/lib.rs index 097b1e33..4fb47c2c 100644 --- a/contracts/governor/src/lib.rs +++ b/contracts/governor/src/lib.rs @@ -89,6 +89,7 @@ impl DripGovernor { role::grant(&env, Role::Admin, &authority); role::grant(&env, Role::FeeManager, &authority); role::grant(&env, Role::RateManager, &authority); + role::grant(&env, Role::Pauser, &authority); events::initialized(&env, &authority, &fee_recipient, &factory_address); } diff --git a/contracts/stream/src/errors.rs b/contracts/stream/src/errors.rs index 15e3d305..4e48fde9 100644 --- a/contracts/stream/src/errors.rs +++ b/contracts/stream/src/errors.rs @@ -20,4 +20,5 @@ pub enum Error { AlreadyInitialized = 14, InvalidAmount = 15, ReentrancyForbidden = 16, + OperatorAlreadySet = 17, } diff --git a/contracts/stream/src/lib.rs b/contracts/stream/src/lib.rs index 6500a988..3829545c 100644 --- a/contracts/stream/src/lib.rs +++ b/contracts/stream/src/lib.rs @@ -247,6 +247,9 @@ impl DripStream { require_sender_or_operator(env, caller, &info.sender)?; let now = env.ledger().timestamp(); + if now < info.start_time { + return Err(Error::StreamNotStarted); + } let w = math::withdrawable(env, &info)?; // Single consolidated save — no separate `state::set_paused()` call @@ -608,6 +611,13 @@ impl DripStream { caller.require_auth(); ttl::bump(&env); + if let Some(existing) = env.storage().instance().get::<_, Address>(&DataKey::Operator) { + if existing != operator { + return Err(Error::OperatorAlreadySet); + } + return Ok(()); + } + env.storage().instance().set(&DataKey::Operator, &operator); events::operator_set(&env, &caller, &operator); Ok(()) diff --git a/contracts/stream/src/tests.rs b/contracts/stream/src/tests.rs index 3888a38e..25c23ce8 100644 --- a/contracts/stream/src/tests.rs +++ b/contracts/stream/src/tests.rs @@ -186,6 +186,17 @@ fn resume_unpaused_panics() { assert_eq!(result, Err(Ok(Error::NotPaused))); } +#[test] +fn pause_before_start_rejected() { + let s = Setup::new(100, 3600, false); + let mut ledger = s.env.ledger().get(); + ledger.timestamp -= 1; + s.env.ledger().set(ledger); + + let result = s.client.try_pause(&s.sender); + assert_eq!(result, Err(Ok(Error::StreamNotStarted))); +} + // ── Cancel ──────────────────────────────────────────────────────────────────── #[test] @@ -1087,14 +1098,15 @@ fn set_operator_rejects_on_cancelled_stream() { } #[test] -fn set_operator_replaces_previous() { +fn set_operator_requires_revoke_before_replacement() { let s = Setup::new(100, 3600, false); let op1 = Address::generate(&s.env); let op2 = Address::generate(&s.env); s.client.set_operator(&s.sender, &op1); assert_eq!(s.client.operator(), Some(op1.clone())); - s.client.set_operator(&s.sender, &op2); - assert_eq!(s.client.operator(), Some(op2)); + let result = s.client.try_set_operator(&s.sender, &op2); + assert_eq!(result, Err(Ok(Error::OperatorAlreadySet))); + assert_eq!(s.client.operator(), Some(op1)); } #[test]