11#![ no_std]
22use soroban_sdk:: contractevent;
33use soroban_sdk:: token:: Client as TokenClient ;
4+ use soroban_sdk:: { contract, contractevent, contractimpl, contracttype, vec, Address , Bytes , Env , Vec } ;
45use soroban_sdk:: { contract, contractimpl, contracttype, vec, Address , Env , Vec } ;
56
67use soroban_sdk:: token:: Client as TokenClient ;
@@ -15,6 +16,11 @@ const FREE_TRIAL_DURATION: u64 = 7 * 24 * 60 * 60;
1516#[ contracttype]
1617#[ derive( Clone , Debug , Eq , PartialEq ) ]
1718pub enum DataKey {
19+ Stream ( Address , Address ) , // (subscriber, stream_id)
20+ TotalStreamed ( Address , Address ) , // (subscriber, creator) - cumulative tokens streamed
21+ CliffThreshold ( Address ) , // creator -> threshold amount for access
22+ CreatorSubscribers ( Address ) , // creator -> Vec<subscriber>
23+ CreatorMetadata ( Address ) , // creator -> IPFS CID bytes
1824 Stream ( Address , Address ) , // (subscriber, creator)
1925 TotalStreamed ( Address , Address ) , // (subscriber, creator)
2026 CliffThreshold ( Address ) , // creator -> threshold amount for access
@@ -208,6 +214,21 @@ fn subscribe_internal(
208214 percentages,
209215 } ;
210216
217+ env. storage ( ) . persistent ( ) . set ( & key, & stream) ;
218+
219+ // Track subscriber under this stream_id for withdraw_all
220+ let creator_key = DataKey :: CreatorSubscribers ( stream_id. clone ( ) ) ;
221+ let mut subs: Vec < Address > = env
222+ . storage ( )
223+ . persistent ( )
224+ . get ( & creator_key)
225+ . unwrap_or ( vec ! [ env] ) ;
226+ subs. push_back ( subscriber. clone ( ) ) ;
227+ env. storage ( ) . persistent ( ) . set ( & creator_key, & subs) ;
228+ }
229+
230+ fn collect_internal ( env : & Env , subscriber : & Address , stream_id : & Address ) {
231+ let key = stream_key ( subscriber, stream_id) ;
211232 set_stream ( env, & key, & stream) ;
212233}
213234
@@ -261,6 +282,8 @@ fn distribute_and_collect(
261282 return 0 ;
262283 }
263284
285+ let time_elapsed = ( current_time - stream. last_collected ) as i128 ;
286+ let mut amount_to_collect = time_elapsed * stream. rate_per_second ;
264287 if let Some ( creator) = total_streamed_creator {
265288 if is_creator_paused ( env, creator) {
266289 // While paused, advance accounting clock so paused time is never billed.
@@ -296,6 +319,35 @@ fn distribute_and_collect(
296319 amount_to_collect = stream. balance ;
297320 }
298321
322+ if amount_to_collect > 0 {
323+ let token_client = TokenClient :: new ( env, & stream. token ) ;
324+ let mut remaining = amount_to_collect;
325+ let creators_len = stream. creators . len ( ) ;
326+
327+ for i in 0 ..creators_len {
328+ let creator = stream. creators . get ( i) . unwrap ( ) ;
329+ let payout = if ( i + 1 ) == creators_len {
330+ remaining
331+ } else {
332+ let percentage = stream. percentages . get ( i) . unwrap ( ) as i128 ;
333+ let amount = ( amount_to_collect * percentage) / 100 ;
334+ remaining -= amount;
335+ amount
336+ } ;
337+
338+ if payout > 0 {
339+ token_client. transfer ( & env. current_contract_address ( ) , & creator, & payout) ;
340+ }
341+ }
342+
343+ stream. balance -= amount_to_collect;
344+ stream. last_collected = current_time;
345+ env. storage ( ) . persistent ( ) . set ( & key, & stream) ;
346+
347+ // Update cumulative streamed for each creator
348+ for i in 0 ..stream. creators . len ( ) {
349+ let creator = stream. creators . get ( i) . unwrap ( ) ;
350+ SubStreamContract :: update_total_streamed ( env, subscriber, & creator, amount_to_collect) ;
299351 if amount_to_collect <= 0 {
300352 return 0 ;
301353 }
@@ -352,6 +404,18 @@ fn cancel_group_internal(env: &Env, subscriber: &Address, stream_id: &Address) {
352404 panic ! ( "stream not found" ) ;
353405 }
354406
407+ // Check minimum flow duration
408+ let stream: Stream = env. storage ( ) . persistent ( ) . get ( & key) . unwrap ( ) ;
409+ let current_time = env. ledger ( ) . timestamp ( ) ;
410+ if current_time < stream. start_time + MINIMUM_FLOW_DURATION {
411+ let remaining_time = stream. start_time + MINIMUM_FLOW_DURATION - current_time;
412+ panic ! (
413+ "cannot cancel stream: minimum duration not met. {} seconds remaining" ,
414+ remaining_time
415+ ) ;
416+ }
417+
418+ collect_internal ( env, subscriber, stream_id) ;
355419 distribute_and_collect ( env, subscriber, stream_id, None ) ;
356420
357421 let stream = get_stream ( env, & key) ;
@@ -360,6 +424,23 @@ fn cancel_group_internal(env: &Env, subscriber: &Address, stream_id: &Address) {
360424 token_client. transfer ( & env. current_contract_address ( ) , subscriber, & stream. balance ) ;
361425 }
362426
427+ env. storage ( ) . persistent ( ) . remove ( & key) ;
428+
429+ // Remove subscriber from stream_id's subscriber list
430+ let creator_key = DataKey :: CreatorSubscribers ( stream_id. clone ( ) ) ;
431+ if let Some ( subs) = env
432+ . storage ( )
433+ . persistent ( )
434+ . get :: < DataKey , Vec < Address > > ( & creator_key)
435+ {
436+ let mut updated: Vec < Address > = vec ! [ env] ;
437+ for s in subs. iter ( ) {
438+ if s != * subscriber {
439+ updated. push_back ( s) ;
440+ }
441+ }
442+ env. storage ( ) . persistent ( ) . set ( & creator_key, & updated) ;
443+ }
363444 remove_stream ( env, & key) ;
364445}
365446
@@ -658,6 +739,16 @@ impl SubStreamContract {
658739 }
659740 }
660741
742+ /// Collect from all active streams for a creator in a single call.
743+ /// `max_count` caps the batch size to avoid hitting ledger instruction limits.
744+ /// Returns the total amount collected across all processed streams.
745+ pub fn withdraw_all( env : Env , creator : Address , max_count : u32 ) -> i128 {
746+ let creator_key = DataKey :: CreatorSubscribers ( creator. clone ( ) ) ;
747+ let subs: Vec < Address > = env
748+ . storage ( )
749+ . persistent ( )
750+ . get ( & creator_key)
751+ . unwrap_or ( vec ! [ & env] ) ;
661752 /// Collect from a creator's active streams in a batch.
662753 pub fn withdraw_all ( env : Env , creator : Address , max_count : u32 ) -> i128 {
663754 let subs_key = DataKey :: CreatorSubscribers ( creator. clone ( ) ) ;
@@ -679,6 +770,25 @@ impl SubStreamContract {
679770 }
680771 }
681772
773+ // Single transfer of the total collected amount to the creator
774+ if total_collected > 0 {
775+ for i in 0 ..limit {
776+ let subscriber = subs. get ( i as u32 ) . unwrap ( ) ;
777+ let stream_key = DataKey :: Stream ( subscriber. clone ( ) , creator. clone ( ) ) ;
778+ if env. storage ( ) . persistent ( ) . has ( & stream_key) {
779+ let stream: Stream = env. storage ( ) . persistent ( ) . get ( & stream_key) . unwrap ( ) ;
780+ let token_client = TokenClient :: new ( & env, & stream. token ) ;
781+ token_client. transfer (
782+ & env. current_contract_address ( ) ,
783+ & creator,
784+ & total_collected,
785+ ) ;
786+ break ;
787+ }
788+ }
789+ }
790+
791+ total_collected
682792 total
683793 }
684794
@@ -701,6 +811,20 @@ impl SubStreamContract {
701811 . unwrap_or ( 0 )
702812 }
703813
814+ /// Store an IPFS CID pointing to the creator's profile, links, and tier descriptions.
815+ /// Only the creator themselves can update their own metadata.
816+ pub fn set_creator_metadata ( env : Env , creator : Address , cid : Bytes ) {
817+ creator. require_auth ( ) ;
818+ let key = DataKey :: CreatorMetadata ( creator. clone ( ) ) ;
819+ env. storage ( ) . persistent ( ) . set ( & key, & cid) ;
820+ }
821+
822+ /// Retrieve the IPFS CID for a creator. Returns None if not set.
823+ pub fn get_creator_metadata ( env : Env , creator : Address ) -> Option < Bytes > {
824+ let key = DataKey :: CreatorMetadata ( creator. clone ( ) ) ;
825+ env. storage ( ) . persistent ( ) . get ( & key)
826+ }
827+
704828 pub fn get_total_streamed ( env : Env , subscriber : Address , creator : Address ) -> i128 {
705829 env. storage ( )
706830 . persistent ( )
@@ -859,6 +983,12 @@ fn cancel_internal(env: &Env, subscriber: &Address, stream_id: &Address) {
859983 env. storage ( ) . persistent ( ) . remove ( & key) ;
860984}
861985
986+ fn update_total_streamed( env : & Env , subscriber : & Address , creator : & Address , amount : i128 ) {
987+ let key = DataKey :: TotalStreamed ( subscriber. clone ( ) , creator. clone ( ) ) ;
988+ let current_total: i128 = env. storage ( ) . persistent ( ) . get ( & key) . unwrap_or ( 0 ) ;
989+ env. storage ( )
990+ . persistent ( )
991+ . set ( & key, & ( current_total + amount) ) ;
862992fn top_up_internal ( env : & Env , subscriber : & Address , stream_id : & Address , amount : i128 ) {
863993 subscriber. require_auth ( ) ;
864994 if amount <= 0 {
0 commit comments