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
5 changes: 5 additions & 0 deletions .github/workflows/deny.yml
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,13 @@ jobs:
runs-on: ubuntu-latest
strategy:
matrix:
# Every check configured in deny.toml must appear here, or a section can
# silently stop being enforced: `[bans]` was configured but never listed
# in the matrix, so duplicate-crate / banned-crate regressions reached
# main unblocked (issue #656).
checks:
- advisories
- bans
- licenses
steps:
- uses: actions/checkout@v4
Expand Down
111 changes: 97 additions & 14 deletions contracts/batch-processor/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,11 @@
#[cfg(test)]
mod tests;

use drip_common::is_zero_address;
use soroban_sdk::{contract, contracterror, contractimpl, token, Address, Env, Symbol, Vec};
use drip_common::{is_zero_address, ttl};
use soroban_sdk::{
contract, contracterror, contractimpl, contracttype, symbol_short, token, Address, BytesN, Env,
Symbol, Vec,
};

/// Maximum number of transfers permitted in a single batch.
const MAX_BATCH_SIZE: u32 = 100;
Expand All @@ -20,19 +23,19 @@ const MAX_BATCH_SIZE: u32 = 100;
/// [`BatchTransferProcessor::version`].
const VERSION: u32 = 1;

/// Errors returned by [`BatchTransferProcessor::process_batch`].
/// Instance-storage key space.
///
/// # Validation order
/// The numeric order of the variants **is** the check order: `process_batch`
/// checks `1`, then `2`, then `3`, and so on, and returns the first failure
/// without evaluating the rest. A client pre-checking a batch before
/// submitting it should apply the same sequence (length → size → per-amount →
/// total → token), mirroring the validation list the README documents for
/// `DripFactory::create_stream`. A batch that is both too large and contains
/// a zero amount always reports `BatchTooLarge`, never `InvalidAmount`.
///
/// All five checks run before `funder.require_auth()` and before any token
/// movement.
/// The processor is stateless with respect to transfers — every call recomputes
/// its batch total from the arguments — but it does hold one durable value: the
/// admin that may replace the contract's own WASM (issue #651). Without a
/// stored authority an `upgrade` entry point would have to be permissionless,
/// which would let anyone replace the code that custodies funds in flight.
#[contracttype]
#[derive(Clone)]
pub enum DataKey {
/// Address allowed to call `upgrade`. Set by `initialize`.
Admin,
}
#[contracterror]
#[derive(Copy, Clone, Debug, Eq, PartialEq, PartialOrd, Ord)]
#[repr(u32)]
Expand All @@ -49,13 +52,57 @@ pub enum Error {
/// Checked fifth (last, immediately before auth): `token` is the
/// all-zero Stellar address, which cannot be a SEP-41 token contract.
InvalidToken = 5,
/// The caller is not the admin stored by `initialize`, so it may not
/// replace the contract's WASM.
NotAuthorized = 6,
/// The WASM hash provided to `upgrade` is all zeros (invalid).
InvalidWasmHash = 7,
/// `initialize` was called on a processor that already has an admin.
AlreadyInitialized = 8,
/// `upgrade` was called before `initialize` has set an admin.
NotInitialized = 9,
}

#[contract]
pub struct BatchTransferProcessor;

/// Emitted by `initialize` when the processor's admin is first set.
fn event_initialized(env: &Env, admin: &Address) {
env.events()
.publish((symbol_short!("init"), admin.clone()), admin.clone());
}

/// Emitted by `upgrade` after the processor's own WASM has been replaced.
///
/// Topics: `("upgraded", caller)` — the admin that authorized the swap.
/// Data: `upgraded_at` — the ledger timestamp at which the swap took effect.
fn event_upgraded(env: &Env, caller: &Address, upgraded_at: u64) {
env.events()
.publish((symbol_short!("upgraded"), caller.clone()), upgraded_at);
}

#[contractimpl]
impl BatchTransferProcessor {
/// One-time setup: record the admin allowed to `upgrade` this contract.
///
/// `process_batch` is permissionless and works without initialization; this
/// call only establishes who may replace the implementation. Guarded
/// against re-initialization so a second call cannot hand the upgrade
/// authority to a different address after the fact.
///
/// # Errors
///
/// - `AlreadyInitialized` — an admin is already recorded.
pub fn initialize(env: Env, admin: Address) -> Result<(), Error> {
if env.storage().instance().has(&DataKey::Admin) {
return Err(Error::AlreadyInitialized);
}
ttl::bump_instance(&env);
env.storage().instance().set(&DataKey::Admin, &admin);
event_initialized(&env, &admin);
Ok(())
}

/// Transfer tokens from `funder` to each address in `recipients`.
///
/// # Auth
Expand Down Expand Up @@ -225,4 +272,40 @@ impl BatchTransferProcessor {
pub fn version(_env: Env) -> u32 {
VERSION
}

// ── Self-upgrade (admin-gated) ──────────────────────────────────────────

/// Replace this contract's own WASM bytecode.
///
/// The new WASM must already be uploaded to the ledger (via
/// `stellar contract upload`); only the hash is passed here. Gated on the
/// admin recorded by [`Self::initialize`], mirroring
/// `DripGovernor::upgrade` and `DripFactory::upgrade_self` — which this
/// contract previously lacked entirely (issue #651). The gate matters here
/// for the same reason it does elsewhere: the processor custodies the whole
/// batch total between the inbound pull and the outbound fan-out, so a
/// replaced implementation runs with funds in flight.
///
/// An all-zero hash is rejected: `update_current_contract_wasm` with it
/// would replace the contract with something that cannot transfer.
pub fn upgrade(env: Env, caller: Address, new_wasm_hash: BytesN<32>) -> Result<(), Error> {
let admin: Address = env
.storage()
.instance()
.get(&DataKey::Admin)
.ok_or(Error::NotInitialized)?;
if caller != admin {
return Err(Error::NotAuthorized);
}
caller.require_auth();

if new_wasm_hash == BytesN::from_array(&env, &[0u8; 32]) {
return Err(Error::InvalidWasmHash);
}

ttl::bump_instance(&env);
env.deployer().update_current_contract_wasm(new_wasm_hash);
event_upgraded(&env, &caller, env.ledger().timestamp());
Ok(())
}
}
2 changes: 2 additions & 0 deletions contracts/common/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,9 @@

//! Shared constants and utilities for the Drip protocol contracts.

pub mod pause;
pub mod rbac;
pub mod ttl;

use soroban_sdk::{Address, Env};

Expand Down
215 changes: 215 additions & 0 deletions contracts/common/src/pause.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,215 @@
//! Shared emergency-pause helpers for the Drip protocol contracts.
//!
//! Every contract in the protocol can be halted by its authority
//! (`pause`/`unpause`/`is_paused`) and every state-mutating entry point is
//! expected to reject the call while halted. The gate itself is the same
//! everywhere — read one instance-storage flag, and refuse to proceed when it
//! is set — but it was previously open-coded per contract, with each copy
//! naming the flag key differently and mapping the failure to its own error
//! variant. A copy that silently loses its guard is a protocol-halt failure
//! mode, so the gate lives here once.
//!
//! Contracts keep their own `Error` enum (every code means something different
//! per contract), so [`PausedError`] is mapped at each call site:
//!
//! ```ignore
//! drip_common::pause::require_not_paused(env, &DataKey::Paused)
//! .map_err(|_| Error::ContractPaused)?;
//! ```
//!
//! The flag is *not* gated here — reading and writing it is the contract's
//! business (only its pause authority may flip it), so [`is_paused`] and
//! [`set_paused`] are low-level storage helpers like the per-contract ones they
//! replace.

use soroban_sdk::Env;

use crate::rbac::StorageKey;

/// Errors the shared pause gate can return.
///
/// Each contract maps this onto its own `Error` variant (typically
/// `ContractPaused`), so no `#[contracterror]` attribute is needed here — the
/// same approach `drip_common::rbac::RbacError` takes.
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
pub enum PausedError {
/// The contract is under an emergency pause, so the call is refused.
ContractPaused,
}

/// Reads the emergency-pause flag stored under `paused_key` in `instance()`
/// storage.
///
/// Defaults to `false` when the key was never written, so a contract deployed
/// before the pause feature existed is treated as running normally rather than
/// halting every call.
pub fn is_paused<K: StorageKey>(env: &Env, paused_key: &K) -> bool {
env.storage().instance().get(paused_key).unwrap_or(false)
}

/// Writes the emergency-pause flag to `instance()` storage.
///
/// Performs no authorization and no state-transition validation: the caller
/// (the contract's own `pause`/`unpause`) owns both, and is responsible for
/// emitting the matching event and bumping TTL.
pub fn set_paused<K: StorageKey>(env: &Env, paused_key: &K, paused: bool) {
env.storage().instance().set(paused_key, &paused);
}

/// The pause gate: `Ok(())` when the contract is running, `Err` while halted.
///
/// Place it at the top of a state-mutating entry point — before validation,
/// storage reads, and TTL payment — so a halted protocol rejects the call
/// immediately and cheaply, and so no state is touched on the rejected path.
pub fn require_not_paused<K: StorageKey>(env: &Env, paused_key: &K) -> Result<(), PausedError> {
require_not_paused_flag(is_paused(env, paused_key))
}

/// [`require_not_paused`] for a flag the caller already has in memory — e.g. a
/// field inside a stored struct rather than a standalone instance key
/// (`DripStream` keeps its pause bit in the consolidated `StreamInfo`).
///
/// Accepting the flag instead of re-reading storage keeps the gate free to use
/// where the contract has already loaded the record it needs to gate on.
pub fn require_not_paused_flag(paused: bool) -> Result<(), PausedError> {
if paused {
Err(PausedError::ContractPaused)
} else {
Ok(())
}
}

#[cfg(test)]
mod tests {
use super::*;
use soroban_sdk::{contract, contractimpl, contracttype};

// ── Contract frame ────────────────────────────────────────────────────────
//
// `pause` is storage-only shared code and owns no `#[contract]` of its own,
// so the suite supplies the frame `soroban_sdk` requires before any
// instance-storage access can run. The host is registered per test purely
// so the frame has instance storage to open.

#[contract]
struct PauseTestHost;

#[contractimpl]
impl PauseTestHost {
/// Never invoked. Present only because registering a contract requires
/// at least one callable entry point.
pub fn noop(_env: Env) {}
}

fn in_contract<T>(env: &Env, f: impl FnOnce() -> T) -> T {
env.as_contract(&env.register_contract(None, PauseTestHost), f)
}

// ── Minimal key space ─────────────────────────────────────────────────────
//
// Mirrors the rbac suite: the helpers are generic over the flag key, so the
// tests supply their own `#[contracttype]` key instead of reaching into a
// consumer contract's `DataKey`.

#[contracttype]
#[derive(Clone, Debug, Eq, PartialEq)]
enum TestKey {
Paused,
StreamAddr(u64),
}

const PAUSED: TestKey = TestKey::Paused;

// ── The flag is readable and writable, defaulting to running ──────────────

#[test]
fn unset_flag_reads_as_not_paused() {
let env = Env::default();

in_contract(&env, || {
// Absent key means "running normally" — a contract deployed before
// the pause feature existed must not be halted by a missing entry.
assert!(!is_paused(&env, &PAUSED));
});
}

#[test]
fn set_paused_round_trips_and_clears() {
let env = Env::default();

in_contract(&env, || {
set_paused(&env, &PAUSED, true);
assert!(is_paused(&env, &PAUSED));

set_paused(&env, &PAUSED, false);
assert!(!is_paused(&env, &PAUSED));
});
}

// ── require_not_paused ────────────────────────────────────────────────────

#[test]
fn require_not_paused_accepts_a_running_contract() {
let env = Env::default();

in_contract(&env, || {
assert_eq!(require_not_paused(&env, &PAUSED), Ok(()));
});
}

#[test]
fn require_not_paused_rejects_a_halted_contract() {
let env = Env::default();

in_contract(&env, || {
set_paused(&env, &PAUSED, true);
assert_eq!(
require_not_paused(&env, &PAUSED),
Err(PausedError::ContractPaused)
);
});
}

/// The two entry points must agree: the key-based gate is a thin wrapper
/// over the flag-based one, and a consumer that already holds the flag
/// (inside a loaded struct) has to reach the same verdict.
#[test]
fn flag_based_gate_matches_the_key_based_gate() {
let env = Env::default();

in_contract(&env, || {
for paused in [false, true] {
set_paused(&env, &PAUSED, paused);
assert_eq!(
require_not_paused(&env, &PAUSED),
require_not_paused_flag(paused)
);
}
});
}

// ── Generic over the flag key ─────────────────────────────────────────────

/// `DripStream` gates on a pause bit held inside its loaded `StreamInfo`
/// rather than a standalone instance key, and the factory/oracle gate on
/// `DataKey::Paused`. Both key shapes must work through the same helper,
/// which is what the generic bound buys.
#[test]
fn gate_is_generic_over_the_flag_key() {
let env = Env::default();
let stream_key = TestKey::StreamAddr(7);

in_contract(&env, || {
assert_eq!(require_not_paused(&env, &stream_key), Ok(()));

set_paused(&env, &stream_key, true);
assert_eq!(
require_not_paused(&env, &stream_key),
Err(PausedError::ContractPaused)
);

// Two distinct keys hold independent flags.
assert!(!is_paused(&env, &PAUSED));
});
}
}
Loading