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
33 changes: 24 additions & 9 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 @@ -61,9 +61,9 @@ use events::{
StreamToppedUpEvent, 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::{
BatchStreamInput, ConditionalMilestone, DataKey, DisputeStatus, OracleAsset, OracleClient,
Expand Down Expand Up @@ -2210,6 +2210,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
25 changes: 18 additions & 7 deletions contracts/stream_contract/src/storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,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 @@ -155,6 +155,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.
///
/// **Key:** `DataKey::Stream(stream_id)` in persistent storage. This is the
Expand Down Expand Up @@ -190,11 +202,7 @@ 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);
}

/// Removes a stream record from persistent storage.
Expand Down Expand Up @@ -231,13 +239,16 @@ pub fn remove_stream(env: &Env, stream_id: u64) {
/// count), or a legacy record failed to decode. Callers that need to tell
/// "no such stream" from "unreadable stream" cannot with this signature.
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
185 changes: 184 additions & 1 deletion contracts/stream_contract/src/test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ use std::string::ToString;

use super::*;
use soroban_sdk::{
testutils::{Address as _, Events, Ledger},
testutils::{storage::Persistent as _, Address as _, Events, Ledger},
token, vec, xdr, Address, Bytes, BytesN, Env, Symbol, TryFromVal, Val, Vec as SorobanVec,
};

Expand Down Expand Up @@ -6315,3 +6315,186 @@ mod emit_helpers {
assert_eq!(payload.refunded_amount, 2_500);
}
}

// ─── Storage TTL Extension & Position Bumping (Issue #1519) ───────────────────

#[test]
fn test_querying_stream_metadata_automatically_extends_ttl() {
let env = Env::default();
env.mock_all_auths();
let (token, _) = create_token(&env);
let client = create_contract(&env);
let sender = Address::generate(&env);
let recipient = Address::generate(&env);
mint(&env, &token, &sender, 10_000);

let id = client.create_stream(&sender, &recipient, &token, &1_000, &1_000);
let contract = client.address.clone();
let key = DataKey::Stream(id);

let initial_ttl = env.as_contract(&contract, || env.storage().persistent().get_ttl(&key));
assert!(initial_ttl >= storage::PERSISTENT_LIFETIME_THRESHOLD);

// Mock ledger sequence increment: advance ledger by 400,000 ledgers
// This simulates infrequent access over a long timeframe (e.g. 4-year vesting)
// reducing the remaining TTL well below PERSISTENT_LIFETIME_THRESHOLD (120,960)
env.ledger().with_mut(|l| {
l.sequence_number += 400_000;
});

let ttl_before_query = env.as_contract(&contract, || env.storage().persistent().get_ttl(&key));
assert!(
ttl_before_query < storage::PERSISTENT_LIFETIME_THRESHOLD,
"TTL before query ({ttl_before_query}) must be below threshold ({})",
storage::PERSISTENT_LIFETIME_THRESHOLD
);

// Querying stream metadata via get_stream automatically triggers bump_position_ttl
let stream = client.get_stream(&id).expect("stream must exist");
assert_eq!(stream.deposited_amount, 1_000);

let ttl_after_query = env.as_contract(&contract, || env.storage().persistent().get_ttl(&key));
assert!(
ttl_after_query >= storage::PERSISTENT_BUMP_AMOUNT - 10,
"TTL after get_stream query ({ttl_after_query}) should be bumped near max ({})",
storage::PERSISTENT_BUMP_AMOUNT
);
assert!(ttl_after_query > ttl_before_query);
}

#[test]
fn test_all_read_only_position_queries_extend_ttl() {
let env = Env::default();
env.mock_all_auths();
let (token, _) = create_token(&env);
let client = create_contract(&env);
let sender = Address::generate(&env);
let recipient = Address::generate(&env);
mint(&env, &token, &sender, 10_000);

let id = client.create_stream(&sender, &recipient, &token, &1_000, &1_000);
let contract = client.address.clone();
let key = DataKey::Stream(id);

let refresh_instance = || {
env.as_contract(&contract, || {
env.storage().instance().extend_ttl(
storage::INSTANCE_LIFETIME_THRESHOLD,
storage::INSTANCE_BUMP_AMOUNT,
);
});
};

// Test get_claimable_amount extends TTL
env.ledger().with_mut(|l| l.sequence_number += 400_000);
refresh_instance();
assert!(
env.as_contract(&contract, || env.storage().persistent().get_ttl(&key))
< storage::PERSISTENT_LIFETIME_THRESHOLD
);
assert!(client.get_claimable_amount(&id).is_some());
let ttl_after_claimable =
env.as_contract(&contract, || env.storage().persistent().get_ttl(&key));
assert!(ttl_after_claimable >= storage::PERSISTENT_BUMP_AMOUNT - 10);

// Test get_vesting_schedule extends TTL
env.ledger().with_mut(|l| l.sequence_number += 400_000);
refresh_instance();
assert!(
env.as_contract(&contract, || env.storage().persistent().get_ttl(&key))
< storage::PERSISTENT_LIFETIME_THRESHOLD
);
assert_eq!(
client.get_vesting_schedule(&id),
Some(VestingSchedule::Linear)
);
let ttl_after_schedule =
env.as_contract(&contract, || env.storage().persistent().get_ttl(&key));
assert!(ttl_after_schedule >= storage::PERSISTENT_BUMP_AMOUNT - 10);

// Test get_projected_end_time extends TTL
env.ledger().with_mut(|l| l.sequence_number += 400_000);
refresh_instance();
assert!(
env.as_contract(&contract, || env.storage().persistent().get_ttl(&key))
< storage::PERSISTENT_LIFETIME_THRESHOLD
);
assert!(client.get_projected_end_time(&id).is_some());
let ttl_after_end_time =
env.as_contract(&contract, || env.storage().persistent().get_ttl(&key));
assert!(ttl_after_end_time >= storage::PERSISTENT_BUMP_AMOUNT - 10);

// Test is_stream_completed extends TTL
env.ledger().with_mut(|l| l.sequence_number += 400_000);
refresh_instance();
assert!(
env.as_contract(&contract, || env.storage().persistent().get_ttl(&key))
< storage::PERSISTENT_LIFETIME_THRESHOLD
);
assert!(!client.is_stream_completed(&id));
let ttl_after_completed =
env.as_contract(&contract, || env.storage().persistent().get_ttl(&key));
assert!(ttl_after_completed >= storage::PERSISTENT_BUMP_AMOUNT - 10);
}

#[test]
fn test_bump_position_ttl_helper_directly() {
let env = Env::default();
env.mock_all_auths();
let (token, _) = create_token(&env);
let client = create_contract(&env);
let sender = Address::generate(&env);
let recipient = Address::generate(&env);
mint(&env, &token, &sender, 10_000);

let id = client.create_stream(&sender, &recipient, &token, &1_000, &1_000);
let contract = client.address.clone();
let key: types::StorageKey = types::DataKey::Stream(id);

let refresh_instance = || {
env.as_contract(&contract, || {
env.storage().instance().extend_ttl(
storage::INSTANCE_LIFETIME_THRESHOLD,
storage::INSTANCE_BUMP_AMOUNT,
);
});
};

// Advance sequence number to drain TTL
env.ledger().with_mut(|l| l.sequence_number += 400_000);
refresh_instance();
let ttl_before = env.as_contract(&contract, || env.storage().persistent().get_ttl(&key));
assert!(ttl_before < storage::PERSISTENT_LIFETIME_THRESHOLD);

// Call bump_position_ttl helper directly
env.as_contract(&contract, || {
storage::bump_position_ttl(&env, &key);
});

let ttl_after = env.as_contract(&contract, || env.storage().persistent().get_ttl(&key));
assert!(ttl_after >= storage::PERSISTENT_BUMP_AMOUNT - 10);

// Calling bump_position_ttl on a non-existent key gracefully does not panic
let non_existent_key = types::DataKey::Stream(999_999);
env.as_contract(&contract, || {
storage::bump_position_ttl(&env, &non_existent_key);
});

// Test explicit contract method bump_stream_ttl
env.ledger().with_mut(|l| l.sequence_number += 400_000);
refresh_instance();
assert!(
env.as_contract(&contract, || env.storage().persistent().get_ttl(&key))
< storage::PERSISTENT_LIFETIME_THRESHOLD
);
client.bump_stream_ttl(&id);
let ttl_after_contract_bump =
env.as_contract(&contract, || env.storage().persistent().get_ttl(&key));
assert!(ttl_after_contract_bump >= storage::PERSISTENT_BUMP_AMOUNT - 10);

// Contract method bump_stream_ttl returns StreamNotFound on non-existent stream
assert_eq!(
client.try_bump_stream_ttl(&999_999),
Err(Ok(StreamError::StreamNotFound))
);
}
3 changes: 3 additions & 0 deletions contracts/stream_contract/src/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -248,6 +248,9 @@ pub struct Stream {
pub is_allowance_based: bool,
}

/// Alias for DataKey representing storage keys in the contract.
pub type StorageKey = DataKey;

/// A single stream to create inside `batch_create_streams`.
///
/// The batch entrypoint takes one struct per stream rather than N parallel
Expand Down
Loading