DistributedCacheEventOutbox Class
The outbox over the host's Microsoft.Extensions.Caching.Distributed.IDistributedCache: one entry per stream holding the queue, so pending events survive a process restart when the store behind the cache does. A stream here is the receiver and the identifier together, as everywhere else. The tier is deliberate, and it is our decision rather than a permission the specification grants: SSF 1.0 Section 8.1.2.1 lets a transmitter drop events held while a stream is PAUSED, and requires transmission for an enabled one. Treating the whole queue as cache-tier follows the delivery protocols' own tolerance for loss over a broken transport, and is why it belongs in the cache tier rather than beside data that earns backups.
public sealed class DistributedCacheEventOutbox : Abblix.SharedSignals.Transmitter.IEventOutbox, System.IDisposableInheritance System.Object → DistributedCacheEventOutbox
Implements IEventOutbox, System.IDisposable
Remarks
Microsoft.Extensions.Caching.Distributed.IDistributedCache reads and writes whole values with no compare-and-set, so queue
mutations are serialized through an in-process gate per stream, taken under the same composed
key the entry lives under - so the gate and the entry cannot disagree about which stream is
being guarded. That gate excludes this
instance's threads from each other and reaches no further, which makes this implementation
correct for a SINGLE transmitter instance and only that: two instances mutating one stream's
queue read the same value, each writes its own edit over the whole entry, and the later write
silently discards the earlier one's. No compare-and-set means the interface cannot express the
fix either - it is not a gap in this class. A transmitter running more than one instance takes
the outbox built on native list operations, AddSharedSignalsRedisOutbox.
Constructors
DistributedCacheEventOutbox(IDistributedCache) Constructor
The outbox over the host's Microsoft.Extensions.Caching.Distributed.IDistributedCache: one entry per stream holding the queue, so pending events survive a process restart when the store behind the cache does. A stream here is the receiver and the identifier together, as everywhere else. The tier is deliberate, and it is our decision rather than a permission the specification grants: SSF 1.0 Section 8.1.2.1 lets a transmitter drop events held while a stream is PAUSED, and requires transmission for an enabled one. Treating the whole queue as cache-tier follows the delivery protocols' own tolerance for loss over a broken transport, and is why it belongs in the cache tier rather than beside data that earns backups.
public DistributedCacheEventOutbox(Microsoft.Extensions.Caching.Distributed.IDistributedCache cache);Parameters
cache Microsoft.Extensions.Caching.Distributed.IDistributedCache
The distributed cache the queues live in; the store is the host's choice.
Remarks
Microsoft.Extensions.Caching.Distributed.IDistributedCache reads and writes whole values with no compare-and-set, so queue
mutations are serialized through an in-process gate per stream, taken under the same composed
key the entry lives under - so the gate and the entry cannot disagree about which stream is
being guarded. That gate excludes this
instance's threads from each other and reaches no further, which makes this implementation
correct for a SINGLE transmitter instance and only that: two instances mutating one stream's
queue read the same value, each writes its own edit over the whole entry, and the later write
silently discards the earlier one's. No compare-and-set means the interface cannot express the
fix either - it is not a gap in this class. A transmitter running more than one instance takes
the outbox built on native list operations, AddSharedSignalsRedisOutbox.
Methods
DistributedCacheEventOutbox.AcknowledgeAsync(string, string, IReadOnlyCollection<string>, CancellationToken) Method
Removes acknowledged SETs from a stream's queue, releasing the transmitter from retaining them (RFC 8936 Section 2.2). Identifiers with nothing to match are ignored - an acknowledgement can only arrive for something that was once here.
public System.Threading.Tasks.Task AcknowledgeAsync(string receiverId, string streamId, System.Collections.Generic.IReadOnlyCollection<string> jwtIds, System.Threading.CancellationToken cancellationToken=default(System.Threading.CancellationToken));Parameters
receiverId System.String
The receiver the stream belongs to.
streamId System.String
The stream whose queue is acknowledged.
jwtIds System.Collections.Generic.IReadOnlyCollection<System.String>
The "jti" values being acknowledged.
cancellationToken System.Threading.CancellationToken
Cancels I/O a durable implementation performs.
Implements AcknowledgeAsync(string, string, IReadOnlyCollection<string>, CancellationToken)
Returns
DistributedCacheEventOutbox.ClearAsync(string, string, CancellationToken) Method
Drops a stream's whole queue - the companion of deleting or disabling the stream, whose events are not held for later (SSF 1.0 Sections 8.1.1.5, 8.1.2.1).
public System.Threading.Tasks.Task ClearAsync(string receiverId, string streamId, System.Threading.CancellationToken cancellationToken=default(System.Threading.CancellationToken));Parameters
receiverId System.String
The receiver the stream belongs to.
streamId System.String
The stream whose queue is dropped.
cancellationToken System.Threading.CancellationToken
Cancels I/O a durable implementation performs.
Implements ClearAsync(string, string, CancellationToken)
Returns
DistributedCacheEventOutbox.Dispose() Method
Performs application-defined tasks associated with freeing, releasing, or resetting unmanaged resources.
public void Dispose();Implements Dispose()
DistributedCacheEventOutbox.EnqueueAsync(string, string, OutboxItem, CancellationToken) Method
Appends a SET to a stream's queue.
public System.Threading.Tasks.Task EnqueueAsync(string receiverId, string streamId, Abblix.SharedSignals.Transmitter.OutboxItem item, System.Threading.CancellationToken cancellationToken=default(System.Threading.CancellationToken));Parameters
receiverId System.String
The receiver the stream belongs to.
streamId System.String
The stream the SET was minted for.
item OutboxItem
The minted SET.
cancellationToken System.Threading.CancellationToken
Cancels I/O a durable implementation performs.
Implements EnqueueAsync(string, string, OutboxItem, CancellationToken)
Returns
DistributedCacheEventOutbox.PendingAsync(string, string, Nullable<int>, CancellationToken) Method
Reads the unacknowledged head of a stream's queue, oldest first, without removing anything - redelivery of the unacknowledged is the delivery protocols' own semantics.
public System.Threading.Tasks.Task<System.Collections.Generic.IReadOnlyList<Abblix.SharedSignals.Transmitter.OutboxItem>> PendingAsync(string receiverId, string streamId, System.Nullable<int> maxCount=null, System.Threading.CancellationToken cancellationToken=default(System.Threading.CancellationToken));Parameters
receiverId System.String
The receiver the stream belongs to.
streamId System.String
The stream whose queue is read.
maxCount System.Nullable<System.Int32>
The most items to return; null returns everything pending, mirroring an absent "maxEvents" (RFC 8936 Section 2.2).
cancellationToken System.Threading.CancellationToken
Cancels I/O a durable implementation performs.
Implements PendingAsync(string, string, Nullable<int>, CancellationToken)
Returns
System.Threading.Tasks.Task<System.Collections.Generic.IReadOnlyList<OutboxItem>>