From 949a67ee827cf4defc730f83d17ea0249e4a9ccd Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Fri, 24 Apr 2026 12:50:58 +0000 Subject: [PATCH 1/8] Add multi-targeting (netstandard2.0 + net10.0) and initial metrics instrumentation Agent-Logs-Url: https://github.com/Abc-Arbitrage/Zebus/sessions/8f197dc4-0b8b-4ce0-a40d-97ea7c647727 Co-authored-by: Kuinox <18743295+Kuinox@users.noreply.github.com> --- src/Abc.Zebus/Abc.Zebus.csproj | 2 +- src/Abc.Zebus/Core/Bus.cs | 9 +++ .../Directory/PeerDirectoryClient.cs | 34 ++++++++- src/Abc.Zebus/DomainException.cs | 2 + src/Abc.Zebus/MessageId.cs | 10 +++ src/Abc.Zebus/MessageProcessingException.cs | 2 + src/Abc.Zebus/Monitoring/ZebusMetrics.cs | 69 +++++++++++++++++++ src/Abc.Zebus/Routing/BindingKey.cs | 3 +- .../Serialization/ProtoBufConvert.cs | 8 +++ .../Protobuf/ProtoBufferReader.cs | 3 + src/Abc.Zebus/Transport/ZmqOutboundSocket.cs | 15 ++++ src/Abc.Zebus/Transport/ZmqTransport.cs | 18 +++++ .../Util/Extensions/ExtendDictionary.cs | 2 + .../Util/Extensions/ExtendIEnumerable.cs | 2 + src/Abc.Zebus/Util/IsExternalInit.cs | 8 +-- src/Abc.Zebus/Util/NullableAnnotations.cs | 2 + src/Abc.Zebus/Util/SkipLocalsInitAttribute.cs | 8 +-- 17 files changed, 180 insertions(+), 17 deletions(-) create mode 100644 src/Abc.Zebus/Monitoring/ZebusMetrics.cs diff --git a/src/Abc.Zebus/Abc.Zebus.csproj b/src/Abc.Zebus/Abc.Zebus.csproj index 6e2d87ee..71427afd 100644 --- a/src/Abc.Zebus/Abc.Zebus.csproj +++ b/src/Abc.Zebus/Abc.Zebus.csproj @@ -1,6 +1,6 @@  - netstandard2.0 + netstandard2.0;net10.0 Zebus $(ZebusVersion) diff --git a/src/Abc.Zebus/Core/Bus.cs b/src/Abc.Zebus/Core/Bus.cs index 68a11fa6..403676ec 100644 --- a/src/Abc.Zebus/Core/Bus.cs +++ b/src/Abc.Zebus/Core/Bus.cs @@ -8,6 +8,9 @@ using Abc.Zebus.Directory; using Abc.Zebus.Dispatch; using Abc.Zebus.Lotus; +#if NET10_0_OR_GREATER +using Abc.Zebus.Monitoring; +#endif using Abc.Zebus.Persistence; using Abc.Zebus.Serialization; using Abc.Zebus.Subscriptions; @@ -135,6 +138,9 @@ public virtual void Start() } Started?.Invoke(); +#if NET10_0_OR_GREATER + ZebusMetrics.ActiveBusCount.Add(1); +#endif } private void PerformStartupSubscribe() @@ -202,6 +208,9 @@ public virtual void Stop() InternalStop(true); Stopped?.Invoke(); +#if NET10_0_OR_GREATER + ZebusMetrics.ActiveBusCount.Add(-1); +#endif } private void InternalStop(bool unregister) diff --git a/src/Abc.Zebus/Directory/PeerDirectoryClient.cs b/src/Abc.Zebus/Directory/PeerDirectoryClient.cs index 2bd4be1b..8283ec30 100644 --- a/src/Abc.Zebus/Directory/PeerDirectoryClient.cs +++ b/src/Abc.Zebus/Directory/PeerDirectoryClient.cs @@ -4,6 +4,9 @@ using System.Diagnostics; using System.Linq; using System.Threading.Tasks; +#if NET10_0_OR_GREATER +using Abc.Zebus.Monitoring; +#endif using Abc.Zebus.Util; using Abc.Zebus.Util.Extensions; using Microsoft.Extensions.Logging; @@ -259,13 +262,24 @@ private void AddOrUpdatePeerEntry(PeerDescriptor peerDescriptor, bool shouldRais peerEntry.SetSubscriptions(subscriptions, peerDescriptor.TimestampUtc); if (shouldRaisePeerUpdated) + { PeerUpdated?.Invoke(peerDescriptor.Peer.Id, PeerUpdateAction.Started); +#if NET10_0_OR_GREATER + ZebusMetrics.PeerUpdates.Add(1); +#endif + } var observedSubscriptions = GetObservedSubscriptions(subscriptions); if (observedSubscriptions.Count > 0) PeerSubscriptionsUpdated?.Invoke(peerDescriptor.PeerId, observedSubscriptions); - PeerEntry CreatePeerEntry() => new(peerDescriptor, _globalSubscriptionsIndex); + PeerEntry CreatePeerEntry() + { +#if NET10_0_OR_GREATER + ZebusMetrics.KnownPeerCount.Add(1); +#endif + return new(peerDescriptor, _globalSubscriptionsIndex); + } PeerEntry UpdatePeerEntry(PeerEntry entry) { @@ -325,6 +339,9 @@ public void Handle(PeerStopped message) peer.Value.TimestampUtc = message.TimestampUtc ?? DateTime.UtcNow; PeerUpdated?.Invoke(message.PeerId, PeerUpdateAction.Stopped); +#if NET10_0_OR_GREATER + ZebusMetrics.PeerUpdates.Add(1); +#endif } public void Handle(PeerDecommissioned message) @@ -336,8 +353,14 @@ public void Handle(PeerDecommissioned message) return; removedPeer.RemoveSubscriptions(); +#if NET10_0_OR_GREATER + ZebusMetrics.KnownPeerCount.Add(-1); +#endif PeerUpdated?.Invoke(message.PeerId, PeerUpdateAction.Decommissioned); +#if NET10_0_OR_GREATER + ZebusMetrics.PeerUpdates.Add(1); +#endif } public void Handle(PeerSubscriptionsUpdated message) @@ -358,6 +381,9 @@ public void Handle(PeerSubscriptionsUpdated message) peer.Value.TimestampUtc = message.PeerDescriptor.TimestampUtc ?? DateTime.UtcNow; PeerUpdated?.Invoke(message.PeerDescriptor.PeerId, PeerUpdateAction.Updated); +#if NET10_0_OR_GREATER + ZebusMetrics.PeerUpdates.Add(1); +#endif var observedSubscriptions = GetObservedSubscriptions(subscriptions); if (observedSubscriptions.Count > 0) @@ -397,6 +423,9 @@ public void Handle(PeerSubscriptionsForTypesUpdated message) peer.Value.SetSubscriptionsForType(subscriptionsForTypes, message.TimestampUtc); PeerUpdated?.Invoke(message.PeerId, PeerUpdateAction.Updated); +#if NET10_0_OR_GREATER + ZebusMetrics.PeerUpdates.Add(1); +#endif var observedSubscriptions = GetObservedSubscriptions(subscriptionsForTypes); if (observedSubscriptions.Count > 0) @@ -446,6 +475,9 @@ private void HandlePeerRespondingChange(PeerId peerId, bool isResponding) peer.Peer.IsResponding = isResponding; PeerUpdated?.Invoke(peerId, PeerUpdateAction.Updated); +#if NET10_0_OR_GREATER + ZebusMetrics.PeerUpdates.Add(1); +#endif } private PeerEntryResult GetPeerCheckTimestamp(PeerId peerId, DateTime? timestampUtc) diff --git a/src/Abc.Zebus/DomainException.cs b/src/Abc.Zebus/DomainException.cs index ba4e69a3..41535c41 100644 --- a/src/Abc.Zebus/DomainException.cs +++ b/src/Abc.Zebus/DomainException.cs @@ -48,10 +48,12 @@ public DomainException(Expression> errorCodeExpression, params object[ { } +#if !NETCOREAPP protected DomainException(SerializationInfo info, StreamingContext context) : base(info, context) { } +#endif private static string ReadDescriptionFromAttribute(Expression> errorCodeExpression) { diff --git a/src/Abc.Zebus/MessageId.cs b/src/Abc.Zebus/MessageId.cs index 4f92a60c..cecf7bf2 100644 --- a/src/Abc.Zebus/MessageId.cs +++ b/src/Abc.Zebus/MessageId.cs @@ -62,7 +62,9 @@ public static IDisposable PauseIdGenerationAtDate(DateTime utcDatetime) private class TimeGuidGenerator { private static readonly long _gregorianCalendarTimeTicks = new DateTime(1582, 10, 15, 0, 0, 0, DateTimeKind.Utc).Ticks; +#if !NETCOREAPP private static readonly RNGCryptoServiceProvider _cryptoServiceProvider = new(); +#endif private readonly uint _nodeIdPart1; private readonly ushort _nodeIdPart2; @@ -139,7 +141,11 @@ public unsafe Guid NewGuid(long absoluteTimestamp) private static byte[] GetRandomNodeId() { var nodeId = new byte[6]; +#if !NETCOREAPP _cryptoServiceProvider.GetBytes(nodeId); +#else + RandomNumberGenerator.Fill(nodeId); +#endif return nodeId; } @@ -147,7 +153,11 @@ private static byte[] GetRandomNodeId() private static ushort GetRandomClockId() { var clockId = new byte[2]; +#if !NETCOREAPP _cryptoServiceProvider.GetBytes(clockId); +#else + RandomNumberGenerator.Fill(clockId); +#endif return BitConverter.ToUInt16(clockId, 0); } diff --git a/src/Abc.Zebus/MessageProcessingException.cs b/src/Abc.Zebus/MessageProcessingException.cs index de19cf95..d3977cbd 100644 --- a/src/Abc.Zebus/MessageProcessingException.cs +++ b/src/Abc.Zebus/MessageProcessingException.cs @@ -32,8 +32,10 @@ public MessageProcessingException(string message, Exception? inner) { } +#if !NETCOREAPP protected MessageProcessingException(SerializationInfo info, StreamingContext context) : base(info, context) { } +#endif } diff --git a/src/Abc.Zebus/Monitoring/ZebusMetrics.cs b/src/Abc.Zebus/Monitoring/ZebusMetrics.cs new file mode 100644 index 00000000..7acf62ed --- /dev/null +++ b/src/Abc.Zebus/Monitoring/ZebusMetrics.cs @@ -0,0 +1,69 @@ +#if NET10_0_OR_GREATER +using System; +using System.Diagnostics.Metrics; + +namespace Abc.Zebus.Monitoring; + +/// +/// Provides instruments for monitoring Zebus connection state and messaging. +/// +internal static class ZebusMetrics +{ + internal static readonly Meter Meter = new("Abc.Zebus", typeof(ZebusMetrics).Assembly.GetName().Version?.ToString()); + + // Transport: messages + internal static readonly Counter MessagesSent = Meter.CreateCounter( + "zebus.transport.messages.sent", + unit: "{message}", + description: "Number of transport messages sent to peers"); + + internal static readonly Counter MessagesReceived = Meter.CreateCounter( + "zebus.transport.messages.received", + unit: "{message}", + description: "Number of transport messages received"); + + internal static readonly Counter MessageSendFailures = Meter.CreateCounter( + "zebus.transport.messages.send_failures", + unit: "{message}", + description: "Number of transport message send failures"); + + // Transport: connections + internal static readonly Counter PeerConnections = Meter.CreateCounter( + "zebus.transport.peer_connections", + unit: "{connection}", + description: "Number of outbound peer socket connections established"); + + internal static readonly Counter PeerDisconnections = Meter.CreateCounter( + "zebus.transport.peer_disconnections", + unit: "{disconnection}", + description: "Number of outbound peer socket disconnections"); + + internal static readonly Counter PeerConnectionFailures = Meter.CreateCounter( + "zebus.transport.peer_connection_failures", + unit: "{failure}", + description: "Number of outbound peer socket connection failures"); + + // Transport: outbound socket state + internal static readonly UpDownCounter OutboundSocketCount = Meter.CreateUpDownCounter( + "zebus.transport.outbound_sockets", + unit: "{socket}", + description: "Current number of active outbound sockets"); + + // Directory: peers + internal static readonly UpDownCounter KnownPeerCount = Meter.CreateUpDownCounter( + "zebus.directory.known_peers", + unit: "{peer}", + description: "Current number of known peers in the directory"); + + internal static readonly Counter PeerUpdates = Meter.CreateCounter( + "zebus.directory.peer_updates", + unit: "{update}", + description: "Number of peer update events received"); + + // Bus: status + internal static readonly UpDownCounter ActiveBusCount = Meter.CreateUpDownCounter( + "zebus.bus.active", + unit: "{bus}", + description: "Current number of active (started) bus instances"); +} +#endif diff --git a/src/Abc.Zebus/Routing/BindingKey.cs b/src/Abc.Zebus/Routing/BindingKey.cs index 53edc272..4f3ee431 100644 --- a/src/Abc.Zebus/Routing/BindingKey.cs +++ b/src/Abc.Zebus/Routing/BindingKey.cs @@ -114,7 +114,8 @@ internal static BindingKey Create(Type messageType, IDictionary var parts = new string[routingMembers.Length]; for (var tokenIndex = 0; tokenIndex < routingMembers.Length; ++tokenIndex) { - parts[tokenIndex] = fieldValues.GetValueOrDefault(routingMembers[tokenIndex].Member.Name, BindingKeyPart.StarToken); + var memberName = routingMembers[tokenIndex].Member.Name; + parts[tokenIndex] = fieldValues.TryGetValue(memberName, out var value) ? value : BindingKeyPart.StarToken; } return new BindingKey(parts); diff --git a/src/Abc.Zebus/Serialization/ProtoBufConvert.cs b/src/Abc.Zebus/Serialization/ProtoBufConvert.cs index 5d5e9a19..7d234d72 100644 --- a/src/Abc.Zebus/Serialization/ProtoBufConvert.cs +++ b/src/Abc.Zebus/Serialization/ProtoBufConvert.cs @@ -3,7 +3,11 @@ using System.Diagnostics.CodeAnalysis; using System.IO; using System.Reflection; +#if !NETCOREAPP using System.Runtime.Serialization; +#else +using System.Runtime.CompilerServices; +#endif using ProtoBuf.Meta; namespace Abc.Zebus.Serialization; @@ -49,7 +53,11 @@ public static object Deserialize(Type messageType, Stream stream) private static object? CreateMessageIfRequired(Type messageType) { if (!HasParameterLessConstructor(messageType) && messageType != typeof(string)) +#if !NETCOREAPP return FormatterServices.GetUninitializedObject(messageType); +#else + return RuntimeHelpers.GetUninitializedObject(messageType); +#endif return null; } diff --git a/src/Abc.Zebus/Serialization/Protobuf/ProtoBufferReader.cs b/src/Abc.Zebus/Serialization/Protobuf/ProtoBufferReader.cs index 2cfed750..dcfff638 100644 --- a/src/Abc.Zebus/Serialization/Protobuf/ProtoBufferReader.cs +++ b/src/Abc.Zebus/Serialization/Protobuf/ProtoBufferReader.cs @@ -104,7 +104,10 @@ public bool TrySkipString() public bool TryReadGuid(out Guid value) { if (!TryReadLength(out var length) || !CanRead(length) || length != ProtoBufferWriter.GuidSize) + { + value = default; return false; + } // Skip tag _buffer.AsSpan(_position + 1, 8).CopyTo(_guidBuffer.AsSpan(0)); diff --git a/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs b/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs index 783ace3d..7f598b71 100644 --- a/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs +++ b/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs @@ -1,5 +1,8 @@ using System; using System.Diagnostics; +#if NET10_0_OR_GREATER +using Abc.Zebus.Monitoring; +#endif using Abc.Zebus.Transport.Zmq; using Microsoft.Extensions.Logging; @@ -42,6 +45,9 @@ public void ConnectFor(TransportMessage message) _socket.Connect(EndPoint); IsConnected = true; +#if NET10_0_OR_GREATER + ZebusMetrics.PeerConnections.Add(1); +#endif _logger.LogInformation($"Socket connected, Peer: {PeerId}, EndPoint: {EndPoint}"); } @@ -53,6 +59,9 @@ public void ConnectFor(TransportMessage message) _logger.LogError(ex, $"Unable to connect socket, Peer: {PeerId}, EndPoint: {EndPoint}"); _errorHandler.OnConnectException(PeerId, EndPoint, ex); +#if NET10_0_OR_GREATER + ZebusMetrics.PeerConnectionFailures.Add(1); +#endif SwitchToClosedState(_options.ClosedStateDurationAfterConnectFailure); } @@ -95,6 +104,9 @@ public void Disconnect() { _socket!.SetOption(ZmqSocketOption.LINGER, 0); _socket!.Dispose(); +#if NET10_0_OR_GREATER + ZebusMetrics.PeerDisconnections.Add(1); +#endif _logger.LogInformation($"Socket disconnected, Peer: {PeerId}"); } @@ -123,6 +135,9 @@ public void Send(byte[] buffer, int length, TransportMessage message) _logger.LogError($"Unable to send message, destination peer: {PeerId}, MessageTypeId: {message.MessageTypeId}, MessageId: {message.Id}, Error: {errorMessage}"); _errorHandler.OnSendFailed(PeerId, EndPoint, message.MessageTypeId, message.Id); +#if NET10_0_OR_GREATER + ZebusMetrics.MessageSendFailures.Add(1); +#endif if (_failedSendCount >= _options.SendRetriesBeforeSwitchingToClosedState) SwitchToClosedState(_options.ClosedStateDurationAfterSendFailure); diff --git a/src/Abc.Zebus/Transport/ZmqTransport.cs b/src/Abc.Zebus/Transport/ZmqTransport.cs index fbd362fe..4c004e36 100644 --- a/src/Abc.Zebus/Transport/ZmqTransport.cs +++ b/src/Abc.Zebus/Transport/ZmqTransport.cs @@ -5,6 +5,9 @@ using System.Linq; using System.Threading; using Abc.Zebus.Directory; +#if NET10_0_OR_GREATER +using Abc.Zebus.Monitoring; +#endif using Abc.Zebus.Serialization.Protobuf; using Abc.Zebus.Transport.Zmq; using Abc.Zebus.Util; @@ -265,7 +268,12 @@ private void DeserializeAndForwardTransportMessage(ProtoBufferReader bufferReade } if (_isListening) + { MessageReceived?.Invoke(transportMessage); +#if NET10_0_OR_GREATER + ZebusMetrics.MessagesReceived.Add(1); +#endif + } } catch (Exception ex) { @@ -397,6 +405,9 @@ private void SendToPeer(TransportMessage transportMessage, ProtoBufferWriter buf try { outboundSocket.Send(bufferWriter.Buffer, bufferWriter.Position, transportMessage); +#if NET10_0_OR_GREATER + ZebusMetrics.MessagesSent.Add(1); +#endif } catch (Exception ex) { @@ -412,6 +423,9 @@ private void DisconnectPeers(IEnumerable peerIds) continue; outboundSocket.Disconnect(); +#if NET10_0_OR_GREATER + ZebusMetrics.OutboundSocketCount.Add(-1); +#endif } } @@ -423,6 +437,10 @@ private ZmqOutboundSocket GetConnectedOutboundSocket(Peer peer, TransportMessage outboundSocket.ConnectFor(transportMessage); _outboundSockets.TryAdd(peer.Id, outboundSocket); +#if NET10_0_OR_GREATER + if (outboundSocket.IsConnected) + ZebusMetrics.OutboundSocketCount.Add(1); +#endif } else if (!string.Equals(outboundSocket.EndPoint, peer.EndPoint, StringComparison.OrdinalIgnoreCase)) { diff --git a/src/Abc.Zebus/Util/Extensions/ExtendDictionary.cs b/src/Abc.Zebus/Util/Extensions/ExtendDictionary.cs index 73e05504..0fb43912 100644 --- a/src/Abc.Zebus/Util/Extensions/ExtendDictionary.cs +++ b/src/Abc.Zebus/Util/Extensions/ExtendDictionary.cs @@ -34,6 +34,7 @@ public static TValue GetValueOrAdd(this IDictionary return dictionary.TryGetValue(key, out var value) ? value : (TValue?)null; } +#if !NETCOREAPP [Pure] [return: MaybeNull] public static TValue GetValueOrDefault(this IDictionary dictionary, TKey key) @@ -50,6 +51,7 @@ public static TValue GetValueOrDefault(this IDictionary(this IDictionary dictionary, TKey key, [InstantHandle] Func defaultValueBuilder) where TKey : notnull diff --git a/src/Abc.Zebus/Util/Extensions/ExtendIEnumerable.cs b/src/Abc.Zebus/Util/Extensions/ExtendIEnumerable.cs index ba33411c..a629dc8d 100644 --- a/src/Abc.Zebus/Util/Extensions/ExtendIEnumerable.cs +++ b/src/Abc.Zebus/Util/Extensions/ExtendIEnumerable.cs @@ -7,6 +7,7 @@ namespace Abc.Zebus.Util.Extensions; internal static class ExtendIEnumerable { +#if !NETCOREAPP [Pure] public static HashSet ToHashSet([InstantHandle] this IEnumerable collection) { @@ -18,6 +19,7 @@ public static HashSet ToHashSet([InstantHandle] this IEnumerable collec { return new HashSet(collection, comparer); } +#endif [Pure] public static IEnumerable EmptyIfNull(this IEnumerable? collection) diff --git a/src/Abc.Zebus/Util/IsExternalInit.cs b/src/Abc.Zebus/Util/IsExternalInit.cs index 0d2a8b2c..591796bf 100644 --- a/src/Abc.Zebus/Util/IsExternalInit.cs +++ b/src/Abc.Zebus/Util/IsExternalInit.cs @@ -1,10 +1,4 @@ -#if NETCOREAPP - -using System.Runtime.CompilerServices; - -[assembly: TypeForwardedTo(typeof(IsExternalInit))] - -#else +#if !NETCOREAPP // ReSharper disable once CheckNamespace namespace System.Runtime.CompilerServices diff --git a/src/Abc.Zebus/Util/NullableAnnotations.cs b/src/Abc.Zebus/Util/NullableAnnotations.cs index 233358cb..7defe46e 100644 --- a/src/Abc.Zebus/Util/NullableAnnotations.cs +++ b/src/Abc.Zebus/Util/NullableAnnotations.cs @@ -1,3 +1,4 @@ +#if !NETCOREAPP // ReSharper disable CheckNamespace namespace System.Diagnostics.CodeAnalysis; @@ -69,3 +70,4 @@ public NotNullWhenAttribute(bool returnValue) public bool ReturnValue { get; } } +#endif diff --git a/src/Abc.Zebus/Util/SkipLocalsInitAttribute.cs b/src/Abc.Zebus/Util/SkipLocalsInitAttribute.cs index 0115a4c3..26b9d60c 100644 --- a/src/Abc.Zebus/Util/SkipLocalsInitAttribute.cs +++ b/src/Abc.Zebus/Util/SkipLocalsInitAttribute.cs @@ -1,12 +1,7 @@ +#if !NETCOREAPP // ReSharper disable once CheckNamespace namespace System.Runtime.CompilerServices; -#if NETCOREAPP - -[assembly: TypeForwardedTo(typeof(SkipLocalsInitAttribute))] - -#else - [AttributeUsage(AttributeTargets.Module | AttributeTargets.Class | AttributeTargets.Struct @@ -19,5 +14,4 @@ namespace System.Runtime.CompilerServices; internal sealed class SkipLocalsInitAttribute : Attribute { } - #endif From 950d97bff7e0e7868dd673dc08573d689f0169b2 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Fri, 24 Apr 2026 12:55:07 +0000 Subject: [PATCH 2/8] Add per-connection metrics with peer ID tags and observable gauge for alive/dead status Agent-Logs-Url: https://github.com/Abc-Arbitrage/Zebus/sessions/8f197dc4-0b8b-4ce0-a40d-97ea7c647727 Co-authored-by: Kuinox <18743295+Kuinox@users.noreply.github.com> --- src/Abc.Zebus/Monitoring/ZebusMetrics.cs | 13 ++++++--- src/Abc.Zebus/Transport/ZmqOutboundSocket.cs | 8 +++--- src/Abc.Zebus/Transport/ZmqTransport.cs | 28 +++++++++++++++++++- 3 files changed, 41 insertions(+), 8 deletions(-) diff --git a/src/Abc.Zebus/Monitoring/ZebusMetrics.cs b/src/Abc.Zebus/Monitoring/ZebusMetrics.cs index 7acf62ed..256d3a33 100644 --- a/src/Abc.Zebus/Monitoring/ZebusMetrics.cs +++ b/src/Abc.Zebus/Monitoring/ZebusMetrics.cs @@ -1,17 +1,21 @@ #if NET10_0_OR_GREATER using System; +using System.Collections.Generic; using System.Diagnostics.Metrics; namespace Abc.Zebus.Monitoring; /// /// Provides instruments for monitoring Zebus connection state and messaging. +/// All per-connection metrics include a zebus.peer.id tag for filtering by individual peer. /// internal static class ZebusMetrics { internal static readonly Meter Meter = new("Abc.Zebus", typeof(ZebusMetrics).Assembly.GetName().Version?.ToString()); - // Transport: messages + private const string PeerIdTag = "zebus.peer.id"; + + // Transport: per-connection message counters internal static readonly Counter MessagesSent = Meter.CreateCounter( "zebus.transport.messages.sent", unit: "{message}", @@ -27,7 +31,7 @@ internal static class ZebusMetrics unit: "{message}", description: "Number of transport message send failures"); - // Transport: connections + // Transport: per-connection state internal static readonly Counter PeerConnections = Meter.CreateCounter( "zebus.transport.peer_connections", unit: "{connection}", @@ -43,7 +47,7 @@ internal static class ZebusMetrics unit: "{failure}", description: "Number of outbound peer socket connection failures"); - // Transport: outbound socket state + // Transport: aggregate outbound socket count internal static readonly UpDownCounter OutboundSocketCount = Meter.CreateUpDownCounter( "zebus.transport.outbound_sockets", unit: "{socket}", @@ -65,5 +69,8 @@ internal static class ZebusMetrics "zebus.bus.active", unit: "{bus}", description: "Current number of active (started) bus instances"); + + internal static KeyValuePair PeerTag(PeerId peerId) + => new(PeerIdTag, peerId.ToString()); } #endif diff --git a/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs b/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs index 7f598b71..cb3541ef 100644 --- a/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs +++ b/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs @@ -46,7 +46,7 @@ public void ConnectFor(TransportMessage message) IsConnected = true; #if NET10_0_OR_GREATER - ZebusMetrics.PeerConnections.Add(1); + ZebusMetrics.PeerConnections.Add(1, ZebusMetrics.PeerTag(PeerId)); #endif _logger.LogInformation($"Socket connected, Peer: {PeerId}, EndPoint: {EndPoint}"); @@ -60,7 +60,7 @@ public void ConnectFor(TransportMessage message) _logger.LogError(ex, $"Unable to connect socket, Peer: {PeerId}, EndPoint: {EndPoint}"); _errorHandler.OnConnectException(PeerId, EndPoint, ex); #if NET10_0_OR_GREATER - ZebusMetrics.PeerConnectionFailures.Add(1); + ZebusMetrics.PeerConnectionFailures.Add(1, ZebusMetrics.PeerTag(PeerId)); #endif SwitchToClosedState(_options.ClosedStateDurationAfterConnectFailure); @@ -105,7 +105,7 @@ public void Disconnect() _socket!.SetOption(ZmqSocketOption.LINGER, 0); _socket!.Dispose(); #if NET10_0_OR_GREATER - ZebusMetrics.PeerDisconnections.Add(1); + ZebusMetrics.PeerDisconnections.Add(1, ZebusMetrics.PeerTag(PeerId)); #endif _logger.LogInformation($"Socket disconnected, Peer: {PeerId}"); @@ -136,7 +136,7 @@ public void Send(byte[] buffer, int length, TransportMessage message) _logger.LogError($"Unable to send message, destination peer: {PeerId}, MessageTypeId: {message.MessageTypeId}, MessageId: {message.Id}, Error: {errorMessage}"); _errorHandler.OnSendFailed(PeerId, EndPoint, message.MessageTypeId, message.Id); #if NET10_0_OR_GREATER - ZebusMetrics.MessageSendFailures.Add(1); + ZebusMetrics.MessageSendFailures.Add(1, ZebusMetrics.PeerTag(PeerId)); #endif if (_failedSendCount >= _options.SendRetriesBeforeSwitchingToClosedState) diff --git a/src/Abc.Zebus/Transport/ZmqTransport.cs b/src/Abc.Zebus/Transport/ZmqTransport.cs index 4c004e36..d5d71c8a 100644 --- a/src/Abc.Zebus/Transport/ZmqTransport.cs +++ b/src/Abc.Zebus/Transport/ZmqTransport.cs @@ -1,6 +1,9 @@ using System; using System.Collections.Concurrent; using System.Collections.Generic; +#if NET10_0_OR_GREATER +using System.Diagnostics.Metrics; +#endif using System.IO; using System.Linq; using System.Threading; @@ -34,6 +37,9 @@ public class ZmqTransport : ITransport private string _environment = string.Empty; private CountdownEvent? _outboundSocketsToStop; private bool _isRunning; +#if NET10_0_OR_GREATER + private ObservableGauge? _connectionStateGauge; +#endif public ZmqTransport(IZmqTransportConfiguration configuration, ZmqSocketOptions socketOptions, IZmqOutboundSocketErrorHandler errorHandler) { @@ -113,6 +119,14 @@ public void Start() startSequenceState.Wait(); _isRunning = true; + +#if NET10_0_OR_GREATER + _connectionStateGauge = ZebusMetrics.Meter.CreateObservableGauge( + "zebus.transport.connection.alive", + observeValues: ObserveConnectionStates, + unit: "{connection}", + description: "Whether a peer connection is alive (1) or dead (0)"); +#endif } public void Stop() @@ -406,7 +420,7 @@ private void SendToPeer(TransportMessage transportMessage, ProtoBufferWriter buf { outboundSocket.Send(bufferWriter.Buffer, bufferWriter.Position, transportMessage); #if NET10_0_OR_GREATER - ZebusMetrics.MessagesSent.Add(1); + ZebusMetrics.MessagesSent.Add(1, ZebusMetrics.PeerTag(target.Id)); #endif } catch (Exception ex) @@ -508,6 +522,18 @@ private void SafeAdd(BlockingCollection collection, T item) } } +#if NET10_0_OR_GREATER + private IEnumerable> ObserveConnectionStates() + { + foreach (var (peerId, socket) in _outboundSockets) + { + yield return new Measurement( + socket.IsConnected ? 1 : 0, + ZebusMetrics.PeerTag(peerId)); + } + } +#endif + private readonly struct OutboundSocketAction { private static readonly TransportMessage _disconnectMessage = new(default, null!, new PeerId(), null!); From 543866e25fb0594726b5fba8dfe491c8c2cf130b Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Fri, 24 Apr 2026 12:57:56 +0000 Subject: [PATCH 3/8] Tag MessagesReceived metric with sender peer ID for per-connection consistency Agent-Logs-Url: https://github.com/Abc-Arbitrage/Zebus/sessions/8f197dc4-0b8b-4ce0-a40d-97ea7c647727 Co-authored-by: Kuinox <18743295+Kuinox@users.noreply.github.com> --- src/Abc.Zebus/Transport/ZmqTransport.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/Abc.Zebus/Transport/ZmqTransport.cs b/src/Abc.Zebus/Transport/ZmqTransport.cs index d5d71c8a..97b48cf6 100644 --- a/src/Abc.Zebus/Transport/ZmqTransport.cs +++ b/src/Abc.Zebus/Transport/ZmqTransport.cs @@ -285,7 +285,7 @@ private void DeserializeAndForwardTransportMessage(ProtoBufferReader bufferReade { MessageReceived?.Invoke(transportMessage); #if NET10_0_OR_GREATER - ZebusMetrics.MessagesReceived.Add(1); + ZebusMetrics.MessagesReceived.Add(1, ZebusMetrics.PeerTag(transportMessage.Originator.SenderId)); #endif } } From e04af38edd9a520cad4303ca0016facfa891a04e Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Tue, 28 Apr 2026 18:06:01 +0000 Subject: [PATCH 4/8] Track bus metrics per instance using a stable Interlocked-incremented ID Agent-Logs-Url: https://github.com/Abc-Arbitrage/Zebus/sessions/ca085f25-a9b5-488a-ad33-858646b84131 Co-authored-by: Kuinox <18743295+Kuinox@users.noreply.github.com> --- src/Abc.Zebus/Core/Bus.cs | 8 ++++++-- src/Abc.Zebus/Monitoring/ZebusMetrics.cs | 11 ++++++++--- 2 files changed, 14 insertions(+), 5 deletions(-) diff --git a/src/Abc.Zebus/Core/Bus.cs b/src/Abc.Zebus/Core/Bus.cs index 403676ec..1bbd5a17 100644 --- a/src/Abc.Zebus/Core/Bus.cs +++ b/src/Abc.Zebus/Core/Bus.cs @@ -26,6 +26,10 @@ public class Bus : IInternalBus, IMessageDispatchFactory { private static readonly BusMessageLogger _messageLogger = new(typeof(Bus)); private static readonly ILogger _logger = ZebusLogManager.GetLogger(typeof(Bus)); +#if NET10_0_OR_GREATER + private static int _nextInstanceId; + private readonly int _instanceId = Interlocked.Increment(ref _nextInstanceId); +#endif private readonly ConcurrentDictionary> _messageIdToTaskCompletionSources = new(); private readonly UniqueTimestampProvider _deserializationFailureTimestampProvider = new(); @@ -139,7 +143,7 @@ public virtual void Start() Started?.Invoke(); #if NET10_0_OR_GREATER - ZebusMetrics.ActiveBusCount.Add(1); + ZebusMetrics.BusActive.Add(1, ZebusMetrics.BusTag(_instanceId)); #endif } @@ -209,7 +213,7 @@ public virtual void Stop() Stopped?.Invoke(); #if NET10_0_OR_GREATER - ZebusMetrics.ActiveBusCount.Add(-1); + ZebusMetrics.BusActive.Add(-1, ZebusMetrics.BusTag(_instanceId)); #endif } diff --git a/src/Abc.Zebus/Monitoring/ZebusMetrics.cs b/src/Abc.Zebus/Monitoring/ZebusMetrics.cs index 256d3a33..8c62cffc 100644 --- a/src/Abc.Zebus/Monitoring/ZebusMetrics.cs +++ b/src/Abc.Zebus/Monitoring/ZebusMetrics.cs @@ -64,13 +64,18 @@ internal static class ZebusMetrics unit: "{update}", description: "Number of peer update events received"); - // Bus: status - internal static readonly UpDownCounter ActiveBusCount = Meter.CreateUpDownCounter( + // Bus: per-instance status + private const string BusIdTag = "zebus.bus.id"; + + internal static readonly UpDownCounter BusActive = Meter.CreateUpDownCounter( "zebus.bus.active", unit: "{bus}", - description: "Current number of active (started) bus instances"); + description: "Whether a bus instance is active (1) or stopped (0)"); internal static KeyValuePair PeerTag(PeerId peerId) => new(PeerIdTag, peerId.ToString()); + + internal static KeyValuePair BusTag(int busInstanceId) + => new(BusIdTag, busInstanceId); } #endif From e1e006bcea24b725ab42a419b6b5935fd9a165b1 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Tue, 28 Apr 2026 18:21:18 +0000 Subject: [PATCH 5/8] Split ZebusMetrics into per-service classes and remove bus active metric Agent-Logs-Url: https://github.com/Abc-Arbitrage/Zebus/sessions/efc82941-92e3-4ee8-b641-de8cf4c1f8ec Co-authored-by: Kuinox <18743295+Kuinox@users.noreply.github.com> --- src/Abc.Zebus/Core/Bus.cs | 13 ---- .../Directory/PeerDirectoryClient.cs | 16 ++--- src/Abc.Zebus/Monitoring/DirectoryMetrics.cs | 21 ++++++ src/Abc.Zebus/Monitoring/TransportMetrics.cs | 50 ++++++++++++++ src/Abc.Zebus/Monitoring/ZebusMetrics.cs | 65 +------------------ src/Abc.Zebus/Transport/ZmqOutboundSocket.cs | 8 +-- src/Abc.Zebus/Transport/ZmqTransport.cs | 8 +-- 7 files changed, 89 insertions(+), 92 deletions(-) create mode 100644 src/Abc.Zebus/Monitoring/DirectoryMetrics.cs create mode 100644 src/Abc.Zebus/Monitoring/TransportMetrics.cs diff --git a/src/Abc.Zebus/Core/Bus.cs b/src/Abc.Zebus/Core/Bus.cs index 1bbd5a17..68a11fa6 100644 --- a/src/Abc.Zebus/Core/Bus.cs +++ b/src/Abc.Zebus/Core/Bus.cs @@ -8,9 +8,6 @@ using Abc.Zebus.Directory; using Abc.Zebus.Dispatch; using Abc.Zebus.Lotus; -#if NET10_0_OR_GREATER -using Abc.Zebus.Monitoring; -#endif using Abc.Zebus.Persistence; using Abc.Zebus.Serialization; using Abc.Zebus.Subscriptions; @@ -26,10 +23,6 @@ public class Bus : IInternalBus, IMessageDispatchFactory { private static readonly BusMessageLogger _messageLogger = new(typeof(Bus)); private static readonly ILogger _logger = ZebusLogManager.GetLogger(typeof(Bus)); -#if NET10_0_OR_GREATER - private static int _nextInstanceId; - private readonly int _instanceId = Interlocked.Increment(ref _nextInstanceId); -#endif private readonly ConcurrentDictionary> _messageIdToTaskCompletionSources = new(); private readonly UniqueTimestampProvider _deserializationFailureTimestampProvider = new(); @@ -142,9 +135,6 @@ public virtual void Start() } Started?.Invoke(); -#if NET10_0_OR_GREATER - ZebusMetrics.BusActive.Add(1, ZebusMetrics.BusTag(_instanceId)); -#endif } private void PerformStartupSubscribe() @@ -212,9 +202,6 @@ public virtual void Stop() InternalStop(true); Stopped?.Invoke(); -#if NET10_0_OR_GREATER - ZebusMetrics.BusActive.Add(-1, ZebusMetrics.BusTag(_instanceId)); -#endif } private void InternalStop(bool unregister) diff --git a/src/Abc.Zebus/Directory/PeerDirectoryClient.cs b/src/Abc.Zebus/Directory/PeerDirectoryClient.cs index 8283ec30..c2723e18 100644 --- a/src/Abc.Zebus/Directory/PeerDirectoryClient.cs +++ b/src/Abc.Zebus/Directory/PeerDirectoryClient.cs @@ -265,7 +265,7 @@ private void AddOrUpdatePeerEntry(PeerDescriptor peerDescriptor, bool shouldRais { PeerUpdated?.Invoke(peerDescriptor.Peer.Id, PeerUpdateAction.Started); #if NET10_0_OR_GREATER - ZebusMetrics.PeerUpdates.Add(1); + DirectoryMetrics.PeerUpdates.Add(1); #endif } @@ -276,7 +276,7 @@ private void AddOrUpdatePeerEntry(PeerDescriptor peerDescriptor, bool shouldRais PeerEntry CreatePeerEntry() { #if NET10_0_OR_GREATER - ZebusMetrics.KnownPeerCount.Add(1); + DirectoryMetrics.KnownPeerCount.Add(1); #endif return new(peerDescriptor, _globalSubscriptionsIndex); } @@ -340,7 +340,7 @@ public void Handle(PeerStopped message) PeerUpdated?.Invoke(message.PeerId, PeerUpdateAction.Stopped); #if NET10_0_OR_GREATER - ZebusMetrics.PeerUpdates.Add(1); + DirectoryMetrics.PeerUpdates.Add(1); #endif } @@ -354,12 +354,12 @@ public void Handle(PeerDecommissioned message) removedPeer.RemoveSubscriptions(); #if NET10_0_OR_GREATER - ZebusMetrics.KnownPeerCount.Add(-1); + DirectoryMetrics.KnownPeerCount.Add(-1); #endif PeerUpdated?.Invoke(message.PeerId, PeerUpdateAction.Decommissioned); #if NET10_0_OR_GREATER - ZebusMetrics.PeerUpdates.Add(1); + DirectoryMetrics.PeerUpdates.Add(1); #endif } @@ -382,7 +382,7 @@ public void Handle(PeerSubscriptionsUpdated message) PeerUpdated?.Invoke(message.PeerDescriptor.PeerId, PeerUpdateAction.Updated); #if NET10_0_OR_GREATER - ZebusMetrics.PeerUpdates.Add(1); + DirectoryMetrics.PeerUpdates.Add(1); #endif var observedSubscriptions = GetObservedSubscriptions(subscriptions); @@ -424,7 +424,7 @@ public void Handle(PeerSubscriptionsForTypesUpdated message) PeerUpdated?.Invoke(message.PeerId, PeerUpdateAction.Updated); #if NET10_0_OR_GREATER - ZebusMetrics.PeerUpdates.Add(1); + DirectoryMetrics.PeerUpdates.Add(1); #endif var observedSubscriptions = GetObservedSubscriptions(subscriptionsForTypes); @@ -476,7 +476,7 @@ private void HandlePeerRespondingChange(PeerId peerId, bool isResponding) PeerUpdated?.Invoke(peerId, PeerUpdateAction.Updated); #if NET10_0_OR_GREATER - ZebusMetrics.PeerUpdates.Add(1); + DirectoryMetrics.PeerUpdates.Add(1); #endif } diff --git a/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs b/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs new file mode 100644 index 00000000..5ddbe8c7 --- /dev/null +++ b/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs @@ -0,0 +1,21 @@ +#if NET10_0_OR_GREATER +using System.Diagnostics.Metrics; + +namespace Abc.Zebus.Monitoring; + +/// +/// Provides instruments for monitoring Zebus peer directory. +/// +internal static class DirectoryMetrics +{ + internal static readonly UpDownCounter KnownPeerCount = ZebusMetrics.Meter.CreateUpDownCounter( + "zebus.directory.known_peers", + unit: "{peer}", + description: "Current number of known peers in the directory"); + + internal static readonly Counter PeerUpdates = ZebusMetrics.Meter.CreateCounter( + "zebus.directory.peer_updates", + unit: "{update}", + description: "Number of peer update events received"); +} +#endif diff --git a/src/Abc.Zebus/Monitoring/TransportMetrics.cs b/src/Abc.Zebus/Monitoring/TransportMetrics.cs new file mode 100644 index 00000000..c5660823 --- /dev/null +++ b/src/Abc.Zebus/Monitoring/TransportMetrics.cs @@ -0,0 +1,50 @@ +#if NET10_0_OR_GREATER +using System.Diagnostics.Metrics; + +namespace Abc.Zebus.Monitoring; + +/// +/// Provides instruments for monitoring Zebus transport layer. +/// All per-connection metrics include a zebus.peer.id tag for filtering by individual peer. +/// +internal static class TransportMetrics +{ + // Per-connection message counters + internal static readonly Counter MessagesSent = ZebusMetrics.Meter.CreateCounter( + "zebus.transport.messages.sent", + unit: "{message}", + description: "Number of transport messages sent to peers"); + + internal static readonly Counter MessagesReceived = ZebusMetrics.Meter.CreateCounter( + "zebus.transport.messages.received", + unit: "{message}", + description: "Number of transport messages received"); + + internal static readonly Counter MessageSendFailures = ZebusMetrics.Meter.CreateCounter( + "zebus.transport.messages.send_failures", + unit: "{message}", + description: "Number of transport message send failures"); + + // Per-connection state + internal static readonly Counter PeerConnections = ZebusMetrics.Meter.CreateCounter( + "zebus.transport.peer_connections", + unit: "{connection}", + description: "Number of outbound peer socket connections established"); + + internal static readonly Counter PeerDisconnections = ZebusMetrics.Meter.CreateCounter( + "zebus.transport.peer_disconnections", + unit: "{disconnection}", + description: "Number of outbound peer socket disconnections"); + + internal static readonly Counter PeerConnectionFailures = ZebusMetrics.Meter.CreateCounter( + "zebus.transport.peer_connection_failures", + unit: "{failure}", + description: "Number of outbound peer socket connection failures"); + + // Aggregate outbound socket count + internal static readonly UpDownCounter OutboundSocketCount = ZebusMetrics.Meter.CreateUpDownCounter( + "zebus.transport.outbound_sockets", + unit: "{socket}", + description: "Current number of active outbound sockets"); +} +#endif diff --git a/src/Abc.Zebus/Monitoring/ZebusMetrics.cs b/src/Abc.Zebus/Monitoring/ZebusMetrics.cs index 8c62cffc..d334e00f 100644 --- a/src/Abc.Zebus/Monitoring/ZebusMetrics.cs +++ b/src/Abc.Zebus/Monitoring/ZebusMetrics.cs @@ -1,13 +1,12 @@ #if NET10_0_OR_GREATER -using System; using System.Collections.Generic; using System.Diagnostics.Metrics; namespace Abc.Zebus.Monitoring; /// -/// Provides instruments for monitoring Zebus connection state and messaging. -/// All per-connection metrics include a zebus.peer.id tag for filtering by individual peer. +/// Shared and tag helpers for Zebus metrics. +/// Service-specific instruments are defined in and . /// internal static class ZebusMetrics { @@ -15,67 +14,7 @@ internal static class ZebusMetrics private const string PeerIdTag = "zebus.peer.id"; - // Transport: per-connection message counters - internal static readonly Counter MessagesSent = Meter.CreateCounter( - "zebus.transport.messages.sent", - unit: "{message}", - description: "Number of transport messages sent to peers"); - - internal static readonly Counter MessagesReceived = Meter.CreateCounter( - "zebus.transport.messages.received", - unit: "{message}", - description: "Number of transport messages received"); - - internal static readonly Counter MessageSendFailures = Meter.CreateCounter( - "zebus.transport.messages.send_failures", - unit: "{message}", - description: "Number of transport message send failures"); - - // Transport: per-connection state - internal static readonly Counter PeerConnections = Meter.CreateCounter( - "zebus.transport.peer_connections", - unit: "{connection}", - description: "Number of outbound peer socket connections established"); - - internal static readonly Counter PeerDisconnections = Meter.CreateCounter( - "zebus.transport.peer_disconnections", - unit: "{disconnection}", - description: "Number of outbound peer socket disconnections"); - - internal static readonly Counter PeerConnectionFailures = Meter.CreateCounter( - "zebus.transport.peer_connection_failures", - unit: "{failure}", - description: "Number of outbound peer socket connection failures"); - - // Transport: aggregate outbound socket count - internal static readonly UpDownCounter OutboundSocketCount = Meter.CreateUpDownCounter( - "zebus.transport.outbound_sockets", - unit: "{socket}", - description: "Current number of active outbound sockets"); - - // Directory: peers - internal static readonly UpDownCounter KnownPeerCount = Meter.CreateUpDownCounter( - "zebus.directory.known_peers", - unit: "{peer}", - description: "Current number of known peers in the directory"); - - internal static readonly Counter PeerUpdates = Meter.CreateCounter( - "zebus.directory.peer_updates", - unit: "{update}", - description: "Number of peer update events received"); - - // Bus: per-instance status - private const string BusIdTag = "zebus.bus.id"; - - internal static readonly UpDownCounter BusActive = Meter.CreateUpDownCounter( - "zebus.bus.active", - unit: "{bus}", - description: "Whether a bus instance is active (1) or stopped (0)"); - internal static KeyValuePair PeerTag(PeerId peerId) => new(PeerIdTag, peerId.ToString()); - - internal static KeyValuePair BusTag(int busInstanceId) - => new(BusIdTag, busInstanceId); } #endif diff --git a/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs b/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs index cb3541ef..c42703a6 100644 --- a/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs +++ b/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs @@ -46,7 +46,7 @@ public void ConnectFor(TransportMessage message) IsConnected = true; #if NET10_0_OR_GREATER - ZebusMetrics.PeerConnections.Add(1, ZebusMetrics.PeerTag(PeerId)); + TransportMetrics.PeerConnections.Add(1, ZebusMetrics.PeerTag(PeerId)); #endif _logger.LogInformation($"Socket connected, Peer: {PeerId}, EndPoint: {EndPoint}"); @@ -60,7 +60,7 @@ public void ConnectFor(TransportMessage message) _logger.LogError(ex, $"Unable to connect socket, Peer: {PeerId}, EndPoint: {EndPoint}"); _errorHandler.OnConnectException(PeerId, EndPoint, ex); #if NET10_0_OR_GREATER - ZebusMetrics.PeerConnectionFailures.Add(1, ZebusMetrics.PeerTag(PeerId)); + TransportMetrics.PeerConnectionFailures.Add(1, ZebusMetrics.PeerTag(PeerId)); #endif SwitchToClosedState(_options.ClosedStateDurationAfterConnectFailure); @@ -105,7 +105,7 @@ public void Disconnect() _socket!.SetOption(ZmqSocketOption.LINGER, 0); _socket!.Dispose(); #if NET10_0_OR_GREATER - ZebusMetrics.PeerDisconnections.Add(1, ZebusMetrics.PeerTag(PeerId)); + TransportMetrics.PeerDisconnections.Add(1, ZebusMetrics.PeerTag(PeerId)); #endif _logger.LogInformation($"Socket disconnected, Peer: {PeerId}"); @@ -136,7 +136,7 @@ public void Send(byte[] buffer, int length, TransportMessage message) _logger.LogError($"Unable to send message, destination peer: {PeerId}, MessageTypeId: {message.MessageTypeId}, MessageId: {message.Id}, Error: {errorMessage}"); _errorHandler.OnSendFailed(PeerId, EndPoint, message.MessageTypeId, message.Id); #if NET10_0_OR_GREATER - ZebusMetrics.MessageSendFailures.Add(1, ZebusMetrics.PeerTag(PeerId)); + TransportMetrics.MessageSendFailures.Add(1, ZebusMetrics.PeerTag(PeerId)); #endif if (_failedSendCount >= _options.SendRetriesBeforeSwitchingToClosedState) diff --git a/src/Abc.Zebus/Transport/ZmqTransport.cs b/src/Abc.Zebus/Transport/ZmqTransport.cs index 97b48cf6..4b7c89c9 100644 --- a/src/Abc.Zebus/Transport/ZmqTransport.cs +++ b/src/Abc.Zebus/Transport/ZmqTransport.cs @@ -285,7 +285,7 @@ private void DeserializeAndForwardTransportMessage(ProtoBufferReader bufferReade { MessageReceived?.Invoke(transportMessage); #if NET10_0_OR_GREATER - ZebusMetrics.MessagesReceived.Add(1, ZebusMetrics.PeerTag(transportMessage.Originator.SenderId)); + TransportMetrics.MessagesReceived.Add(1, ZebusMetrics.PeerTag(transportMessage.Originator.SenderId)); #endif } } @@ -420,7 +420,7 @@ private void SendToPeer(TransportMessage transportMessage, ProtoBufferWriter buf { outboundSocket.Send(bufferWriter.Buffer, bufferWriter.Position, transportMessage); #if NET10_0_OR_GREATER - ZebusMetrics.MessagesSent.Add(1, ZebusMetrics.PeerTag(target.Id)); + TransportMetrics.MessagesSent.Add(1, ZebusMetrics.PeerTag(target.Id)); #endif } catch (Exception ex) @@ -438,7 +438,7 @@ private void DisconnectPeers(IEnumerable peerIds) outboundSocket.Disconnect(); #if NET10_0_OR_GREATER - ZebusMetrics.OutboundSocketCount.Add(-1); + TransportMetrics.OutboundSocketCount.Add(-1); #endif } } @@ -453,7 +453,7 @@ private ZmqOutboundSocket GetConnectedOutboundSocket(Peer peer, TransportMessage _outboundSockets.TryAdd(peer.Id, outboundSocket); #if NET10_0_OR_GREATER if (outboundSocket.IsConnected) - ZebusMetrics.OutboundSocketCount.Add(1); + TransportMetrics.OutboundSocketCount.Add(1); #endif } else if (!string.Equals(outboundSocket.EndPoint, peer.EndPoint, StringComparison.OrdinalIgnoreCase)) From 4c7958dd6d810b6ab9bd559f939c7efdbda231c7 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Tue, 28 Apr 2026 18:23:45 +0000 Subject: [PATCH 6/8] Remove trailing newlines from new metric files Agent-Logs-Url: https://github.com/Abc-Arbitrage/Zebus/sessions/efc82941-92e3-4ee8-b641-de8cf4c1f8ec Co-authored-by: Kuinox <18743295+Kuinox@users.noreply.github.com> --- src/Abc.Zebus/Monitoring/DirectoryMetrics.cs | 2 +- src/Abc.Zebus/Monitoring/TransportMetrics.cs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs b/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs index 5ddbe8c7..47562d78 100644 --- a/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs +++ b/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs @@ -18,4 +18,4 @@ internal static class DirectoryMetrics unit: "{update}", description: "Number of peer update events received"); } -#endif +#endif \ No newline at end of file diff --git a/src/Abc.Zebus/Monitoring/TransportMetrics.cs b/src/Abc.Zebus/Monitoring/TransportMetrics.cs index c5660823..56c42592 100644 --- a/src/Abc.Zebus/Monitoring/TransportMetrics.cs +++ b/src/Abc.Zebus/Monitoring/TransportMetrics.cs @@ -47,4 +47,4 @@ internal static class TransportMetrics unit: "{socket}", description: "Current number of active outbound sockets"); } -#endif +#endif \ No newline at end of file From c3c154c6dd5a962b7a4a07a8b87629076d63f10a Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Tue, 28 Apr 2026 18:47:06 +0000 Subject: [PATCH 7/8] Restructure #ifdef so metric classes are declared but empty on netstandard, removing #ifdef on usings Agent-Logs-Url: https://github.com/Abc-Arbitrage/Zebus/sessions/81a35d95-9277-416c-9f67-15fb83ff98cf Co-authored-by: Kuinox <18743295+Kuinox@users.noreply.github.com> --- src/Abc.Zebus/Directory/PeerDirectoryClient.cs | 2 -- src/Abc.Zebus/Monitoring/DirectoryMetrics.cs | 6 ++++-- src/Abc.Zebus/Monitoring/TransportMetrics.cs | 6 ++++-- src/Abc.Zebus/Monitoring/ZebusMetrics.cs | 6 ++++-- src/Abc.Zebus/Transport/ZmqOutboundSocket.cs | 2 -- src/Abc.Zebus/Transport/ZmqTransport.cs | 2 -- 6 files changed, 12 insertions(+), 12 deletions(-) diff --git a/src/Abc.Zebus/Directory/PeerDirectoryClient.cs b/src/Abc.Zebus/Directory/PeerDirectoryClient.cs index c2723e18..e8957b9b 100644 --- a/src/Abc.Zebus/Directory/PeerDirectoryClient.cs +++ b/src/Abc.Zebus/Directory/PeerDirectoryClient.cs @@ -4,9 +4,7 @@ using System.Diagnostics; using System.Linq; using System.Threading.Tasks; -#if NET10_0_OR_GREATER using Abc.Zebus.Monitoring; -#endif using Abc.Zebus.Util; using Abc.Zebus.Util.Extensions; using Microsoft.Extensions.Logging; diff --git a/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs b/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs index 47562d78..1b86d5ff 100644 --- a/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs +++ b/src/Abc.Zebus/Monitoring/DirectoryMetrics.cs @@ -1,5 +1,6 @@ #if NET10_0_OR_GREATER using System.Diagnostics.Metrics; +#endif namespace Abc.Zebus.Monitoring; @@ -8,6 +9,7 @@ namespace Abc.Zebus.Monitoring; /// internal static class DirectoryMetrics { +#if NET10_0_OR_GREATER internal static readonly UpDownCounter KnownPeerCount = ZebusMetrics.Meter.CreateUpDownCounter( "zebus.directory.known_peers", unit: "{peer}", @@ -17,5 +19,5 @@ internal static class DirectoryMetrics "zebus.directory.peer_updates", unit: "{update}", description: "Number of peer update events received"); -} -#endif \ No newline at end of file +#endif +} \ No newline at end of file diff --git a/src/Abc.Zebus/Monitoring/TransportMetrics.cs b/src/Abc.Zebus/Monitoring/TransportMetrics.cs index 56c42592..24f1add5 100644 --- a/src/Abc.Zebus/Monitoring/TransportMetrics.cs +++ b/src/Abc.Zebus/Monitoring/TransportMetrics.cs @@ -1,5 +1,6 @@ #if NET10_0_OR_GREATER using System.Diagnostics.Metrics; +#endif namespace Abc.Zebus.Monitoring; @@ -9,6 +10,7 @@ namespace Abc.Zebus.Monitoring; /// internal static class TransportMetrics { +#if NET10_0_OR_GREATER // Per-connection message counters internal static readonly Counter MessagesSent = ZebusMetrics.Meter.CreateCounter( "zebus.transport.messages.sent", @@ -46,5 +48,5 @@ internal static class TransportMetrics "zebus.transport.outbound_sockets", unit: "{socket}", description: "Current number of active outbound sockets"); -} -#endif \ No newline at end of file +#endif +} \ No newline at end of file diff --git a/src/Abc.Zebus/Monitoring/ZebusMetrics.cs b/src/Abc.Zebus/Monitoring/ZebusMetrics.cs index d334e00f..06631349 100644 --- a/src/Abc.Zebus/Monitoring/ZebusMetrics.cs +++ b/src/Abc.Zebus/Monitoring/ZebusMetrics.cs @@ -1,20 +1,22 @@ #if NET10_0_OR_GREATER using System.Collections.Generic; using System.Diagnostics.Metrics; +#endif namespace Abc.Zebus.Monitoring; /// -/// Shared and tag helpers for Zebus metrics. +/// Shared and tag helpers for Zebus metrics. /// Service-specific instruments are defined in and . /// internal static class ZebusMetrics { +#if NET10_0_OR_GREATER internal static readonly Meter Meter = new("Abc.Zebus", typeof(ZebusMetrics).Assembly.GetName().Version?.ToString()); private const string PeerIdTag = "zebus.peer.id"; internal static KeyValuePair PeerTag(PeerId peerId) => new(PeerIdTag, peerId.ToString()); -} #endif +} diff --git a/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs b/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs index c42703a6..83d626dd 100644 --- a/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs +++ b/src/Abc.Zebus/Transport/ZmqOutboundSocket.cs @@ -1,8 +1,6 @@ using System; using System.Diagnostics; -#if NET10_0_OR_GREATER using Abc.Zebus.Monitoring; -#endif using Abc.Zebus.Transport.Zmq; using Microsoft.Extensions.Logging; diff --git a/src/Abc.Zebus/Transport/ZmqTransport.cs b/src/Abc.Zebus/Transport/ZmqTransport.cs index 4b7c89c9..d8ce1ab2 100644 --- a/src/Abc.Zebus/Transport/ZmqTransport.cs +++ b/src/Abc.Zebus/Transport/ZmqTransport.cs @@ -8,9 +8,7 @@ using System.Linq; using System.Threading; using Abc.Zebus.Directory; -#if NET10_0_OR_GREATER using Abc.Zebus.Monitoring; -#endif using Abc.Zebus.Serialization.Protobuf; using Abc.Zebus.Transport.Zmq; using Abc.Zebus.Util; From 4dd924db88bea6183e1540dd87fb4d7d470f1a75 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Thu, 30 Apr 2026 14:20:22 +0000 Subject: [PATCH 8/8] Add byte count per message metrics (zebus.transport.bytes.sent/received) Agent-Logs-Url: https://github.com/Abc-Arbitrage/Zebus/sessions/f0604d23-5c82-4f7f-baaa-b80b93aa52cf Co-authored-by: Kuinox <18743295+Kuinox@users.noreply.github.com> --- src/Abc.Zebus/Monitoring/TransportMetrics.cs | 10 ++++++++++ src/Abc.Zebus/Transport/ZmqTransport.cs | 8 ++++++-- 2 files changed, 16 insertions(+), 2 deletions(-) diff --git a/src/Abc.Zebus/Monitoring/TransportMetrics.cs b/src/Abc.Zebus/Monitoring/TransportMetrics.cs index 24f1add5..0883e866 100644 --- a/src/Abc.Zebus/Monitoring/TransportMetrics.cs +++ b/src/Abc.Zebus/Monitoring/TransportMetrics.cs @@ -22,6 +22,16 @@ internal static class TransportMetrics unit: "{message}", description: "Number of transport messages received"); + internal static readonly Counter BytesSent = ZebusMetrics.Meter.CreateCounter( + "zebus.transport.bytes.sent", + unit: "By", + description: "Number of bytes sent to peers"); + + internal static readonly Counter BytesReceived = ZebusMetrics.Meter.CreateCounter( + "zebus.transport.bytes.received", + unit: "By", + description: "Number of bytes received from peers"); + internal static readonly Counter MessageSendFailures = ZebusMetrics.Meter.CreateCounter( "zebus.transport.messages.send_failures", unit: "{message}", diff --git a/src/Abc.Zebus/Transport/ZmqTransport.cs b/src/Abc.Zebus/Transport/ZmqTransport.cs index d8ce1ab2..3fb7de7e 100644 --- a/src/Abc.Zebus/Transport/ZmqTransport.cs +++ b/src/Abc.Zebus/Transport/ZmqTransport.cs @@ -283,7 +283,9 @@ private void DeserializeAndForwardTransportMessage(ProtoBufferReader bufferReade { MessageReceived?.Invoke(transportMessage); #if NET10_0_OR_GREATER - TransportMetrics.MessagesReceived.Add(1, ZebusMetrics.PeerTag(transportMessage.Originator.SenderId)); + var peerTag = ZebusMetrics.PeerTag(transportMessage.Originator.SenderId); + TransportMetrics.MessagesReceived.Add(1, peerTag); + TransportMetrics.BytesReceived.Add(bufferReader.Length, peerTag); #endif } } @@ -418,7 +420,9 @@ private void SendToPeer(TransportMessage transportMessage, ProtoBufferWriter buf { outboundSocket.Send(bufferWriter.Buffer, bufferWriter.Position, transportMessage); #if NET10_0_OR_GREATER - TransportMetrics.MessagesSent.Add(1, ZebusMetrics.PeerTag(target.Id)); + var peerTag = ZebusMetrics.PeerTag(target.Id); + TransportMetrics.MessagesSent.Add(1, peerTag); + TransportMetrics.BytesSent.Add(bufferWriter.Position, peerTag); #endif } catch (Exception ex)