Zongsoft.Messaging.ZeroMQ 2.0.0

There is a newer version of this package available.
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
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="Zongsoft.Messaging.ZeroMQ" Version="2.0.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Zongsoft.Messaging.ZeroMQ" Version="2.0.0" />
                    
Directory.Packages.props
<PackageReference Include="Zongsoft.Messaging.ZeroMQ" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add Zongsoft.Messaging.ZeroMQ --version 2.0.0
                    
#r "nuget: Zongsoft.Messaging.ZeroMQ, 2.0.0"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package Zongsoft.Messaging.ZeroMQ@2.0.0
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=Zongsoft.Messaging.ZeroMQ&version=2.0.0
                    
Install as a Cake Addin
#tool nuget:?package=Zongsoft.Messaging.ZeroMQ&version=2.0.0
                    
Install as a Cake Tool

Zongsoft.Messaging.ZeroMQ Message Queue Plugin

License NuGet Version NuGet Downloads GitHub Stars

English | 简体中文


<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, and IEventChannel abstractions;
  • 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:

  • MostOnce is transient broadcast. If the application XPUB sees no matching subscription at send time, ProduceAsync returns null without 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.
  • LeastOnce returns null without 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.
  • LeastOnce permits 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:

  1. Verify that discovery, control, incoming, and outgoing ports are reachable in the required directions;
  2. Verify that publisher and subscriber use the same Group and compatible topic prefixes;
  3. Check the Filter setting—self-produced messages are excluded by default;
  4. If ProduceAsync returns null, verify that a matching subscription was visible to the Broker at that instant and apply the application's retry policy if appropriate;
  5. For LeastOnce, verify Server Storage, the Control endpoint, explicit acknowledgement, and expiration;
  6. 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 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. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

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
Loading failed