Zongsoft.Messaging.ZeroMQ
2.0.0
See the version list below for details.
dotnet add package Zongsoft.Messaging.ZeroMQ --version 2.0.0
NuGet\Install-Package Zongsoft.Messaging.ZeroMQ -Version 2.0.0
<PackageReference Include="Zongsoft.Messaging.ZeroMQ" Version="2.0.0" />
<PackageVersion Include="Zongsoft.Messaging.ZeroMQ" Version="2.0.0" />
<PackageReference Include="Zongsoft.Messaging.ZeroMQ" />
paket add Zongsoft.Messaging.ZeroMQ --version 2.0.0
#r "nuget: Zongsoft.Messaging.ZeroMQ, 2.0.0"
#:package Zongsoft.Messaging.ZeroMQ@2.0.0
#addin nuget:?package=Zongsoft.Messaging.ZeroMQ&version=2.0.0
#tool nuget:?package=Zongsoft.Messaging.ZeroMQ&version=2.0.0
Zongsoft.Messaging.ZeroMQ Message Queue Plugin
<a name="abstract"></a>
Abstract
Zongsoft.Messaging.ZeroMQ is a NetMQ-based adapter for the messaging and communication abstractions in Zongsoft.Core. It provides topic publishing and subscription through IMessageQueue, and also supplies request/response and event-channel adapters.
The included ZeroQueueServer combines an XPUB/XSUB exchange for at-most-once traffic with a durable acknowledgement path for at-least-once traffic. Clients discover the current Broker epoch and all runtime endpoints automatically.
<a name="features"></a>
Features
- Implements the Zongsoft
IMessageQueue,IRequester,IResponder, andIEventChannelabstractions; - Supports multiple publishers and subscribers through an XPUB/XSUB exchange;
- Supports topic prefixes, optional message groups, instance filtering, and heartbeats;
- Supports Brotli, GZip, ZLib, or Deflate payload compression above a configurable threshold;
- Supports immediate at-most-once broadcast and Broker-persisted, explicitly acknowledged competing at-least-once delivery;
- Supports standalone use and Zongsoft plugin-based hosting;
- Targets .NET 8, .NET 9, and .NET 10.
<a name="installation"></a>
Installation
Install the NuGet package:
dotnet add package Zongsoft.Messaging.ZeroMQ
To build from this repository, build Zongsoft.Core first:
dotnet build Zongsoft.Core/src/Zongsoft.Core.csproj
dotnet build messaging/zero/Zongsoft.Messaging.ZeroMQ.slnx
<a name="topology"></a>
Exchange Topology
| Endpoint | Default in packaged configuration | Purpose |
|---|---|---|
| Discovery | 7969 |
Clients request the Broker epoch and runtime endpoint ports. |
| Reliability control | 32100 |
Reliable subscription registration, delivery, and acknowledgement. |
| Publisher incoming | 32101 |
Application publishers connect here. |
| Subscriber outgoing | 32102 |
Application subscribers connect here. |
7969 is the built-in discovery-port default. Omitted runtime ports are selected dynamically and rediscovered after a Broker restart. Fixed control, incoming, and outgoing ports are still recommended for predictable firewall and operations configuration.
The server binds TCP endpoints on all network interfaces and does not configure authentication or encryption. Restrict access at the host or network boundary, or add an authenticated transport before using it across an untrusted network.
<a name="configuration"></a>
Configuration
Server
The packaged daemon plugin starts ZeroQueueServer automatically. Configure its data endpoints under /Messaging/ZeroMQ/Servers:
<configuration>
<option path="/Messaging/ZeroMQ">
<servers port="32100,32101,32102">
<server server.name="unnamed" />
</servers>
</option>
</configuration>
The three values are reliability control, publisher incoming, and subscriber outgoing. A two-value configuration is Incoming,Outgoing; when Storage is available, Control is selected dynamically. Control starts only when the Server has an IMessageStorage. Port precedence is: explicit startup arguments, the named server's own Port, the collection-level Servers.Port, then dynamic selection. * explicitly requests dynamic ports.
The actual server sample starts a broadcast server without persistent storage. Its startup excerpt is below; see the plugin/storage sections for reliable broker composition:
using var server = new ZeroQueueServer();
await server.StartAsync(["--incoming:32101", "--outgoing:32102"]);
Client Connection
Define a ZeroMQ connection under /Messaging/ConnectionSettings:
The client name and group come from the sample client, expressed here in host option format:
<configuration>
<option path="/Messaging">
<connectionSettings default="ZeroMQ">
<connectionSetting connectionSetting.name="ZeroMQ"
driver="ZeroMQ"
value="server=127.0.0.1;port=7969;group=Demo;client=Zongsoft.Messaging.ZeroMQ.Sample;" />
</connectionSettings>
</option>
</configuration>
| Setting | Default | Description |
|---|---|---|
Server |
required | Host name or IP address of the discovery endpoint. Do not include tcp://. |
Port |
7969 |
Discovery endpoint port. |
Topic |
empty | Topic used when ProduceAsync or SubscribeAsync omits a topic. |
Group |
empty | Prefix added as Group:Topic to isolate applications sharing an exchange. |
Client |
empty | Stable client name used as part of an automatically generated instance identifier. |
Instance |
generated | Explicit producer instance identifier. Empty or * generates a unique identifier. |
Filter |
excludes self | Comma-separated instance filter controlling which producers are accepted. |
Timeout |
10s |
Discovery and subscription-synchronization timeout. |
Heartbeat |
10s |
Heartbeat interval. A value less than or equal to zero disables heartbeats. |
ReconnectInterval |
1s |
Minimum interval between endpoint rediscovery attempts. |
The default filter excludes messages produced by the same queue instance. Use Filter=* to accept every instance, Filter=. (or ~) to accept only the current instance, ordinary identifiers as an allow list, and !identifier entries as exclusions.
<a name="usage"></a>
Usage
Publish and Subscribe
In a host with the ZeroMQ plugin deployed, obtain the shared queue by provider and connection name. This uses the ZeroMQ connection above; for this in-process loopback example set filter=* on that connection and start a compatible broker beforehand. The consumer module references only Core:
using System.Text;
using Zongsoft.Services;
using Zongsoft.Messaging;
var provider = ApplicationContext.Current.Services
.FindRequired<IMessageQueueProvider>("ZeroMQ");
var queue = provider.Queue("ZeroMQ");
var received = new TaskCompletionSource<string>(TaskCreationOptions.RunContinuationsAsynchronously);
var topic = $"docs/{Guid.NewGuid():N}";
var consumer = await queue.SubscribeAsync(topic, message =>
received.TrySetResult(Encoding.UTF8.GetString(message.Data.Span)));
try
{
var identifier = await queue.ProduceAsync(topic, "demo".AsMemory());
if(identifier == null)
throw new InvalidOperationException("No matching subscription was visible.");
Console.WriteLine(await received.Task.WaitAsync(TimeSpan.FromSeconds(10)));
}
finally
{
await consumer.UnsubscribeAsync();
}
The provider reuses named queues; business operations must not dispose a shared queue. This example cancels only its own unique-topic subscription and waits for its handler before exiting. A non-null publish identifier still does not mean business processing completed. A standalone tool that constructs ZeroQueue directly owns its entire lifetime.
Subscriptions use prefix matching. One ZeroQueue keeps one consumer for each logical topic; subscribing to the same topic again returns the existing consumer and does not replace its handler or options. With Group=Demo, the physical wire topic is Demo:topic/reliable, while handlers receive the logical Message.Topic value topic/reliable.
Each subscriber invokes its handler sequentially in receive order. When its bounded pending queue reaches capacity, that subscriber pauses Poller reads and resumes after the handler frees capacity; other sockets remain responsive.
Compression
Set MessageEnqueueOptions.Compression to an algorithm and the minimum payload size that enables compression. MessageCompression.Value is an integer byte threshold; its text format is <algorithm>:<threshold>. Zero compresses every non-empty payload and the default value disables compression:
var options = new MessageEnqueueOptions()
{
Compression = MessageCompression.Parse("Brotli:4096"),
};
await queue.ProduceAsync("documents/updated", payload, options);
The equivalent strongly typed construction is new MessageCompression("Brotli", 4096).
At-least-once delivery
Set LeastOnce on both the subscription and publication. Only the Broker requires an IMessageStorage; publishers do not persist messages. A handler must call AcknowledgeAsync; returning normally is not an acknowledgement.
The following is from ZeroQueueReliabilityTests.LeastOnceCompletesAfterBrokerAcceptanceBeforeAcknowledge. ReliableServerScope, CreateQueue and AcknowledgingHandler are real fixtures in that test file. Read or run them in the original test project; they are not application plugin types:
await using var scope = await ReliableServerScope.StartAsync();
using var publisher = CreateQueue(scope.Port, "publisher", "publisher");
using var subscriber = CreateQueue(scope.Port, "subscriber", "subscriber");
var handler = new AcknowledgingHandler(1, false);
await subscriber.SubscribeAsync("topic/reliable", handler, ReliableSubscribeOptions());
var identifier = await publisher.ProduceAsync("topic/reliable", Encoding.UTF8.GetBytes("reliable"), ReliableEnqueueOptions()).AsTask().WaitAsync(TimeSpan.FromSeconds(5));
var message = await handler.ReceiveAsync(TimeSpan.FromSeconds(5));
Assert.False(string.IsNullOrWhiteSpace(identifier));
Assert.Equal(identifier, message.Identifier);
Assert.Single(await GetPendingAsync(scope));
await message.AcknowledgeAsync();
Assert.True(await WaitForPendingCountAsync(scope, 0, TimeSpan.FromSeconds(5)));
The fixture checks that Pending remains before acknowledgement and is removed afterwards. It is not an order-storage implementation or proof of idempotent business effects. Its in-process test storage does not establish production restart durability.
The Broker accepts a publication only when an online matching subscription exists. No match returns null without writing Storage. With a match, the Broker persists Pending first and then returns the identifier. Delivery competes among online subscribers; any one acknowledgement removes Pending. Retries reuse Message.Identifier and may choose another consumer, so handlers must be idempotent.
Assign ZeroQueueServer.Storages only while the Server is stopped. On first start, the Server calls Storages.Create(Name) and owns the returned storage. Ordinary Stop retains that storage for restart; replacing the factory, a failed start, or disposing the Server releases it, preferring IAsyncDisposable. A Broker without a storage factory still serves MostOnce Broadcast, but does not start Control and returns only Incoming,Outgoing in discovery Ports; LeastOnce operations then fail.
| Messaging option | Support |
|---|---|
Compression |
Supported by both MostOnce and LeastOnce with Brotli, GZip, ZLib, or Deflate; only Message.Data is compressed. |
| Tags and identity | Both modes carry Identifier, Identity, and Tags as independent metadata. |
Delay |
Unsupported; Core checks Features and rejects a positive delay before entering the driver. |
| Expiration | Supported by LeastOnce; zero means no expiration. |
| Priority | Not implemented. |
MostOnce |
Supported; returns null when no subscription is visible at send time, otherwise sends locally once. |
LeastOnce |
Supported with Broker persistence, competing consumers, explicit acknowledgement, and same-identifier retry. |
ExactlyOnce |
Not supported and fails before transport state is created. |
| Subscription fallback | Not implemented by the current handler dispatcher. |
Request and Response
ZeroRequester and ZeroResponder exchange requests through logical URL topics; the default reply topic is <url>/reply. The real ZeroRequesterTests.RequesterReceivesImmediateResponses starts separate requester/responder queues, registers the same file's EchoHandler, sends to rpc/echo, and checks payload and request-identifier correlation. It stops/disposes the responder in finally and also disposes each request token.
ZeroResponderTests additionally checks subscription rollback on startup failure. No nonexistent PingHandler is supplied here. Applications use shared communication contracts and host composition; standalone tests own their directly constructed instances.
Event Channel
ZeroQueueEventChannel connects an EventExchanger to queue topics under Events/...:
await using var channel = new ZeroQueueEventChannel(queue);
await channel.OpenAsync(exchanger);
await channel.SendAsync(eventContext);
The plugin manifest registers this channel automatically for hosted applications. Request/response and event channels support Group; the prefix is applied only at the network boundary and adapters always use logical topics.
<a name="semantics"></a>
Delivery Semantics
The selected MessageReliability determines the contract:
MostOnceis transient broadcast. If the application XPUB sees no matching subscription at send time,ProduceAsyncreturnsnullwithout sending. Otherwise it sends locally once and returns a unique identifier. Every broadcast recipient sees that identifier, but it is not remote or Handler acknowledgement.LeastOncereturnsnullwithout persistence when no online matching subscription exists. Otherwise the Broker persists Pending first and immediately returns the unique identifier without waiting for a Handler acknowledgement.- The Broker selects one online consumer per attempt. It retries the same identifier until any valid acknowledgement removes Pending. If all consumers disconnect after acceptance, Pending remains until a subscription returns.
LeastOncepermits duplicate handler invocations. It does not deduplicate business effects and does not provide exactly-once delivery.- A Control timeout, caller cancellation, or disconnect may leave the publisher unable to determine whether the Broker accepted the message. These stop only local waiting and cannot revoke acceptance already in progress; a business retry can still produce a duplicate.
- Expired reliable messages are removed from Broker Pending with a diagnostic record.
- A queue snapshots its connection, ports, group, filter, timeout, and heartbeat settings at construction; mutating the original settings object does not reconfigure a running queue;
- Empty business payloads are supported.
Message storage is an independent plugin concept, not part of the ZeroMQ driver. This package provides no default file store. Applications using LeastOnce inject an IMessageStorageFactory; the factory looks up a connection with the same name as the Broker and creates its exclusive storage. Storage supports exact-topic reads and clears. An implementation must hold a message snapshot before SetAsync returns and provide the required restart durability.
See the ZeroMQ 1.0 protocol for the complete Discovery, Broadcast, and Control frame definitions.
<a name="samples"></a>
Samples and Troubleshooting
The .NET 10 samples contain an interactive exchange server and client. Start the server first, then run one client as a subscriber and another as a publisher. See the sample guide for commands.
If messages are not received:
- Verify that discovery, control, incoming, and outgoing ports are reachable in the required directions;
- Verify that publisher and subscriber use the same
Groupand compatible topic prefixes; - Check the
Filtersetting—self-produced messages are excluded by default; - If
ProduceAsyncreturnsnull, verify that a matching subscription was visible to the Broker at that instant and apply the application's retry policy if appropriate; - For
LeastOnce, verify ServerStorage, the Control endpoint, explicit acknowledgement, and expiration; - Inspect Broker Pending data, online subscriptions, and consumer idempotency when reliable delivery remains unresolved.
Plugin-Based Integration
Compose this feature through the host; a package reference supplies compile-time APIs, while plugin loading also requires deployed manifests and runtime assets. See the complete plugin workflow.
The deployment manifest selects Zongsoft.Messaging.ZeroMQ-$(site).plugin. site=daemon includes the broker startup contribution; client-only and broker hosts must be deliberately distinguished. Configure the endpoint and storage factory before starting the broker.
| Runtime artifact | Source of truth |
|---|---|
Zongsoft.Messaging.ZeroMQ.Daemon |
Zongsoft.Messaging.ZeroMQ-daemon.plugin |
Zongsoft.Messaging.ZeroMQ |
Zongsoft.Messaging.ZeroMQ.plugin |
Zongsoft.Messaging.ZeroMQ.Storage |
Zongsoft.Messaging.ZeroMQ.Storage.plugin |
| File copying and dependencies | Zongsoft.Messaging.ZeroMQ.deploy |
Add this fragment to an existing host .deploy (retain Main and the host’s other base manifests; do not replace the whole file):
[plugins zongsoft messaging zeromq]
nuget:Zongsoft.Messaging.ZeroMQ
Run dotnet deploy against a test deployment as explained in the workflow, with the host's framework, platform, architecture and, where needed, site. Pin compatible versions in real deployments; application dependencies such as databases, caches or commercial runtimes are still separate prerequisites.
Additional artifacts listed by the deployment manifest include Zongsoft.Messaging.ZeroMQ.option, Zongsoft.Messaging.ZeroMQ.plugin, Zongsoft.Messaging.ZeroMQ-$(site).plugin, Zongsoft.Messaging.ZeroMQ.Storage.plugin. Retain assemblies, dependencies and satellite resource directories as well. Restart the host after deployment, check plugin loading and service/driver registration, then verify the workflow above; copied files alone do not prove that the feature is active.
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net8.0 is compatible. net8.0-android was computed. net8.0-browser was computed. net8.0-ios was computed. net8.0-maccatalyst was computed. net8.0-macos was computed. net8.0-tvos was computed. net8.0-windows was computed. net9.0 is compatible. net9.0-android was computed. net9.0-browser was computed. net9.0-ios was computed. net9.0-maccatalyst was computed. net9.0-macos was computed. net9.0-tvos was computed. net9.0-windows was computed. net10.0 is compatible. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
-
net10.0
- Microsoft.Extensions.Configuration.Abstractions (>= 10.0.0)
- NetMQ (>= 4.0.4.3)
- Zongsoft.Core (>= 7.59.0)
-
net8.0
- Microsoft.Extensions.Configuration.Abstractions (>= 8.0.0)
- NetMQ (>= 4.0.4.3)
- Zongsoft.Core (>= 7.59.0)
-
net9.0
- Microsoft.Extensions.Configuration.Abstractions (>= 9.0.2)
- NetMQ (>= 4.0.4.3)
- Zongsoft.Core (>= 7.59.0)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 2.1.0 | 75 | 9/26/2026 |
| 2.0.0 | 92 | 9/18/2026 |
| 1.8.1 | 139 | 7/23/2026 |
| 1.8.0 | 126 | 6/30/2026 |
| 1.7.0 | 127 | 6/23/2026 |
| 1.6.1 | 118 | 4/29/2026 |
| 1.6.0 | 125 | 4/1/2026 |
| 1.5.1 | 118 | 4/1/2026 |
| 1.5.0 | 119 | 3/31/2026 |
| 1.4.1 | 156 | 1/27/2026 |
| 1.4.0 | 325 | 5/6/2025 |
| 1.3.0 | 250 | 4/28/2025 |
| 1.2.0 | 240 | 4/25/2025 |
| 1.1.0 | 241 | 3/27/2025 |
| 1.0.0 | 228 | 2/24/2025 |