dotnet/reactive · error · ArgumentNullException

observer

Error message

observer

What it means

AsyncObserver.Merge's sink factory throws ArgumentNullException when the downstream observer is null. The sink needs a valid observer to forward merged inner-stream events to. This surfaces when composing operator sinks manually rather than through the extension method.

Solutions

  1. Pass a valid IAsyncObserver<TSource> created earlier in the pipeline
  2. Use the public extension source.Merge() which supplies the observer automatically
  3. Make the observer producer non-null-returning (throw earlier with a clearer message)

Example fix

// before
var (sink, cancel) = AsyncObserver.Merge<int>(null);
// after
var (sink, cancel) = AsyncObserver.Merge<int>(observer);
Defensive patterns

Strategy: validation

Validate before calling

if (observer == null) throw new InvalidOperationException("Merge sink requires an observer");

Type guard

bool HasObserver<T>(IAsyncObserver<T> o) => o is not null;

Try / catch

try { var (sink, cancel) = AsyncObserver.Merge(observer); } catch (ArgumentNullException ex) when (ex.ParamName == "observer") { /* reconstruct observer */ }

Prevention

When it happens

Trigger: Calling AsyncObserver.Merge<TSource>(null) while wiring a custom async-observable pipeline.

Common situations: Hand-rolled operator chains where an upstream Create(...) returned null; passing an observer before it is constructed.

Related errors


AI-assisted analysis of dotnet/reactive@94b5d5ab91 (2026-09-15). Data as JSON: /api/errors/4d2085019f8b0c48. Report an issue: GitHub.

Appendix: source

Thrown at AsyncRx.NET/System.Reactive.Async/Linq/Operators/Merge.cs:36

                throw new ArgumentNullException(nameof(source));

            return Create<TSource>(async observer =>
            {
                var (sink, cancel) = AsyncObserver.Merge(observer);

                var subscription = await source.SubscribeSafeAsync(sink).ConfigureAwait(false);

                return StableCompositeAsyncDisposable.Create(subscription, cancel);
            });
        }
    }

    public partial class AsyncObserver
    {
        public static (IAsyncObserver<IAsyncObservable<TSource>>, IAsyncDisposable) Merge<TSource>(IAsyncObserver<TSource> observer)
        {
            if (observer == null)
                throw new ArgumentNullException(nameof(observer));

            var gate = new AsyncGate();

            var count = 1;

            var disposable = new CompositeAsyncDisposable();

            async ValueTask OnErrorAsync(Exception ex)
            {
                using (await gate.LockAsync().ConfigureAwait(false))
                {
                    await observer.OnErrorAsync(ex).ConfigureAwait(false);
                }
            };

            async ValueTask OnCompletedAsync()
            {
                using (await gate.LockAsync().ConfigureAwait(false))

View on GitHub (pinned to 94b5d5ab91)