import { merge, Observable, OperatorFunction } from 'rxjs'; import { IChangeSet } from '../IChangeSet'; import { IQuery } from '../IQuery'; import { map, publish, scan } from 'rxjs/operators'; import { Cache } from '../Cache'; import { AnonymousQuery } from '../AnonymousQuery'; import { mergeMany } from './mergeMany'; /** * The latest copy of the cache is exposed for querying i) after each modification to the underlying data ii) on subscription * @typeparam TObject The type of the object. * @typeparam TKey The type of the key. * @typeparam TValue The type of the value. * @param itemChangedTrigger Should the query be triggered for observables on individual items */ export function queryWhenChanged(itemChangedTrigger?: (value: TObject) => Observable): OperatorFunction, IQuery> { return function queryWhenChangedBaseOperator(source) { if (itemChangedTrigger == undefined) { return source.pipe( scan((cache, changes) => { cache.clone(changes); return cache; }, new Cache()), map(list => new AnonymousQuery(list)), ); } return source.pipe( publish(shared => { const state = new Cache(); const inlineChange = shared.pipe( mergeMany(itemChangedTrigger), map(_ => new AnonymousQuery(state)), ); const sourceChanged = shared.pipe( scan((list, changes) => { list.clone(changes); return list; }, state), map(list => new AnonymousQuery(list)), ); return merge(sourceChanged, inlineChange); }), ); }; }