Skip to content

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.

C#
public sealed class DistributedCacheEventOutbox : Abblix.SharedSignals.Transmitter.IEventOutbox, System.IDisposable

Inheritance 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.

C#
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.

C#
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

System.Threading.Tasks.Task

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).

C#
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

System.Threading.Tasks.Task

DistributedCacheEventOutbox.Dispose() Method

Performs application-defined tasks associated with freeing, releasing, or resetting unmanaged resources.

C#
public void Dispose();

Implements Dispose()

DistributedCacheEventOutbox.EnqueueAsync(string, string, OutboxItem, CancellationToken) Method

Appends a SET to a stream's queue.

C#
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

System.Threading.Tasks.Task

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.

C#
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>>