Skip to content

Commit 49e2515

Browse files
authored
feat!: deliver notifications outside the lock, one sink per operator, no runtime-async (#217)
* feat(scheduling)!: deliver notifications outside the lock - Route every synchronous operator, sink and subject through DeliveryGate and SerializedDelivery, so no observer or user callback runs while a lock is held. - Replace the public abstract witness, coordinator and scheduler bases with interfaces, embedded state structs and static helpers. - Add the Serialize operator, SerializeSignal, SerializeWitness, SerializedSignal and Signal.Serialized, which serialize notifications without a lock. - Add CurrentValueSignal, CurrentValueWitness, CurrentValueDelivery and ICurrentValueReader for latest-value sources that are read on subscribe. - Add SyncLatestSlot and the slot subscription extensions, so a custom coordinator wires its sources without computing completion bits. - Add DisposableSet, an inline set that holds a few disposables without allocating a list. - Fix Switch never completing when an inner sequence finishes during its own subscribe. - Run the handler and callbacks of the async subscribe sink outside its queue lock. - Remove the unused Handle, SequencerWorkItem, SequencerWorkItemDisposal and Extensions ConcurrencyRaceHelpers helpers. - Make CoalescingDispatchState, DispatchWorkState and IDispatchHost public API of ReactiveUI.Primitives.Reactive, so a custom UI sequencer can be written outside this repository. Each platform package now references them instead of compiling its own copy. - Move PendingNotifications and PendingDelivery into an internal namespace, since they are storage details of SerializedDelivery rather than something a caller builds against. - Update Microsoft.NET.Test.Sdk to 18.10.1 and TUnit to 1.68.0. - Drop the .NET Framework leg from both benchmark projects, which cannot compile there: the generator harness needs AppContext.GetData and a generic Enum.GetValues, and the comparison benchmarks need GeneratedRegex. - Resolve the .NET Framework polyfill sources from the props file's own directory, so a project nested below src finds them. - Adopt the shared analyzer set in .editorconfig and update the NuGet packages that are not deliberately pinned. - Add benchmarks covering every touched production file, including a diagnostics benchmark for the observable-event generator. - Give the Avalonia and Blazor benchmarks their own assemblies, which target .NET only, so the main benchmark project keeps its .NET Framework leg. - Add dispatchable workflows that run the benchmark projects, one for a suite run and one comparing two commits. BREAKING CHANGE: Six public types are removed. Implement an interface and hold the matching state struct instead of deriving: WitnessAsync becomes IWitnessAsync plus WitnessAsyncState, SyncLatestCoordinatorBase becomes ISyncLatestCoordinator plus SyncLatestLifecycle, TaskResultWitnessAsyncBase becomes ITaskSignalJob plus TaskSignalState, ForwardingWitnessAsync becomes IWitnessAsync, CoalescingDispatchScheduler is replaced by the platform sequencers, and IReentrantAsyncDisposable is gone. * refactor(async): share the operator subscribe and terminal paths - Add WitnessSubscription, which subscribes an operator's witness to its source and links the two teardowns, so every single-source operator shares one subscribe path. - Add TaskResultCompletionSource.CompleteAndDisposeAsync, so a terminal witness publishes its result or its exception in one call and no longer carries its own set-and-dispose wrappers. - Name the CatchSignal advance handoff, so the walker reports teardown from the gate rather than through a local flag. * docs: rewrite the readme around what the library offers - Rewrite README.md for a reader new to the library: the problem first, a six-step first-signal walkthrough, then the reference material. - List every operator in a table with what it does and its LINQ or System.Reactive name, and add a worked example for the ones in common use. - Add reference sections for the extension helpers, the async operators, the subjects, the sequencers and the disposables. - Move the advanced types to the end, one or two sentences each, for a reader who wants the concrete type instead of the extension method. - Add a Writing Docs section to CLAUDE.md covering who the docs are written for, sentence shape, word choice and structure. * feat(operators): expose the fused operator signal types - Move the signal types behind Fold, Reduce, Unique, Zip, CombineLatest, Calm, Shift, Probe, Latch, KeepNotNull, KeepType, Reattempt and absolute-time Expire into Advanced as public types. - Give the Latch, CombineLatest, Reattempt and Calm coordinators their own files as internal types. - Update the public API baselines for the core and System.Reactive shim packages. * build: compile the net11 targets without runtime-async - Runtime-async is unsupported on Mono, so a net11.0 package asset built with it fails for Blazor WebAssembly consumers. - The shared framework can enable it because it ships a separate Mono build; a NuGet package resolves one asset for every runtime. * docs: describe the fused operator types as constructible - The readme lists every fused operator type as public, with an example that builds one directly. * docs: contrast the composed and dedicated operator designs - The comparison section explains that System.Reactive and R3 compose a general Synchronize operator while this library builds the behaviour into each sink. - States the trade: more classes for speed and fewer allocations. * ci: raise the copy/paste detection threshold to 200 tokens - Below 200 tokens the detector reports the interface members every sink must declare, whose bodies already delegate to shared static helpers. - At 300 tokens the detector reports nothing at all for this repository. * build(deps): update PublicApiSharp.Analyzers to 2.0.1 - The baseline lookup now keys on the generic constraint, so overloads that differ only by a constraint match their own baseline entry. - Both SubscribeSafe overloads drop their PAS0003 suppressions and are tracked against the baseline again. * refactor(sinks): forward values through a shared delivery helper - SinkDelivery.Next forwards a value to the downstream observer and disposes the sink when that observer throws. - Sixteen sinks call it instead of repeating the try/catch that tore themselves down on throw. * refactor(operators): collapse duplicated sink and witness pairs - OnDisposeWitness holds the synchronous action and the asynchronous callback, so the sync and async dispose overloads share one sink. - FirstTaskWitness carries the default value and a flag for whether an empty sequence yields it or fails, replacing FirstOrDefaultTaskWitness. - Regenerate the public API baselines for every target framework. * docs: show the array form of the collection CombineLatest - The readme row passes an array, because listing sources individually binds to the tuple overload and returns a tuple rather than a list. * fix(operators): correct buffering, absolute timeouts, cancellation and rethrow - Buffer(count, skip) opens a window every skip values, so a skip below the count overlaps windows and a skip above it leaves a gap; completion flushes every window still filling. - Timeout(DateTimeOffset) and its sequencer overload arm one window at subscription, so arriving values no longer push the deadline back. - HandleCancellation on an observable passes the token into the wait, so cancelling part-way through ends it instead of waiting forever. - Exception.Throw and Exception.Rethrow go through ExceptionDispatchInfo on every target, keeping the stack trace from the original throw site. - AsObservable returns a read-only view, so a caller cannot cast it back to the subject and push values in. - Correct the Expire, Timeout, Probe and Sample summaries to describe an inactivity window rather than a fixed schedule, and the readme rows for OnErrorResumeNext, Repeat, Reattempt and AsObservable. * docs: correct the readme rows that misdescribe members - ShareLatest says it shares one live subscription and does not replay to a late subscriber. - ToReadOnlyState describes Changed sending the current value on subscribe and repeating unchanged projections, and qualifies the ToProperty equivalence. - Probe records that a value still waiting when the source completes is dropped, in the readme and in the XML remarks. - Synchronize splits the Lock overload onto its own row, since it exists only on net9.0 and later. - Signal, StateSignal, AsyncSignal, PrioritySemaphoreSignal and CommandSignal list the members their public surface actually carries. * chore(api): regenerate public API baselines for the read-only observable view * docs: add the ITaskSignal row to the signals table - The row names what Signal.FromTask hands back and the cancellation members it carries. * docs: name the helper types and correct two wrong names - Add a table for the extension, option, enum and collection types that sit beside the operators, with the namespace each lives in. - State the naming convention that gives every operator a public type behind it. - The migration step names AsyncSignal rather than a FinalSignal that does not exist. - The blocking-helper note names WaitForValue, WaitForCompletion and WaitForError rather than a WaitFor that does not exist. * fix(disposables): run the slot action exactly once - A single-assignment slot runs its action on disposal whether or not a value was ever assigned, matching what its Dispose documents. - A replaceable slot disposes a value assigned after disposal without running the action a second time. * docs: record the slot dispose order and the event-pattern constraint - A table shows which single-value slots run the constructor action before disposing their value and which run it after, since swapping one family for the other reverses the order. - FromEventPattern documents that TEventArgs must derive from EventArgs, and points at FromEvent for an event whose argument type does not. * docs: explain why some signatures are narrower than their Rx counterparts - A comparison section records that several System.Reactive entry points build delegates or look events up by name at run time, and that reflection does not survive trimming or ahead-of-time compilation. - FromEventPattern states that its EventArgs constraint is what lets it bind the handler at compile time, and points at FromEvent for the argument types it excludes. * docs: correct the mapping tables and trim the reflection section - The subject mapping names AsyncSignal, and the disposable mapping names Scope.Create and Scope.Empty, which are the types this library exports. - The reflection section states the constraint and the reason without restating them. * fix(operators)!: count Retry as total runs and flush Probe on completion - Retry counts total runs, matching the System.Reactive operator of that name, so Retry(3) runs the source three times and Retry(0) completes without running it. Reattempt keeps counting extra tries. - Probe sends a value still waiting when the source completes, ahead of the completion, as Calm and the time-based Buffer already do. BREAKING CHANGE: Retry(n) now runs the source n times in total rather than n times after the first, so a call that relied on the extra attempt should use Reattempt(n). Probe and Sample now deliver a pending value on completion instead of dropping it. * feat(async): add Probe and Sample to the async library - Probe forwards the newest element once each sampling period elapses, with an optional TimeProvider. - Sample is the Rx name for the same operator. - An element held when the source completes is forwarded ahead of the completion. * docs: correct the R3Async observer mapping - The async observer row names IWitnessAsync and the WitnessAsyncState field a custom observer holds. * fix(operators)!: correct disposal, retry filtering and async timer parity - A replaceable slot that must dispose what it displaces now uses SwapDisposable, so Heartbeat stops stacking a live periodic timer per value and the retry, switch-if-empty and while operators stop leaking a subscription per attempt. - The retry operators hold the pending retry timer in its own slot, so a source that fails during re-subscribe no longer displaces the timer that was about to fire. - OnErrorRetry retries only the exception type it was given; any other failure goes straight downstream instead of being retried forever. - SelectAsync passes a subscription-scoped cancellation token to the selector, and cancels it on disposal. - SelectManyThen delivers through one fused coordinator that counts both projection stages, so completion arrives once rather than once per stage. - SignalAsync.Use disposes its resource exactly once. - Async SwitchTo ignores a superseded inner sequence's outcome, so an inner that is still running when the next arrives no longer deadlocks the producer. - Async Interval counts from zero, matching Every, Pulse, Timer and the synchronous operator of that name. - Async TakeUntil with a predicate emits the element that matched before completing, matching the synchronous helper. - Async LogErrors reports a terminal failure to the logger, not only resumable errors. - AnyAsync takes a predicate without a cancellation token, matching the other terminals. - Awaiting an empty AsyncSignal reports the same message as ToTask, FirstAsync and LastAsync. - Readme rows for async Unique, UniqueBy, Retry, Reattempt, Interval and the R3Async observer mapping describe what the code does. - Comments and editorconfig rule descriptions are ASCII throughout. BREAKING CHANGE: OnErrorRetry no longer retries exceptions that are not TException; those now terminate the sequence. Async Interval starts at 0 rather than 1. Async TakeUntil(predicate) includes the matching element. SelectManyThen completes once instead of once per projection stage. * more work * feat(advanced): make the remaining operator sinks constructible - Eleven sinks under Advanced are public, so callers can build them directly the way they already can with BufferSignal and UniqueSignal: CreateSignal, CreateSignal with state, CreateSafeSignal, DeferSignal, WitnessOnSignal, CatchSignal, CallbackSignalAsync, LatchCoordinator, CombineLatestCoordinator, ReattemptCoordinator and CalmCoordinator. - Each coordinator exposes the constructor and Run entry point a caller needs, rather than a public type with no way in. - The multi-source CombineLatestCoordinator exposes Attach and the slot it returns, which is the path the factory passed to CombineLatestSignal has to use. - The readme links the detailed documentation at reactiveui.net, lists the newly constructible sinks, and states that an operator is one sink and never builds its behaviour from other operators. * feat(advanced): let callers construct OnErrorResumeNextSignal - The signal behind OnErrorResumeNext takes its sources through a public constructor, matching every other fused operator type. * docs(async): correct two summaries that describe superseded behaviour - IntervalSignal numbers its ticks from zero. - LogErrorsSignal reports terminal failures to the logger as well as resumable errors. * build(android)!: keep the generated resource designer out of the public surface - The Android targets no longer generate the resource designer, which was published as a public Resource type in the library's root namespace. BREAKING CHANGE: the generated Resource type is no longer part of the Android public surface. * docs(operators): describe the SynchronizeAsync handle as the acknowledgement it is - The handles are independent: the producer does not wait on one, and one value's handle does not gate the next. - A subscriber that ignores the handle still receives every value and the terminal notification. * docs(operators): state that Conflate emits at the end of each window - Each window emits its newest value at the end of that window, and the clock starts at subscription, so the first value is held for a full period. - A value arriving after a quiet gap longer than the period goes out at once. * feat(operators)!: name the value-shaped Schedule overloads apart - Scheduling a single value is ScheduleValue, so a concrete signal type can no longer bind the value overload and emit the signal object itself instead of scheduling its values. BREAKING CHANGE: Schedule on a non-observable value is now ScheduleValue. * fix(operators): deliver a cold source to the first side of a Partition - The first subscriber is registered before the source is subscribed, so a cold source that runs to completion during subscribe reaches it instead of an empty observer list. - The source is subscribed outside the gate, so it never runs while the gate is held. - The remarks state that both sides share one subscription, so over a cold source the first side to subscribe consumes it. * docs(operators): describe ReplayLastOnSubscribe and Partition as they behave - ReplayLastOnSubscribe gives each subscriber its own subscription and the initial value, so a late subscriber does not receive the newest source value; the async operator of that name shares one subscription and does replay the newest. - The readme and the Partition example state that both sides share one subscription, so over a cold source the first side to subscribe consumes it. * docs(operators): state that WhereIsNotNull keeps the source element type - The synchronous operator filters nulls without narrowing the element type, so a nullable source stays nullable downstream and a handler taking a non-nullable parameter still warns. - The async operator of that name narrows a T? source to T, and the readme names the difference. * chore(api): regenerate the Apple public API baselines - The iOS, tvOS, macOS and Mac Catalyst baselines carry this branch's public surface: the constructible Advanced sinks, the OnErrorResumeNextSignal constructor and the ScheduleValue rename.
1 parent 8aba4af commit 49e2515

810 files changed

Lines changed: 68774 additions & 19283 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.editorconfig‎

Lines changed: 908 additions & 1185 deletions
Large diffs are not rendered by default.
Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,28 @@
1+
name: Benchmarks A/B
2+
3+
on:
4+
workflow_dispatch:
5+
inputs:
6+
base:
7+
description: 'Baseline commit SHA.'
8+
type: string
9+
required: true
10+
head:
11+
description: 'Commit SHA compared against the baseline.'
12+
type: string
13+
required: true
14+
15+
permissions:
16+
contents: read
17+
18+
jobs:
19+
benchmark:
20+
uses: reactiveui/actions-common/.github/workflows/workflow-common-benchmarks-ab.yml@main
21+
with:
22+
base: ${{ inputs.base }}
23+
head: ${{ inputs.head }}
24+
projects: |
25+
benchmarks/ReactiveUI.Primitives.Benchmarks/ReactiveUI.Primitives.Benchmarks.csproj
26+
benchmarks/ReactiveUI.Primitives.ObservableEvents.Benchmarks/ReactiveUI.Primitives.ObservableEvents.Benchmarks.csproj
27+
solutionFile: ReactiveUI.Primitives.slnx
28+
installWorkloads: true

‎.github/workflows/benchmarks.yml‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,23 @@
1+
name: Benchmarks
2+
3+
on:
4+
workflow_dispatch:
5+
inputs:
6+
sha:
7+
description: 'Commit to benchmark. Empty benchmarks the latest commit of the chosen branch.'
8+
type: string
9+
default: ''
10+
11+
permissions:
12+
contents: read
13+
14+
jobs:
15+
benchmark:
16+
uses: reactiveui/actions-common/.github/workflows/workflow-common-benchmarks.yml@main
17+
with:
18+
revision: ${{ inputs.sha }}
19+
projects: |
20+
benchmarks/ReactiveUI.Primitives.Benchmarks/ReactiveUI.Primitives.Benchmarks.csproj
21+
benchmarks/ReactiveUI.Primitives.ObservableEvents.Benchmarks/ReactiveUI.Primitives.ObservableEvents.Benchmarks.csproj
22+
solutionFile: ReactiveUI.Primitives.slnx
23+
installWorkloads: true

‎.github/workflows/sonarcloud.yml‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,9 @@ jobs:
1818
minverMinimumMajorMinor: '0.1'
1919
sonarProjectKey: reactiveui_Primitives
2020
sonarOrganization: reactiveui
21+
# Below 200 tokens copy/paste detection reports the interface members every sink must
22+
# declare, whose bodies already delegate to shared static helpers. At 300 it reports nothing.
23+
sonarExtraBeginArgs: '/d:sonar.cpd.cs.minimumTokens=200'
2124
sonarExclusions: '**/tests/**,**/tools/**,**/benchmarks/**,**/TestResults/**'
2225
sonarCoverageExclusions: '**/tests/**,**/tools/**,**/benchmarks/**,**/*Tests/**,**/*Tests.cs,**/Generated/**'
2326
sonarCpdExclusions: '**/tests/**,**/tools/**,**/benchmarks/**,**/Operators/SyncLatest*.cs,**/SyncLatest*Coordinator*.cs,**/SyncLatest*Signal*.cs,**/SyncLatest*State*.cs,**/SyncLatestWitness*.cs,**/SynchronizeWitness.cs,**/SynchronizeObjectWitness*.cs,**/SequencerSchedulingExtensions.cs,**/Signal*RxAliases*.cs,**/ReactiveUI.Primitives.Blazor/**,**/ReactiveUI.Primitives.Maui/**,**/ReactiveUI.Primitives.WinUI/**,**/ReactiveUI.Primitives.WinForms/**,**/ReactiveUI.Primitives.Wpf/**,**/SignalOperatorMixins.CombineLatest.cs,**/SignalOperatorMixins.SyncLatest.MultiSource.cs,**/SignalOperatorParityMixins.RxNames.cs,**/SignalOperatorParityMixins.RxNames.CombineLatest.cs,**/SignalOperatorMixins.CombineLatest.WideArity.cs,**/SignalOperatorMixins.SyncLatest.WideArity.cs,**/SignalOperatorParityMixins.RxNames.CombineLatest.WideArity.cs'

‎CLAUDE.md‎

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -197,6 +197,47 @@ These rules are authoritative for everything under `src/tests/`.
197197

198198
---
199199

200+
## Writing Docs
201+
202+
These rules cover `README.md`, `CLAUDE.md` and every other doc in the repository. They do not cover XML doc comments.
203+
204+
### Who you write for
205+
206+
Write for a reader at a grade 8 level who knows basic C#. They know what a class, a property and an event are. They do
207+
not know this library.
208+
209+
### Sentences
210+
211+
- Put the main point first.
212+
- Give each sentence one subject. Use two only when they are tightly coupled.
213+
- Keep sentences short. Split a sentence that needs a dash, a semicolon or a "which" to hold together.
214+
- Use the active voice. Say who does what: "the operator delivers the value", not "the value is delivered".
215+
- Use verbs, not nouns made from verbs. Write "decide", not "make a decision".
216+
- Say what is true. Avoid double negatives.
217+
- Cut words that add nothing. Do not restate a point in the next sentence.
218+
219+
### Words
220+
221+
- Use everyday words. When you need a technical term, define it the first time you use it.
222+
- Define each term once. After that, use it without explaining it again.
223+
- Use the same word for the same thing every time. Do not swap in a synonym for variety.
224+
- Use "you" for the reader.
225+
226+
### Structure
227+
228+
- Use headings so a reader can find a topic.
229+
- Use a list for steps or for separate items. Use a table to compare items across the same columns.
230+
- Show a short code example when it explains faster than words.
231+
232+
### Scope
233+
234+
- Each section says what this library does, on its own terms.
235+
- A comparison with System.Reactive, R3 or R3Async goes only under the comparison headers: "Why not System.Reactive or
236+
R3?" and the migration guides in `README.md`. Do not compare with them anywhere else.
237+
- Describe the code as it is. Do not describe what it used to do.
238+
239+
---
240+
200241
## Agent Compatibility
201242

202243
If another agent entrypoint file exists, it should defer to this file.

0 commit comments

Comments
 (0)