Skip to content

Channels

Channels are message pipelines that carry messages between framework components.


What is a Channel?

+------------------+                      +------------------+
|     Producer     |                      |     Consumer     |
| (writes messages)|--->[  Channel  ]---->|(reads messages)  |
+------------------+                      +------------------+

Channel = async queue that connects producers to consumers

Channels are:

  • Pipelines - messages flow through them
  • Async - non-blocking writes and reads
  • Buffered - can hold messages temporarily
  • Typed - each channel carries specific message type

The Three Queues

                         Adapter
+----------------------------------------------------------+
|                                                           |
|  INBOUND (from PQ)              OUTBOUND (to PQ)         |
|                                                           |
|  +------------------+           +------------------+      |
|  | CommandChannel   |           | EventPublisher   |      |
|  +------------------+           +------------------+      |
|         |                              ^                  |
|         v                              |                  |
|  [Command Handlers]             +------------------+      |
|                                 | StatusPublisher  |      |
|                                 +------------------+      |
|                                        ^                  |
|                                        |                  |
|                                 [Your device code]        |
|                                                           |
+----------------------------------------------------------+
Queue Direction Carries Purpose
CommandChannel PQ -> Adapter Commands Every request from PQ (unlock, sync, restart)
EventPublisher Adapter -> PQ Events Something happened (access granted)
StatusPublisher Adapter -> PQ Status Current state (online, door locked)

All three are unbounded in-memory queues with a background reader. NatsCommunication owns them and is the single Start/Stop entry point.


Message Flow

Inbound (Commands)

PQ System
    |
    v (NATS message)
+-------------------+
| NatsCommand       |
| Receiver          |
+-------------------+
    |
    v (writes to channel)
+-------------------+
| Command           |
| Channel           |
+-------------------+
    |
    v (executor reads)
+-------------------+
| Command           |
| Executor          |
+-------------------+
    |
    v (calls handler)
+-------------------+
| Your Command      |
| Handler           |
+-------------------+

Outbound (Events)

Your Code
    |
    | thing.PublishEvent()
    v
+-------------------+
| Event             |
| Publisher queue   |
+-------------------+
    |
    v (background loop reads)
+-------------------+
| Event             |
| Publisher         |
+-------------------+
    |
    v (NATS message)
PQ System

Why Separate Queues?

Events vs Status

+------------------+     +------------------+
| EventPublisher   |     | StatusPublisher  |
+------------------+     +------------------+
        |                        |
        v                        v
  High volume              Low volume
  Critical delivery        Latest-only OK
  (access.granted)         (heartbeat)

Events must all be delivered. Status only needs latest value.


Queue Depth

A queue grows when the reader falls behind the writer:

+--------------------------------------------------+
|                    Queue                          |
|  [msg][msg][msg][msg][msg][msg][msg][msg]        |
|   ^                                      ^        |
|   |                                      |        |
|   Write pointer                   Read pointer    |
+--------------------------------------------------+

If the queue keeps growing:
- The device is producing faster than NATS accepts
- Memory grows with it, so the depth warning is your alarm

ChannelConfiguration sets the depth at which the framework logs a warning:

var channels = new ChannelConfiguration
{
    DataCommandDepthWarningThreshold = 1000,   // default 1000
    EventChannelDepthWarningThreshold = 10000, // default 10000
    StatusChannelDepthWarningThreshold = 10000 // default 10000
};

Reaching the Queues from Your Code

You never touch a queue directly. The Thing that owns the device is the publishing surface.

Publishing Events

// inside the Thing that owns the door
public async Task OnDeviceEvent(DeviceAccessLog log, CancellationToken ct)
{
    var deviceEvent = this.AccessGranted()
        .At(log.Timestamp, this)
        .WithParameter("credential", log.CardNumber)
        .Build();

    await PublishEvent(deviceEvent, ct);
}

Thing.PublishEvent hands the event to the adapter, which writes it to the event queue. The framework builds the NATS subject and publishes it.

Publishing Status

Status is not an event. Collect the current state in a batch and commit it:

await thing.StatusBatch()
    .Set(ConnectionState.Online)
    .Set(LockState.Locked)
    .Commit(ct);

Who Owns the Queues

NatsCommunication is the facade. The framework host builds it, starts it and stops it; your adapter code does not register it:

public sealed class NatsCommunication
{
    public Task Start(/* adapter id, category, dispatcher */ CancellationToken ct);
    public Task Stop();

    public ValueTask<bool> Publish(EventMessage message, CancellationToken ct = default);
    public ValueTask<bool> Publish(StatusMessage message, CancellationToken ct = default);
}

Behind it sit NatsCommandReceiver → CommandChannel on the way in, and EventPublisher / StatusPublisher on the way out. Both publishers implement IPublisher<TMessage> (TryWrite, Start, Stop).


Command Processing

Commands flow through processor:

+-------------------+
| Command           |
| Channel           |
+-------------------+
         |
         v
+-------------------+
| NatsCommand       |
| Receiver          |
+-------------------+
         |
         | for each command:
         v
+-------------------+
| CommandExecutor   |  <-- your handler
| (delegate)        |
+-------------------+
         |
         v
+-------------------+
| ResponseMessage   |
+-------------------+
         |
         v
+-------------------+
| Return via NATS   |
+-------------------+

CommandExecutor is the delegate that runs one command. The receiver reads the channel and calls it with the configured parallelism. Your own command handlers are reached through the source-generated CommandDispatcher, so nothing here needs wiring by hand.


Graceful Shutdown

On shutdown, channels drain:

Stop signal received
        |
        v
+-------------------+
| Stop accepting    |
| new messages      |
+-------------------+
        |
        v
+-------------------+
| Process remaining |
| messages in queue |
+-------------------+
        |
        | (up to 30 second timeout)
        v
+-------------------+
| Force close       |
+-------------------+

Events especially need time to flush - don't lose audit data.


Backpressure

When channel fills up:

[Producer] --write--> [FULL Channel] --X-- blocked

Producer waits until:
1. Consumer reads some messages (space available)
2. Timeout expires
3. Cancellation requested
// publishing can be cancelled during shutdown
try
{
    await PublishEvent(deviceEvent, ct);
}
catch (OperationCanceledException)
{
    _logger.LogWarning("Event publish cancelled (shutdown)");
}

Channel Depth Monitoring

Monitor for problems:

Normal:     [msg][msg][   ][   ][   ]     depth: 2
Warning:    [msg][msg][msg][msg][msg]     depth: 5 (threshold hit)
Critical:   [msg][msg][msg][msg][msg]...  depth: 100+ (backlog)

When depth warning triggers:

  • Something is slow (consumer not keeping up)
  • Or burst of events (temporary spike)

Check:

  • Is protocol blocking?
  • Is NATS connection slow?
  • Is device flooding events?

Common Mistakes

1. Reaching Past the Thing

// WRONG - writing to the publisher queue yourself
_eventPublisher.TryWrite(message, ct);

// RIGHT - publish from the Thing that owns the device
await PublishEvent(deviceEvent, ct);

2. Fire and Forget

// WRONG - ignoring result
PublishEvent(deviceEvent, ct);  // no await!

// RIGHT - await, so a failed write is visible
await PublishEvent(deviceEvent, ct);

3. Not Handling Cancellation

// WRONG - ignores shutdown
while (true)
{
    await PublishHeartbeat();
    await Task.Delay(30000);  // ignores ct
}

// RIGHT - respects cancellation
while (!ct.IsCancellationRequested)
{
    await PublishHeartbeat(ct);
    await Task.Delay(30000, ct);
}

Summary

+--------------------------------------------------------------+
|                        Adapter                                |
|                                                               |
|   PQ --> [CommandChannel] --> [Receiver] --> Handler         |
|                                                               |
|   Thing --> [EventPublisher] --> PQ                          |
|   Thing --> [StatusPublisher] --> PQ                         |
|                                                               |
+--------------------------------------------------------------+

The queues are the plumbing. NatsCommunication owns them, and your code reaches them through the Thing.


Done with Concepts!

Now you understand the building blocks. Continue to: