materialize, timestamped, sequenceEqual

Reifica notificaciones como Next / Err / Done, sella eventos con tiempo y comprueba si dos secuencias son iguales.

sealed class StreamEvent<T> — Next(T value) | Err(Object error, [StackTrace]) | Done() FxEvents<StreamEvent<T>> FxEvents<T>.materialize() // chain (events) FxEvents<T> FxEvents<StreamEvent<T>>.dematerialize() FxEvents<(DateTime at, T value)> FxEvents<T>.timestamped({DateTime Function()? now}) FxEvents<(Duration dt, T value)> FxEvents<T>.intervals({DateTime Function()? now}) (FxEvents<T> matches, FxEvents<T> rest) FxEvents<T>.partition(bool Function(T) test) Future<bool> FxEvents<T>.sequenceEqual(Stream<T> other, {bool Function(T, T)? eq}) bool sequenceEqual<T>(Iterable<T> a, Iterable<T> b, [bool Function(T, T)? eq]) Future<bool> sequenceEqualAsync<T>(FxAsyncIterable<T> a, FxAsyncIterable<T> b, [bool Function(T, T)? eq]) bool Fx.sequenceEqual(Iterable<T> other, [bool Function(T, T)? eq]) Future<bool> FxAsync.sequenceEqual(FxAsyncIterable<T> other, [bool Function(T, T)? eq])

Lección

Los tres terminales de un Stream de Dart — un valor, un error, un cierre — normalmente abandonan la tubería. materialize convierte cada uno en un valor StreamEvent que puede viajar a través de la cadena: un evento de datos se convierte en Next(value), un error se convierte en Err y entonces el resultado completa — no falla — y un cierre se convierte en Done y entonces el resultado completa. Ese es el punto de reificarlos: toList puede recolectar un error en vez de fallar, un log puede imprimir Err(boom) junto a Next(1), un test puede afirmar la secuencia exacta de notificaciones. dematerialize es la inversa — Next se convierte en un valor, Err se convierte en Stream.addError, Done cierra el resultado, y todo lo que venga después de Done se ignora. Ida y vuelta: materialize().dematerialize() es el stream original de valores. Capa de eventos de fxdart, siguiendo a materialize / dematerialize de Rx.

El tiempo es el otro metadato que una notificación puede llevar. timestamped empareja cada evento con la hora de reloj de pared en la que llegó ((DateTime at, T value)); intervals lo empareja con el tiempo desde el anterior ((Duration dt, T value)), y el primer evento es siempre Duration.zero. Ambos aceptan now: para que un test pueda pasar un reloj falso — now: () => DateTime.utc(2020) — en vez de DateTime.now. Los errores y el cierre pasan sin cambios. Siguiendo a timestamp y timeInterval de Rx.

partition(test) en eventos no es el partition del lado pull (que recorre una vez y devuelve dos listas). Divide una cadena viva en (matches, rest) compartiendo una ejecución de la fuente. El record se devuelve de inmediato; escuchar cualquiera de los lados arranca la fuente; un valor que pertenece a un lado al que nadie escucha se descarta, no se almacena en búfer. Escucha ambos antes de que la fuente dispare si quieres ambas mitades.

sequenceEqual pregunta si dos secuencias contienen los mismos valores en el mismo orden y se detienen juntas. En la capa de eventos es un terminal: fxEvents(a).sequenceEqual(b) devuelve Future<bool>, false ante el primer valor o desajuste de longitud, y un error de cualquiera de los lados hace fallar el future. La misma pregunta existe en pull: sequenceEqual / Fx.sequenceEqual para iterables, sequenceEqualAsync / FxAsync.sequenceEqual para FxAsyncIterables. Siguiendo a sequenceEqual de Rx.

Demo 1 · Next, Err, Done

Un cierre limpio se convierte en Done. Un error se convierte en Err y entonces la cadena completa — así toList devuelve la lista de StreamEvent en vez de fallar:

Demo 2 · Un reloj que puedes pasar

now: () => DateTime.utc(2020) mantiene los sellos deterministas. intervals usa el mismo gancho, con un reloj que avanza a pasos para que los huecos sean exactos:

Demo 3 · sequenceEqual, y partition

La grafía pull es simplemente fx([1, 2]).sequenceEqual([1, 2]). El terminal de eventos toma un Stream. Escucha ambos lados de partition antes de que la fuente corra, o la mitad no escuchada se descarta:

Relacionado: fxEvents — la cadena sobre la que se asientan estos operadores · Puentes de Stream — cuatro formas de tirar de un Stream hacia FxAsync · partition — el original del lado pull, dos listas de un solo recorrido · FxSubscriptions — cancela un saco de listeners juntos