Synced Query Data Source

A synced query data source is a live collection that lives inside a hub's workspace and stays in sync with a mesh query. The framework subscribes to IMeshQueryProvider.Query<T> when the hub starts, seeds the workspace's EntityStore with the initial result set, and continuously folds Added / Updated / Removed deltas into that same store via IDataChangeNotifier.

The payoff: hub-internal code — validators, layout areas, compile pipelines, access checks — reads its source data through the standard workspace.GetStream<T>() / workspace.GetStream(new CollectionReference("name")) surface. No Observe round-trip, no CQRS staleness lag, no Observable.FromAsync at every leaf. By the time the hub handles its first message, the synced collection is already populated.

IMeshQuery Query<MeshNode> IMeshQuery Query<MeshNode> one per query string Scan / Fold ImmutableDictionary keyed by Path Initial Gate suppress until all queries fire Initial Hub Workspace Added / Updated / Removed deltas re-emits full snapshot on every change

Synced query data source pipeline: mesh queries fold into a path-keyed dictionary, gated until all Initial events arrive, then emit live snapshots into the hub workspace.


When to use it

Use a synced query data source whenever a hub needs a local view of nodes that live elsewhere in the mesh.

Use case What you sync
NodeType hubs Sources and Tests Code nodes — the compile pipeline reads these collections directly, zero round-trips
Per-data-node hubs AccessAssignments — the access pipeline checks permissions synchronously without touching the security service
Aggregator hubs Cross-namespace dashboards built from nodeType:Order status:Open, rendered straight from the workspace stream
System defaults + extensions Agent, Model, Role — static built-ins appear on first subscribe via IStaticNodeProvider; user-created instances stream in as Added / Updated / Removed deltas

For system-default / mesh-extension composition, see Extensible Defaults.

Tip: Use workspace.GetMeshNodeStream(path) for one-shot single-node reads and for nodes you never keep in a collection. The synced collection shines for sets you read repeatedly — and it is bidirectional, so writes work through the same surface.


How it works

The data source is built on VirtualDataSource.WithVirtualType<T> — the framework primitive for "this collection comes from an IObservable stream." The mesh adds a thin extension method WithMeshQuery that composes three pieces:

1. Subscribe to each mesh query

One IMeshQueryCore.Query<MeshNode> subscription per query string. Multi-query collections are fine — the result is their union. Each QueryResultChange<MeshNode> carries Initial / Reset / Added / Updated / Removed deltas together with the matching MeshNode payloads.

2. Fold deltas into a path-keyed dictionary

A single Scan accumulates an ImmutableDictionary<string, MeshNode> keyed by MeshNode.Path, spanning every change event from every query. It also tracks how many Initial / Reset events each query has produced. The result is always the union of every query's current matches, deduplicated by path.

3. Gate the first emission until every query has sent its Initial

A Where(...) clause suppresses emission until each query's Initial count reaches the configured provider count. The first .Take(1) consumer therefore sees a complete snapshot, not a partial one. After the gate opens, every change re-emits the full IEnumerable<MeshNode> — the dictionary's values.

A companion BehaviorSubject<ImmutableHashSet<string>> tracks the live path set in parallel. The AddSyncedQuery reducer reads it synchronously to decide whether this source Owns a given path before opening the per-node remote stream for a write.

Note: There is no per-path read subscription. Read content comes from the MeshNode payloads carried by the query events. The synchronization protocol re-pushes those payloads on every owning-hub change, so the synced collection re-emits as soon as the query layer notices.


Configuration

config.AddData(data => data
    .WithVirtualDataSource("$mesh-sources", vs => vs.WithMeshQuery(
        query: $"namespace:{hubPath}/Source scope:subtree nodeType:Code",
        collectionName: "Sources")));

Multiple synced collections per hub are fine — each goes into its own virtual data source. The canonical example is MeshDataSource registration in MeshDataSourceExtensions.AddMeshDataSource: every per-node hub automatically gets Sources and Tests synced collections derived from the hub's path.

API reference — extension methods

public static VirtualDataSource WithMeshQuery(
    this VirtualDataSource ds,
    string query,
    string? collectionName = null);
Parameter Description
query Mesh query string in the standard syntax (see Query Syntax). Common shapes: namespace:X scope:subtree nodeType:Y, path:X, path:X scope:descendants, namespace:X scope:nextLevel (the next populated level — graph navigation).
collectionName Workspace collection name. Defaults to nameof(MeshNode). Required when several synced collections would otherwise collide — for example, Sources and Tests both hold MeshNode.

There is one overload, and it yields a collection of MeshNode. To work with a single content shape, filter on the node's Content when you read the collection.


Reading and writing

The synced collection has two distinct surfaces. Reads come from the dict-of-MeshNode snapshot the data source emits. Writes route through the per-node hub via the workspace's (address, reference) remote-stream cache.

Read — subscribe to the snapshot observable

var workspace = hub.GetWorkspace();
var collection = workspace.GetQuery(
    "my-collection-id",
    "namespace:Agent nodeType:Agent");

var sub = collection.Subscribe(snapshot =>
{
    // snapshot is IEnumerable<MeshNode> — the COMPLETE current set,
    // path-keyed and deduplicated. Rebuild your view from this each time.
});

Every emission is the full current collection, not individual deltas — the Scan inside SyncedQueryMeshNodes already merged them. The observable is Replay(1).RefCount(), so a late subscriber gets the cached latest snapshot immediately, and upstream subscriptions are shared across every consumer of the same id.

For a single-node read by path, use workspace.GetMeshNodeStream(path). GetQuery is for collections, not for known-path lookups.

Write — Update on the per-node remote stream

// AddSyncedQuery's reducer routes any
// workspace.GetStream(new MeshNodeReference(path)) call to
// GetRemoteStream<MeshNode, MeshNodeReference>(new Address(path), ref)
// when `path` is in the source's live path set (Owns(path)).
var stream = workspace.GetRemoteStream<MeshNode, MeshNodeReference>(
    new Address(path), new MeshNodeReference());

stream.Update(current => current is null
    ? null
    : new ChangeItem<MeshNode>(
        current with { Name = "New Name" },
        changedBy: hub.Address.ToString(),
        stream.StreamId,
        ChangeType.Full,
        stream.Hub.Version,
        Updates: null));

Because workspace.GetRemoteStream<MeshNode, MeshNodeReference>(addr, ref) caches per (address, reference), every caller in the hub gets the same instance — write-paths share one upstream subscription per node. A write through this stream propagates to the owning per-node hub via the synchronization protocol, and the owning hub's update echoes back through the query layer to the next synced-collection emission.


Live updates without polling

When any node is updated anywhere in the mesh, the change arrives in the synced collection through the upstream IMeshQueryProvider.Query subscription. The query layer already mirrors per-hub change notifications into Updated / Removed events; the Scan folds those into the dictionary and re-emits the snapshot on the next tick.

There is no manual cache to invalidate — the query stream is the cache.

In practice, a synced collection of nodeType:Code nodes re-emits within tens of milliseconds of a developer's save reaching persistence. Consumers such as the compile pipeline and side-menu listings re-render off the new snapshot automatically.


Why this is safe — the actor model

A hub is a single-threaded actor. At any moment exactly one message is in flight, and the data source's stream subscription, the workspace's per-node remote-stream cache, and all user code that reads or writes through them run on that same thread. There are no concurrent callers, no torn reads, no partial state.

The actor model is the integrity guarantee. That is why the cache can be plain in-memory state — no locks, no compare-and-swap, no reconciliation logic.

Concretely: when your handler calls workspace.GetRemoteStream<MeshNode, MeshNodeReference>(addr, ref).Update(...), the synchronization protocol routes the patch to the owning per-node hub, which is itself a single-threaded actor. It processes the update in order with every other write to that node. No two writes collide, and no reader ever observes a half-applied state.

Workspace remote-stream cache — one stream per (address, reference)

The mechanism that makes the write path share a single subscription per node across all writers in the same hub is Workspace._remoteStreamCache (src/MeshWeaver.Data/Workspace.cs). It is a plain ConcurrentDictionary<(Address, WorkspaceReference), ISynchronizationStream> keyed by the inputs to workspace.GetRemoteStream<TReduced, TReference>(addr, ref). The first call for a given key opens an external client subscription and stores the stream; every subsequent call returns the same instance.

This is exactly what the MeshNodeReference(path) reducer registered by AddSyncedQuery relies on: when it returns workspace.GetRemoteStream<MeshNode, MeshNodeReference>(new Address(path), new MeshNodeReference()), every writer through the synced collection — and any other code in the workspace asking for the same (addr, ref) — gets back the same stream instance. One upstream pump per node, no matter how many writers share it.

Cache eviction is lazy and happens in three cases:

That third case is not cosmetic. A stream's store is a ReplaySubject, so an OnError on it is permanent under the Rx grammar: the stream can never emit again, and every later subscriber replays that same error the instant it subscribes. CreateExternalClient's terminal arm faults the stream and tears down only its keep-alive, leaving the object errored but undisposed — so while liveness was judged on disposal alone, the cache kept serving the corpse for the whole process lifetime. Every later write to that node then failed instantly while reporting the 30-second initial-state bound, which is what made the symptom unreadable: a timeout announced after a fraction of a millisecond. On boot it cost each replica one of its default packages, and the installer's own fall-back-to-a-full-install repair re-entered the same dead mirror and could not possibly succeed — #2387.

Dropping a faulted stream is not a retry: nothing re-attempts on a timer and nothing loops. The next natural caller opens a fresh mirror, and a stream that is born dead (a construction-time fault — a same-process owner NACKing the subscribe inline) is handed back rather than re-created, so create → fault → create can never spin. This is the same contract PromiseCache states for pooled I/O: cache a success, evict a fault, and let the next caller re-attempt.

An evicted stream is closed as soon as its last declared holder lets go. Eviction removes the stream from the cache but cannot dispose it at the eviction site — it does not know whether anyone is still reading. Rx subscriber counts cannot answer that either: the reduce chain CreateExternalClient builds subscribes to the stream itself, so an evicted stream sits at 2–3 subscribers forever. So holders declare themselves: anything that keeps a remote stream past the call that resolved it takes one through Workspace.AcquireRemoteStreamUnchecked and disposes the returned lease when it is done — the shared mesh-node cache's hydration for its entry's whole life, each cross-hub write for the duration of its Observable.Create subscription. When the last lease on an already-evicted stream goes, it is disposed at once: UnsubscribeRequest reaches the owner and both sync/ hubs die. A stream nobody leased keeps the conservative parking (it may still have undeclared readers) and is reclaimed by the mesh-node cache's idle release or at workspace disposal. Without the leases every write to a subscribed path left a client sync/ hub and its owner-side twin behind for the process lifetime — #1324.

The full integrity story is a plain dictionary, single-threaded access, and shared-by-default semantics — made safe entirely by the actor model.

The RemoteStreamCacheTest (test/MeshWeaver.Query.Test) pins this contract: two calls for the same key are reference-equal, and a disposed or faulted stream is evicted before the next caller receives it. Two details the test makes explicit and this page should not blur:


Caveat — RAM footprint

The only real cost is memory. Every synced collection holds its full result-set dictionary in the hub's address space. A query matching 100 nodes pins 100 MeshNodes' worth of memory on every subscriber hub. Keep queries narrow enough that the live set genuinely belongs in RAM.


Testing

The synced query data source is a production wiring, not a side-channel. Tests must exercise it as it ships, not through test-only handlers. The pattern is fixed and short.

1. Define a NodeType that uses WithMeshQuery

Register a static MeshNode via AddMeshNodes whose HubConfiguration adds the synced data source — exactly the way a production NodeType (Sources / Tests / AccessAssignments) does it.

protected override MeshBuilder ConfigureMesh(MeshBuilder builder)
    => base.ConfigureMesh(builder)
        .AddMeshNodes(new MeshNode("Subscriber", TestPartition)
        {
            Name      = "Subscriber",
            NodeType  = "Markdown",
            State     = MeshNodeState.Active,
            HubConfiguration = config => config.AddData(data =>
                data.WithVirtualDataSource("$mesh-subjects",
                    vs => vs.WithMeshQuery(
                        query: $"namespace:{SubjectsNamespace} scope:subtree nodeType:Markdown")))
        });

The subscriber's per-node hub now mirrors every matching node via the upstream mesh-query subscription — same wiring as production.

2. Drive the test from the client, never from Mesh

Mesh is a virtual coordinator. Real players are hubs created via Mesh.ServiceProvider.CreateMessageHub(...) (the test base does this through GetClient()). Post messages from the client so routing, serialization, and the synchronization protocol all run end to end.

3. Use standard data-layer messages — not custom request handlers

Reads and writes go through the framework's data-layer messages.

Do not write GetMyThingRequest / WriteMyThingRequest test-only handlers. They route around the contract you are trying to test.

4. Verify with the existing test base helpers

MonolithMeshTestBase exposes:

Verification follows an observe-until-condition rhythm (mirrors ObservableQueryTests). The synced collection stays subscribed the whole time — no Take(1), no draining.

End-to-end test sketch

[Fact]
public async Task DataChangeRequestOnSubscriber_PropagatesToOwningHub()
{
    var path = $"{SubjectsNamespace}/alpha";

    // Source-side state. The write is a COLD observable, so the assertion's
    // Subscribe IS the write — there is nothing to bridge.
    await NodeFactory.CreateNode(MakeSubject("alpha", "Original")).Should().Emit();

    // Wait for the CONDITION, never a fixed sleep: a Task.Delay races CI load —
    // too short and it flakes, too long and it wastes the suite's budget.
    // 🚨 The assertion OWNS the wait. No FirstAsync + ToTask bridge — forbidden
    // repo-wide (2026-08-30, "no ToTask ever") — and no bare `await observable`
    // either: Rx's awaiter resumes the test inline on the signalling thread.
    // `.Within(t)` is the deadline, which is why there is no CancellationToken here.
    var current = await ReadNode(path).Should().Within(10.Seconds()).Match(n => n is not null);

    // Write at the SUBSCRIBER via DataChangeRequest. The subscriber's
    // synced data source routes the update through its cached per-node
    // remote stream → owning per-node hub.
    await client.Observe(
            new DataChangeRequest { Updates = [current! with { Name = "Updated" }] },
            o => o.WithTarget(new Address(SubscriberPath)))
        .Should().Emit();

    // Source side reflects the write — again, observe until the condition holds.
    await Observable.Interval(TimeSpan.FromMilliseconds(50)).StartWith(0L)
        .SelectMany(_ => ReadNode(path))
        .Should().Within(10.Seconds()).Match(n => n?.Name == "Updated");
}

Reconnecting…
The server was updated. Reloading the page to pick up the latest version.