Skip to content
Merged
Show file tree
Hide file tree
Changes from 13 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
280 changes: 254 additions & 26 deletions PowerKit.Tests/ObservableTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,21 +6,10 @@

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 +23,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 +44,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 +62,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 +81,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 +90,139 @@ 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 =>
{
new Thread(() =>

Copy link
Copy Markdown
Owner

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why not Task.Run instead? @copilot

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Switched to Task.Run in the latest commit. No reason to spin up a dedicated thread here — threadpool is fine for this pattern.

{
startSignal.Wait(TestTimeout);
for (var i = 1; i <= 5; i++)
observer.OnNext(i);
})
{
IsBackground = true,
}.Start();
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 =>
{
new Thread(() =>
{
startSignal.Wait(TestTimeout);
for (var i = 1; i <= 3; i++)
observer.OnNext(i);
observer.OnCompleted();
})
{
IsBackground = true,
}.Start();
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 =>
{
new Thread(() =>
{
startSignal.Wait(TestTimeout);
for (var i = 1; i <= 5; i++)
observer.OnNext(i);
observer.OnCompleted();
})
{
IsBackground = true,
}.Start();
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 +240,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 +279,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