watchWithTree method
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;
}