-
Notifications
You must be signed in to change notification settings - Fork 2
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Fix a bunch of typing problems in Buffer, Sample, Fold and Reduce ope…
…rator.
- Loading branch information
Showing
7 changed files
with
188 additions
and
151 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,47 @@ | ||
import '../core/observable.dart'; | ||
import '../core/observer.dart'; | ||
import '../core/subscriber.dart'; | ||
import '../disposables/disposable.dart'; | ||
import '../events/event.dart'; | ||
import '../shared/functions.dart'; | ||
|
||
extension FoldOperator<T> on Observable<T> { | ||
/// Combines a sequence of values by repeatedly applying [transform], starting | ||
/// with the provided [initialValue]. | ||
Observable<R> fold<R>(R initialValue, Map2<R, T, R> transform) => | ||
FoldObservable<T, R>(this, transform, initialValue); | ||
} | ||
|
||
class FoldObservable<T, R> extends Observable<R> { | ||
final Observable<T> delegate; | ||
final Map2<R, T, R> transform; | ||
final R seedValue; | ||
|
||
FoldObservable(this.delegate, this.transform, this.seedValue); | ||
|
||
@override | ||
Disposable subscribe(Observer<R> observer) { | ||
final subscriber = FoldSubscriber<T, R>(observer, transform, seedValue); | ||
subscriber.add(delegate.subscribe(subscriber)); | ||
return subscriber; | ||
} | ||
} | ||
|
||
class FoldSubscriber<T, R> extends Subscriber<T> { | ||
final Map2<R, T, R> transform; | ||
R seedValue; | ||
|
||
FoldSubscriber(Observer<R> destination, this.transform, this.seedValue) | ||
: super(destination); | ||
|
||
@override | ||
void onNext(T value) { | ||
final transformEvent = Event.map2(transform, seedValue, value); | ||
if (transformEvent.isError) { | ||
doError(transformEvent.error, transformEvent.stackTrace); | ||
} else { | ||
seedValue = transformEvent.value; | ||
} | ||
doNext(seedValue); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,51 @@ | ||
import '../core/observable.dart'; | ||
import '../core/observer.dart'; | ||
import '../core/subscriber.dart'; | ||
import '../disposables/disposable.dart'; | ||
import '../events/event.dart'; | ||
import '../shared/functions.dart'; | ||
|
||
extension ReduceOperator<T> on Observable<T> { | ||
/// Combines a sequence of values by repeatedly applying [transform]. | ||
Observable<T> reduce(Map2<T, T, T> transform) => | ||
ReduceObservable<T>(this, transform); | ||
} | ||
|
||
class ReduceObservable<T> extends Observable<T> { | ||
final Observable<T> delegate; | ||
final Map2<T, T, T> transform; | ||
|
||
ReduceObservable(this.delegate, this.transform); | ||
|
||
@override | ||
Disposable subscribe(Observer<T> observer) { | ||
final subscriber = ReduceSubscriber<T>(observer, transform); | ||
subscriber.add(delegate.subscribe(subscriber)); | ||
return subscriber; | ||
} | ||
} | ||
|
||
class ReduceSubscriber<T> extends Subscriber<T> { | ||
final Map2<T, T, T> transform; | ||
bool hasSeed = false; | ||
late T seedValue; | ||
|
||
ReduceSubscriber(Observer<T> destination, this.transform) | ||
: super(destination); | ||
|
||
@override | ||
void onNext(T value) { | ||
if (hasSeed) { | ||
final transformEvent = Event.map2(transform, seedValue, value); | ||
if (transformEvent.isError) { | ||
doError(transformEvent.error, transformEvent.stackTrace); | ||
} else { | ||
seedValue = transformEvent.value; | ||
} | ||
} else { | ||
seedValue = value; | ||
hasSeed = true; | ||
} | ||
doNext(seedValue); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file was deleted.
Oops, something went wrong.
Oops, something went wrong.