@@ -55,6 +55,35 @@ fn stream_key(subscriber: &Address, stream_id: &Address) -> DataKey {
5555 DataKey :: Stream ( subscriber. clone ( ) , stream_id. clone ( ) )
5656}
5757
58+ fn stream_exists ( env : & Env , key : & DataKey ) -> bool {
59+ env. storage ( ) . persistent ( ) . has ( key) || env. storage ( ) . temporary ( ) . has ( key)
60+ }
61+
62+ fn get_stream ( env : & Env , key : & DataKey ) -> Stream {
63+ if env. storage ( ) . persistent ( ) . has ( key) {
64+ env. storage ( ) . persistent ( ) . get ( key) . unwrap ( )
65+ } else if env. storage ( ) . temporary ( ) . has ( key) {
66+ env. storage ( ) . temporary ( ) . get ( key) . unwrap ( )
67+ } else {
68+ panic ! ( "stream not found" )
69+ }
70+ }
71+
72+ fn set_stream ( env : & Env , key : & DataKey , stream : & Stream ) {
73+ if stream. balance > 0 {
74+ env. storage ( ) . persistent ( ) . set ( key, stream) ;
75+ env. storage ( ) . temporary ( ) . remove ( key) ;
76+ } else {
77+ env. storage ( ) . temporary ( ) . set ( key, stream) ;
78+ env. storage ( ) . persistent ( ) . remove ( key) ;
79+ }
80+ }
81+
82+ fn remove_stream ( env : & Env , key : & DataKey ) {
83+ env. storage ( ) . persistent ( ) . remove ( key) ;
84+ env. storage ( ) . temporary ( ) . remove ( key) ;
85+ }
86+
5887fn validate_distribution (
5988 creators : & Vec < Address > ,
6089 percentages : & Vec < u32 > ,
@@ -145,7 +174,7 @@ fn subscribe_internal(
145174 }
146175
147176 let key = stream_key ( subscriber, stream_id) ;
148- if env . storage ( ) . persistent ( ) . has ( & key) {
177+ if stream_exists ( env , & key) {
149178 panic ! ( "stream already exists" ) ;
150179 }
151180
@@ -166,7 +195,7 @@ fn subscribe_internal(
166195 percentages,
167196 } ;
168197
169- env . storage ( ) . persistent ( ) . set ( & key, & stream) ;
198+ set_stream ( env , & key, & stream) ;
170199}
171200
172201fn distribute_and_collect (
@@ -176,6 +205,38 @@ fn distribute_and_collect(
176205 total_streamed_creator : Option < & Address > ,
177206) -> i128 {
178207 let key = stream_key ( subscriber, stream_id) ;
208+ let mut stream = get_stream ( env, & key) ;
209+ let now = env. ledger ( ) . timestamp ( ) ;
210+
211+ if now <= stream. last_collected {
212+ return 0 ;
213+ }
214+
215+ if let Some ( creator) = total_streamed_creator {
216+ if is_creator_paused ( env, creator) {
217+ // While paused, advance accounting clock so paused time is never billed.
218+ stream. last_collected = now;
219+ set_stream ( env, & key, & stream) ;
220+ return 0 ;
221+ }
222+ }
223+
224+ let trial_end = stream
225+ . start_time
226+ . saturating_add ( stream. tier . trial_duration ) ;
227+ let charge_start = if stream. last_collected > trial_end {
228+ stream. last_collected
229+ } else {
230+ trial_end
231+ } ;
232+
233+ if now <= charge_start {
234+ return 0 ;
235+ }
236+
237+ let elapsed = ( now - charge_start) as i128 ;
238+ let mut amount_to_collect = elapsed
239+ . checked_mul ( stream. tier . rate_per_second )
179240 if !env. storage ( ) . persistent ( ) . has ( & key) {
180241 panic ! ( "stream not found" ) ;
181242 }
@@ -226,6 +287,11 @@ fn distribute_and_collect(
226287 return 0 ;
227288 }
228289
290+ let token_client = TokenClient :: new ( env, & stream. token ) ;
291+ let mut remaining = amount_to_collect;
292+ let len = stream. creators . len ( ) ;
293+
294+
229295 let token_client = TokenClient :: new ( env, & stream. token ) ;
230296 let mut remaining = amount_to_collect;
231297 let len = stream. creators . len ( ) ;
@@ -252,6 +318,7 @@ fn distribute_and_collect(
252318
253319 stream. balance -= amount_to_collect;
254320 stream. last_collected = now;
321+ set_stream ( env, & key, & stream) ;
255322 env. storage ( ) . persistent ( ) . set ( & key, & stream) ;
256323
257324 if let Some ( creator) = total_streamed_creator {
@@ -268,19 +335,19 @@ fn cancel_group_internal(env: &Env, subscriber: &Address, stream_id: &Address) {
268335 subscriber. require_auth ( ) ;
269336
270337 let key = stream_key ( subscriber, stream_id) ;
271- if !env . storage ( ) . persistent ( ) . has ( & key) {
338+ if !stream_exists ( env , & key) {
272339 panic ! ( "stream not found" ) ;
273340 }
274341
275342 distribute_and_collect ( env, subscriber, stream_id, None ) ;
276343
277- let stream: Stream = env . storage ( ) . persistent ( ) . get ( & key) . unwrap ( ) ;
344+ let stream = get_stream ( env , & key) ;
278345 if stream. balance > 0 {
279346 let token_client = TokenClient :: new ( env, & stream. token ) ;
280347 token_client. transfer ( & env. current_contract_address ( ) , subscriber, & stream. balance ) ;
281348 }
282349
283- env . storage ( ) . persistent ( ) . remove ( & key) ;
350+ remove_stream ( env , & key) ;
284351}
285352
286353fn top_up_internal ( env : & Env , subscriber : & Address , stream_id : & Address , amount : i128 ) {
@@ -291,15 +358,16 @@ fn top_up_internal(env: &Env, subscriber: &Address, stream_id: &Address, amount:
291358 }
292359
293360 let key = stream_key ( subscriber, stream_id) ;
294- if !env . storage ( ) . persistent ( ) . has ( & key) {
361+ if !stream_exists ( env , & key) {
295362 panic ! ( "stream not found" ) ;
296363 }
297364
298- let mut stream: Stream = env . storage ( ) . persistent ( ) . get ( & key) . unwrap ( ) ;
365+ let mut stream = get_stream ( env , & key) ;
299366 let token_client = TokenClient :: new ( env, & stream. token ) ;
300367 token_client. transfer ( subscriber, & env. current_contract_address ( ) , & amount) ;
301368
302369 stream. balance = stream. balance . checked_add ( amount) . expect ( "overflow" ) ;
370+ set_stream ( env, & key, & stream) ;
303371 env. storage ( ) . persistent ( ) . set ( & key, & stream) ;
304372}
305373
@@ -339,6 +407,11 @@ impl SubStreamContract {
339407 subscriber. require_auth ( ) ;
340408
341409 let key = stream_key ( & subscriber, & creator) ;
410+ if !stream_exists ( & env, & key) {
411+ panic ! ( "stream not found" ) ;
412+ }
413+
414+ let stream = get_stream ( & env, & key) ;
342415 if !env. storage ( ) . persistent ( ) . has ( & key) {
343416 panic ! ( "stream not found" ) ;
344417 }
@@ -351,6 +424,7 @@ impl SubStreamContract {
351424
352425 distribute_and_collect ( & env, & subscriber, & creator, Some ( & creator) ) ;
353426
427+ let stream_after = get_stream ( & env, & key) ;
354428 let stream_after: Stream = env. storage ( ) . persistent ( ) . get ( & key) . unwrap ( ) ;
355429 if stream_after. balance > 0 {
356430 let token_client = TokenClient :: new ( & env, & stream_after. token ) ;
@@ -361,6 +435,7 @@ impl SubStreamContract {
361435 ) ;
362436 }
363437
438+ remove_stream ( & env, & key) ;
364439 env. storage ( ) . persistent ( ) . remove ( & key) ;
365440 remove_subscriber_from_creator ( & env, & creator, & subscriber) ;
366441 }
@@ -417,6 +492,8 @@ impl SubStreamContract {
417492
418493 // Settle all streams up to pause timestamp, then freeze charging.
419494 for subscriber in subs. iter ( ) {
495+ let s_key = stream_key ( & subscriber, & creator) ;
496+ if stream_exists ( & env, & s_key) {
420497 if env
421498 . storage ( )
422499 . persistent ( )
@@ -445,6 +522,10 @@ impl SubStreamContract {
445522 // Resume billing from now so paused window is never charged.
446523 for subscriber in subs. iter ( ) {
447524 let s_key = stream_key ( & subscriber, & creator) ;
525+ if stream_exists ( & env, & s_key) {
526+ let mut stream = get_stream ( & env, & s_key) ;
527+ stream. last_collected = now;
528+ set_stream ( & env, & s_key, & stream) ;
448529 if env. storage ( ) . persistent ( ) . has ( & s_key) {
449530 let mut stream: Stream = env. storage ( ) . persistent ( ) . get ( & s_key) . unwrap ( ) ;
450531 stream. last_collected = now;
@@ -478,6 +559,11 @@ impl SubStreamContract {
478559 }
479560
480561 let key = stream_key ( & subscriber, & creator) ;
562+ if !stream_exists ( & env, & key) {
563+ panic ! ( "stream not found" ) ;
564+ }
565+
566+ let stream_before = get_stream ( & env, & key) ;
481567 if !env. storage ( ) . persistent ( ) . has ( & key) {
482568 panic ! ( "stream not found" ) ;
483569 }
@@ -487,7 +573,7 @@ impl SubStreamContract {
487573
488574 distribute_and_collect ( & env, & subscriber, & creator, Some ( & creator) ) ;
489575
490- let mut stream: Stream = env . storage ( ) . persistent ( ) . get ( & key) . unwrap ( ) ;
576+ let mut stream = get_stream ( & env , & key) ;
491577 let mut balance = stream. balance ;
492578
493579 if new_rate_per_second < old_rate && balance > 0 {
@@ -520,7 +606,7 @@ impl SubStreamContract {
520606 . expect ( "overflow" ) ;
521607 }
522608
523- env . storage ( ) . persistent ( ) . set ( & key, & stream) ;
609+ set_stream ( & env , & key, & stream) ;
524610
525611 if old_rate != new_rate_per_second {
526612 TierChanged {
@@ -543,6 +629,8 @@ impl SubStreamContract {
543629
544630 for i in 0 ..limit {
545631 let subscriber = subs. get ( i) . unwrap ( ) ;
632+ let s_key = stream_key ( & subscriber, & creator) ;
633+ if stream_exists ( & env, & s_key) {
546634 if env
547635 . storage ( )
548636 . persistent ( )
0 commit comments