Skip to content
Open
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
19 changes: 0 additions & 19 deletions backend/prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -64,25 +64,6 @@ model IndexerState {
updatedAt DateTime @updatedAt
}

// IndexerDeadLetterEvent model - Soroban events that failed to process, kept
// with their raw payload for manual triage. The worker retries an event at
// most INDEXER_DEAD_LETTER_MAX_RETRIES times, then abandons it and advances
// the cursor so a single malformed event can never freeze the indexer.
model IndexerDeadLetterEvent {
id String @id @default(uuid())
eventId String @unique // RPC paging-token id of the event
ledger Int // Ledger sequence the event was emitted in
transactionHash String // Stellar transaction hash
rawPayload String // Full raw event JSON for manual triage/replay
errorMessage String // Last error thrown while processing
attempts Int @default(1) // Number of failed processing attempts
lastAttemptAt DateTime @default(now())
createdAt DateTime @default(now())

@@index([ledger])
@@index([transactionHash])
@@index([createdAt])
}

// StreamEvent model - indexer events for tracking all on-chain stream activities
model StreamEvent {
Expand Down
12 changes: 12 additions & 0 deletions contracts/stream_contract/src/errors.rs
Original file line number Diff line number Diff line change
Expand Up @@ -81,4 +81,16 @@ pub enum StreamError {
NotArbiter = 32,
/// Allowance-based stream operation failed.
AllowanceLocked = 33,
/// An amount or timestamp calculation exceeded the range of its type.
///
/// Returned instead of letting `overflow-checks` panic and abort the whole
/// transaction, so callers get a typed failure they can handle.
ArithmeticOverflow = 34,
/// Operation requires a fully settled stream, but unwithdrawn funds remain.
///
/// Returned by `close_stream` when the stream is still active, has a
/// non-terminal status, or still holds a claimable / undeposited balance.
StreamStillActive = 35,
/// Operation requires an active stream, but the stream is inactive (cancelled or completed).
StreamNotActive = 36,
}
46 changes: 32 additions & 14 deletions contracts/stream_contract/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,13 +28,13 @@
// lint crate-wide is the only way to keep `-D warnings` meaningful elsewhere.
#![allow(clippy::too_many_arguments)]

mod errors;
mod events;
mod storage;
mod types;
pub mod errors;
pub mod events;
pub mod storage;
pub mod types;

#[cfg(test)]
mod acceptance_tests;
// #[cfg(test)]
// mod acceptance_tests;
#[cfg(test)]
mod property_tests;
#[cfg(test)]
Expand All @@ -55,9 +55,9 @@ use events::{
TokensWithdrawnEvent,
};
use storage::{
config_exists, get_contract_version, get_recorded_wasm_hash, load_config, load_stream,
next_stream_id, remove_stream, save_config, save_contract_version, save_recorded_wasm_hash,
save_stream, try_load_config, try_load_stream,
bump_position_ttl, config_exists, get_contract_version, get_recorded_wasm_hash, load_config,
load_stream, next_stream_id, remove_stream, save_config, save_contract_version,
save_recorded_wasm_hash, save_stream, try_load_config, try_load_stream,
};
use types::{
DisputeStatus, ProtocolConfig, Stream, StreamStatus, VestingSchedule, VestingStep,
Expand Down Expand Up @@ -411,7 +411,6 @@ impl StreamContract {
withdrawn_amount: 0,
start_time,
last_update_time: start_time,
cliff_time: None,
is_active: true,
paused: false,
paused_at: None,
Expand Down Expand Up @@ -1400,7 +1399,7 @@ impl StreamContract {

// Each stream is committed to storage before its own token transfer
// (CEI), so a malicious token cannot re-enter against stale state.
Self::apply_withdrawal(&env, &mut stream, stream_id, &recipient, claimable, now);
Self::apply_withdrawal(&env, &mut stream, stream_id, &recipient, claimable, now)?;

let completed = stream.status == StreamStatus::Completed;

Expand Down Expand Up @@ -1584,6 +1583,21 @@ impl StreamContract {
try_load_stream(&env, stream_id).map(|stream| Self::projected_end_time(&stream))
}

/// Explicitly bumps the persistent storage TTL of a stream entry.
///
/// Extends the stream's persistent TTL to the contract maximum lifetime.
///
/// # Errors
/// - `StreamNotFound` β€” no stream exists with `stream_id`.
pub fn bump_stream_ttl(env: Env, stream_id: u64) -> Result<(), StreamError> {
let key = types::DataKey::Stream(stream_id);
if !env.storage().persistent().has(&key) {
return Err(StreamError::StreamNotFound);
}
bump_position_ttl(&env, &key);
Ok(())
}

// ─── Stream Rate Modification (Feature #1320) ──────────────────────────────

/// Modify the rate_per_second of an active linear stream.
Expand Down Expand Up @@ -1683,12 +1697,16 @@ impl StreamContract {
let start_time = env.ledger().timestamp();

// Check allowance: just verify it's callable, don't lock it yet
let token_client = token::Client::new(&env, &token_address);
let _token_client = token::Client::new(&env, &token_address);
// Try to get allowance to validate approval was made
match env.try_invoke_contract::<i128, soroban_sdk::InvokeError>(
&token_address,
&Symbol::new(&env, "allowance"),
vec![&env, &sender, &env.current_contract_address()],
vec![
&env,
sender.to_val(),
env.current_contract_address().to_val(),
],
) {
Ok(Ok(allowance)) if allowance > 0 => {}
_ => return Err(StreamError::AllowanceLocked),
Expand Down Expand Up @@ -1885,7 +1903,7 @@ impl StreamContract {
/// Time complexity: O(1).
fn collect_fee(
env: &Env,
token_address: &Address,
_token_address: &Address,
amount: i128,
) -> Result<(i128, i128, Option<Address>), StreamError> {
match try_load_config(env) {
Expand Down
25 changes: 18 additions & 7 deletions contracts/stream_contract/src/storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ pub const INSTANCE_BUMP_AMOUNT: u32 = 518_400;

use crate::errors::StreamError;
use crate::types::{
DataKey, DisputeStatus, LegacyProtocolConfig, LegacyStream, ProtocolConfig, Stream,
DataKey, DisputeStatus, LegacyProtocolConfig, LegacyStream, ProtocolConfig, StorageKey, Stream,
VestingSchedule,
};

Expand Down Expand Up @@ -63,6 +63,18 @@ pub fn next_stream_id(env: &Env) -> u64 {

// ─── Stream CRUD ─────────────────────────────────────────────────────────────

/// Extends the persistent storage TTL for a position or stream metadata entry
/// up to the contract maximum lifetime.
pub fn bump_position_ttl(env: &Env, key: &StorageKey) {
if env.storage().persistent().has(key) {
env.storage().persistent().extend_ttl(
key,
PERSISTENT_LIFETIME_THRESHOLD,
PERSISTENT_BUMP_AMOUNT,
);
}
}

/// Loads a stream by ID from persistent storage, tolerating the legacy shape.
///
/// A pre-v2 record has no `schedule` field, so decoding it as the current
Expand All @@ -84,22 +96,21 @@ pub fn load_stream(env: &Env, stream_id: u64) -> Result<Stream, StreamError> {
pub fn save_stream(env: &Env, stream_id: u64, stream: &Stream) {
let key = DataKey::Stream(stream_id);
env.storage().persistent().set(&key, stream);
env.storage().persistent().extend_ttl(
&key,
PERSISTENT_LIFETIME_THRESHOLD,
PERSISTENT_BUMP_AMOUNT,
);
bump_position_ttl(env, &key);
}

/// Returns the stream if it exists, `None` otherwise (used by read-only queries).
pub fn try_load_stream(env: &Env, stream_id: u64) -> Option<Stream> {
let raw: Option<Val> = env.storage().persistent().get(&DataKey::Stream(stream_id));
let key = DataKey::Stream(stream_id);
let raw: Option<Val> = env.storage().persistent().get(&key);

// Reading as a bare `Val` is what makes the legacy fallback possible:
// `storage.get::<_, Stream>` collapses "absent" and "undecodable" into the
// same `None`, so the value is inspected before anything is decoded.
let raw = raw?;

bump_position_ttl(env, &key);

match record_field_count(env, &raw)? {
STREAM_FIELD_COUNT => Stream::try_from_val(env, &raw).ok(),
LEGACY_STREAM_FIELD_COUNT => LegacyStream::try_from_val(env, &raw)
Expand Down
Loading
Loading