Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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
222 changes: 198 additions & 24 deletions PowerKit.Tests/ObservableTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,19 +6,6 @@

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
{
[Fact]
Expand All @@ -34,7 +21,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 +42,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 +60,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 +79,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 @@ -103,15 +90,85 @@ public void Observable_Create_Dispose_Test()
{
// Arrange
var disposed = false;
var observable = Observable.Create<int>(_ => Disposable.Create(() => disposed = true));
IObserver<int>? producer = null;
var received = new List<int>();

// Act
disposed.Should().BeFalse();
var subscription = observable.Subscribe(new FakeObserver<int>());
var observable = Observable.Create<int>(observer =>
{
producer = observer;
return Disposable.Create(() => disposed = true);
});

var subscription = observable.Subscribe(Observer.Create<int>(received.Add));

producer!.OnNext(1);
producer.OnNext(2);
producer.OnNext(3);
subscription.Dispose();
producer.OnNext(4);
producer.OnNext(5);

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

[Fact]
public void Observable_Create_Dispose_AfterOnCompleted_NoDoubleDispose_Test()
{
// Arrange
var disposeCount = 0;
IObserver<int>? producer = null;
var received = new List<int>();

var observable = Observable.Create<int>(observer =>
{
producer = observer;
return Disposable.Create(() => disposeCount++);
});

var subscription = observable.Subscribe(Observer.Create<int>(received.Add));

producer!.OnNext(1);
producer.OnNext(2);
producer.OnNext(3);
producer.OnCompleted();
subscription.Dispose();

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

[Fact]
public void Observable_Create_Dispose_NoEventAfterDispose_Test()
{
// Arrange
var completedCalled = false;
IObserver<int>? producer = null;
var received = new List<int>();

var observable = Observable.Create<int>(observer =>
{
producer = observer;
return Disposable.Null;
});

var subscription = observable.Subscribe(
Observer.Create<int>(onNext: received.Add, onCompleted: () => completedCalled = true)
);

producer!.OnNext(1);
producer.OnNext(2);
producer.OnNext(3);
subscription.Dispose();
producer.OnNext(4);
producer.OnNext(5);
producer.OnCompleted();

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

[Fact]
Expand All @@ -129,7 +186,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 +225,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