Skip to content
Merged
Show file tree
Hide file tree
Changes from 14 commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
272 changes: 246 additions & 26 deletions PowerKit.Tests/ObservableTests.cs
Original file line number Diff line number Diff line change
@@ -1,26 +1,16 @@
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using FluentAssertions;
using Xunit;

namespace PowerKit.Tests;

file class FakeObserver<T>(
Action<T>? onNext = null,
Action<Exception>? onError = null,
Action? onCompleted = null
) : IObserver<T>
{
public void OnNext(T value) => onNext?.Invoke(value);

public void OnError(Exception error) => onError?.Invoke(error);

public void OnCompleted() => onCompleted?.Invoke();
}

public class ObservableTests
{
private static readonly TimeSpan TestTimeout = TimeSpan.FromSeconds(5);

[Fact]
public void Observable_Create_Subscribe_Test()
{
Expand All @@ -34,7 +24,7 @@ public void Observable_Create_Subscribe_Test()

// Act
subscribed.Should().BeFalse();
observable.Subscribe(new FakeObserver<int>());
observable.Subscribe(Observer.Create<int>());

// Assert
subscribed.Should().BeTrue();
Expand All @@ -55,7 +45,7 @@ public void Observable_Create_OnNext_Test()
});

// Act
observable.Subscribe(new FakeObserver<int>(received.Add));
observable.Subscribe(Observer.Create<int>(onNext: received.Add));

// Assert
received.Should().Equal(1, 2, 3);
Expand All @@ -73,11 +63,11 @@ public void Observable_Create_OnError_Test()
});

// Act
observable.Subscribe(new FakeObserver<int>(onError: ex => receivedError = ex));
observable.Subscribe(Observer.Create<int>(onError: ex => receivedError = ex));

// Assert
receivedError.Should().BeOfType<InvalidOperationException>();
receivedError.Message.Should().Be("test error");
receivedError!.Message.Should().Be("test error");
}

[Fact]
Expand All @@ -92,7 +82,7 @@ public void Observable_Create_OnCompleted_Test()
});

// Act
observable.Subscribe(new FakeObserver<int>(onCompleted: () => completed = true));
observable.Subscribe(Observer.Create<int>(onCompleted: () => completed = true));

// Assert
completed.Should().BeTrue();
Expand All @@ -101,17 +91,130 @@ public void Observable_Create_OnCompleted_Test()
[Fact]
public void Observable_Create_Dispose_Test()
Comment thread
Tyrrrz marked this conversation as resolved.
{
// Arrange
// Arrange: an observable that emits 5 items from a background thread.
// A start signal ensures the thread doesn't emit until the subscription
// is fully assigned; the onNext callback disposes after 3 items.
var received = new List<int>();
var disposed = false;
var observable = Observable.Create<int>(_ => Disposable.Create(() => disposed = true));
var startSignal = new ManualResetEventSlim(false);
var disposedEvent = new ManualResetEventSlim(false);
IDisposable? subscription = null;

// Act
disposed.Should().BeFalse();
var subscription = observable.Subscribe(new FakeObserver<int>());
subscription.Dispose();
var observable = Observable.Create<int>(observer =>
{
Task.Run(() =>
{
startSignal.Wait(TestTimeout);
for (var i = 1; i <= 5; i++)
observer.OnNext(i);
});
return Disposable.Create(() =>
{
disposed = true;
disposedEvent.Set();
});
});

subscription = observable.Subscribe(
Observer.Create<int>(v =>
{
received.Add(v);
if (received.Count >= 3)
subscription!.Dispose();
})
);
startSignal.Set();

disposedEvent.Wait(TestTimeout).Should().BeTrue();

// Assert
disposed.Should().BeTrue();
received.Should().Equal(1, 2, 3);
}

[Fact]
public void Observable_Create_Dispose_AfterOnCompleted_NoDoubleDispose_Test()
{
// Arrange: an observable that emits 3 items and then completes from a
// background thread. OnCompleted auto-disposes; the subsequent external
// Dispose must be a no-op (source disposable fires exactly once).
var received = new List<int>();
var disposeCount = 0;
var startSignal = new ManualResetEventSlim(false);
var disposedEvent = new ManualResetEventSlim(false);
IDisposable? subscription = null;

var observable = Observable.Create<int>(observer =>
{
Task.Run(() =>
{
startSignal.Wait(TestTimeout);
for (var i = 1; i <= 3; i++)
observer.OnNext(i);
observer.OnCompleted();
});
return Disposable.Create(() =>
{
disposeCount++;
disposedEvent.Set();
});
});

subscription = observable.Subscribe(Observer.Create<int>(received.Add));
startSignal.Set();
disposedEvent.Wait(TestTimeout).Should().BeTrue();
subscription.Dispose();

// Assert
received.Should().Equal(1, 2, 3);
disposeCount.Should().Be(1);
}

[Fact]
public void Observable_Create_Dispose_NoEventAfterDispose_Test()
{
// Arrange: an observable that emits 5 items and then completes from a
// background thread. The onNext callback disposes after 3; items 4-5
// and the OnCompleted must be silently dropped.
var received = new List<int>();
var completedCalled = false;
var startSignal = new ManualResetEventSlim(false);
var disposedEvent = new ManualResetEventSlim(false);
IDisposable? subscription = null;

var observable = Observable.Create<int>(observer =>
{
Task.Run(() =>
{
startSignal.Wait(TestTimeout);
for (var i = 1; i <= 5; i++)
observer.OnNext(i);
observer.OnCompleted();
});
return Disposable.Null;
});

subscription = observable.Subscribe(
Observer.Create<int>(
onNext: v =>
{
received.Add(v);
if (received.Count >= 3)
{
subscription!.Dispose();
disposedEvent.Set();
}
},
onCompleted: () => completedCalled = true
)
);
startSignal.Set();

disposedEvent.Wait(TestTimeout).Should().BeTrue();

// Assert
received.Should().Equal(1, 2, 3);
completedCalled.Should().BeFalse();
}

[Fact]
Expand All @@ -129,7 +232,7 @@ public void Observable_CreateSynchronized_OnNext_Test()
});

// Act
observable.Subscribe(new FakeObserver<int>(received.Add));
observable.Subscribe(Observer.Create<int>(onNext: received.Add));

// Assert
received.Should().Equal(1, 2, 3);
Expand Down Expand Up @@ -168,9 +271,126 @@ public void Observable_CreateSynchronized_ThreadSafe_Test()
});

// Act
observable.Subscribe(new FakeObserver<int>(v => received.Add(v)));
observable.Subscribe(Observer.Create<int>(onNext: v => received.Add(v)));

// Assert
received.Should().HaveCount(threadCount * valuesPerThread);
}

[Fact]
public void Observable_Create_AutoDetach_OnNext_Throw_DisposesSource_Test()
{
// Arrange: the subscribe callback emits 5 items synchronously; the observer
// collects them and throws on the third. Only 2 items are received and the
// source disposable must be disposed once subscribe returns.
var received = new List<int>();
var disposed = false;
var observable = Observable.Create<int>(observer =>
{
for (var i = 1; i <= 5; i++)
{
try
{
observer.OnNext(i);
}
catch
{
break;
}
}
return Disposable.Create(() => disposed = true);
});

// Act
observable.Subscribe(
Observer.Create<int>(onNext: v =>
{
if (v == 3)
throw new InvalidOperationException();
received.Add(v);
})
);

// Assert
disposed.Should().BeTrue();
received.Should().Equal(1, 2);
}

[Fact]
public void Observable_Create_AutoDetach_OnError_Throw_DisposesSource_Test()
{
// Arrange
IObserver<int>? capturedObserver = null;
var disposed = false;
var observable = Observable.Create<int>(observer =>
{
capturedObserver = observer;
return Disposable.Create(() => disposed = true);
});
observable.Subscribe(
Observer.Create<int>(onError: _ => throw new InvalidOperationException("boom"))
);

// Act & Assert
var act = () => capturedObserver!.OnError(new Exception("source error"));
act.Should().Throw<InvalidOperationException>().WithMessage("boom");
disposed.Should().BeTrue();
}

[Fact]
public void Observable_Create_AutoDetach_OnCompleted_Throw_DisposesSource_Test()
{
// Arrange
IObserver<int>? capturedObserver = null;
var disposed = false;
var observable = Observable.Create<int>(observer =>
{
capturedObserver = observer;
return Disposable.Create(() => disposed = true);
});
observable.Subscribe(
Observer.Create<int>(onCompleted: () => throw new InvalidOperationException("boom"))
);

// Act & Assert
var act = () => capturedObserver!.OnCompleted();
act.Should().Throw<InvalidOperationException>().WithMessage("boom");
disposed.Should().BeTrue();
}

[Fact]
public void Observable_Create_AutoDetach_SuccessfulOnError_DisposesSource_Test()
{
// Arrange
var disposed = false;
var observable = Observable.Create<int>(observer =>
{
observer.OnError(new Exception("source error"));
return Disposable.Create(() => disposed = true);
});

// Act
observable.Subscribe(Observer.Create<int>(onError: _ => { }));

// Assert
disposed.Should().BeTrue();
}

[Fact]
public void Observable_Create_AutoDetach_SuccessfulOnCompleted_DisposesSource_Test()
{
// Arrange
var disposed = false;
var observable = Observable.Create<int>(observer =>
{
observer.OnCompleted();
return Disposable.Create(() => disposed = true);
});

// Act
observable.Subscribe(Observer.Create<int>());

// Assert
disposed.Should().BeTrue();
}
}
Loading