materialize, timestamped, sequenceEqual

알림을 Next / Err / Done으로 재화하고, 이벤트에 시각을 찍고, 두 수열이 같은지 묻습니다.

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])

강의

Dart Stream의 세 종단 — 값, 에러, 닫힘 — 은 보통 파이프를 떠납니다. materialize는 각각을 StreamEvent 값으로 바꿔 체인을 통과할 수 있게 합니다: 데이터 이벤트는 Next(value)가 되고, 에러는 Err가 된 뒤 결과가 완료됩니다 — 에러를 내지 않습니다 — 그리고 닫힘은 Done이 된 뒤 결과가 완료됩니다. 재화하는 이유입니다: toList가 실패하는 대신 에러를 모을 수 있고, 로그가 Err(boom)Next(1) 옆에 찍을 수 있고, 테스트가 정확한 알림 수열을 단언할 수 있습니다. dematerialize는 역입니다 — Next는 값이 되고, ErrStream.addError가 되고, Done은 결과를 닫고, Done 이후는 무시됩니다. 왕복: materialize().dematerialize()는 원래 값의 스트림입니다. fxdart 이벤트 레이어, Rx의 materialize / dematerialize를 따랐습니다.

시간은 알림이 실을 수 있는 다른 메타데이터입니다. timestamped는 각 이벤트를 도착한 벽시계 시각과 짝짓고 ((DateTime at, T value)), intervals는 이전 이벤트 이후의 시간과 짝짓습니다 ((Duration dt, T value)). 첫 이벤트의 dt는 언제나 Duration.zero입니다. 둘 다 now:를 받으므로 테스트가 가짜 시계 — now: () => DateTime.utc(2020) — 를 DateTime.now 대신 넘길 수 있습니다. 에러와 닫힘은 그대로 통과합니다. Rx의 timestamptimeInterval을 따랐습니다.

이벤트의 partition(test)는 풀 쪽 partition이 아닙니다 (그건 한 번 걸으며 리스트 둘을 돌려줍니다). 살아 있는 체인 하나를 (matches, rest)로 나누되 소스 번의 실행을 공유합니다. 레코드는 즉시 반환되고, 어느 쪽이든 듣기 시작하면 소스가 시작되며, 아무도 듣지 않는 쪽에 속하는 값은 버퍼에 담기지 않고 버려집니다. 양쪽 반쪽을 모두 원하면 소스가 발화하기 전에 둘 다 들으세요.

sequenceEqual은 두 수열이 같은 값을 같은 순서로 담고 함께 끝나는지 묻습니다. 이벤트 레이어에서는 터미널입니다: fxEvents(a).sequenceEqual(b)Future<bool>을 돌려주고, 첫 값 또는 길이 불일치에서 false이며, 어느 쪽의 에러든 퓨처를 실패시킵니다. 같은 질문이 pull에도 있습니다: iterable에는 sequenceEqual / Fx.sequenceEqual, sequenceEqualAsync / FxAsync.sequenceEqualFxAsyncIterable용입니다. Rx의 sequenceEqual을 따랐습니다.

데모 1 · Next, Err, Done

깨끗한 닫힘은 Done이 됩니다. 에러는 Err가 된 뒤 체인이 완료되므로, toList는 실패하는 대신 StreamEvent 리스트를 돌려줍니다:

데모 2 · 넘길 수 있는 시계

now: () => DateTime.utc(2020)가 스탬프를 결정적으로 만듭니다. intervals도 같은 훅을 쓰며, 간격이 정확하도록 한 칸씩 나아가는 시계를 넘깁니다:

데모 3 · sequenceEqual, 그리고 partition

pull 표기는 그저 fx([1, 2]).sequenceEqual([1, 2])입니다. 이벤트 터미널은 Stream을 받습니다. partition의 양쪽을 소스가 돌기 전에 들으세요. 듣지 않은 반쪽은 버려집니다:

관련 항목: fxEvents — 이 연산자들이 앉는 체인 · Stream 브리지StreamFxAsync로 당기는 네 가지 방법 · partition — 한 번의 순회로 리스트 둘을 얻는 풀 쪽 원본 · FxSubscriptions — 리스너 자루를 함께 취소