From 823918c1fb58c9d47ea36b04c42e07a7e7a78544 Mon Sep 17 00:00:00 2001 From: Oyebisi Ayobami Date: Tue, 29 Sep 2026 11:51:37 +0100 Subject: [PATCH 01/10] feat(stream): add prospective rate checkpoint --- contracts/stream/src/storage.rs | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/contracts/stream/src/storage.rs b/contracts/stream/src/storage.rs index 25e590ee..e1ba24af 100644 --- a/contracts/stream/src/storage.rs +++ b/contracts/stream/src/storage.rs @@ -55,6 +55,11 @@ pub enum DataKey { MinWithdrawalInterval, /// Timestamp of the last withdrawal. LastWithdrawalTime, + /// Optional accrual checkpoint created by `change_rate`. + /// + /// Kept outside `StreamInfo` so adding rate changes does not alter the + /// serialized `Config` layout of already-deployed streams. + RateCheckpoint, } #[contracttype] @@ -72,6 +77,19 @@ pub struct CliffConfig { pub cliff_unlock_amount: i128, } +/// Accrued-value checkpoint used to make rate changes prospective. +/// +/// `accrued` is the total amount earned at `at`. From that point onward +/// streaming uses the current `StreamInfo::rate_per_second`. For a pending +/// stream (or a stream still before its cliff), `at` may be in the future; +/// no value is considered streamed until that timestamp is reached. +#[contracttype] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct RateCheckpoint { + pub accrued: i128, + pub at: u64, +} + /// Alias for SplitConfig to support split recipients. pub type StreamConfig = SplitConfig; From 410bb59f6db81bc5f29da8461aeb443308734773 Mon Sep 17 00:00:00 2001 From: Oyebisi Ayobami Date: Tue, 29 Sep 2026 11:51:40 +0100 Subject: [PATCH 02/10] feat(stream): preserve accrual across rate changes --- contracts/stream/src/math.rs | 26 +++++++++++++++++++++++++- 1 file changed, 25 insertions(+), 1 deletion(-) diff --git a/contracts/stream/src/math.rs b/contracts/stream/src/math.rs index a0546067..cc73e130 100644 --- a/contracts/stream/src/math.rs +++ b/contracts/stream/src/math.rs @@ -1,7 +1,7 @@ use soroban_sdk::Env; use crate::errors::Error; -use crate::storage::{DataKey, CliffConfig, StreamInfo}; +use crate::storage::{DataKey, CliffConfig, RateCheckpoint, StreamInfo}; /// Returns the total tokens that have streamed up to `now`, /// excluding any paused time. Does not account for withdrawals. @@ -43,6 +43,30 @@ pub fn streamed_amount(env: &Env, info: &StreamInfo) -> Result { now }; + // A rate change checkpoints all value accrued under the previous rate. + // This makes the new rate prospective: previously earned value never gets + // retroactively repriced. The checkpoint is optional so pre-upgrade streams + // continue to use the original start/cliff calculation until their first + // rate change. + let checkpoint: Option = + env.storage().instance().get(&DataKey::RateCheckpoint); + if let Some(checkpoint) = checkpoint { + if effective_now < checkpoint.at { + return Ok(0); + } + let elapsed = effective_now + .checked_sub(checkpoint.at) + .ok_or(Error::ArithmeticOverflow)?; + let linear = info + .rate_per_second + .checked_mul(elapsed as i128) + .ok_or(Error::ArithmeticOverflow)?; + return checkpoint + .accrued + .checked_add(linear) + .ok_or(Error::ArithmeticOverflow); + } + // Check for optional cliff config (Issue #719) let cliff_cfg: Option = env.storage().instance().get(&DataKey::CliffConfig); if let Some(cfg) = cliff_cfg { From 14e6f4b825d6520583c7143a4c502cc4ad3440df Mon Sep 17 00:00:00 2001 From: Oyebisi Ayobami Date: Tue, 29 Sep 2026 11:51:43 +0100 Subject: [PATCH 03/10] feat(stream): add combined and rate-change events --- contracts/stream/src/events.rs | 35 ++++++++++++++++++++++++++++++++++ 1 file changed, 35 insertions(+) diff --git a/contracts/stream/src/events.rs b/contracts/stream/src/events.rs index 27948e97..ac292291 100644 --- a/contracts/stream/src/events.rs +++ b/contracts/stream/src/events.rs @@ -139,6 +139,41 @@ pub fn topped_up(env: &Env, sender: &Address, amount: i128, new_balance: i128) { ); } + +/// Publish a single event for the atomic `top_up_and_extend` transition. +/// +/// A dedicated event prevents indexers from having to correlate independent +/// top-up and duration-extension records from the same transaction. +pub fn topped_up_and_extended( + env: &Env, + caller: &Address, + old_balance: i128, + new_balance: i128, + old_end_time: u64, + new_end_time: u64, +) { + assert_non_negative_amount(env, old_balance); + assert_non_negative_amount(env, new_balance); + + let sequence = next_sequence(env); + env.events().publish( + (symbol_short!("top_ext"), caller.clone(), sequence), + (old_balance, new_balance, old_end_time, new_end_time), + ); +} + +/// Publish a prospective stream-rate change. +pub fn rate_changed(env: &Env, caller: &Address, old_rate: i128, new_rate: i128) { + assert_non_negative_amount(env, old_rate); + assert_non_negative_amount(env, new_rate); + + let sequence = next_sequence(env); + env.events().publish( + (symbol_short!("rate_chg"), caller.clone(), sequence), + (old_rate, new_rate), + ); +} + pub fn clawback(env: &Env, sender: &Address, amount: i128) { assert_non_negative_amount(env, amount); From 32e508addbc7f45f78f1c70de30b61cf5d16477d Mon Sep 17 00:00:00 2001 From: Oyebisi Ayobami Date: Tue, 29 Sep 2026 11:52:25 +0100 Subject: [PATCH 04/10] feat(stream): support prospective rate changes and combined extension events --- contracts/stream/src/lib.rs | 128 +++++++++++++++++++++++++++++++++++- 1 file changed, 126 insertions(+), 2 deletions(-) diff --git a/contracts/stream/src/lib.rs b/contracts/stream/src/lib.rs index fcdff6fa..ca08e95c 100644 --- a/contracts/stream/src/lib.rs +++ b/contracts/stream/src/lib.rs @@ -14,7 +14,7 @@ use soroban_sdk::{contract, contractimpl, panic_with_error, token, Address, Env} use drip_common::{is_zero_address, pause}; pub use errors::Error; -use storage::{DataKey, StreamInfo, FLAG_CLAWBACK_ENABLED, FLAG_PAUSED}; +use storage::{DataKey, RateCheckpoint, StreamInfo, FLAG_CLAWBACK_ENABLED, FLAG_PAUSED}; pub use storage::{CliffConfig, SplitConfig, StreamConfig, StreamStatus, StreamSummary}; #[contract] @@ -458,6 +458,26 @@ impl DripStream { .checked_sub(info.paused_at) .ok_or(Error::ArithmeticOverflow)?; + // If a rate change has already become active, move its accrual + // checkpoint forward by the paused duration just like start/end time. + // A checkpoint that is still in the future (for example at a pending + // cliff) is left untouched, matching the existing cliff semantics. + if let Some(mut checkpoint) = env + .storage() + .instance() + .get::<_, RateCheckpoint>(&DataKey::RateCheckpoint) + { + if checkpoint.at <= info.paused_at { + checkpoint.at = checkpoint + .at + .checked_add(paused_duration) + .ok_or(Error::ArithmeticOverflow)?; + env.storage() + .instance() + .set(&DataKey::RateCheckpoint, &checkpoint); + } + } + // A stream that stays paused beyond the protocol's safe grace window is // at risk of instance-storage archival; reject the resume before the host // turns this into the opaque "entry archived" error path. The same @@ -493,6 +513,102 @@ impl DripStream { Ok(()) } + /// Sender or delegated operator changes the stream's payout rate. + /// + /// The new rate is prospective: value accrued under the old rate is + /// checkpointed before the rate is updated, so already-earned value is + /// never retroactively repriced. The caller may change the rate on a + /// pending, active, or paused stream; a completed/cancelled stream is + /// rejected. + pub fn change_rate( + env: Env, + caller: Address, + new_rate_per_second: i128, + ) -> Result<(), Error> { + state::with_guard(&env, |env| { + Self::_change_rate(env, &caller, new_rate_per_second) + }) + } + + fn _change_rate( + env: &Env, + caller: &Address, + new_rate_per_second: i128, + ) -> Result<(), Error> { + if new_rate_per_second <= 0 { + return Err(Error::InvalidAmount); + } + + let mut info = state::load(env); + state::assert_not_cancelled(&info)?; + require_sender_or_operator(env, caller, &info.sender)?; + + let now = env.ledger().timestamp(); + if info.end_time > 0 && now >= info.end_time { + return Err(Error::StreamEnded); + } + + // Compute the value earned under the old rate before changing it. + let accrued = math::streamed_amount(env, &info)?; + + // If currently paused, the new rate starts when the stream resumes, + // not while wall-clock time is frozen at paused_at. + let effective_change_time = if info.is_paused() { + info.paused_at + } else { + now + }; + + // Before start/cliff, the checkpoint lives at the first timestamp at + // which value can accrue. For a cliffed stream the upfront unlock is + // included exactly at the cliff boundary. + let cliff_cfg: Option = + env.storage().instance().get(&DataKey::CliffConfig); + let (checkpoint_at, checkpoint_accrued) = if let Some(cfg) = cliff_cfg { + if effective_change_time < cfg.cliff_time { + (cfg.cliff_time, cfg.cliff_unlock_amount) + } else { + (effective_change_time, accrued) + } + } else if effective_change_time < info.start_time { + (info.start_time, 0) + } else { + (effective_change_time, accrued) + }; + + // Re-run the bounded-stream overflow invariant with the new rate. + // This is the same class of guard initialize performs, but includes + // any value already accrued before this rate change. + if info.end_time > 0 && checkpoint_at < info.end_time { + let remaining = info + .end_time + .checked_sub(checkpoint_at) + .ok_or(Error::ArithmeticOverflow)? as i128; + let future = new_rate_per_second + .checked_mul(remaining) + .ok_or(Error::ArithmeticOverflow)?; + checkpoint_accrued + .checked_add(future) + .ok_or(Error::ArithmeticOverflow)?; + } + + ttl::bump(env); + + let old_rate = info.rate_per_second; + info.rate_per_second = new_rate_per_second; + state::save(env, &info); + env.storage().instance().set( + &DataKey::RateCheckpoint, + &RateCheckpoint { + accrued: checkpoint_accrued, + at: checkpoint_at, + }, + ); + + events::rate_changed(env, caller, old_rate, new_rate_per_second); + Ok(()) + } + /// Sender or delegated operator deposits additional tokens into the stream. /// /// Auth is checked immediately after the minimal state load needed to @@ -672,6 +788,7 @@ impl DripStream { let tk = token::Client::new(env, &info.token); let contract_addr = env.current_contract_address(); + let old_balance = tk.balance(&contract_addr); // Transfer funds from the caller (sender or operator) into the // contract. See `_top_up` for why this isn't always `info.sender`. @@ -688,7 +805,14 @@ impl DripStream { state::save(env, &updated); let new_balance = tk.balance(&contract_addr); - events::topped_up(env, caller, amount, new_balance); + events::topped_up_and_extended( + env, + caller, + old_balance, + new_balance, + info.end_time, + new_end_time, + ); Ok(()) } From 007de91cd17f7834920d3d1aba62655d7fdd7c32 Mon Sep 17 00:00:00 2001 From: Oyebisi Ayobami Date: Tue, 29 Sep 2026 11:52:54 +0100 Subject: [PATCH 05/10] test(stream): cover rate changes and combined top-up extension event --- contracts/stream/src/tests.rs | 85 ++++++++++++++++++++++++++++++++++- 1 file changed, 84 insertions(+), 1 deletion(-) diff --git a/contracts/stream/src/tests.rs b/contracts/stream/src/tests.rs index b1bf629a..c2178c41 100644 --- a/contracts/stream/src/tests.rs +++ b/contracts/stream/src/tests.rs @@ -1288,6 +1288,68 @@ fn cancelled_flag_is_durable_across_invocations() { assert_eq!(s.client.streamed_total(), 0); } +// ── Issue #619: mutable prospective stream rate ────────────────────────────── + +#[test] +fn change_rate_preserves_accrued_value_and_applies_new_rate_prospectively() { + let s = Setup::new(100, 3_600, false); + s.advance_secs(100); + assert_eq!(s.client.streamed_total(), 10_000); + + s.client.change_rate(&s.sender, &200); + assert_eq!(s.client.info().rate_per_second, 200); + assert_eq!( + s.client.streamed_total(), + 10_000, + "changing the rate must not reprice value already accrued" + ); + + s.advance_secs(50); + assert_eq!(s.client.streamed_total(), 20_000); +} + +#[test] +fn change_rate_rejects_invalid_unauthorized_and_overflowing_rates() { + let s = Setup::new(1, 2, false); + + assert_eq!( + s.client.try_change_rate(&s.sender, &0), + Err(Ok(Error::InvalidAmount)) + ); + + let stranger = Address::generate(&s.env); + assert_eq!( + s.client.try_change_rate(&stranger, &2), + Err(Ok(Error::NotAuthorized)) + ); + + assert_eq!( + s.client.try_change_rate(&s.sender, &i128::MAX), + Err(Ok(Error::ArithmeticOverflow)) + ); +} + +#[test] +fn change_rate_emits_rate_changed_event() { + let s = Setup::new(100, 3_600, false); + s.client.change_rate(&s.sender, &250); + + assert_eq!(s.client.event_sequence(), 2); + let all_events = s.env.events().all(); + let stream_events: std::vec::Vec<_> = all_events + .iter() + .filter(|(contract, _, _)| contract == &s.client.address) + .collect(); + let last = stream_events.last().unwrap(); + + assert_eq!( + last.1, + (symbol_short!("rate_chg"), s.sender.clone(), 2_u64).into_val(&s.env) + ); + let data: (i128, i128) = last.2.try_into_val(&s.env).unwrap(); + assert_eq!(data, (100, 250)); +} + // ── Issue #205: top_up_and_extend convenience ──────────────────────────────── #[test] @@ -1303,7 +1365,28 @@ fn top_up_and_extend_updates_balance_and_end_time() { s.client.top_up_and_extend(&s.sender, &20_000, &200); assert_eq!(s.client.info().end_time, before_end + 200); - assert_eq!(s.token.balance(&s.client.address), contract_before + 20_000); + let contract_after = s.token.balance(&s.client.address); + assert_eq!(contract_after, contract_before + 20_000); + + // The atomic operation emits one dedicated event rather than a generic + // top-up event plus a separate duration event that indexers must correlate. + assert_eq!(s.client.event_sequence(), 2); + let all_events = s.env.events().all(); + let stream_events: std::vec::Vec<_> = all_events + .iter() + .filter(|(contract, _, _)| contract == &s.client.address) + .collect(); + assert_eq!(stream_events.len(), 2); + let last = stream_events.last().unwrap(); + assert_eq!( + last.1, + (symbol_short!("top_ext"), s.sender.clone(), 2_u64).into_val(&s.env) + ); + let data: (i128, i128, u64, u64) = last.2.try_into_val(&s.env).unwrap(); + assert_eq!( + data, + (contract_before, contract_after, before_end, before_end + 200) + ); } #[test] From 4ed89c592a9b70a5aa4587e3fcb9f4888e6bd0eb Mon Sep 17 00:00:00 2001 From: Oyebisi Ayobami Date: Tue, 29 Sep 2026 11:54:03 +0100 Subject: [PATCH 06/10] feat(oracle): aggregate pair prices across fresh reporters --- contracts/oracle/src/lib.rs | 385 +++++++++++++++++++++++++++++------- 1 file changed, 315 insertions(+), 70 deletions(-) diff --git a/contracts/oracle/src/lib.rs b/contracts/oracle/src/lib.rs index 1439cd4f..85e52666 100644 --- a/contracts/oracle/src/lib.rs +++ b/contracts/oracle/src/lib.rs @@ -21,6 +21,11 @@ pub const PRICE_PRECISION: u128 = 100_000_000; const MAX_SUBMITTERS: u32 = 32; +/// Maximum number of reporters retained for any one pair feed. +const MAX_PAIR_REPORTERS: u32 = 5; +/// Pair reports older than 15 minutes never participate in the median. +const PAIR_REPORT_FRESHNESS_SECS: u64 = 15 * 60; + /// Extends the instance storage TTL so the oracle's entries /// (`DataKey::Admin`, `DataKey::Config`, `DataKey::Price`, etc.) never silently /// archive during idle periods. Called from every state-mutating entry point. @@ -98,10 +103,20 @@ pub enum DataKey { /// `DripGovernor::DataKey::RoleMembers`, closing the gap noted in the /// off-chain tooling audit. RoleMembers(Role), - /// Direct pair price by token Address (base, quote). + /// Legacy/latest direct pair price by token Address (base, quote). + /// New reads aggregate timestamped per-reporter submissions, but this + /// scalar remains as an upgrade fallback for pre-aggregation deployments. PairPrice(Address, Address), - /// Direct pair price by Symbol (base, quote). + /// Timestamped pair report keyed by (base, quote, reporter). + PairSubmission(Address, Address, Address), + /// Bounded reporter index for one Address pair. + PairReporters(Address, Address), + /// Legacy/latest direct pair price by Symbol (base, quote). PairPriceSymbol(Symbol, Symbol), + /// Timestamped symbol-pair report keyed by (base, quote, reporter). + PairSymbolSubmission(Symbol, Symbol, Address), + /// Bounded reporter index for one Symbol pair. + PairSymbolReporters(Symbol, Symbol), } /// Configuration parameters for the TWAP oracle. @@ -179,6 +194,13 @@ pub struct PriceData { pub updated_at: u64, } +#[contracttype] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct PairPriceData { + pub price: u128, + pub updated_at: u64, +} + /// Aggregated price health, computed from the same per-feeder submission set /// that [`get_twap_price`](TwapOracle::get_twap_price) aggregates, so these /// metrics can never disagree with the median it returns. @@ -266,6 +288,8 @@ pub enum Error { InvalidWasmHash = 1021, /// Price feed not found for requested pair. PriceNotFound = 1022, + /// A pair already has the maximum five active reporters. + TooManyPairReporters = 1023, } #[contract] @@ -818,10 +842,11 @@ impl TwapOracle { Ok(value as u64) } - /// Sets the price for a token pair (base/quote) identified by `Address`. + /// Submit a price report for a token pair (base/quote) identified by `Address`. /// - /// Requires that caller is an authorized PriceFeeder or Admin, and that - /// the oracle is not paused. Rejects zero price with [`Error::InvalidPrice`]. + /// Up to five authorized reporters are tracked per pair. `get_price` + /// returns the median of reports no older than 15 minutes, preventing one + /// reporter from unilaterally setting the pair price. pub fn set_price( env: Env, caller: Address, @@ -834,59 +859,72 @@ impl TwapOracle { if price == 0 { return Err(Error::InvalidPrice); } + bump_instance(&env); + add_pair_reporter(&env, &base, "e, &caller)?; + + let data = PairPriceData { + price, + updated_at: env.ledger().timestamp(), + }; + let submission_key = + DataKey::PairSubmission(base.clone(), quote.clone(), caller.clone()); + env.storage().persistent().set(&submission_key, &data); + ttl::bump_persistent(&env, &submission_key); + + // Keep the historical scalar populated for upgrade compatibility. New + // reads prefer the timestamped reporter set above. env.storage() .instance() .set(&DataKey::PairPrice(base, quote), &price); Ok(()) } - /// Fetches the direct price for a token pair (base/quote) identified by `Address`. + /// Returns the decentralized median price for an Address pair. /// - /// Returns [`Error::PriceNotFound`] if no direct feed exists. + /// Only reports submitted within the last 15 minutes by reporters that + /// remain authorized are included. If a contract predates pair-report + /// aggregation and has no reporter index yet, the legacy scalar is + /// returned as a migration fallback. pub fn get_price(env: Env, base: Address, quote: Address) -> Result { - env.storage() - .instance() - .get(&DataKey::PairPrice(base, quote)) - .ok_or(Error::PriceNotFound) + aggregate_pair_price(&env, &base, "e) } - /// Fetches the inverted price for a token pair (base/quote) identified by `Address`. - /// - /// Returns the direct feed price if found; otherwise, if the reciprocal - /// feed (quote/base) is found, calculates and returns: - /// `(PRICE_PRECISION^2 / price)` - /// Returns [`Error::PriceNotFound`] if neither feed exists. + /// Fetches the direct pair median, or the reciprocal median for the reverse pair. pub fn get_price_inverted(env: Env, base: Address, quote: Address) -> Result { - if let Some(price) = env - .storage() - .instance() - .get::<_, u128>(&DataKey::PairPrice(base.clone(), quote.clone())) - { - return Ok(price); - } + let direct_err = match aggregate_pair_price(&env, &base, "e) { + Ok(price) => return Ok(price), + Err(err @ Error::PriceNotFound) | Err(err @ Error::OracleStalePrice) => err, + Err(err) => return Err(err), + }; - if let Some(price) = env - .storage() - .instance() - .get::<_, u128>(&DataKey::PairPrice(quote, base)) - { - if price == 0 { - return Err(Error::InvalidPrice); + match aggregate_pair_price(&env, "e, &base) { + Ok(price) => { + if price == 0 { + return Err(Error::InvalidPrice); + } + let precision_sq = PRICE_PRECISION + .checked_mul(PRICE_PRECISION) + .ok_or(Error::ArithmeticOverflow)?; + precision_sq + .checked_div(price) + .ok_or(Error::ArithmeticOverflow) + } + Err(reverse_err) => { + if direct_err == Error::OracleStalePrice + || reverse_err == Error::OracleStalePrice + { + Err(Error::OracleStalePrice) + } else { + Err(Error::PriceNotFound) + } } - let precision_sq = PRICE_PRECISION - .checked_mul(PRICE_PRECISION) - .ok_or(Error::ArithmeticOverflow)?; - let inverted = precision_sq - .checked_div(price) - .ok_or(Error::ArithmeticOverflow)?; - return Ok(inverted); } - - Err(Error::PriceNotFound) } - /// Sets the price for a pair identified by `Symbol` (e.g. XLM/USDC). + /// Submit a price report for a symbol pair (e.g. XLM/USDC). + /// + /// Uses the same five-reporter / 15-minute median policy as Address pairs. pub fn set_price_symbol( env: Env, caller: Address, @@ -899,49 +937,60 @@ impl TwapOracle { if price == 0 { return Err(Error::InvalidPrice); } + bump_instance(&env); + add_symbol_pair_reporter(&env, &base, "e, &caller)?; + + let data = PairPriceData { + price, + updated_at: env.ledger().timestamp(), + }; + let submission_key = + DataKey::PairSymbolSubmission(base.clone(), quote.clone(), caller.clone()); + env.storage().persistent().set(&submission_key, &data); + ttl::bump_persistent(&env, &submission_key); + env.storage() .instance() .set(&DataKey::PairPriceSymbol(base, quote), &price); Ok(()) } - /// Fetches the direct price for a symbol pair. + /// Returns the decentralized median price for a symbol pair. pub fn get_price_symbol(env: Env, base: Symbol, quote: Symbol) -> Result { - env.storage() - .instance() - .get(&DataKey::PairPriceSymbol(base, quote)) - .ok_or(Error::PriceNotFound) + aggregate_symbol_pair_price(&env, &base, "e) } - /// Fetches the inverted price for a symbol pair. + /// Fetches the direct symbol-pair median, or the reciprocal reverse-pair median. pub fn get_price_inverted_symbol(env: Env, base: Symbol, quote: Symbol) -> Result { - if let Some(price) = env - .storage() - .instance() - .get::<_, u128>(&DataKey::PairPriceSymbol(base.clone(), quote.clone())) - { - return Ok(price); - } + let direct_err = match aggregate_symbol_pair_price(&env, &base, "e) { + Ok(price) => return Ok(price), + Err(err @ Error::PriceNotFound) | Err(err @ Error::OracleStalePrice) => err, + Err(err) => return Err(err), + }; - if let Some(price) = env - .storage() - .instance() - .get::<_, u128>(&DataKey::PairPriceSymbol(quote, base)) - { - if price == 0 { - return Err(Error::InvalidPrice); + match aggregate_symbol_pair_price(&env, "e, &base) { + Ok(price) => { + if price == 0 { + return Err(Error::InvalidPrice); + } + let precision_sq = PRICE_PRECISION + .checked_mul(PRICE_PRECISION) + .ok_or(Error::ArithmeticOverflow)?; + precision_sq + .checked_div(price) + .ok_or(Error::ArithmeticOverflow) + } + Err(reverse_err) => { + if direct_err == Error::OracleStalePrice + || reverse_err == Error::OracleStalePrice + { + Err(Error::OracleStalePrice) + } else { + Err(Error::PriceNotFound) + } } - let precision_sq = PRICE_PRECISION - .checked_mul(PRICE_PRECISION) - .ok_or(Error::ArithmeticOverflow)?; - let inverted = precision_sq - .checked_div(price) - .ok_or(Error::ArithmeticOverflow)?; - return Ok(inverted); } - - Err(Error::PriceNotFound) } // ── Emergency pause (Pauser or Admin-gated) ────────────────────────── @@ -1061,6 +1110,202 @@ fn store_oracle_config(env: &Env, caller: &Address, config: OracleConfig) -> Res Ok(()) } + +/// True while a reporter may contribute to pair-price aggregation. +fn is_authorized_pair_reporter(env: &Env, reporter: &Address) -> bool { + has_role(env, Role::PriceFeeder, reporter) || has_role(env, Role::Admin, reporter) +} + +fn add_pair_reporter( + env: &Env, + base: &Address, + quote: &Address, + reporter: &Address, +) -> Result<(), Error> { + let key = DataKey::PairReporters(base.clone(), quote.clone()); + let existing: Vec
= env + .storage() + .persistent() + .get(&key) + .unwrap_or(Vec::new(env)); + + let mut active = Vec::new(env); + let mut already_present = false; + for account in existing.iter() { + if is_authorized_pair_reporter(env, &account) { + if account == *reporter { + already_present = true; + } + active.push_back(account); + } + } + + if !already_present { + if active.len() >= MAX_PAIR_REPORTERS { + return Err(Error::TooManyPairReporters); + } + active.push_back(reporter.clone()); + } + + env.storage().persistent().set(&key, &active); + ttl::bump_persistent(env, &key); + Ok(()) +} + +fn add_symbol_pair_reporter( + env: &Env, + base: &Symbol, + quote: &Symbol, + reporter: &Address, +) -> Result<(), Error> { + let key = DataKey::PairSymbolReporters(base.clone(), quote.clone()); + let existing: Vec
= env + .storage() + .persistent() + .get(&key) + .unwrap_or(Vec::new(env)); + + let mut active = Vec::new(env); + let mut already_present = false; + for account in existing.iter() { + if is_authorized_pair_reporter(env, &account) { + if account == *reporter { + already_present = true; + } + active.push_back(account); + } + } + + if !already_present { + if active.len() >= MAX_PAIR_REPORTERS { + return Err(Error::TooManyPairReporters); + } + active.push_back(reporter.clone()); + } + + env.storage().persistent().set(&key, &active); + ttl::bump_persistent(env, &key); + Ok(()) +} + +fn aggregate_pair_price(env: &Env, base: &Address, quote: &Address) -> Result { + let reporters_key = DataKey::PairReporters(base.clone(), quote.clone()); + let reporters: Vec
= env + .storage() + .persistent() + .get(&reporters_key) + .unwrap_or(Vec::new(env)); + + // Upgrade compatibility for feeds written before pair aggregation existed. + if reporters.is_empty() { + return env + .storage() + .instance() + .get(&DataKey::PairPrice(base.clone(), quote.clone())) + .ok_or(Error::PriceNotFound); + } + + let now = env.ledger().timestamp(); + let mut fresh = Vec::new(env); + let mut saw_authorized_submission = false; + + for reporter in reporters.iter() { + if !is_authorized_pair_reporter(env, &reporter) { + continue; + } + let key = DataKey::PairSubmission(base.clone(), quote.clone(), reporter); + if let Some(data) = env.storage().persistent().get::<_, PairPriceData>(&key) { + saw_authorized_submission = true; + if now.saturating_sub(data.updated_at) <= PAIR_REPORT_FRESHNESS_SECS { + fresh.push_back(data.price); + } + } + } + + if fresh.is_empty() { + return if saw_authorized_submission { + Err(Error::OracleStalePrice) + } else { + Err(Error::PriceNotFound) + }; + } + + Ok(median_u128(fresh)) +} + +fn aggregate_symbol_pair_price(env: &Env, base: &Symbol, quote: &Symbol) -> Result { + let reporters_key = DataKey::PairSymbolReporters(base.clone(), quote.clone()); + let reporters: Vec
= env + .storage() + .persistent() + .get(&reporters_key) + .unwrap_or(Vec::new(env)); + + if reporters.is_empty() { + return env + .storage() + .instance() + .get(&DataKey::PairPriceSymbol(base.clone(), quote.clone())) + .ok_or(Error::PriceNotFound); + } + + let now = env.ledger().timestamp(); + let mut fresh = Vec::new(env); + let mut saw_authorized_submission = false; + + for reporter in reporters.iter() { + if !is_authorized_pair_reporter(env, &reporter) { + continue; + } + let key = DataKey::PairSymbolSubmission(base.clone(), quote.clone(), reporter); + if let Some(data) = env.storage().persistent().get::<_, PairPriceData>(&key) { + saw_authorized_submission = true; + if now.saturating_sub(data.updated_at) <= PAIR_REPORT_FRESHNESS_SECS { + fresh.push_back(data.price); + } + } + } + + if fresh.is_empty() { + return if saw_authorized_submission { + Err(Error::OracleStalePrice) + } else { + Err(Error::PriceNotFound) + }; + } + + Ok(median_u128(fresh)) +} + +fn median_u128(values: Vec) -> u128 { + let len = values.len(); + let mut sorted = values.clone(); + let mut i: u32 = 1; + while i < len { + let key = sorted.get(i).unwrap(); + let mut j = i; + while j > 0 { + let prev = sorted.get(j - 1).unwrap(); + if prev <= key { + break; + } + sorted.set(j, prev); + j -= 1; + } + sorted.set(j, key); + i += 1; + } + + let mid = len / 2; + if len & 1 == 0 { + let a = sorted.get(mid - 1).unwrap(); + let b = sorted.get(mid).unwrap(); + a.saturating_add(b) / 2 + } else { + sorted.get(mid).unwrap() + } +} + // ── Internal RBAC helpers (delegate to drip_common::rbac) ───────────────── // // These thin wrappers translate the oracle's DataKey / Role types into the From e31a772637ba95939089ff88995157ba47e6318f Mon Sep 17 00:00:00 2001 From: Oyebisi Ayobami Date: Tue, 29 Sep 2026 11:56:00 +0100 Subject: [PATCH 07/10] test(oracle): cover pair median freshness and reporter cap --- contracts/oracle/src/lib.rs | 134 +++++++++++++++++++++++++++++++++++- 1 file changed, 133 insertions(+), 1 deletion(-) diff --git a/contracts/oracle/src/lib.rs b/contracts/oracle/src/lib.rs index 85e52666..c75e5b12 100644 --- a/contracts/oracle/src/lib.rs +++ b/contracts/oracle/src/lib.rs @@ -1300,7 +1300,11 @@ fn median_u128(values: Vec) -> u128 { if len & 1 == 0 { let a = sorted.get(mid - 1).unwrap(); let b = sorted.get(mid).unwrap(); - a.saturating_add(b) / 2 + // Overflow-safe integer average. Avoid `a + b`, which can overflow + // for valid u128 prices even though their median is representable. + (a / 2) + .saturating_add(b / 2) + .saturating_add(((a % 2) + (b % 2)) / 2) } else { sorted.get(mid).unwrap() } @@ -3305,6 +3309,134 @@ mod tests { ); } + // ── Issue #695: decentralized pair-price aggregation ─────────────────── + + #[test] + fn pair_price_uses_median_of_fresh_authorized_reporters() { + let (env, client, admin) = setup(); + client.initialize(&admin); + + let base = Address::generate(&env); + let quote = Address::generate(&env); + let f1 = Address::generate(&env); + let f2 = Address::generate(&env); + let f3 = Address::generate(&env); + + for feeder in [f1.clone(), f2.clone(), f3.clone()] { + client.grant_role(&admin, &Role::PriceFeeder, &feeder); + } + + client.set_price(&f1, &base, "e, &100u128); + client.set_price(&f2, &base, "e, &300u128); + client.set_price(&f3, &base, "e, &200u128); + + assert_eq!(client.get_price(&base, "e), 200u128); + } + + #[test] + fn pair_price_ignores_stale_reports_and_errors_when_all_are_stale() { + let (env, client, admin) = setup(); + client.initialize(&admin); + + let base = Address::generate(&env); + let quote = Address::generate(&env); + let f1 = Address::generate(&env); + let f2 = Address::generate(&env); + client.grant_role(&admin, &Role::PriceFeeder, &f1); + client.grant_role(&admin, &Role::PriceFeeder, &f2); + + env.ledger().set(LedgerInfo { + timestamp: 1_000_000, + protocol_version: 21, + sequence_number: 1, + network_id: Default::default(), + base_reserve: 10, + min_temp_entry_ttl: 16, + min_persistent_entry_ttl: 4096, + max_entry_ttl: 6_312_000, + }); + client.set_price(&f1, &base, "e, &100u128); + + // f1 is now older than the 15-minute pair-report window. + env.ledger().set(LedgerInfo { + timestamp: 1_000_901, + protocol_version: 21, + sequence_number: 2, + network_id: Default::default(), + base_reserve: 10, + min_temp_entry_ttl: 16, + min_persistent_entry_ttl: 4096, + max_entry_ttl: 6_312_000, + }); + client.set_price(&f2, &base, "e, &250u128); + assert_eq!(client.get_price(&base, "e), 250u128); + + env.ledger().set(LedgerInfo { + timestamp: 1_001_802, + protocol_version: 21, + sequence_number: 3, + network_id: Default::default(), + base_reserve: 10, + min_temp_entry_ttl: 16, + min_persistent_entry_ttl: 4096, + max_entry_ttl: 6_312_000, + }); + assert_eq!( + client.try_get_price(&base, "e), + Err(Ok(Error::OracleStalePrice)) + ); + } + + #[test] + fn pair_price_caps_active_reporters_at_five() { + let (env, client, admin) = setup(); + client.initialize(&admin); + + let base = Address::generate(&env); + let quote = Address::generate(&env); + let mut feeders = std::vec::Vec::new(); + + for _ in 0..6 { + let feeder = Address::generate(&env); + client.grant_role(&admin, &Role::PriceFeeder, &feeder); + feeders.push(feeder); + } + + for feeder in feeders.iter().take(5) { + client.set_price(feeder, &base, "e, &100u128); + } + + assert_eq!( + client.try_set_price(&feeders[5], &base, "e, &100u128), + Err(Ok(Error::TooManyPairReporters)) + ); + } + + #[test] + fn symbol_pair_price_uses_the_same_median_policy() { + let (env, client, admin) = setup(); + client.initialize(&admin); + + let base = symbol_short!("XLM"); + let quote = symbol_short!("USDC"); + let f1 = Address::generate(&env); + let f2 = Address::generate(&env); + let f3 = Address::generate(&env); + + for feeder in [f1.clone(), f2.clone(), f3.clone()] { + client.grant_role(&admin, &Role::PriceFeeder, &feeder); + } + + client.set_price_symbol(&f1, &base, "e, &20_000_000u128); + client.set_price_symbol(&f2, &base, "e, &21_000_000u128); + client.set_price_symbol(&f3, &base, "e, &19_000_000u128); + + assert_eq!( + client.get_price_symbol(&base, "e), + 20_000_000u128 + ); + } + // ── Pair price and inverted price calculation tests ─────────────────── #[test] From e040e0261324eea573da944a14ed43b9011ca0b4 Mon Sep 17 00:00:00 2001 From: Oyebisi Ayobami Date: Tue, 29 Sep 2026 11:56:38 +0100 Subject: [PATCH 08/10] docs: document rate changes and atomic extension event --- README.md | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index d78dfd93..dafa009c 100644 --- a/README.md +++ b/README.md @@ -75,6 +75,7 @@ fn cancel(env: Env, caller: Address) -> Result<(), Error> fn pause(env: Env, caller: Address) -> Result<(), Error> fn resume(env: Env, caller: Address) -> Result<(), Error> fn top_up(env: Env, caller: Address, amount: i128) -> Result<(), Error> +fn change_rate(env: Env, caller: Address, new_rate_per_second: i128) -> Result<(), Error> fn clawback(env: Env, caller: Address) -> Result // rejected while paused; resume() first // Extend end_time by extra_time_seconds, pulling the exact rate-implied deposit from the sender @@ -109,7 +110,7 @@ fn storage_version(env: Env) -> u32 ``` > **Not yet in the SDK.** `force_cancel`, `transfer_recipient`, `streamed_total`, `extend_duration`, -> `top_up_and_extend`, `set_operator`, `revoke_operator`, `operator`, `event_sequence`, and +> `top_up_and_extend`, `change_rate`, `set_operator`, `revoke_operator`, `operator`, `event_sequence`, and > `storage_version` exist in the contract but aren't wrapped by `conduit-sdk` yet — callers need > to invoke them directly until the SDK catches up. @@ -123,6 +124,7 @@ A sender can delegate a subset of sender-level actions to another address via `s | `resume` | sender **or** operator | | `cancel` | sender **or** operator | | `top_up` | sender **or** operator (funds come from the caller) | +| `change_rate` | sender **or** operator; new rate applies prospectively without repricing accrued value | | `extend_duration` | sender **or** operator (funds come from the caller) | | `top_up_and_extend` | sender **or** operator (funds come from the caller) | | `clawback` | sender **or** operator | @@ -140,6 +142,8 @@ A sender can delegate a subset of sender-level actions to another address via `s | `stream_paused` | `[sender]` | `{ paused_at, withdrawable }` | | `stream_resumed` | `[sender]` | `{ resumed_at }` | | `stream_topped_up` | `[sender]` | `{ amount, new_balance }` | +| `top_ext` | `[caller, sequence]` | `{ old_balance, new_balance, old_end_time, new_end_time }` | +| `rate_chg` | `[caller, sequence]` | `{ old_rate, new_rate }` | | `stream_clawback` | `[sender]` | `{ amount }` | | `xfer_rec` | `[old_recipient]` | `new_recipient` | From 280af4f4a501d24837b43c9946f882c318c9fa90 Mon Sep 17 00:00:00 2001 From: Ayobami Date: Tue, 29 Sep 2026 12:08:20 +0100 Subject: [PATCH 09/10] style: format stream and oracle changes --- contracts/oracle/src/lib.rs | 16 ++++------------ contracts/stream/src/events.rs | 2 -- contracts/stream/src/lib.rs | 18 ++++-------------- contracts/stream/src/math.rs | 10 ++++++---- contracts/stream/src/tests.rs | 17 +++++++++++------ 5 files changed, 25 insertions(+), 38 deletions(-) diff --git a/contracts/oracle/src/lib.rs b/contracts/oracle/src/lib.rs index c75e5b12..af30e1ed 100644 --- a/contracts/oracle/src/lib.rs +++ b/contracts/oracle/src/lib.rs @@ -867,8 +867,7 @@ impl TwapOracle { price, updated_at: env.ledger().timestamp(), }; - let submission_key = - DataKey::PairSubmission(base.clone(), quote.clone(), caller.clone()); + let submission_key = DataKey::PairSubmission(base.clone(), quote.clone(), caller.clone()); env.storage().persistent().set(&submission_key, &data); ttl::bump_persistent(&env, &submission_key); @@ -911,9 +910,7 @@ impl TwapOracle { .ok_or(Error::ArithmeticOverflow) } Err(reverse_err) => { - if direct_err == Error::OracleStalePrice - || reverse_err == Error::OracleStalePrice - { + if direct_err == Error::OracleStalePrice || reverse_err == Error::OracleStalePrice { Err(Error::OracleStalePrice) } else { Err(Error::PriceNotFound) @@ -982,9 +979,7 @@ impl TwapOracle { .ok_or(Error::ArithmeticOverflow) } Err(reverse_err) => { - if direct_err == Error::OracleStalePrice - || reverse_err == Error::OracleStalePrice - { + if direct_err == Error::OracleStalePrice || reverse_err == Error::OracleStalePrice { Err(Error::OracleStalePrice) } else { Err(Error::PriceNotFound) @@ -3431,10 +3426,7 @@ mod tests { client.set_price_symbol(&f2, &base, "e, &21_000_000u128); client.set_price_symbol(&f3, &base, "e, &19_000_000u128); - assert_eq!( - client.get_price_symbol(&base, "e), - 20_000_000u128 - ); + assert_eq!(client.get_price_symbol(&base, "e), 20_000_000u128); } // ── Pair price and inverted price calculation tests ─────────────────── diff --git a/contracts/stream/src/events.rs b/contracts/stream/src/events.rs index ac292291..bb99124d 100644 --- a/contracts/stream/src/events.rs +++ b/contracts/stream/src/events.rs @@ -139,7 +139,6 @@ pub fn topped_up(env: &Env, sender: &Address, amount: i128, new_balance: i128) { ); } - /// Publish a single event for the atomic `top_up_and_extend` transition. /// /// A dedicated event prevents indexers from having to correlate independent @@ -251,4 +250,3 @@ pub fn min_withdrawal_interval_set(env: &Env, sender: &Address, interval_seconds interval_seconds, ); } - diff --git a/contracts/stream/src/lib.rs b/contracts/stream/src/lib.rs index ca08e95c..7d11b4df 100644 --- a/contracts/stream/src/lib.rs +++ b/contracts/stream/src/lib.rs @@ -14,8 +14,8 @@ use soroban_sdk::{contract, contractimpl, panic_with_error, token, Address, Env} use drip_common::{is_zero_address, pause}; pub use errors::Error; -use storage::{DataKey, RateCheckpoint, StreamInfo, FLAG_CLAWBACK_ENABLED, FLAG_PAUSED}; pub use storage::{CliffConfig, SplitConfig, StreamConfig, StreamStatus, StreamSummary}; +use storage::{DataKey, RateCheckpoint, StreamInfo, FLAG_CLAWBACK_ENABLED, FLAG_PAUSED}; #[contract] pub struct DripStream; @@ -520,21 +520,13 @@ impl DripStream { /// never retroactively repriced. The caller may change the rate on a /// pending, active, or paused stream; a completed/cancelled stream is /// rejected. - pub fn change_rate( - env: Env, - caller: Address, - new_rate_per_second: i128, - ) -> Result<(), Error> { + pub fn change_rate(env: Env, caller: Address, new_rate_per_second: i128) -> Result<(), Error> { state::with_guard(&env, |env| { Self::_change_rate(env, &caller, new_rate_per_second) }) } - fn _change_rate( - env: &Env, - caller: &Address, - new_rate_per_second: i128, - ) -> Result<(), Error> { + fn _change_rate(env: &Env, caller: &Address, new_rate_per_second: i128) -> Result<(), Error> { if new_rate_per_second <= 0 { return Err(Error::InvalidAmount); } @@ -562,8 +554,7 @@ impl DripStream { // Before start/cliff, the checkpoint lives at the first timestamp at // which value can accrue. For a cliffed stream the upfront unlock is // included exactly at the cliff boundary. - let cliff_cfg: Option = - env.storage().instance().get(&DataKey::CliffConfig); + let cliff_cfg: Option = env.storage().instance().get(&DataKey::CliffConfig); let (checkpoint_at, checkpoint_accrued) = if let Some(cfg) = cliff_cfg { if effective_change_time < cfg.cliff_time { (cfg.cliff_time, cfg.cliff_unlock_amount) @@ -1395,4 +1386,3 @@ impl DripStream { .unwrap_or(0) } } - diff --git a/contracts/stream/src/math.rs b/contracts/stream/src/math.rs index cc73e130..93bdaee7 100644 --- a/contracts/stream/src/math.rs +++ b/contracts/stream/src/math.rs @@ -1,7 +1,7 @@ use soroban_sdk::Env; use crate::errors::Error; -use crate::storage::{DataKey, CliffConfig, RateCheckpoint, StreamInfo}; +use crate::storage::{CliffConfig, DataKey, RateCheckpoint, StreamInfo}; /// Returns the total tokens that have streamed up to `now`, /// excluding any paused time. Does not account for withdrawals. @@ -48,8 +48,7 @@ pub fn streamed_amount(env: &Env, info: &StreamInfo) -> Result { // retroactively repriced. The checkpoint is optional so pre-upgrade streams // continue to use the original start/cliff calculation until their first // rate change. - let checkpoint: Option = - env.storage().instance().get(&DataKey::RateCheckpoint); + let checkpoint: Option = env.storage().instance().get(&DataKey::RateCheckpoint); if let Some(checkpoint) = checkpoint { if effective_now < checkpoint.at { return Ok(0); @@ -79,7 +78,10 @@ pub fn streamed_amount(env: &Env, info: &StreamInfo) -> Result { let linear = (info.rate_per_second) .checked_mul(elapsed as i128) .ok_or(Error::ArithmeticOverflow)?; - return cfg.cliff_unlock_amount.checked_add(linear).ok_or(Error::ArithmeticOverflow); + return cfg + .cliff_unlock_amount + .checked_add(linear) + .ok_or(Error::ArithmeticOverflow); } let elapsed = effective_now diff --git a/contracts/stream/src/tests.rs b/contracts/stream/src/tests.rs index c2178c41..9c3f93bb 100644 --- a/contracts/stream/src/tests.rs +++ b/contracts/stream/src/tests.rs @@ -1385,7 +1385,12 @@ fn top_up_and_extend_updates_balance_and_end_time() { let data: (i128, i128, u64, u64) = last.2.try_into_val(&s.env).unwrap(); assert_eq!( data, - (contract_before, contract_after, before_end, before_end + 200) + ( + contract_before, + contract_after, + before_end, + before_end + 200 + ) ); } @@ -2576,7 +2581,10 @@ fn test_cliff_unlock_percentage_and_linear_streaming() { let upfront_unlock = 72_000; // 20% of 360,000 // Set cliff - assert!(s.client.try_set_cliff(&s.sender, &cliff_time, &upfront_unlock).is_ok()); + assert!(s + .client + .try_set_cliff(&s.sender, &cliff_time, &upfront_unlock) + .is_ok()); // Before cliff: 0 withdrawable s.advance_secs(500); @@ -2624,16 +2632,13 @@ fn test_set_min_withdrawal_interval_auth() { let s = Setup::new(100, 3600, false); let unauthorized = Address::generate(&s.env); - let res = s - .client - .try_set_min_withdrawal_interval(&unauthorized, &60); + let res = s.client.try_set_min_withdrawal_interval(&unauthorized, &60); assert_eq!( res.err().unwrap().unwrap(), crate::errors::Error::NotAuthorized ); } - // ── Issue #621: validate_withdraw dry-run query ───────────────────────────── #[test] From 84f024a85e42397d0c06cd738b0475dced5b631f Mon Sep 17 00:00:00 2001 From: Ayobami Date: Tue, 29 Sep 2026 23:54:36 +0100 Subject: [PATCH 10/10] fix(oracle): remove redundant quorum cast --- contracts/oracle/src/lib.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/contracts/oracle/src/lib.rs b/contracts/oracle/src/lib.rs index af30e1ed..a3cf1088 100644 --- a/contracts/oracle/src/lib.rs +++ b/contracts/oracle/src/lib.rs @@ -692,7 +692,7 @@ impl TwapOracle { // from fewer fresh feeders than `min_submitters` is not reliable // enough to return. `PriceStatus` exposes both counts so callers can // observe how far below quorum the set is. - if (fresh_prices.len() as u32) < config.min_submitters { + if fresh_prices.len() < config.min_submitters { return Err(Error::InsufficientQuorum); }