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:
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:
- Hello World Tutorial - build your first adapter
- Access Synchronization Tutorial - sync credentials to devices