diff --git a/Chickensoft.Sync.Tests/badges/branch_coverage.svg b/Chickensoft.Sync.Tests/badges/branch_coverage.svg index 689b6fd..69d0b7f 100644 --- a/Chickensoft.Sync.Tests/badges/branch_coverage.svg +++ b/Chickensoft.Sync.Tests/badges/branch_coverage.svg @@ -117,7 +117,7 @@ - Generated by: ReportGenerator 5.4.8.0 + Generated by: ReportGenerator 5.5.11.0 diff --git a/Chickensoft.Sync.Tests/badges/line_coverage.svg b/Chickensoft.Sync.Tests/badges/line_coverage.svg index 9067937..2a07312 100644 --- a/Chickensoft.Sync.Tests/badges/line_coverage.svg +++ b/Chickensoft.Sync.Tests/badges/line_coverage.svg @@ -117,13 +117,13 @@ - Generated by: ReportGenerator 5.4.8.0 + Generated by: ReportGenerator 5.5.11.0 Coverage Coverage - 99.2%99.2% + 99.3%99.3% diff --git a/Chickensoft.Sync.Tests/src/primitives/AutoQueueTest.cs b/Chickensoft.Sync.Tests/src/primitives/AutoQueueTest.cs new file mode 100644 index 0000000..c578d7d --- /dev/null +++ b/Chickensoft.Sync.Tests/src/primitives/AutoQueueTest.cs @@ -0,0 +1,399 @@ +namespace Chickensoft.Sync.Tests.Primitives; + +using Chickensoft.Sync.Primitives; +using Shouldly; + +public sealed class AutoQueueTest +{ + private readonly record struct TestValue(int Value); + private readonly record struct TestMessage(string Text); + + [Fact] + public void Initializes() + { + var queue = new AutoQueue(); + var enqueueReceived = false; + var dequeueReceived = false; + var discardReceived = false; + var clearReceived = false; + var modifyReceived = false; + + queue.Bind() + .OnEnqueue((in int v) => enqueueReceived = true) + .OnDequeue((in int v) => dequeueReceived = true) + .OnDiscard(count => discardReceived = true) + .OnClear(() => clearReceived = true) + .OnModify(() => modifyReceived = true); + + queue.Count.ShouldBe(0); + queue.HasValues.ShouldBeFalse(); + enqueueReceived.ShouldBeFalse(); + dequeueReceived.ShouldBeFalse(); + discardReceived.ShouldBeFalse(); + clearReceived.ShouldBeFalse(); + modifyReceived.ShouldBeFalse(); + } + + [Fact] + public void NoReentrancy() + { + var queue = new AutoQueue(); + var inEnqueueCallback = false; + var inDequeueCallback = false; + var enqueueReentered = false; + var dequeueReentered = false; + + queue.Bind() + .OnEnqueue((in int v) => + { + if (inEnqueueCallback) + { + enqueueReentered = true; // signals immediate (re-entrant) delivery + } + + inEnqueueCallback = true; + + if (v == 1) + { + queue.Enqueue(2); // attempt re-entrant enqueue + } + + inEnqueueCallback = false; + }) + .OnDequeue((in int v) => + { + if (inDequeueCallback) + { + dequeueReentered = true; // signals immediate (re-entrant) delivery + } + + inDequeueCallback = true; + + if (v == 1) + { + queue.Dequeue(); // attempt re-entrant dequeue + } + + inDequeueCallback = false; + }); + + queue.Enqueue(1); + queue.Dequeue(); + + enqueueReentered.ShouldBeFalse(); + dequeueReentered.ShouldBeFalse(); + } + + [Fact] + public void BroadcastsAllEnqueuedDequeuedValues() + { + var queue = new AutoQueue(); + var enqueueValues = new List(); + var dequeueValues = new List(); + var modifyCount = 0; + + queue.Bind() + .OnEnqueue((in int v) => enqueueValues.Add(v)) + .OnDequeue((in int v) => dequeueValues.Add(v)) + .OnModify(() => modifyCount++); + + queue.Enqueue(5); + queue.Enqueue(10); + queue.Enqueue(10); + queue.Enqueue(15); + + queue.Count.ShouldBe(4); + enqueueValues.ShouldBe([5, 10, 10, 15]); + dequeueValues.ShouldBe([]); + modifyCount.ShouldBe(4); + + while (queue.HasValues) + { + queue.Dequeue(); + } + + queue.Count.ShouldBe(0); + dequeueValues.ShouldBe([5, 10, 10, 15]); + modifyCount.ShouldBe(8); + } + + [Fact] + public void BroadcastsMultipleTypes() + { + var queue = new AutoQueue(); + var enqueueValues = new List(); + var dequeueValues = new List(); + var modifyCount = 0; + + queue.Bind() + .OnEnqueue((in int v) => enqueueValues.Add(v)) + .OnEnqueue((in double v) => enqueueValues.Add(v)) + .OnEnqueue((in TestValue v) => enqueueValues.Add(v)) + .OnEnqueue((in TestMessage v) => enqueueValues.Add(v)) + .OnDequeue((in int v) => dequeueValues.Add(v)) + .OnDequeue((in double v) => dequeueValues.Add(v)) + .OnDequeue((in TestValue v) => dequeueValues.Add(v)) + .OnDequeue((in TestMessage v) => dequeueValues.Add(v)) + .OnModify(() => modifyCount++); + + queue.Enqueue(42); + queue.Enqueue(3.14); + queue.Enqueue(new TestValue(100)); + queue.Enqueue(new TestMessage("hello")); + + queue.Count.ShouldBe(4); + enqueueValues + .ShouldBe([42, 3.14, new TestValue(100), new TestMessage("hello")]); + dequeueValues.ShouldBe([]); + modifyCount.ShouldBe(4); + + while (queue.HasValues) + { + queue.Dequeue(); + } + + queue.Count.ShouldBe(0); + dequeueValues + .ShouldBe([42, 3.14, new TestValue(100), new TestMessage("hello")]); + modifyCount.ShouldBe(8); + } + + [Fact] + public void ConditionalCallbacks() + { + var queue = new AutoQueue(); + var enqueueValues = new List(); + var evenEnqueueValues = new List(); + var dequeueValues = new List(); + var evenDequeueValues = new List(); + var modifyCount = 0; + + queue.Bind() + .OnEnqueue((in int v) => enqueueValues.Add(v)) + .OnEnqueue( + (in int v) => evenEnqueueValues.Add(v), + condition: v => v % 2 == 0 + ) + .OnDequeue((in int v) => dequeueValues.Add(v)) + .OnDequeue( + (in int v) => evenDequeueValues.Add(v), + condition: v => v % 2 == 0 + ) + .OnModify(() => modifyCount++); + + queue.Enqueue(1); + queue.Enqueue(2); + queue.Enqueue(3); + queue.Enqueue(4); + queue.Enqueue(5); + + enqueueValues.ShouldBe([1, 2, 3, 4, 5]); + evenEnqueueValues.ShouldBe([2, 4]); + dequeueValues.ShouldBe([]); + evenDequeueValues.ShouldBe([]); + modifyCount.ShouldBe(5); + + while (queue.HasValues) + { + queue.Dequeue(); + } + + dequeueValues.ShouldBe([1, 2, 3, 4, 5]); + evenDequeueValues.ShouldBe([2, 4]); + modifyCount.ShouldBe(10); + } + + [Fact] + public void DequeueOnEmptyDoesNotBroadcast() + { + var queue = new AutoQueue(); + var dequeueValues = new List(); + var modifyCount = 0; + + queue.Bind() + .OnDequeue((in int v) => dequeueValues.Add(v)) + .OnModify(() => modifyCount++); + + queue.Dequeue(); + + queue.Count.ShouldBe(0); + dequeueValues.ShouldBe([]); + modifyCount.ShouldBe(0); + } + + [Fact] + public void MultipleBindings() + { + var queue = new AutoQueue(); + var enqueueValues1 = new List(); + var enqueueValues2 = new List(); + var enqueueValues3 = new List(); + var dequeueValues1 = new List(); + var dequeueValues2 = new List(); + var dequeueValues3 = new List(); + var modifyCount1 = 0; + var modifyCount2 = 0; + var modifyCount3 = 0; + + queue.Bind() + .OnEnqueue((in int v) => enqueueValues1.Add(v)) + .OnDequeue((in int v) => dequeueValues1.Add(v)) + .OnModify(() => modifyCount1++); + + queue.Bind() + .OnEnqueue((in int v) => enqueueValues2.Add(v)) + .OnDequeue((in int v) => dequeueValues2.Add(v)) + .OnModify(() => modifyCount2++); + + queue.Bind() + .OnEnqueue((in int v) => enqueueValues3.Add(v)) + .OnDequeue((in int v) => dequeueValues3.Add(v)) + .OnModify(() => modifyCount3++); + + queue.Enqueue(7); + queue.Enqueue(14); + + enqueueValues1.ShouldBe([7, 14]); + enqueueValues2.ShouldBe([7, 14]); + enqueueValues3.ShouldBe([7, 14]); + dequeueValues1.ShouldBe([]); + dequeueValues2.ShouldBe([]); + dequeueValues3.ShouldBe([]); + modifyCount1.ShouldBe(2); + modifyCount2.ShouldBe(2); + modifyCount3.ShouldBe(2); + + while (queue.HasValues) + { + queue.Dequeue(); + } + + dequeueValues1.ShouldBe([7, 14]); + dequeueValues2.ShouldBe([7, 14]); + dequeueValues3.ShouldBe([7, 14]); + modifyCount1.ShouldBe(4); + modifyCount2.ShouldBe(4); + modifyCount3.ShouldBe(4); + } + + [Fact] + public void DiscardBroadcastsOnlyOnNonEmpty() + { + var queue = new AutoQueue(); + var enqueueValues = new List(); + var dequeueValues = new List(); + var discardCount = 0; + var modifyCount = 0; + + queue.Bind() + .OnEnqueue((in int v) => enqueueValues.Add(v)) + .OnDequeue((in int v) => dequeueValues.Add(v)) + .OnDiscard(count => discardCount += count) + .OnModify(() => modifyCount++); + + queue.Enqueue(1); + queue.Enqueue(2); + queue.Enqueue(3); + queue.Enqueue(4); + queue.Enqueue(5); + + queue.Dequeue(); // dequeue 1 + queue.Discard(2); // discard 2, 3 + + // dequeue 4, 5 + while (queue.HasValues) + { + queue.Dequeue(); + } + + queue.Enqueue(6); + queue.Discard(3); // discard 6 + queue.Discard(5); // no-op + + enqueueValues.ShouldBe([1, 2, 3, 4, 5, 6]); + dequeueValues.ShouldBe([1, 4, 5]); + discardCount.ShouldBe(3); + modifyCount.ShouldBe(11); + } + + [Fact] + public void ClearBroadcastsOnlyOnNonEmpty() + { + var queue = new AutoQueue(); + var clearCount = 0; + var modifyCount = 0; + + var binding = queue.Bind() + .OnClear(() => clearCount++) + .OnModify(() => modifyCount++); + + queue.Clear(); + + clearCount.ShouldBe(0); + modifyCount.ShouldBe(0); + clearCount = 0; + modifyCount = 0; + + queue.Enqueue(1); + queue.Enqueue(2); + queue.Clear(); + + clearCount.ShouldBe(1); + modifyCount.ShouldBe(3); + clearCount = 0; + modifyCount = 0; + + queue.Clear(); + + clearCount.ShouldBe(0); + modifyCount.ShouldBe(0); + } + + [Fact] + public void DisposedBindingStopsReceiving() + { + var queue = new AutoQueue(); + var modifyCount = 0; + + var binding = queue.Bind() + .OnModify(() => modifyCount++); + + queue.Enqueue(1); + modifyCount.ShouldBe(1); + + binding.Dispose(); + + queue.Enqueue(2); + modifyCount.ShouldBe(1); // Should not receive the second value + } + + [Fact] + public void ClearsBindings() + { + var queue = new AutoQueue(); + var modifyCount = 0; + + var binding = queue.Bind() + .OnModify(() => modifyCount++); + + queue.Enqueue(1); + queue.Enqueue(2); + + modifyCount.ShouldBe(2); + modifyCount = 0; + + queue.ClearBindings(); + queue.Enqueue(3); + modifyCount.ShouldBe(0); + } + + [Fact] + public void Disposes() + { + var queue = new AutoQueue(); + + queue.Dispose(); + + Should.Throw(() => queue.Enqueue(2)); + } +} diff --git a/Chickensoft.Sync/src/primitives/AutoQueue.cs b/Chickensoft.Sync/src/primitives/AutoQueue.cs new file mode 100644 index 0000000..8c26e3d --- /dev/null +++ b/Chickensoft.Sync/src/primitives/AutoQueue.cs @@ -0,0 +1,260 @@ +namespace Chickensoft.Sync.Primitives; + +using System; +using System.Diagnostics.CodeAnalysis; +using Chickensoft.Collections; + +/// +/// +/// A channel which queues value types, and broadcasts them on demand, on a +/// first-in-first-out basis +/// +/// +public interface IAutoQueue : IAutoObject; + +/// +public sealed class AutoQueue + : IAutoQueue, + IPerformAnyOperation, + IPerform, + IPerform, + IPerform +{ + // Atomic operations + private readonly record struct DequeueOp; + private readonly record struct DiscardOp(int Count); + private readonly record struct ClearOp; + + // Broadcasts + private readonly record struct EnqueueBroadcast(T Value) + where T : struct; + private readonly record struct DequeueBroadcast(T Value) + where T : struct; + private readonly record struct DiscardBroadcast(int Count); + private readonly record struct ClearBroadcast; + private readonly record struct ModifyBroadcast; + + /// + /// A binding to an . + /// + public sealed class Binding : SyncBinding + { + internal Binding(ISyncSubject subject) + : base(subject) { } + + /// + /// Registers a callback that is invoked whenever a value is enqueued of a + /// specific value type + /// + /// Callback to invoke. + /// Optional condition that must be true for the + /// callback to be invoked. + /// Value Type to Listen For + /// This binding (for chaining). + [SuppressMessage( + "Style", + "IDE0350", + Justification = "Implicit lambda with ref type won't compile" + )] + public Binding OnEnqueue(Callback callback, Func? condition = null) + where T : struct + { + bool predicate(T value) => condition?.Invoke(value) ?? true; + + AddCallback( + (in EnqueueBroadcast broadcast) => callback(broadcast.Value), + (in EnqueueBroadcast broadcast) => predicate(broadcast.Value) + ); + + return this; + } + + /// + /// Registers a callback that is invoked whenever a value is dequeued of a + /// specific value type + /// + /// Callback to invoke. + /// Optional condition that must be true for the + /// callback to be invoked. + /// Value Type to Listen For + /// This binding (for chaining). + [SuppressMessage( + "Style", + "IDE0350", + Justification = "Implicit lambda with ref type won't compile" + )] + public Binding OnDequeue(Callback callback, Func? condition = null) + where T : struct + { + bool predicate(T value) => condition?.Invoke(value) ?? true; + + AddCallback( + (in DequeueBroadcast broadcast) => callback(broadcast.Value), + (in DequeueBroadcast broadcast) => predicate(broadcast.Value) + ); + + return this; + } + + /// + /// Registers a callback to be invoked when items are discarded from the queue. + /// + /// + /// The value passed to the callback is the number of items actually discarded, + /// and not the number requested. + /// + /// Callback to be invoked. + /// This binding (for chaining). + public Binding OnDiscard(Action callback) + { + AddCallback((in DiscardBroadcast broadcast) => callback(broadcast.Count)); + return this; + } + + /// + /// Registers a callback to be invoked when the queue is cleared. + /// + /// Callback to be invoked. + /// This binding (for chaining). + public Binding OnClear(Action callback) + { + AddCallback((in ClearBroadcast _) => callback()); + return this; + } + + /// + /// Registers a callback to be invoked on any modification to the queue. + /// + /// Callback to be invoked. + /// This binding (for chaining). + public Binding OnModify(Action callback) + { + AddCallback((in ModifyBroadcast _) => callback()); + return this; + } + } + + private readonly struct MessageHandler : IBoxlessValueHandler + { + private readonly SyncSubject _subject; + + public MessageHandler(SyncSubject subject) + { + _subject = subject; + } + + public void HandleValue(in TValue value) + where TValue : struct => _subject.Broadcast(new DequeueBroadcast(value)); + } + + /// + /// Total number of values in the queue. + /// + public int Count => _queue.Count; + + /// + /// Returns whether the auto queue has any values. + /// + public bool HasValues => _queue.HasValues; + + private readonly SyncSubject _subject; + private readonly BoxlessQueue _queue; + private readonly MessageHandler _handler; + + /// + /// + /// Creates a new auto queue. + /// + /// + /// + /// + /// + public AutoQueue() + { + _subject = new(this); + _queue = new(); + _handler = new(_subject); + } + + /// + public Binding Bind() => new(_subject); + + /// + public void ClearBindings() => _subject.ClearBindings(); + + /// + public void Dispose() => _subject.Dispose(); + + /// + /// Add a message to the queue without boxing it. + /// + /// The type of message to enqueue. + public void Enqueue(in T value) + where T : struct => _subject.Perform(value); + + /// + /// Dequeues and broadcasts the next value in the queue, if any. + /// + public void Dequeue() => _subject.Perform(new DequeueOp()); + + /// + /// Discard a number of values from the queue without handling them. + /// + /// + /// The number of values to discard. If this is greater than the number of + /// values in the queue, all values will be discarded. + /// + public void Discard(int count = 1) => _subject.Perform(new DiscardOp(count)); + + /// Clear all values from the queue. + public void Clear() => _subject.Perform(new ClearOp()); + + void IPerformAnyOperation.Perform(in TOp op) + where TOp : struct + { + if (op is not DequeueOp and not DiscardOp and not ClearOp) + { + _queue.Enqueue(op); + _subject.Broadcast(new EnqueueBroadcast(op)); + _subject.Broadcast(new ModifyBroadcast()); + } + } + + void IPerform.Perform(in DequeueOp op) + { + if (!_queue.HasValues) + { + return; + } + + _queue.Dequeue(_handler); + _subject.Broadcast(new ModifyBroadcast()); + } + + void IPerform.Perform(in DiscardOp op) + { + var count = op.Count; + if (!_queue.HasValues || count <= 0) + { + return; + } + + var discardCount = count > _queue.Count ? _queue.Count : count; + + _queue.Discard(count); + _subject.Broadcast(new DiscardBroadcast(discardCount)); + _subject.Broadcast(new ModifyBroadcast()); + } + + void IPerform.Perform(in ClearOp op) + { + if (!_queue.HasValues) + { + return; + } + + _queue.Clear(); + _subject.Broadcast(new ClearBroadcast()); + _subject.Broadcast(new ModifyBroadcast()); + } +}