Server-side realtime plan: change streams, SignalR and the subscription model (October 2026)¶
This plan answers two questions from Chris: is Rx.NET the right way to turn MongoDB change streams into
SignalR notifications, and what should the server-side subscription model be? It maps the current
pipeline, compares the options against current .NET 10 / ASP.NET Core 10 facts, recommends a target
architecture, and fixes a hub contract v2 for the parallel client plan to consume. It is planning
only; nothing starts until Chris approves and answers the decisions in §9. Code facts are from main
at 8f73bcbc9 (4 Oct 2026). Markers: VERIFIED (code path read), CORRECTED (an earlier
finding was wrong), NOT VERIFIED (reason given). The client rewrite is out of scope (see
state-management plan §5b).
Chris decisions recorded 2026-10-04:
- Protocol: the API-wide Newtonsoft → System.Text.Json migration is a prerequisite, planned
separately in system-text-json-migration-plan-2026-10.md.
Hub v2 uses the STJ hub protocol.
- Full-stats push: add it back on v2 as a follow-up release (R4), staging only; production is
not planned for now. The legacy SubscribeToProjectFullStats is still retired in R0.1 (D5).
1. Current state¶
1.1 Pipeline map (VERIFIED unless marked)¶
| Stage | How it works today | Evidence |
|---|---|---|
| Change streams | One cached hot stream per aggregate type per API pod, opened lazily on the first subscriber: Observable.Create loop over WatchAsync, FullDocument = UpdateLookup, op-type $match only (no projection, no per-project filter), then RetryWithExponentialBackoff(15, 2 s…180 s) → Publish().RefCount() |
MongoContext.cs:206-314; MongoRepositoryBase.cs:38-39 |
| Streams opened in practice | Project (project details + listings), DataExportJob, Study (presence or full stats only). Presence is off everywhere (activeReviewerTrackingEnabled default false; no cluster-gitops override), and full stats has no web caller (row RT3), so production most likely runs 2 streams per pod. NOT VERIFIED live (no metrics read) |
env-mapping.yaml:1596-1598; FeatureFlags.cs:35-36; grep of cluster-gitops |
| Resume tokens | Stored per event after OnNext, in memory only (InMemoryResumePointRepository, MongoLamarRegistry.cs:50). They survive an in-process retry, not a pod restart. Invalid token → cleared, then OnError |
MongoContext.cs:220-235,267-288 |
| Per-subscriber fan-out | Every subscription builds its own Rx chain over the whole collection stream and filters by id in memory, with three LogInformation calls per event per subscriber |
MongoRepositoryBase.cs:165-212 |
| Subscription registry | ConnectionKeyedSubDictionaryCache → per type, per connection, per key (("ENTITY", id), "COLLECTION", ("FULL-STATS", p, u), ("STUDY-PRESENCE", p, s)) → IDisposable |
AggregateRootEntitySubscriptionManager.cs:29-47,187-262,279-282,417-418 |
| Subscriber retry | A second, unfiltered retry (MyUtils.cs:588-640, 24 h budget) wraps each subscriber chain. Two retry operators with different semantics coexist |
AggregateRootEntitySubscriptionManager.cs:60-71 |
| Delivery | _hubContext.Clients.Client(connectionId).X(dto) from inside the Rx callback. The legacy path (authority mode Off, the deployed default) maps DTOs inline on the change-stream loop, including a synchronous Investigators.GetMany per subscriber per event |
NotificationHub.cs:1155-1190,1255-1296; ProjectStandardDto.cs:137-144 |
| Groups | User Group: {id} (joined on connect), Project Group: {id} (joined on SubscribeToProject). Nothing is sent to either group on main; #3932 adds the first user-group send |
NotificationHub.cs:92-104,1148,1194; gh pr diff 3932 |
| Authorization at subscribe | Hub policies via SignalRAuthorizationHandler (arg 0 = project id, or job id resolved to project). Enforced mode (M4b) re-decides with RealtimeDeliveryAuthority and, for payload streams, again at delivery (late-query barrier) |
SignalRAuthorizationHandler.cs:95-145; NotificationHub.cs:1394-1408; AggregateRootEntitySubscriptionManager.cs:476-537 |
| Revocation | TakeUntil(ProjectAccessRevoked…) per stream, plus RevokeProjectStreams on a disabled-member update |
AggregateRootEntitySubscriptionManager.cs:596-640; NotificationHub.cs:1159-1166,1212-1231 |
| Disconnect | Unsubscribes project and study managers, then review-session cleanup (Mongo ReviewSessionConnection + Quartz timers). Not the DataExportJob manager |
NotificationHub.cs:112-131,143-320 |
| Non-change-stream pushes | FEAT-024 statistics: Mongo outbox → ProjectStatisticsDispatchService (1 s) → bus publish → per-pod temporary queue → per-connection authority recheck → invalidation-only ProjectStatisticsChanged. Claim revocation: same pattern → Clients.User. Flag revision: Clients.All on the handling pod only. #3932 (open): a raw WatchAsync loop on the inbox → InboxChanged() to the user group |
Statistics/*.cs; ClaimRevocations/*.cs; RuntimeFeatureFlagsController.cs:153 |
| Scale-out | No backplane (no AddStackExchangeRedis for SignalR). Production API HPA 3–10 pods. Every pod opens its own streams; cookie session affinity is on every SyRF ingress |
Program.cs:379-400; cluster-gitops production/api/values.yaml:46-50; _ingress.tpl:20-23 |
| Reconnect | Server: no AllowStatefulReconnects. A reconnect gets a new connection ID and no subscriptions (client bug #3976). A pod restart loses everything; nothing replays |
Program.cs:822; signal-r.service.ts:633 |
1.2 Hub surface (VERIFIED)¶
NotificationHub.cs is 1,650 lines with 15 client-callable methods:
- 5 review-session commands: JoinStudyReview, LeaveStudyReview, StartedAnnotating, StoppedAnnotating, Heartbeat;
- 10 subscription methods: subscribe/unsubscribe for StudyPresence, DataExportJob, Project, ProjectSummaries and ProjectFullStats.
Server→client messages:
| Message | Trigger | Payload | Audience on main |
|---|---|---|---|
ProjectNotification |
Project change stream | Full ProjectWithRelatedInvestigatorsDto or delete |
Subscriber of that project |
ProjectSummariesNotification |
Every Project write in the database | Summary DTO + investigators, or delete when not visible | Every listing subscriber, mapped per connection |
DataExportJobNotification |
DataExportJob change stream | DataExportJobDto |
Job subscriber |
ProjectStatsNotification |
Every Study write in the project, or a Project write | Per-investigator FullStats |
Full-stats subscriber (none in the web) |
StudyPresenceUpdated |
Study change stream (presence fields distinct) | StudyReviewPresenceSnapshot |
Presence subscriber (flag off) |
ProjectStatisticsChanged |
FEAT-024 outbox → bus | Revision strings (invalidation) | FEAT-024 registry |
ActivityClaimRevoked |
Claim-revocation outbox → bus | Own claim DTO | Clients.User |
RuntimeFeatureFlagsRevision |
Admin PUT | long revision (invalidation) |
Clients.All, local pod only |
UiVersionCheck |
On connect | Min UI version | Caller |
InboxChanged (#3932, open) |
Inbox change stream | Empty (invalidation) | User group |
2. Scope¶
| ID | Problem | Impact | Evidence (file:line) | Issue |
|---|---|---|---|---|
| RT1 | Subscribers run inline on the cursor loop; a throwing subscriber reaches the loop's catch, and shouldRetry rejects non-Mongo exceptions |
Live updates stop on that pod. Mechanism VERIFIED; whether Publish() then stays terminated for good is NOT VERIFIED (needs a test) |
MongoContext.cs:262-265,289-306 |
#3974 (backend PR-5) |
| RT2 | Health records success only on an event, the retry budget is cumulative, and the check is tagged live |
Liveness kills pods after the database recovers; after 15 failures with no event in between, the stream is gone | MongoContext.cs:302-310,397-405; ChangeStreamHealthTracker.cs:57-82; Program.cs:167-170 |
#3974 |
| RT3 | CORRECTED. The review says every Study save re-runs the full-stats aggregate for every viewer. The server path is real: CombineLatest → sync Projects.Get + GetFullStatsForInvestigatorAsync per emission, no throttle, unbounded .Merge() (out of order), DateTime.Now. But the web has not invoked SubscribeToProjectFullStats since 2026-04-17 (commit 0cea246dc7: "no client dispatches it"). It is dead code that any member can still call, so it is a cost amplifier |
Latent DoS vector, ~400 lines of server code and tests | AggregateRootEntitySubscriptionManager.cs:301-367,476-537 (:332,340,348,507,516,523); NotificationHub.cs:1305-1349; web grep returns 0 callers |
new |
| RT4 | Head-of-line blocking: the legacy DTO mapping does a sync GetMany per subscriber per event on the cursor thread |
One slow read delays every subscriber of that collection on the pod | NotificationHub.cs:1178,1285; ProjectStandardDto.cs:140 |
new |
| RT5 | Listing fan-out: every Project write in the database is mapped once per listing viewer | Cost grows with writes × viewers × pods | NotificationHub.cs:1236-1303 |
new |
| RT6 | Full 118–465 KB Project documents (UpdateLookup) reach every pod on every Project write | Network and deserialisation on 3–10 pods | MongoContext.cs:221-225; size from repository-cache.md |
new |
| RT7 | Log volume: 3 Info lines per event per subscriber | Log cost; noise | MongoRepositoryBase.cs:169-212 |
new |
| RT8 | Export-job subscriptions are not removed on disconnect | Leak per connection | NotificationHub.cs:112-122 |
#3978 (backend PR-9) |
| RT9 | Reconnect loses subscriptions; no missed-event strategy; DbTime has one-second resolution and no DTO carries Audit.Version |
Silent staleness; clients cannot order two changes inside one second | MongoUtils.cs:18-21; Models/ProjectDto/*.cs (no version field) |
#3976 (client) |
| RT10 | Pushes not driven by change streams reach only the local pod (Clients.All for flag revisions) |
Other pods' clients wait for the 5-minute poll (non-production only) | RuntimeFeatureFlagsController.cs:153 |
new |
| RT11 | Three change-stream idioms: Rx cached stream, the dead blocking variant, and #3932's raw loop with no resume or health | Each fixes the same problems differently | MongoContext.cs:173-201; #3932 InboxInvalidationWorker.cs |
new |
| RT12 | The hub mixes review-session commands, 10 subscription methods and two authority modes in 1,650 lines | Hard to change safely; four open PRs touch it | NotificationHub.cs |
new |
| RT13 | SignalR detailed errors are always on | Exception text reaches clients | Program.cs:381 |
#3978 (backend PR-8) |
3. Options evaluated (web facts checked 2026-10-04)¶
| Option | Facts | Verdict |
|---|---|---|
| Keep Rx.NET, cleaned up | Rx.NET is maintained: 6.0.2 (2025-08-29), 6.1.0 (2025-10-03), 7.0.0 (2026-07-17), targeting net8.0/netstandard2.0 (NuGet). SyRF pins 6.0.0 (SyRF.Mongo.Common.csproj:20, SyRF.SharedKernel.csproj:36) and uses it in 4 production files. |
Viable, but the defects RT1/RT3/RT4 come from Rx idioms (hot shared subjects, inline observers, FromAsync+Merge, two retry operators) that few maintainers can reason about. Cleaning up keeps the per-subscriber pipeline model (RT5/RT6). Not recommended as the target; used for the interim #3974 fix only |
BackgroundService + System.Threading.Channels + IAsyncEnumerable |
All in-box. System.Linq.AsyncEnumerable is in the .NET 10 shared framework (VERIFIED: ~/.dotnet/shared/Microsoft.NETCore.App/10.0.0/System.Linq.AsyncEnumerable.dll; package 10.0.0 2025-11-11). The FEAT-024 and claim-revocation dispatchers and #3932's worker already use BackgroundService loops. TimeProvider makes windows testable |
Recommended. Straight-line await foreach with try/catch, one reader per collection, explicit bounded queues, per-send isolation |
SignalR server-to-client streaming (IAsyncEnumerable/ChannelReader hub methods) |
One producer per invocation; it ends on disconnect and must be re-invoked; no group fan-out | Not adopted: it re-creates per-connection pipelines |
| Notify-then-fetch (invalidation) | Already the house pattern for FEAT-024 (ProjectStatisticsChanged.cs:5-10), flag revisions (feature-flags.md) and #3932. Fits the client plan's resources (reload on event) |
Recommended for project details and listings; small data payloads stay where a refetch would cost more than the payload (§4.4) |
| Stateful reconnect | .NET 8+: server AllowStatefulReconnects on MapHub, StatefulReconnectBufferSize default 100,000 bytes; client withStatefulReconnect() (docs). The buffer lives in the pod, so it needs affinity, and a pod restart still loses it. Protocol support VERIFIED (orchestrator, 2026-10-04): dotnet/aspnetcore main NewtonsoftJsonHubProtocol.cs declares ProtocolVersion = 2 and handles AckMessage/SequenceMessage, so both JSON protocols support it. WebSockets-only is reported by secondary sources; NOT VERIFIED in primary docs |
Recommended on v2, as a latency optimisation, not as the correctness mechanism (versions + resync are) |
| Domain events via the outbox instead of change streams | The outbox is persistence-plan PR-10 (flagged, needs #3973 and PR-5). Events cover only unit-of-work saves: partial updates, raw writes, PM and Quartz writes bypass them. Wolverine switch is targeted for Q3 2027 (ADR-021) | Hybrid: change streams stay the trigger for entity invalidation (they see every committed write); outbox → bus fan-out carries semantic events (statistics, claim revocation, later inbox/attention) |
| Wolverine SignalR transport | WolverineFx.SignalR 6.45.0 (2026-10-02) exists, but hubs must inherit WolverineHub and messages use a CloudEvents envelope (docs) |
Not adopted: incompatible with the typed hub and its client. After S1, Wolverine handlers call the same IRealtimeClients facade |
| Backplane: none / Redis-Valkey / Azure SignalR | Redis backplane (docs) forwards group/user/all sends across servers; sticky sessions are still required. valkey-production exists for the (disabled) BFF. Azure SignalR doesn't apply on GKE |
None. Every source reaches every pod already (per-pod streams, per-pod temporary queues). A backplane adds a hot-path dependency for one gap (RT10), which a per-pod bus fan-out fixes. Revisit triggers in D4 |
4. Target architecture¶
flowchart LR
subgraph pod[Each API pod]
R1[ChangeFeedReader<Project>] --> Q[(bounded Channel<EntityChange>)]
R2[ChangeFeedReader<Study>] --> Q
R3[ChangeFeedReader<DataExportJob>] --> Q
R4[ChangeFeedReader<InboxNotification>] --> Q
B[per-pod bus consumers: statistics, claim revocation, flag revision] --> C
Q --> D[RealtimeDispatcher: coalesce, revoke, authorise]
D --> C[IRealtimeClients facade]
REG[(RealtimeRegistry: topic → connections)] <--> D
H[RealtimeHub /notifications/v2] --> REG
C --> H
end
4.1 Change feed readers¶
- One
ChangeFeedReader<T>BackgroundServiceper watched collection per pod, started at host start when the v2 flag is on (so health is meaningful) and running regardless of subscribers. - Pipeline: op-type
$matchplus a$projectthat keeps only the routing fields:_id,ProjectId,Audit.Version, plus for Study the fieldsStudyReviewPresenceSnapshot.FromStudyreads, and for DataExportJob its status fields. Saves are replaces (MongoExtensionsReplaceOneAsync), which always carryfullDocument, so the projection trims them.UpdateLookupstays only for partialupdateops. - Loop:
await foreachover the cursor,TryWriteinto a bounded channel. Retry with uncapped backoff (2 s → 180 s ceiling, jitter). Healthy when the cursor opens. - Resume token kept in memory.
- An invalid or lost token, or any restart without a token, emits a
ResyncRequiredfor that collection, because events may have been missed. - No subscriber code runs on the reader thread. A full channel coalesces rather than blocking the cursor.
4.2 Registry and dispatcher¶
RealtimeRegistry(singleton):topic → connectionId → Subscription(subject, impersonationActor, hubKind, subscribedAt), plus a reverse index per connection for disconnect cleanup and per-user lookup. It is the only subscription store for v2. SignalR groups are not needed without a backplane.RealtimeDispatcher(oneBackgroundServicereading the channel):- Map each change to topic keys.
- Coalesce per topic, latest wins, with a trailing window (§4.5).
- Skip topics with no local subscriber.
- For a Project change with local subscribers, run one uncached Project read per window per pod and evaluate revocation for every subscriber in memory (
RevokedBy/IsDisabled, which exist today). Drop revoked subscriptions and sendSubscriptionRevoked. - Send through
IRealtimeClients. Each send has its own try/catch, and a failed socket never affects others. IRealtimeClientsfacade:Connection(id),Connections(ids),User(userId),All(). During the transition it fans out to both hub contexts (v1NotificationHub, v2RealtimeHub), so existing producers (statistics consumer, claim revocation, flag revision, #3932 inbox worker) reach clients on either path with a one-line change each.
4.3 Authorization¶
| Moment | Rule |
|---|---|
| Subscribe | Unchanged policies on typed hub methods (SignalRAuthorizationHandler); under enforced mode, RequireCurrentAuthorityAsync as today |
| Revocation | Driven by the Project change (4.2 step 4) and by subject suspension: deny on the next invocation (SubjectAdmissionHubFilter, exists), and D8: abort the subject's local connections on a suspension event |
| Delivery, invalidation messages | No per-delivery read. They carry ids and versions only, and the refetch is authorised over HTTP. Same stance as FEAT-024 and #3932 |
| Delivery, payload messages (presence snapshot, export status, claim revocation) | Per-delivery decision kept under enforced mode (M4b parity), batched per event |
4.4 Payload versus invalidation per message¶
| v2 message | Kind | Why |
|---|---|---|
ProjectChanged |
Invalidation | Project DTO is large and per-user; the client reloads its route resource |
ProjectListingChanged |
Invalidation, filtered by visibility (public, active member, pending invite or join request) so private project ids never leak | Today's per-user mapping cost (RT5) |
ExportJobChanged |
Status data (small) | Progress ticks would double the load as refetches; the client fetches the full job on a terminal status (GET api/projects/{p}/data-export/{j} exists, DataExportController.cs:75) |
StudyPresenceChanged |
Data (existing snapshot) | Small and latency-sensitive; no presence GET exists |
ProjectStatisticsChanged, ActivityClaimRevoked, InboxChanged, RuntimeFeatureFlagsRevision, UiVersionCheck |
Unchanged | Already right-sized |
4.5 Throttling, delivery semantics, reconnect¶
- Coalescing windows (PROPOSAL): Project 1 s, listing 2 s, export 1 s, presence 250 ms. Each send carries the newest version seen.
- Delivery is at-most-once per connection. Correctness comes from three rules:
- (a) every subscribe returns a
SubscriptionAckcarrying the current version where one exists; - (b) the client compares versions and reloads on any doubt;
- ©
ResyncRequiredafter server-side loss. versionisAudit.Version. It is advisory until persistence-plan #3985 makes direct writes bump it; until then clients treat any invalidation as "reload".atis the change's cluster time as a"t.i"string, for diagnostics and ordering inside one second.- Reconnect:
- v2 allows stateful reconnect, so a blip on the same pod replays the buffer;
- a new connection (or a pod restart) re-declares its topics (client plan
RealtimeStore), and the acks drive reloads; - review sessions keep their own rejoin protocol.
4.6 Observability and health¶
- Meter
SyRF.Api.Realtimewith bounded tags only (collection, message, outcome; never ids): realtime.change_feed.events,.restarts{reason},.lag_seconds(now − cluster time);realtime.channel.depth,realtime.coalesced;realtime.subscriptions{topic_kind},realtime.sends{message,outcome};realtime.resyncs.- Health check
realtime-change-feeds, taggedreadyonly: Degraded while retrying, Unhealthy after 60 s of retrying (PROPOSAL), Healthy once the cursor opens. Neverlive. - Logging through
[LoggerMessage], Debug per event, Information per reader state change.
5. Hub contract v2 (for the client plan)¶
Endpoint: RealtimeHub : Hub<IRealtimeClient> mapped at /notifications/v2 with AllowStatefulReconnects = true. This path reuses the existing ingress prefixes (staging's Web-host /notifications rule) and query-token auth (DirectBearerSchemeSelector.cs:82 uses StartsWithSegments("/notifications")), both VERIFIED. It also doesn't collide with MapHub("/notifications") (to be proved by an R1.3 test).
Protocol (decided by Chris, 2026-10-04): the System.Text.Json hub protocol (AddJsonProtocol), with serializer options identical to the API's migrated MVC options (camelCase, the same enum, null and converter policy). This makes the API-wide Newtonsoft → STJ migration (system-text-json-migration-plan-2026-10.md) a prerequisite of R1.3b.
- Today the hub uses AddNewtonsoftJsonProtocol with three custom converters and a camelCase resolver (Program.cs:390-400), VERIFIED.
- Both protocols register under the protocol name json, so one host cannot offer Newtonsoft on v1 and STJ on v2 side by side. VERIFIED 2026-10-04 against dotnet/aspnetcore main: NewtonsoftJsonHubProtocol.cs:39 and JsonHubProtocol.cs:49 both declare ProtocolName = "json". The migration therefore switches the v1 hub's protocol too, as part of the same API-wide change.
- Every v2 payload type (StudyReviewPresenceSnapshot, StudyReviewAccessResult, ActivityClaimRevokedDto, ProjectStatisticsChanged, and the R4 ProjectStatisticsSnapshot sections) gets an STJ round-trip test, either in the STJ plan or in R1.3b, whichever lands first.
- Stateful reconnect works on either protocol (§3, VERIFIED), so the protocol choice does not gate D7.
Versioning:
- contractVersion = 2 is returned in every ack.
- Additive changes only within v2 (new messages, new optional fields). Clients ignore unknown messages and fields.
- Any rename, removal or type change means /notifications/v3.
- The contract is pinned by a .NET snapshot test of method signatures and message JSON, committed as docs/architecture/realtime-contract-v2.md plus a JSON sample file that a web spec reads (D9).
5.1 Client → server methods¶
| Method | Policy | Returns | Notes |
|---|---|---|---|
SubscribeProject(Guid projectId) |
ProjectViewSignalRPolicy |
SubscriptionAck |
Also registers the FEAT-024 statistics subscription (filter updated to the v2 name) |
UnsubscribeProject(Guid projectId) |
none | void | Idempotent; never throws |
SubscribeProjectListing() |
authenticated | SubscriptionAck |
One per connection |
UnsubscribeProjectListing() |
none | void | Idempotent |
SubscribeExportJob(Guid exportJobId) |
ProjectExportDataSignalRPolicy (job → project) |
SubscriptionAck with the current status |
|
UnsubscribeExportJob(Guid exportJobId) |
none | void | Idempotent |
SubscribeStudyPresence(Guid projectId, Guid stageId, Guid studyId) |
ProjectViewStudiesSignalRPolicy |
PresenceAck (ack + current snapshot) |
No-op ack when tracking is unavailable |
UnsubscribeStudyPresence(Guid projectId, Guid stageId, Guid studyId) |
none (today it needs View, NotificationHub.cs:1076) |
void | Idempotent |
SubscribeProjectStatistics(Guid projectId, Guid? stageId) (R4, additive) |
ProjectViewStudiesSignalRPolicy (member content, as the legacy full stats), plus the per-delivery member-content check in every mode (R4 item 7) |
SubscriptionAck (+ a snapshot read fresh for this caller when served) |
Denied with reason: "featureUnavailable" (an additive value on top of the CR-3 reasons) when realtimeStatisticsPush is off or the project is not allowlisted. stageId selects the per-stage (Stage Overview) read |
UnsubscribeProjectStatistics(Guid projectId, Guid? stageId) (R4) |
none | void | Idempotent |
JoinStudyReview, LeaveStudyReview, StartedAnnotating, StoppedAnnotating, Heartbeat |
Unchanged | Unchanged (StudyReviewAccessResult, bool) |
Same signatures; both hubs delegate to one extracted ReviewSessionCommands service |
Not carried over: SubscribeToProjectFullStats (RT3); R4's SubscribeProjectStatistics replaces it on a different design. Groups are not part of the contract; topics are server-internal (project:{id}, listing, export:{id}, presence:{p}:{s}, user:{id}).
5.2 Server → client messages¶
| Message | Shape (camelCase) | Kind |
|---|---|---|
ProjectChanged |
{ projectId, change: "updated" \| "deleted", version?: number, at: string } |
Invalidation |
ProjectListingChanged |
{ projectId, change: "updated" \| "removed", version?: number, at: string } |
Invalidation |
ExportJobChanged |
{ exportJobId, projectId, version?: number, status: string, progress?: number, terminal: boolean, at: string }; the exact status fields are fixed in R1.3 from DataExportJobDto |
Data (status) |
StudyPresenceChanged |
{ projectId, studyId, version: number, snapshot: StudyReviewPresenceSnapshot } |
Data |
SubscriptionRevoked |
{ topic: "project" \| "listing" \| "export" \| "presence", id?: string, reason: "accessRevoked" \| "deleted" } |
Control |
ResyncRequired |
{ scope: "all" \| "project" \| "listing" \| "export" \| "presence", id?: string, reason: "feedRestarted" \| "resumeLost" } |
Control |
ProjectStatisticsChanged |
Unchanged (ProjectStatisticsChanged.cs) |
Invalidation (FEAT-024-owned) |
ActivityClaimRevoked |
Unchanged | Data |
InboxChanged |
Unchanged from #3932 (no arguments) | Invalidation |
RuntimeFeatureFlagsRevision |
Unchanged (long) |
Invalidation |
UiVersionCheck |
Unchanged (string) |
Control |
ProjectStatisticsSnapshot (R4, additive) |
{ projectId, stageId?: string, revision: string, globalRevision: string, sections: { projectScreening?, stageAnnotation?, membershipScreening?, membershipAnnotation? }, at: string }. Each section is the HTTP bundle section for this caller (byte-equal, R4 4.1.3) plus provenance { servedFrom: "materialized", checkpointId, pendingFingerprint, watermarks }. revision is the Int64 project ClientInvalidationRevision and globalRevision the Int64 GlobalClientInvalidationRevision, both read in the same snapshot, sent as strings and compared as integers (the client's pollCoherentStatistics tracks both). A section is present only when its bundle is Fresh and the HTTP consumer flag for it is on, and only for sections backed by the two coherent DTOs (StageOverviewStatistics, ScreeningOverviewStatistics; R4 item 10). Exact field names are fixed in R4.1 with the statistics session |
Data (shaped per subscriber, as the materialised HTTP read) |
SubscriptionAck = { contractVersion: 2, topic, id?: string, version?: number, at: string }.
Client obligations: subscribe before fetching, or compare the ack version with the fetched one. Re-declare every topic on a new connection ID. Reload on ResyncRequired. Never treat an absent message as "unchanged".
5.3 Reconciliation with the client plan (orchestrator, 2026-10-04)¶
Each client requirement (CR-n in realtime-client-plan-2026-10.md §9) is resolved below. These resolutions are part of contract v2 from its first release.
| CR | Resolution |
|---|---|
| CR-1 invalidation shape | Adopt the server's per-family messages (ProjectChanged, ProjectListingChanged) rather than a generic { entityType, … }. The client's message typing maps families to resources. |
CR-2 per-subscription seq |
Deferred (additive later). Gaps are covered by stateful reconnect, ResyncRequired (feed restart or lost resume token), refetch-on-new-connection and the version > check. If staging shows missed updates, add an optional seq to v2 messages; this is additive. The client's RV.2 is deferred with it. |
| CR-3 ack or denial as a return value | Adopt. SubscriptionAck gains granted: boolean, reason?: "notAMember" \| "notFound" \| "disabled" \| "authorityUnavailable" and retryable: boolean. Subscribe methods return a denial instead of throwing a HubException, and policy failures are mapped to reason. Turning off EnableDetailedErrors in production (backend plan #3978) is independent. |
| CR-4 batch resubscribe and full cleanup | Adopt. Add Resubscribe(SubscriptionRequest[] requests) → SubscriptionAck[] (SubscriptionRequest = { topic, id?, stageId?, studyId? }). The registry drops every topic for a connection on disconnect, including export jobs; R1.2 gets an acceptance criterion for it. |
| CR-5 revocation notice | Covered by SubscriptionRevoked. |
| CR-6 stateful reconnect | Covered (AllowStatefulReconnects = true). Closed, 2026-10-04: NewtonsoftJsonHubProtocol declares ProtocolVersion = 2 and handles AckMessage/SequenceMessage (VERIFIED by the orchestrator in dotnet/aspnetcore main), so stateful reconnect does not depend on the protocol. v2 nonetheless uses the STJ protocol by Chris's decision (§5 "Protocol"). R1.3b keeps one stateful-reconnect test on STJ. |
| CR-7 delete versions | ProjectChanged { change: "deleted" } carries version when known. It stays advisory until #3985 makes every write bump Audit.Version. |
| CR-8 contract manifest | Covered by the snapshot test plus docs/architecture/realtime-contract-v2.md and the JSON sample the web spec reads. contractVersion arrives in every ack, so no separate on-connect message is needed. |
| CR-9 ordered delivery | Covered: one dispatcher per pod and per-topic coalescing; the full-stats .Merge() path is removed in R0. The client still applies version > checks. |
| CR-10 user-group delivery without a backplane | Covered by D4 provided the source reaches every pod. #3932's InboxChanged reader must run on every API pod (it does today as a per-pod change stream). Any future user-group message sent from a single consumer needs a fan-out endpoint, like ProjectStatisticsChanged. |
| CR-11 keep-alive values | The server uses the ASP.NET Core defaults today (AddSignalR sets no intervals: Program.cs:379-389): keep-alive 15 s, client timeout 30 s. The contract document states them, and the client sets withKeepAliveInterval(15 s) and withServerTimeout(30 s). Changing them later is a contract change on both sides. |
Flags: keep both flags, with a catalogued dependency. realtimeHubV2 (api + web) turns on the v2 pipeline and the /notifications/v2 hub per environment. realtimeClientV2 (web, page reload) moves browsers to it, and the flag dependency catalog records that it requires realtimeHubV2. Rollout order per environment: server flag on, then the dark-run checks, then the client flag on. Rollback is the reverse.
6. MVP, flag decision and common criteria¶
MVP boundary: R0 + R1 + R2. The v2 server path runs dark behind the flag, the client plan's pilot consumes it, and staging runs on v2 for a soak. Out of scope (tracked where noted): - client code (state-management plan §5b and its follow-up client plan); - #3974's v1 fix (backend PR-5); - export-disconnect leak and detailed errors (backend PR-9/PR-8, #3978); - outbox (persistence PR-10); - Wolverine (ADR-021); - the API-wide Newtonsoft → STJ migration (a prerequisite, planned in system-text-json-migration-plan-2026-10.md); MessagePack is not planned; - R4 (full-stats push, staging only) is a follow-up release outside the MVP; - cross-pod flag revisions (RT10; follow-up issue, per-pod bus fan-out).
Flag decision: flagged. The change is a hot-path rewrite with changed observable semantics, so it needs an environment-scoped kill switch. One catalogued flag realtimeHubV2 (services api + web tsName, default false, via env-mapping.yaml and pnpm run generate:flags):
- The API reads the deployed value at start-up to start the readers and dispatcher; no runtime override starts or stops them, as with the fold flag.
- The web uses the runtime value to choose /notifications/v2 or /notifications. Both hubs are always mapped, so either path is selectable per environment and rollback is a flag flip.
- R0 is not flagged: it removes dead or broken behaviour with no web caller (D5).
- R4 sub-flag realtimeStatisticsPush (api + web, default false). Its catalogue entry requires realtimeHubV2, materializedProjectStatisticsSignalR, materializedProjectStatisticsServing and the per-section consumer flags the HTTP endpoints check (for example …StageOverviewBundle, …ScreeningInfo, …ProjectOverview); R4 also checks those flags at run time on every pass (statistics session review, 2026-10-05). Why a sub-flag rather than reusing realtimeHubV2:
- it has a different owner (FEAT-024) and a different risk (a payload carrying reviewer data);
- it is staging-only while realtimeHubV2 is meant to reach production;
- it needs its own kill switch so the statistics session can turn it off without taking the v2 hub down.
Production values stay false, and enabling it in production is outside this plan.
| # | Common acceptance criterion (every PR) | Verification |
|---|---|---|
| C1 | With realtimeHubV2 off, v1 behaviour is unchanged: existing SyRF.API.Endpoint.Tests/SignalR/* pass untranslated (except files R0/R3 delete) |
dotnet test --filter FullyQualifiedName~SignalR |
| C2 | No subscriber or send can throw into a reader or the dispatcher loop | Unit test with a throwing fake sender |
| C3 | Metrics carry only bounded tags; no project, user or connection id | Meter tag-allowlist test (pattern from ProjectStatisticsMetrics) |
| C4 | New code is in BackgroundService/Channel/TimeProvider idioms; no new System.Reactive use |
Review; grep in the PR |
| C5 | Docs in the same PR: docs/architecture/realtime-contract-v2.md and docs/features/signalr-active-reviewer-tracking.md where touched; flag decision stated |
Docs review |
| C6 | Narrow, niced test runs only (this machine is the CI host); settled review on the head; green checks | PR body; pr-review-settled.sh |
7. Releases → PRs¶
R0: remove dead and hot-path waste (parallel; not flagged)¶
PR-R0.1 Retire the full-stats push (RT3; effort S–M). The statistics session confirmed on 2026-10-05 that
there is no web caller and that FEAT-024 does not depend on this path (ProjectStatisticsSubscriptionFilter records
only SubscribeToProject/UnsubscribeFromProject; the outbox, ProjectStatisticsDispatchService and the sink never
touch AggregateRootEntitySubscriptionManager). MembershipScreeningVisibilityShaper survives (it is still used by
ProjectController.cs:1187-1190, ReviewController.cs:1205-1210, ReviewerScreeningHistoryController.cs:70 and
StageOverviewStatisticsQuery.cs:220,324), and so does GetFullStatsForInvestigatorAsync (ProjectController.cs:1182).
The file list below includes the gaps that review found (statistics session review, 2026-10-05). Files:
- NotificationHub.cs:1305-1349;
- AggregateRootEntitySubscriptionManager.cs:279-410,463-537 (the full-stats parts only);
- INotificationHubClient.ProjectStatsNotification;
- RealtimeSurface.FullStatistics (SignalR/Authority/RealtimeDeliveryAuthority.cs:50), which has no other
production use after R0.1;
- the dead blocking stream GetCollectionChangeStream with its only caller MongoChangeStream (MongoContext.cs:173-201, MongoRepositoryBase.cs:41-46), and the unused DebugSubscribe/DebugSubscribeObservable (MongoContext.cs:366-370,435-461); 0 callers each (grep).
Tests to delete or translate (every reference to the full-stats method, message or surface):
- SignalR/FullStatsPushVisibilityTests.cs (deleted, after the port below);
- SignalR/NotificationHubCurrentAuthorityTests.cs (153-294, 373);
- SignalR/NotificationHubDisabledMembershipRevocationTests.cs (137-249, 385);
- Authorization/LegacyStatisticsAuthorizationTests.cs:187-191,228-229;
- Authorization/EndpointAuthorizationCatalogTests.cs:210;
- Authorization/EndpointAuthorizationExceptions.cs:269;
- ApiSensitiveLogTests.cs:345-348;
- SignalR/RealtimeDeliveryAuthorityTests.cs:232.
Visibility-test port inside R0.1. The shaper rules that FullStatsPushVisibilityTests exercises are ported
in this PR, before the file is deleted, into the shaper's own tests (for example
SyRF.ProjectManagement.Core.Tests/ProjectStatistics/Visibility/MembershipScreeningVisibilityShaperTests.cs), so
the surviving HTTP callers keep their coverage. R0.1 ships releases before R4, so the port cannot wait for R4.
These cases encode the legacy shaper's rules, which disclose more than the materialised read does (the
materialised filter withholds peer rows without ViewMemberships, where the legacy shaper names or anonymises
them; #3642). R4 must not inherit them as its expected output; R4's expected output is the materialised HTTP read
for the same caller (R4 AC 4.1.3).
E2E: e2e/authority/tests/authority-m1a-application-roles.spec.ts:697-799 (calls at 697, 734, 771, 781 and 799)
invokes SubscribeToProjectFullStats on two replicas and asserts ProjectStatsNotification pushes. It is the
cross-replica M4b proof for a member-content payload stream, and it fails after R0.1. R0.1 rewrites it, for
example to prove the same cross-replica late-query barrier on the ProjectNotification stream only, with the
application-authority owner's agreement recorded in the PR.
Docs updated in the same PR: docs/features/materialized-project-statistics/STATUS.md:417 ("old push") and
:939-942 (coordinated with the statistics session, which owns STATUS); phase0-calculation-consumer-catalogue.md:102-103,
910, 935, 1081-1094 (including option "© Wire SubscribeToProjectFullStats up"); docs/how-to/application-authority-mode.md:232.
Server only; the harmless web listener (ProjectStatsNotification → stageActions.receivedFullStats) goes with
the client plan.
| # | Criterion | Verification |
|---|---|---|
| 0.1.1 | Invoking SubscribeToProjectFullStats → HubException "method does not exist"; nothing aggregates |
Hub unit test |
| 0.1.2 | Every other hub method and message still works | Existing SignalR tests (C1), with the 7 test files above translated or deleted |
| 0.1.3 | Web build and specs unchanged | CI Test Web (Angular) |
| 0.1.4 | A grep of src (*.cs) for SubscribeToProjectFullStats, ProjectStatsNotification and RealtimeSurface.FullStatistics → no match |
grep in the PR body |
| 0.1.5 | Every visibility case from FullStatsPushVisibilityTests → has a ported equivalent against MembershipScreeningVisibilityShaper before the file is deleted; the PR lists the mapping |
Unit (MembershipScreeningVisibilityShaperTests) plus the mapping table in the PR body |
| 0.1.6 | The rewritten authority-m1a-application-roles cross-replica test → passes, and still proves a member-content stream is refused at delivery after revocation on the other replica; the application-authority owner has approved the rewrite |
bash e2e/authority/run.sh (hermetic); approval in the PR |
| 0.1.7 | The listed docs → no longer describe the legacy push as live or as an option to wire up | Docs review; validate-docs --skip-indexes |
PR-R0.2 Per-event logging (RT7; effort S): MongoRepositoryBase.cs:169-212 to Debug via [LoggerMessage].
| # | Criterion | Verification |
|---|---|---|
| 0.2.1 | One change event with 10 subscribers → 0 Information-level log lines from the fan-out | Log-capture unit test |
R1: v2 server foundation (dark; flag off)¶
PR-R1.1 Change feed readers (RT1, RT2, RT6, RT11; effort M). Files:
- new src/services/api/SyRF.API.Endpoint/Realtime/ChangeFeed/* (reader, EntityChange, projections, health, metrics);
- API Program.cs registration;
- env-mapping.yaml flag.
| # | Criterion | Verification |
|---|---|---|
| 1.1.1 | Flag on; a Study save → exactly one EntityChange with id, project id, version, and nothing outside the projection |
Testcontainers replica set |
| 1.1.2 | Primary stepdown mid-stream → the reader resumes from its token, and no change committed in between is missed | Testcontainers (replSetStepDown) |
| 1.1.3 | Resume token rejected (code 286 or 260) → the stream restarts at now, and ResyncRequired(collection) is emitted once |
Unit (fake cursor) |
| 1.1.4 | 20 failures, then the cursor opens with no events → health Healthy within one poll; it never reports to live |
Unit (FakeTimeProvider); TestServer /health/live stays 200 |
| 1.1.5 | Flag off → no reader starts and no cursor opens | Unit over the registration |
| 1.1.6 | Wire size per Study change ≤ 1 KB (PROPOSAL) | Testcontainers command monitor |
PR-R1.2 Registry, dispatcher and IRealtimeClients (RT4, RT5, RT10 partly; effort M). Files:
- new Realtime/RealtimeRegistry.cs, RealtimeDispatcher.cs, IRealtimeClients.cs;
- one-line producer switches in ProjectStatisticsChangedConsumer.cs, ActivityClaimRevokedConsumer.cs, RuntimeFeatureFlagsController.cs (and the #3932 worker if merged), coordinated with their owners.
| # | Criterion | Verification |
|---|---|---|
| 1.2.1 | 50 writes to project P within 1 s → each P subscriber receives 1–2 ProjectChanged, the last carrying the newest version |
Unit (FakeTimeProvider) |
| 1.2.2 | One Project change with 100 local subscribers → exactly 1 uncached Project read on that pod | Unit (read counter) |
| 1.2.3 | Membership disabled → that subscriber gets SubscriptionRevoked(project, accessRevoked) and nothing later for P; other members are unaffected |
Unit + Testcontainers |
| 1.2.4 | Private project Q changes → no listing subscriber without visibility receives Q's id | Unit (visibility matrix) |
| 1.2.5 | One socket's send throws → every other recipient is still sent to, and the dispatcher keeps running | Unit (C2) |
| 1.2.6 | Statistics, claim revocation and flag revision → delivered to a v1 connection and to a v2 connection | Unit over the facade with both fake hub contexts |
| 1.2.7 | Disconnect → every registry entry for the connection is removed | Unit |
PR-R1.3a Extract review-session commands (RT12; effort M; behaviour-preserving, not flagged). Move NotificationHub.cs:143-1040 into ReviewSessionCommands (an application-level service) and keep thin hub methods.
| # | Criterion | Verification |
|---|---|---|
| 1.3a.1 | NotificationHubStudyPresenceTests (1,999 lines) and ReviewSessionConnectionTests pass unchanged against the thin hub |
dotnet test --filter |
| 1.3a.2 | Hub file ≤ 700 lines (PROPOSAL) | wc -l in the PR |
PR-R1.3b RealtimeHub and contract v2 (§5; effort M). Files:
- new Realtime/RealtimeHub.cs and IRealtimeClient.cs;
- MapHub("/notifications/v2", o => o.AllowStatefulReconnects = true), on the STJ protocol the migration has configured (prerequisite);
- ProjectStatisticsSubscriptionFilter.cs (accept SubscribeProject, record the hub kind; statistics-session sign-off);
- contract doc and snapshot test.
| # | Criterion | Verification |
|---|---|---|
| 1.3b.1 | Negotiate and connect at /notifications/v2 with a query token → connected; /notifications still works on the same host |
WebApplicationFactory integration |
| 1.3b.2 | Snapshot of method names, parameters and message JSON → matches the committed contract; any change fails the test | Contract test |
| 1.3b.3 | A non-member calls SubscribeStudyPresence → HubException, and no registry entry is created |
Hub unit test |
| 1.3b.4 | Unsubscribe calls for unknown topics → no exception | Hub unit test |
| 1.3b.5 | SubscribeProject → the ack carries contractVersion 2 and the project's current Audit.Version |
Integration |
| 1.3b.6 | A v2 connection on an allowlisted pilot project receives ProjectStatisticsChanged |
Integration (FEAT-024 harness) |
| 1.3b.7 | A v2 connection negotiates the json protocol → payloads are serialised by STJ with the migrated MVC options (camelCase; same enum and null handling) |
Integration, comparing a hub message with the same DTO from an MVC response |
| 1.3b.8 | A client with stateful reconnect drops its socket for 5 s on the same pod → buffered messages arrive once, in order, with no new connection ID | Integration (HubConnection with WithStatefulReconnect) |
R2: parity proof and staging (with the client plan)¶
PR-R2.1 Hermetic E2E parity (effort M). The e2e/ stack runs the realtime specs twice, flag off and on (two-context project edit, listing update, export completion, presence, reconnect).
| # | Criterion | Verification |
|---|---|---|
| 2.1.1 | Both flag values → the same specs pass | E2E (hermetic) |
| 2.1.2 | Context B offline 10 s (PROPOSAL) while A edits P, then online → B shows A's edit without a page reload | E2E |
| 2.1.3 | API process restarted mid-session → the client reconnects, re-subscribes and converges | E2E (restart the API process in the stack) |
PR-R2.2 Staging switch (cluster-gitops staging/api/values.yaml flag; effort S). Staging is not mission critical, so no window.
| # | Criterion | Verification |
|---|---|---|
| 2.2.1 | 7 days on v2 (PROPOSAL) → zero reader terminations; lag_seconds p95 ≤ 2 s; no API restarts attributed to realtime |
Staging rollout check (GMP metrics) |
| 2.2.2 | Flag flipped back → clients use v1 after reload with no errors | Staging rollout check |
Production rollout = a separate Chris-approved step (cluster-gitops production values), after R2.2.
R3: retire v1 (after production has run on v2 for a Chris-approved period)¶
PR-R3.1 (effort L). Delete:
- AggregateRootEntitySubscriptionManager.cs, ConnectionKeyedSubDictionaryCache.cs;
- the v1 subscription methods (the v1 hub is reduced to review-session delegation, or removed once MinUiVersion forces new bundles);
- GetCachedCollectionChangeStream, ChangeStreamDictionaryByType, the Rx RetryWithExponentialBackoff in both libraries;
- ICrudRepository.Get*NotificationStream (ICrudRepository.cs:36-37);
- the System.Reactive package references.
Translate or delete the Rx-based tests (CurrentAuthority, DisabledMembershipRevocation, StudyPresence, Logging, ApiSensitiveLogTests, and the Mongo.Common stream and retry tests).
| # | Criterion | Verification |
|---|---|---|
| 3.1.1 | System.Reactive is gone from every csproj, and no file references System.Reactive |
grep; build |
| 3.1.2 | The behaviours those tests asserted (revocation, delivery-time authority for payloads, presence) are asserted against v2 | Test diff review |
| 3.1.3 | The flag is removed from env-mapping.yaml, and the generator output is updated |
Generator diff |
R4: full-stats push on v2 (follow-up; staging only; sub-flag realtimeStatisticsPush)¶
Production rollout is not in scope of this plan (Chris, 2026-10-04: "leave prod for now"). R4 needs the FEAT-024 statistics session's agreement before it starts, because it reads their rows through a reader entry point they own and hooks their consumer. The statistics session reviews the R4 PR before it merges.
The design was rewritten on 2026-10-05 to meet the ten correctness conditions in the statistics session's review
(marked S1–S10 below). The earlier design (stored rows through ProjectStatisticsServingGate, shaped once per
Named/Anonymised/None visibility class) is superseded: it would have pushed values older than the HTTP
read in fold mode, and the three classes do not partition what the materialised read discloses.
PR-R4.1 ProjectStatisticsSnapshot push (effort M–L; part of the reader work is the statistics session's). Design:
1. Trigger, with no hot-path cost (S6). The existing per-pod ProjectStatisticsChangedConsumer, after its
invalidation send, calls IRealtimeStatisticsPush.Enqueue(projectId, clientInvalidationRevision). Enqueue is
O(1) and non-blocking (it updates an in-memory per-project entry and never awaits), because the consumer's
concurrency is bounded to about 1 message per pod (ProjectStatisticsChangedConsumerDefinition). Nothing is added
to the save path. The push fires only after a materialised update has committed and been published through the
FEAT-024 outbox; no change stream and no Study-save trigger.
- Global changes (S5). A ProjectStatisticsChanged with ProjectId == null (mode or kill-switch
transitions) never becomes a per-project read storm: the pod stops pushing and may send featureUnavailable to
its statistics subscribers. Pushing resumes on the next per-project change whose fresh flags allow it.
2. Coalesce per project with a trailing window of 2 s (PROPOSAL). The pod skips the pass when the
coalesced maximum trigger revision is ≤ the revision it last read for that project, because that earlier read
already covered it (no data is lost: new pending entries always arrive with the next fold, at a higher revision).
Coalescing keeps the maximum revision
seen, because consumer delivery can be out of order (S8).
3. Per-pod reads only with subscribers (S7). A pod reads only when it has at least one statistics:{projectId}
subscriber, so N API pods make at most N reads per coalesced window. The registry is in memory per pod and is lost
on restart, so the client resubscribes on reconnect (#3976, CR-4).
4. Flags per pass (S5, S9). Each pass resolves flags in a fresh scope, as the dispatcher does ("each pass gets
fresh scoped flags"), so runtime overrides and the global kill switch apply at once. A section is included only
when the same consumer flag the HTTP endpoint checks is on (for example
materializedProjectStatisticsStageOverviewBundle, …ScreeningInfo, …ProjectOverview), together with …Pages,
…Writes, …Serving, the family flags and the allowlist. Otherwise HTTP would serve legacy values while the push
served materialised ones.
5. One read through the HTTP reader, fold overlay included (S1). The pass reads through the FEAT-024 bundle
reader (ProjectStatisticsBundleReader, or IStageOverviewStatisticsQuery, which already serves the four Stage
Overview families from materialised rows) in one pinned snapshot. In fold mode that serves stored row +
pending entries, exactly as the HTTP read does (ReadScopesAsync overlay; async-point-fold-design.md §4). A
stored-row read through ProjectStatisticsServingGate alone is not used.
- Coherent-bundle rule: request exactly the selection sets the HTTP consumers request. If a bundle falls
back, omit every section from it. Never mix a materialised section with a section HTTP serves
authoritatively.
- CORRECTED (statistics session review): "reviewer and annotation sections still come from the authoritative
aggregation" is true only for FullStats (ProjectScreeningFullStatsSubstitution.cs:8-18). The Stage Overview
bundle and ScreeningInfo consumers already serve materialised stageAnnotation and membership sections, and they
are on in staging. R4 never calls the authoritative aggregation.
- The read is bounded by the reader's caps (128 Studies / 1,024 pending entries, response budgets).
6. Per-caller authorisation and shaping through a new FEAT-024-owned reader entry point (S2). The materialised
HTTP read is per caller: ReadCurrentAsync loads an authorization context for request.Caller inside the
snapshot, applies ProjectStatisticsAuthorizationFilter.FilterCurrent and honours RequireActiveMembership
(ProjectStatisticsBundleReader.cs:146-161). That filter withholds peer rows without ViewMemberships, where
the legacy shaper names or anonymises them (#3642). Stage Overview adds the per-stage
StageActivity.ViewAnnotationProgressGraph, and the project-wide read never carries annotation
(StageOverviewStatisticsQuery.cs:217-226). So the three shaper classes are not a sufficient partition, and the
per-class topic groups are dropped. Instead:
- the statistics session adds a reader entry point (no existing method fits): one system-scoped snapshot read
per pass, then the same FilterCurrent and shaper decision computed per subscriber, in memory, from
that snapshot's Project;
- a differential test proves the entry point's per-caller output is byte-equal to
ReadCurrentAsync/the HTTP endpoint for the same caller and snapshot;
- the payload has a stage dimension: a subscription is for a project and, optionally, a stage, because the
Stage Overview read is per stage;
- anonymised labels follow the HTTP contract (per response), so the client replaces a section wholesale.
7. Per-subscriber delivery authority in every mode (S3). After the read and before each send (the late-query
barrier), each subscriber is checked for member-content scope (ProjectViewStudies/ActiveMembers) against
the fresh Project facts, including the impersonation actor (RealtimeSubscriber.ImpersonationActorOf). The check
applies in Off and shadow modes as well as enforced mode. The invalidation consumer's View-scope recheck
(ProjectStatisticsChangedConsumer.cs:30-31,47-49) is too weak for a payload and is not reused. A failed check
drops the subscription with SubscriptionRevoked.
8. Never push a non-Fresh or fallback value (S4). A section is omitted whenever its bundle is Fallback or
Unavailable (disabled, fenced, inclusion or definition rewrite in progress, durable-mode disagreement,
PendingOverlayUnavailable, Incompatible); the client keeps its HTTP read for it. Each section carries
provenance: servedFrom: "materialized", the bundle watermarks/CheckpointId and PendingFingerprint, so a
parity check can compare push and HTTP output.
9. No writes; clean callbacks (S6). The pass writes nothing: no receipts, outbox, control or cache rows. The
snapshot transaction is closed before any SignalR send (invariant 14, "clean callbacks"). Metric tags stay
bounded (C3).
10. Revisions and the equal-revision rule (S8). revision is the Int64 control ClientInvalidationRevision
and globalRevision the Int64 GlobalClientInvalidationRevision, both read in the same pinned snapshot as the
sections, not taken from the triggering message (the read can be up to 2 s later and include newer folds). They
travel as strings and are always compared as integers, never as strings. The snapshot carries both because the
client's pollCoherentStatistics tracks both (coherent-statistics-snapshot.ts:126-166).
- Why the rule is safe (VERIFIED by the statistics session, 2026-10-05). A fold-path save only appends a
pending entry; it moves no revision and publishes no ProjectStatisticsChanged. Only the fold worker removes
entries, and it advances the clocks in the same transaction (ProjectStatisticsFoldWorker.cs:412,
AdvanceClocks always increments ClientInvalidationRevision), then publishes. So at a fixed revision the
CheckpointId is fixed and the pending set only grows, and new data always reaches clients at a strictly
higher revision with a trigger.
- Pod: per (project, stage, subscriber), drop a push whose revision is lower than the last one sent. At an
equal revision, send only when the provenance (PendingFingerprint) differs; with the pass-skip in item 2
this almost never happens.
- Which HTTP responses can be compared. Only the two coherent DTOs, StageOverviewStatistics and
ScreeningOverviewStatistics, carry checkpoint, clientInvalidationRevision and
globalClientInvalidationRevision. Neither carries a pending fingerprint, and the reviewer-progress
(ProjectReviewerProgressMetadata: readSource only) and Project Overview screening-stats responses carry
no revision at all. R4.1 therefore pushes only the sections those two coherent DTOs serve; the
reviewer-progress and Project Overview stores keep invalidate-and-refetch unless a later, statistics-owned PR
adds revisions to their responses.
- Decision (recommended): a value read over HTTP has unknown provenance. An equal-revision push for an
HTTP-held value is never applied and triggers one HTTP refetch. The alternative, adding pendingFingerprint
additively to the two coherent DTOs, is a statistics-owned wire change (regenerated client, fold design §4 "the
wire DTOs are unchanged") for a case the pass-skip already makes rare; the measured cost of the refetch is at
most one GET per push per client, below today's invalidate→refetch rate, and a refetch publishes nothing, so it
cannot feed back. Revisit only if staging shows equal-revision refetches are frequent.
- Client (the rule the client plan implements, RS.1): apply a push only when its project revision is
strictly greater than the one held and its global revision is not lower than the one held, whether the held
values came from HTTP or from a push; a push with a lower global revision is dropped. At an equal project revision a push never overwrites the held value; if the held value
came from HTTP (unknown provenance) or from a push with a different pendingFingerprint, the client makes one
HTTP refetch. Pages keep accepting an equal revision from HTTP, and an HTTP response with a lower
revision than held is dropped and the held value kept (today it throws "Statistics revision regressed" and the
page goes unavailable). Accepted transient: an HTTP GET issued before a push can return at the same revision
with a smaller pending set and overwrite the fresher push; the next fold (steady-state lag under 1 s) heals it
at revision+1.
- While subscribed, an invalidation does not refetch the sections the last push carried (the push for that
revision follows); sections the push omitted still refetch on invalidation.
11. Subscribe ack (S7). The SubscribeProjectStatistics ack carries a snapshot read fresh by that pod for
that caller, through the same entry point and checks; never a cached last-sent value.
12. Staging only by construction (S9). The sub-flag defaults to false and the production values are asserted
absent (4.1.8). The catalogue entry requires realtimeHubV2, materializedProjectStatisticsSignalR,
materializedProjectStatisticsServing and the per-section consumer flags (§6).
| # | Criterion | Verification |
|---|---|---|
| 4.1.1 | Flag on, allowlisted project; one materialised update → each subscriber receives exactly one ProjectStatisticsSnapshot whose revision equals the control ClientInvalidationRevision read in the same snapshot as its sections |
Integration (FEAT-024 harness + v2 hub) |
| 4.1.2 | 30 invalidations for project P within 2 s, delivered out of order → 1 snapshot read per pod and 1–2 pushes per subscriber, the last carrying the highest revision | Unit (FakeTimeProvider, read counter) |
| 4.1.3 | For each caller class (a member with ViewMemberships; a member without it; a member without the stage's annotation-graph permission; an impersonating administrator), project-wide and per stage → the entry point's output is byte-equal to ReadCurrentAsync/the HTTP endpoint for that caller on the same snapshot. The legacy shaper cases ported in R0.1 are not used as expected output |
Differential unit + Testcontainers test (statistics-session-owned entry point) |
| 4.1.4 | An older revision arrives after a newer one → no push. A coalesced trigger revision ≤ the last revision read → no read, no push. Equal revision with the same pendingFingerprint → no push. Equal revision with a different fingerprint → a push is sent, and the client refetches over HTTP without overwriting its held value (also when the held value came from HTTP). Every push carries both revision and globalRevision from the same snapshot |
Unit (pod) + client store specs (RS.1.2, RS.1.8) |
| 4.1.5 | A bundle in each Fallback/Unavailable state listed in item 8 → every section from that bundle is absent; no authoritative aggregation runs; no push mixes a materialised section with an authoritatively served one | Unit (aggregation spy) + Testcontainers |
| 4.1.6 | A member loses ProjectViewStudies, is disabled, or an impersonation actor lacks scope between the read and the send, in Off, shadow and enforced modes → no snapshot is sent and the subscriber receives SubscriptionRevoked; a permission grant → the next push matches the HTTP read for the new permissions |
Unit |
| 4.1.7 | realtimeStatisticsPush off, or the project not allowlisted → SubscribeProjectStatistics returns a denial and nothing is read |
Hub unit test |
| 4.1.8 | The production values file → realtimeStatisticsPush absent or false; production rollout is not in scope of this plan |
Governance (cluster-gitops review) |
| 4.1.9 | (S10) Staging, 7 days with the sub-flag on (PROPOSAL), checked per caller class, including a member without ViewMemberships and a member without the stage annotation-graph permission → zero pushes with a revision lower than one already delivered, and the Stage Overview figures match the HTTP read for that caller after each push |
Staging rollout check |
| 4.1.10 | Fold mode with pending entries → the pushed values equal the HTTP read including the overlay, and the provenance carries PendingFingerprint and CheckpointId |
Testcontainers (FEAT-024 harness) |
| 4.1.11 | A consumer flag the HTTP endpoint checks (or Serving, Pages, a family flag or the allowlist) is turned off by runtime override → the next pass omits that section, or sends nothing, with no restart. A global ProjectStatisticsChanged (ProjectId == null) → no per-project reads and no further pushes |
Unit |
| 4.1.12 | A pass → only read commands (no writes); Enqueue returns without awaiting; the snapshot session is closed before the first send; metric tags are bounded |
Unit + Testcontainers command monitor + meter tag test (C3) |
| 4.1.13 | A pod with no subscriber for P → 0 reads for P. A subscribe after a change that was not pushed → the ack carries the fresh value | Unit |
| 4.1.15 | A push never includes a reviewer-progress or Project Overview section (their HTTP responses carry no revision) | Unit (shape test) |
| 4.1.14 | The flag catalogue → realtimeStatisticsPush requires realtimeHubV2, materializedProjectStatisticsSignalR, materializedProjectStatisticsServing and the per-section consumer flags |
Flag dependency-catalogue validation test |
Later, optional (R5): once persistence PR-10's outbox is on and ADR-021 S1 has landed, semantic notifications (inbox, attention) may move from change streams to outbox events through IRealtimeClients. Not scheduled; D2.
8. Order, coordination and risks¶
flowchart LR
P5[backend PR-5 #3974] --> R11[R1.1 readers]
STJ[STJ migration plan] --> R13b
R01[R0.1] & R02[R0.2] --> R13a
N3932[#3932 merged] --> R13a[R1.3a extract] --> R13b[R1.3b v2 hub]
R11 --> R12[R1.2 dispatcher] --> R13b --> R21[R2.1 E2E] --> R22[R2.2 staging] --> PROD{{Chris: production}} --> R31[R3.1 retire v1]
CL[client plan RealtimeStore] --> R21
R22 --> R41[R4.1 statistics push: staging only]
F24{{FEAT-024 session agrees}} --> R41
- R0.1, R0.2 and R1.1 can run in parallel; they share no files except
MongoContext.cs(R0.1 dead code, R1.1 none). - R1.2 and R1.3a run in parallel.
- Critical path: PR-5 → R1.1 → R1.2 → R1.3b → R2.1 → R2.2, and the STJ migration → R1.3b. Whichever of R1.2 and the STJ migration finishes last gates R1.3b. R1.1, R1.2 and R1.3a do not depend on the protocol.
- R4.1 follows R2.2 (the v2 hub is proven on staging), the statistics session's agreement and their reader entry point; they review the R4 PR before merge. It is off the MVP path.
| Stream / PR | Shared files | Sequencing |
|---|---|---|
| Backend plan PR-5 (#3974) | MongoContext.cs, ChangeStreamHealthTracker.cs, API Program.cs:170 |
Lands first and fixes v1 for the transition; R1.1 adopts its ready/live split. R3 deletes what PR-5 touched |
| Backend plan PR-8/PR-9 (#3978) | Program.cs:381, NotificationHub.cs:112-131 |
Independent; R1.3a rebases on PR-9 |
Notifications stack: #3932 (adds InboxInvalidationWorker, InboxChanged, edits INotificationHubClient and signal-r.service.ts); #3938, #3941–#3945, #3965 touch Program.cs and NotificationsController only |
NotificationHub.cs, Program.cs |
R0.1 and R1.3a wait for #3932 to merge (and rebase on #2469, which also edits the hub). The inbox worker moves onto ChangeFeedReader in a follow-up owned with that stream |
| FEAT-024 (statistics session) | Statistics/ProjectStatisticsChangedConsumer.cs, ProjectStatisticsSubscriptionFilter.cs, the bundle reader and IStageOverviewStatisticsQuery (read only), and a new FEAT-024-owned per-caller reader entry point for R4 |
The contract is unchanged. R1.2/R1.3b need the session's review of the one-line sender switch and the v2 method name. R0.1 removes the legacy full-stats push; they confirmed on 2026-10-05 that FEAT-024 does not depend on it, and R0.1 edits their STATUS.md lines only in coordination with them. R4 needs their agreement on: the O(1) enqueue hook in their consumer; the new FEAT-024-owned per-caller reader entry point and its differential test; the coherent-bundle and omitted-section rules; provenance; the equal-revision rule; and the sub-flag dependencies (materializedProjectStatisticsSignalR, Serving, the per-section consumer flags). They review the R4 PR before it merges |
| Application-authority owner | e2e/authority/tests/authority-m1a-application-roles.spec.ts:697-799 (the cross-replica M4b proof on the full-stats stream) |
R0.1 rewrites that test (for example onto ProjectNotification) only with their agreement, recorded in the PR (0.1.6) |
| STJ migration (system-text-json-migration-plan-2026-10.md, parallel plan) | API Program.cs:390-400 (hub protocol), MVC JSON options, converters |
Prerequisite of R1.3b. It also moves the v1 hub to STJ (same protocol name). Its payload round-trip tests cover the hub DTOs listed in §5 |
| Persistence plan (PR-10 outbox; #3985 version bumps) | none | version stays advisory until #3985; R5 depends on PR-10 |
| ADR-021 (Wolverine, S1 Jan–May 2027 PROPOSAL) | the two per-pod temporary consumers | v2 adds no bus consumer; S0.12 already covers per-pod fan-out; IRealtimeClients is bus-agnostic |
| Authentication migration session | SubjectAdmissionHubFilter, query-token auth |
D8 suspension abort needs their sign-off; /notifications/v2 reuses DirectBearerSchemeSelector unchanged |
| Client plan (parallel) | contract §5 | Consumes v2; its pilot can start on v1 behind RealtimeStore and switch at R2 |
| Risk | Mitigation |
|---|---|
| Missed events during the dual-path period or a pod restart | At-most-once is explicit; acks with versions, ResyncRequired, client re-declaration; E2E 2.1.2–2.1.3 |
| Two pipelines double the change-stream load during transition | v1 streams open lazily, only for v1 subscribers; v2 projections are small; staging metrics in 2.2.1 |
version unreliable before #3985 |
Clients reload on any invalidation; patch-in-place only after #3985 (client plan rule) |
| Listing visibility filter leaks a private id | 1.2.4 matrix test, built from the legacy predicate at NotificationHub.cs:1260-1267 |
| The revocation stance for invalidations differs from M4b's per-delivery barrier | Payload messages keep the barrier; invalidations carry ids only; D8 records the auth-session sign-off |
| Stateful reconnect stalls sends when the buffer fills | Small messages; buffer size measured in R2.2; behaviour when full NOT VERIFIED, so R1.3b adds a test |
| Affinity cookie missing on a cross-origin negotiate, so connect lands on another pod | Already required by v1 and working in staging with 2 pods; R2.2 checks ingress logs (NOT VERIFIED for production) |
| Hub-file conflicts with in-flight PRs | R1.3a after #3932/#2469; gh pr list re-checked before each PR |
| The STJ prerequisite slips and delays R1.3b | R1.1, R1.2 and R1.3a proceed dark meanwhile; the client pilot runs on v1 behind RealtimeStore |
| The STJ switch changes v1 payloads (converter gaps, casing, enums, Int64) | Owned by the STJ plan's round-trip and parity tests; 1.3b.7 compares hub and MVC serialisation |
| R4 serves a statistics figure the HTTP read would not (stale overlay, wrong flags, fallback mixed in) | It reads through the same bundle reader as HTTP, overlay included, with the coherent-bundle rule and the same consumer flags per pass; non-Fresh bundles are omitted (4.1.5, 4.1.10, 4.1.11); staging check 4.1.9 per caller class |
| R4 leaks peer rows or annotation progress a caller cannot see over HTTP | Per-caller entry point with a byte-equality differential test (4.1.3); per-subscriber member-content check in every mode (4.1.6); the legacy shaper rules are not R4's expected output |
| A push overwrites a newer HTTP read at the same revision | Strict > on the client; at an equal revision a differing provenance triggers an HTTP refetch, never an overwrite (4.1.4) |
| Anonymised labels differ between pushes, so the client cannot diff rows | Same as the HTTP contract (labels are per response by design); the client replaces the section wholesale |
9. Decisions needed from Chris¶
- D1. Paradigm. Replace Rx.NET with
BackgroundServicereaders, a boundedChannel, a registry and a dispatcher. Rx stays only until R3. Recommended. - D2. Trigger. Hybrid: change streams for entity invalidation; outbox and bus fan-out for semantic events. Moving everything to domain events waits for PR-10 coverage and Wolverine (R5, unscheduled). Recommended.
- D3. Message kinds. Invalidation for project details and listings; small data for presence, export status and claim revocation (§4.4). Recommended.
- D4. Backplane. None. Revisit if API pods exceed 10, if a producer must reach sockets without a per-pod source, or if per-pod change-stream cost shows in metrics. Recommended.
- D5. Retire
SubscribeToProjectFullStatsnow, unflagged (no web caller since 2026-04-17). Stays as recommended: retire it. The statistics session confirmed on 2026-10-05 that FEAT-024 does not depend on it. The rewrite of the cross-replica E2E proof needs the application-authority owner's agreement (R0.1). Its replacement is R4, built on a different design. - D6. Versioning. New typed hub at
/notifications/v2, additive within v2, v3 path for breaks. Recommended. - D7. Stateful reconnect on v2 (the client opts in through the client plan), treated as an optimisation only. Recommended.
- D8. Authorization for invalidations. Revocation-driven, with no per-delivery read. Abort a suspended subject's connections. Payloads keep the M4b barrier. Needs the authentication session's agreement. Recommended.
- D9. Contract sharing. Committed doc plus .NET snapshot plus a web spec over a JSON sample (recommended), versus TS code generation (TypedSignalR) as a later option.
- D10. PROPOSAL values:
- coalescing windows 1 s / 2 s / 1 s / 250 ms;
- readiness Unhealthy after 60 s of retrying;
- ≤ 1 KB per Study change;
- hub ≤ 700 lines after extraction;
- E2E outage 10 s;
- staging soak 7 days with lag p95 ≤ 2 s.
- D11. v1 retirement gate. How long production runs on v2 before R3 deletes v1 and Rx. Recommended: 4 weeks without a rollback.
- D12. Hub protocol: DECIDED (Chris, 2026-10-04). System.Text.Json across the API, as a prerequisite planned separately; v2 uses
AddJsonProtocolwith the migrated MVC options. - D13. Full-stats push: DECIDED (Chris, 2026-10-04). R4, staging only; production not planned. Still open: approve the
realtimeStatisticsPushsub-flag, which requiresrealtimeHubV2,materializedProjectStatisticsSignalR,materializedProjectStatisticsServingand the per-section consumer flags (recommended over reusingrealtimeHubV2; §6), and the PROPOSAL 2 s coalescing window. The R4 design was rewritten on 2026-10-05 to the statistics session's ten conditions (per-caller reads, fold overlay, equal-revision rule).
10. Not verified¶
- Production change-stream count, event rates and lag: no production metrics were read, and the API metrics PodMonitoring is staging-only until promotion.
- Whether
Publish().RefCount()stays terminated after a non-retryable error (RT1): needs the PR-5 test. - Stateful reconnect being WebSockets-only, and its behaviour when the buffer is full: secondary sources only. (Protocol support is now VERIFIED; see CR-6.)
- The exact signature of the new FEAT-024-owned per-caller reader entry point for R4 (no existing method fits; the statistics session owns it), and the exact
ProjectStatisticsSnapshotfield names. - Production affinity on cross-origin negotiate: inferred from the ingress template, not observed.
- The Cluster0 MongoDB server version (it matters only if
$changeStreamSplitLargeEventwere needed; the projections make it unnecessary).