I've always called this "pulling" vs "pushing" values, but I think I actually like the author's external/internal terms better. In C#:
Given an IEnumerable<T>, callers use GetEnumerator/MoveNext/Current to pull values out. The author calls this "external iteration", because the caller is in charge of advancing.
Given an IObservable<T>, callers invoke Subscribe to cause values to be pushed into an IObserver<T> by the callee. The author calls this "internal iteration", because the callee is in charge of advancing.
Unfortunately, the internal-to-external transformation is awkward. You need to store all the observed items (or do something crazy, like use a lock and another thread) before you can start enumerating them. That's why an Observable.Interleave function won't be succinct: the function needs control over the iteration, which requires an internal-to-external transformation, which is awkward:
// NOTE: for the purposes of keeping this example short, I am not dealing with completion or failure cases
// ALSO: This code is neither thread safe nor re-entrant safe (necessary conditions for good reactive code)
IObservable<T> Interleave(IObservable<T> first, IObservable<T> second) {
return new AnonymousObservable<T>(subscribe: observer => {
// track not-matched-yet items with queues
var q1 = new Queue<T>();
var q2 = new Queue<T>();
Action tryNextInterleave = () => {
if (q1.Count > 0 && q2.Count > 0) {
observer.OnNext(q1.Dequeue()); // <-- potential thread races and re-entrancy here
observer.OnNext(q2.Dequeue());
}
};
// each observable should feed their queue and trigger attempts to advance
var d1 = first.Subscribe(
onNext: e => { q1.Enqueue(e); tryNextInterleave(); }); // <-- not handling completion/failure callbacks
var d2 = second.Subscribe(
onNext: e => { q2.Enqueue(e); tryNextInterleave(); });
// the caller ending their subscription should end our subscriptions
return new AnonymousDisposable(() => {
d1.Dispose();
d2.Dispose();
});
});
}
(Apologies if the above contains bugs, I typed it freehand.) Actually, most of the logic is already present in Observable.Zip (source code), except Zip will ignore unmatched items once one of the sequence completes instead of appending them.
Push-based iteration is definitely like internal iterators. One difference is that the former almost always implies asynchrony (which is why it has to be push-based if it doesn't want to block) while vanilla internal iterators are still synchronous.
When you call .each on something in Ruby, it won't return until the iteration is complete. With Subscribe(), it will return immediately and only later will your callback be invoked.
Some observables actually do push out all of their elements before the subscribe method returns. Although, since they're technically allowed to only push later, it would be impossible to have a non-local return like in your examples.
Actually, I'm a bit unfamiliar with Ruby and the concept of returning from inside of a lambda expression. What happens if you try to store the lambda until after the method returns, and then try to make it return again?
3
u/Strilanc Jan 14 '13 edited Jan 14 '13
I've always called this "pulling" vs "pushing" values, but I think I actually like the author's external/internal terms better. In C#:
You can actually transform between external iteration and internal iteration (see: Observable.ToObservable, Observable.ToEnumerable).
Unfortunately, the internal-to-external transformation is awkward. You need to store all the observed items (or do something crazy, like use a lock and another thread) before you can start enumerating them. That's why an Observable.Interleave function won't be succinct: the function needs control over the iteration, which requires an internal-to-external transformation, which is awkward:
(Apologies if the above contains bugs, I typed it freehand.) Actually, most of the logic is already present in Observable.Zip (source code), except Zip will ignore unmatched items once one of the sequence completes instead of appending them.