Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions contracts/factory/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ crate-type = ["cdylib", "rlib"]
soroban-sdk = { workspace = true }
drip-governor = { path = "../governor" }
drip-common = { path = "../common" }
drip-stream = { path = "../stream" }

[dev-dependencies]
soroban-sdk = { workspace = true, features = ["testutils"] }
Expand Down
41 changes: 41 additions & 0 deletions contracts/factory/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -336,6 +336,41 @@ impl DripFactory {
.get(&DataKey::StreamAddr(stream_id))
}

/// Cancel multiple streams in one transaction, all authorized by the
/// same `sender`.
///
/// Mirrors the bulk-creation ergonomics of [`create_batch_streams`](Self::create_batch_streams)
/// on the cancellation side. Each stream address in `stream_addresses`
/// is cancelled via a cross-contract call to `DripStream::cancel`,
/// reusing the per-stream validation, settlement, and event emission.
///
/// Atomicity: Soroban transactions are all-or-nothing at the host
/// level. If any cancellation fails (e.g. stream already cancelled,
/// sender mismatch), the `?` below propagates that error immediately,
/// and every state change already made earlier in this same call is
/// rolled back by the host — no partial-batch state is ever left behind.
pub fn cancel_batch_streams(
env: Env,
sender: Address,
stream_addresses: Vec<Address>,
) -> Result<(), Error> {
sender.require_auth();

if stream_addresses.is_empty() {
return Err(Error::EmptyBatch);
}
if stream_addresses.len() > MAX_BATCH_SIZE {
return Err(Error::BatchTooLarge);
}

for stream_addr in stream_addresses.iter() {
let stream_client = drip_stream::DripStreamClient::new(&env, &stream_addr);
stream_client.cancel(&sender);
}

Ok(())
}

/// Batch-resolve stream IDs to their deployed contract addresses.
///
/// Pairs with `streams_by_sender`/`streams_by_recipient`: a page of IDs
Expand Down Expand Up @@ -370,6 +405,12 @@ impl DripFactory {
query::paginate(&env, all, offset, limit)
}

/// Paginated list of stream IDs where `recipient` is the beneficiary.
///
/// Returns at most `limit` IDs starting at `offset`. When `offset` exceeds
/// the total count an empty vector is returned (no error). `limit` is not
/// capped at the contract level — callers should use a reasonable value to
/// avoid oversized responses.
pub fn streams_by_recipient(env: Env, recipient: Address, offset: u32, limit: u32) -> Vec<u64> {
let all: Vec<u64> = env
.storage()
Expand Down
24 changes: 24 additions & 0 deletions contracts/factory/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -217,4 +217,28 @@ fn bump_persistent_extends_ttl_of_persistent_entry() {
let retrieved: Address = s.env.storage().persistent().get(&key).unwrap();
assert_eq!(retrieved, dummy);
});
}

// ── Issue #204: cancel_batch_streams ─────────────────────────────────────────

#[test]
fn cancel_batch_rejects_empty_list() {
let s = Setup::new();
let sender = Address::generate(&s.env);
let addresses: soroban_sdk::Vec<Address> = soroban_sdk::Vec::new(&s.env);

let result = s.client.try_cancel_batch_streams(&sender, &addresses);
assert_eq!(result, Err(Ok(Error::EmptyBatch)));
}

#[test]
fn cancel_batch_rejects_oversized_list() {
let s = Setup::new();
let sender = Address::generate(&s.env);
let mut addresses: soroban_sdk::Vec<Address> = soroban_sdk::Vec::new(&s.env);
for _ in 0..101 {
addresses.push_back(Address::generate(&s.env));
}
let result = s.client.try_cancel_batch_streams(&sender, &addresses);
assert_eq!(result, Err(Ok(Error::BatchTooLarge)));
}
134 changes: 134 additions & 0 deletions contracts/oracle/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,17 @@ impl TwapOracle {

// ── Reads ────────────────────────────────────────────────────────────

/// Reconfigure the oracle parameters. Admin-gated.
///
/// When `decimals` or `asset_peg` changes relative to the currently stored
/// config, all existing price data (`DataKey::Price`, per-feeder
/// `DataKey::Submission` entries, and the `DataKey::Submitters` list) is
/// cleared. This prevents stale prices submitted under the old config from
/// being silently misinterpreted under the new parameters — the next
/// `get_twap_price` call will return `NoPriceAvailable` until a fresh
/// `submit_price` is made. Changes to `max_staleness` or `oracle_address`
/// alone do not clear price data, as those do not affect price magnitude
/// interpretation.
pub fn configure_oracle(env: Env, caller: Address, config: OracleConfig) -> Result<(), Error> {
require_role_or_admin(&env, &caller, Role::Admin)?;

Expand All @@ -179,6 +190,30 @@ impl TwapOracle {
}

bump_instance(&env);

// Check if decimals or asset_peg changed relative to existing config.
// If so, clear all stored price data to prevent magnitude misinterpretation.
let existing: Option<OracleConfig> = env.storage().instance().get(&DataKey::Config);
if let Some(old) = existing {
if old.decimals != config.decimals || old.asset_peg != config.asset_peg {
// Clear the legacy single-value price slot.
env.storage().instance().remove(&DataKey::Price);

// Clear every per-feeder submission and the submitter list itself.
let submitters: Vec<Address> = env
.storage()
.instance()
.get(&DataKey::Submitters)
.unwrap_or(Vec::new(&env));
for feeder in submitters.iter() {
env.storage()
.instance()
.remove(&DataKey::Submission(feeder));
}
env.storage().instance().remove(&DataKey::Submitters);
}
}

env.storage().instance().set(&DataKey::Config, &config);
events::oracle_configured(&env, &caller, config);
Ok(())
Expand Down Expand Up @@ -1313,4 +1348,103 @@ mod tests {
.as_contract(&client.address, || env.storage().instance().get_ttl());
assert!(ttl >= 100_000, "instance TTL after submit_price: {ttl}");
}

// ── Issue #206: configure_oracle clears stale price on decimals change ────

#[test]
fn configure_oracle_clears_price_when_decimals_change() {
let (env, client, admin) = setup();
client.initialize(&admin);

let oracle_addr = Address::generate(&env);
let config = OracleConfig {
oracle_address: oracle_addr.clone(),
decimals: 8,
asset_peg: 1,
max_staleness: 300,
};
client.configure_oracle(&admin, &config);
client.submit_price(&admin, &50_000_000);

// Price exists and is fresh
let price = client.get_twap_price();
assert_eq!(price, 50_000_000);

// Reconfigure with different decimals
let new_config = OracleConfig {
oracle_address: oracle_addr.clone(),
decimals: 6,
asset_peg: 1,
max_staleness: 300,
};
client.configure_oracle(&admin, &new_config);

// After decimals change, old price data should be cleared
let result = client.try_get_twap_price();
assert_eq!(result, Err(Ok(Error::NoPriceAvailable)));
}

#[test]
fn configure_oracle_clears_price_when_asset_peg_changes() {
let (env, client, admin) = setup();
client.initialize(&admin);

let oracle_addr = Address::generate(&env);
let config = OracleConfig {
oracle_address: oracle_addr.clone(),
decimals: 8,
asset_peg: 1,
max_staleness: 300,
};
client.configure_oracle(&admin, &config);
client.submit_price(&admin, &50_000_000);

let price = client.get_twap_price();
assert_eq!(price, 50_000_000);

// Reconfigure with different asset_peg
let new_config = OracleConfig {
oracle_address: oracle_addr.clone(),
decimals: 8,
asset_peg: 2,
max_staleness: 300,
};
client.configure_oracle(&admin, &new_config);

// After asset_peg change, old price data should be cleared
let result = client.try_get_twap_price();
assert_eq!(result, Err(Ok(Error::NoPriceAvailable)));
}

#[test]
fn configure_oracle_preserves_price_when_only_staleness_changes() {
let (env, client, admin) = setup();
client.initialize(&admin);

let oracle_addr = Address::generate(&env);
let config = OracleConfig {
oracle_address: oracle_addr.clone(),
decimals: 8,
asset_peg: 1,
max_staleness: 300,
};
client.configure_oracle(&admin, &config);
client.submit_price(&admin, &50_000_000);

let price = client.get_twap_price();
assert_eq!(price, 50_000_000);

// Reconfigure with only max_staleness changed
let new_config = OracleConfig {
oracle_address: oracle_addr.clone(),
decimals: 8,
asset_peg: 1,
max_staleness: 600,
};
client.configure_oracle(&admin, &new_config);

// Price should still be available — only staleness window changed
let price_after = client.get_twap_price();
assert_eq!(price_after, 50_000_000);
}
}
67 changes: 67 additions & 0 deletions contracts/stream/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -399,6 +399,73 @@ impl DripStream {
Ok(())
}

/// Sender (or operator) tops up and extends the stream in a single call.
///
/// Combines [`top_up`](Self::top_up) and [`extend_duration`](Self::extend_duration)
/// into one authorized transaction, reducing round-trips and the risk of a
/// sender performing only one half of the pair (which would leave the
/// stream either underfunded for the extended duration or with idle funds
/// past the original `end_time`).
///
/// `amount` is deposited into the stream and `extra_time_seconds` is added
/// to `end_time`. Both must be non-zero. Open-ended streams (`end_time == 0`)
/// cannot be extended — use `top_up` alone instead.
pub fn top_up_and_extend(
env: Env,
caller: Address,
amount: i128,
extra_time_seconds: u64,
) -> Result<(), Error> {
state::with_guard(&env, |env| {
Self::_top_up_and_extend(env, &caller, amount, extra_time_seconds)
})
}

fn _top_up_and_extend(
env: &Env,
caller: &Address,
amount: i128,
extra_time_seconds: u64,
) -> Result<(), Error> {
if amount <= 0 {
return Err(Error::InvalidAmount);
}
if extra_time_seconds == 0 {
return Err(Error::InvalidTimeRange);
}

let info = state::load(env);
require_sender_or_operator(env, caller, &info.sender)?;

ttl::bump(env);
state::assert_not_cancelled(&info)?;

if info.end_time == 0 {
return Err(Error::InvalidTimeRange);
}

let tk = token::Client::new(env, &info.token);
let contract_addr = env.current_contract_address();

// Transfer funds from sender into the contract
tk.transfer(&info.sender, &contract_addr, &amount);

// Update end_time with overflow check
let new_end_time = info
.end_time
.checked_add(extra_time_seconds)
.ok_or(Error::ArithmeticOverflow)?;

let mut updated = info.clone();
updated.end_time = new_end_time;
state::save(env, &updated);

let new_balance = tk.balance(&contract_addr);
events::topped_up(env, caller, amount, new_balance);

Ok(())
}

/// Sender reclaims unstreamed tokens (only if clawback was enabled).
pub fn clawback(env: Env, caller: Address) -> Result<i128, Error> {
state::with_guard(&env, |env| Self::_clawback(env, &caller))
Expand Down
66 changes: 66 additions & 0 deletions contracts/stream/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -985,3 +985,69 @@ fn cancelled_flag_is_durable_across_invocations() {
assert_eq!(s.client.withdrawable(), 0);
assert_eq!(s.client.streamed_total(), 0);
}

// ── Issue #205: top_up_and_extend convenience ────────────────────────────────

#[test]
fn top_up_and_extend_updates_balance_and_end_time() {
let s = Setup::new(100, 3_600, false);
let before_end = s.client.info().end_time;

// Mint exact deposit needed: 100 rate × 200s = 20_000
let token_admin = token::StellarAssetClient::new(&s.env, &s.token.address);
token_admin.mint(&s.sender, &20_000);

let contract_before = s.token.balance(&s.client.address);
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);
}

#[test]
fn top_up_and_extend_rejects_zero_amount() {
let s = Setup::new(100, 3_600, false);
let result = s.client.try_top_up_and_extend(&s.sender, &0, &100);
assert_eq!(result, Err(Ok(Error::InvalidAmount)));
}

#[test]
fn top_up_and_extend_rejects_zero_extra_time() {
let s = Setup::new(100, 3_600, false);
let result = s.client.try_top_up_and_extend(&s.sender, &10_000, &0);
assert_eq!(result, Err(Ok(Error::InvalidTimeRange)));
}

#[test]
fn top_up_and_extend_rejected_on_cancelled_stream() {
let s = Setup::new(100, 3_600, false);
s.client.cancel(&s.sender);

let token_admin = token::StellarAssetClient::new(&s.env, &s.token.address);
token_admin.mint(&s.sender, &10_000);

let result = s.client.try_top_up_and_extend(&s.sender, &10_000, &100);
assert!(result.is_err());
}

#[test]
fn top_up_and_extend_rejected_for_open_ended_stream() {
let env = Env::default();
env.mock_all_auths();

let sender = Address::generate(&env);
let recipient = Address::generate(&env);
let token_admin = Address::generate(&env);
let token_addr = env
.register_stellar_asset_contract_v2(token_admin.clone())
.address();

let now: u64 = 1_000_000;
let stream_id = env.register_contract(None, DripStream);
let client = DripStreamClient::new(&env, &stream_id);

client.initialize(&sender, &recipient, &token_addr, &100, &now, &0, &false);

let result = client.try_top_up_and_extend(&sender, &10_000, &100);
assert_eq!(result, Err(Ok(Error::InvalidTimeRange)));
}