watchWithTree method

Stream<List<TreeNode<T>>> watchWithTree(
  1. List<SubSpec> subSpecs
)

Live reactive terminal that joins each parent matching this query with multiple levels of sub-collections, each level described by a SubSpec. Generalises watchWithSub to arbitrary depth.

final stream = posts.query()
    .watchWithTree([
      SubSpec<Comment>(
        subName: 'comments',
        subDefaultExpiration: const Duration(days: 30),
        subFromJson: Comment.fromJson,
        subTypeTag: 'Comment',
        children: [
          SubSpec<Reply>(
            subName: 'replies',
            subDefaultExpiration: const Duration(days: 30),
            subFromJson: Reply.fromJson,
            subTypeTag: 'Reply',
          ),
        ],
      ),
    ]);
// stream: Stream<List<TreeNode<Post>>>
// tree[i].parent              → CItem<Post>
// tree[i].branches['comments'] → List<TreeNode<dynamic>> (one per comment)
// tree[i].branches['comments'][j].branches['replies']
//                              → List<TreeNode<dynamic>> (one per reply)

Re-emits on any change at any level: parent updates/deletes that affect the result set, plus child / grand-child / ... events within any open sub-tree. When a parent leaves the result set, the library cascade-cancels every descendant subscription rooted at that parent — including transitively. Same on outer-stream cancellation.

The returned stream is single-subscription. Wrap with Stream.asBroadcastStream for multi-listener UIs.

Implementation

Stream<List<TreeNode<T>>> watchWithTree(List<SubSpec<dynamic>> subSpecs) {
  // Per-parent state, keyed by `<owner>:<id>`:
  //   childSubs[parentKey][subName] = stream subscription on that
  //     sub-collection's recursive watchWithTree (one entry per spec).
  //   childLatest[parentKey][subName] = latest emitted children list.
  final childSubs =
      <String, Map<String, StreamSubscription<List<TreeNode<dynamic>>>>>{};
  final childLatest = <String, Map<String, List<TreeNode<dynamic>>>>{};
  List<CItem<T>> latestParents = const [];
  late final StreamController<List<TreeNode<T>>> ctrl;
  StreamSubscription<List<CItem<T>>>? parentSub;

  String keyOf(CItem<dynamic> p) => '${p.owner}:${p.id}';

  void emit() {
    if (ctrl.isClosed) return;
    ctrl.add([
      for (final p in latestParents)
        TreeNode<T>(
          parent: p,
          branches: Map<String, List<TreeNode<dynamic>>>.from(
            childLatest[keyOf(p)] ?? const {},
          ),
        ),
    ]);
  }

  Future<void> onParents(List<CItem<T>> parents) async {
    latestParents = parents;
    final currentKeys = parents.map(keyOf).toSet();
    // Cascade-cancel subs for parents that left the result set —
    // their entire sub-tree is gone.
    final leavers =
        childSubs.keys.where((k) => !currentKeys.contains(k)).toList();
    for (final k in leavers) {
      final perParent = childSubs.remove(k);
      if (perParent != null) {
        for (final s in perParent.values) {
          await s.cancel();
        }
      }
      childLatest.remove(k);
    }
    // Open subs for newly-arrived parents — one stream per declared
    // [SubSpec].
    for (final p in parents) {
      final k = keyOf(p);
      if (childSubs.containsKey(k)) continue;
      childSubs[k] = {};
      childLatest[k] = {};
      for (final spec in subSpecs) {
        // Use SubSpec._openOnForTest<T>(...) so this spec's own U
        // survives the loop iteration over List<SubSpec<dynamic>> —
        // without it Dart would erase U to dynamic and the
        // constructor's implicit (Type, typeTag) registration would
        // clash with the already-registered (RealType, typeTag)
        // entry.
        final subColl = spec._openOnForTest(_collection, p);
        // Recurse: each child's own children come from a nested
        // watchWithTree on the sub-collection's query. Empty
        // [spec.children] makes the recursion a single-level scan.
        final stream = subColl.query().watchWithTree(spec.children);
        childSubs[k]![spec.subName] = stream.listen(
          (children) {
            childLatest[k]![spec.subName] = children;
            emit();
          },
          onError: (Object e, StackTrace st) {
            if (!ctrl.isClosed) ctrl.addError(e, st);
          },
        );
      }
    }
    emit();
  }

  ctrl = StreamController<List<TreeNode<T>>>(
    onListen: () {
      parentSub = watch().listen(
        (parents) => unawaited(onParents(parents)),
        onError: (Object e, StackTrace st) {
          if (!ctrl.isClosed) ctrl.addError(e, st);
        },
      );
    },
    onCancel: () async {
      await parentSub?.cancel();
      for (final perParent in childSubs.values) {
        for (final s in perParent.values) {
          await s.cancel();
        }
      }
      childSubs.clear();
      childLatest.clear();
    },
  );
  return ctrl.stream;
}