dotnet/reactive · error · ArgumentNullException

observer

Error message

observer

What it means

The TakeLast<TSource>(observer, count) operator validates its observer argument first and throws ArgumentNullException when null is passed. AsyncRx operators build an observer pipeline, so a null downstream observer cannot be wired up. This is a fail-fast guard thrown synchronously at operator construction.

Solutions

  1. Pass a non-null IAsyncObserver<TSource> instance as the first argument.
  2. Check where the observer value originates and fix the upstream null producer.
  3. Add a null check before calling TakeLast to fail with a clearer stack trace at the call site.

Example fix

// before
IAsyncObserver<int> obs = GetObserver(); // may return null
var op = AsyncObserver.TakeLast(obs, 3);
// after
var obs = GetObserver() ?? throw new InvalidOperationException("observer required");
var op = AsyncObserver.TakeLast(obs, 3);
Defensive patterns

Strategy: validation

Validate before calling

if (observer is null) throw new ArgumentNullException(nameof(observer));
AsyncObserver.TakeLast(observer, count);

Type guard

bool IsValidObs<TSource>(IAsyncObserver<TSource> o) => o is not null;

Try / catch

try { var op = AsyncObserver.TakeLast(observer, count); }
catch (ArgumentNullException ex) when (ex.ParamName == "observer") { /* supply default observer */ }

Prevention

When it happens

Trigger: Calling AsyncObserver.TakeLast(null, 5) or TakeLast(null, 5, scheduler), e.g. when the observer comes from a failed factory method or an uninitialized field.

Common situations: Chaining operators where an earlier factory returned null; passing a result of a conditional expression that evaluated to null; refactoring where the observer variable lost its assignment.

Related errors


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

Appendix: source

Thrown at AsyncRx.NET/System.Reactive.Async/Linq/Operators/TakeLast.cs:157

                    var (sink, drain) = AsyncObserver.TakeLast(observer, state.duration, state.clock, state.scheduler);

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

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

        public static IAsyncObservable<TSource> TakeLast<TSource>(this IAsyncObservable<TSource> source, TimeSpan duration, IAsyncScheduler scheduler) => TakeLast(source, duration, scheduler, scheduler);
    }

    public partial class AsyncObserver
    {
        public static (IAsyncObserver<TSource>, IAsyncDisposable) TakeLast<TSource>(IAsyncObserver<TSource> observer, int count) => TakeLast(observer, count, TaskPoolAsyncScheduler.Default);

        public static (IAsyncObserver<TSource>, IAsyncDisposable) TakeLast<TSource>(IAsyncObserver<TSource> observer, int count, IAsyncScheduler scheduler)
        {
            if (observer == null)
                throw new ArgumentNullException(nameof(observer));
            if (count <= 0)
                throw new ArgumentOutOfRangeException(nameof(count));
            if (scheduler == null)
                throw new ArgumentNullException(nameof(scheduler));

            var sad = new SingleAssignmentAsyncDisposable();

            var queue = new Queue<TSource>();

            return
                (
                    Create<TSource>(
                        x =>
                        {
                            queue.Enqueue(x);

                            if (queue.Count > count)
                            {

View on GitHub (pinned to 94b5d5ab91)