Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
131 changes: 74 additions & 57 deletions src/Vulthil.Messaging.RabbitMq/Consumers/MessageContextFactory.cs
Original file line number Diff line number Diff line change
Expand Up @@ -66,46 +66,20 @@ public static MessageContext<TMessage> CreateContext<TMessage>(
/// </summary>
public static MessageContextSnapshot CreateSnapshot(BasicDeliverEventArgs ea)
{
var context = BuildMetadata(ea);
var metadata = ExtractMetadata(ea);
return new MessageContextSnapshot
{
MessageId = context.MessageId,
RequestId = context.RequestId,
CorrelationId = context.CorrelationId,
ConversationId = context.ConversationId,
InitiatorId = context.InitiatorId,
SourceAddress = context.SourceAddress,
DestinationAddress = context.DestinationAddress,
ResponseAddress = context.ResponseAddress,
FaultAddress = context.FaultAddress,
RoutingKey = context.RoutingKey,
RetryCount = context.RetryCount,
};
}

private static MessageContext BuildMetadata(BasicDeliverEventArgs ea)
{
var props = ea.BasicProperties;
var headers = props.Headers ?? new Dictionary<string, object?>();
var sentTime = GetSentTime(props);
return new MessageContext
{
MessageId = props.MessageId,
CorrelationId = props.CorrelationId,
RequestId = props.CorrelationId,
RoutingKey = ea.RoutingKey,
Headers = AmqpHeaderValueNormalizer.Normalize(headers),
Redelivered = ea.Redelivered,
RetryCount = RabbitMqConstants.GetRetryCount(headers),
ConversationId = RabbitMqConstants.GetHeaderString(headers, "ConversationId"),
InitiatorId = RabbitMqConstants.GetHeaderString(headers, "InitiatorId"),
SourceAddress = RabbitMqConstants.GetHeaderUri(headers, "SourceAddress"),
DestinationAddress = RabbitMqConstants.GetHeaderUri(headers, "DestinationAddress"),
ResponseAddress = RabbitMqConstants.GetHeaderUri(headers, "ResponseAddress")
?? (string.IsNullOrEmpty(props.ReplyTo) ? null : new Uri($"queue:{props.ReplyTo}")),
FaultAddress = RabbitMqConstants.GetHeaderUri(headers, "FaultAddress"),
SentTime = sentTime,
ExpirationTime = RabbitMqConstants.TryParseExpiration(props.Expiration, sentTime)
MessageId = metadata.MessageId,
RequestId = metadata.RequestId,
CorrelationId = metadata.CorrelationId,
ConversationId = metadata.ConversationId,
InitiatorId = metadata.InitiatorId,
SourceAddress = metadata.SourceAddress,
DestinationAddress = metadata.DestinationAddress,
ResponseAddress = metadata.ResponseAddress,
FaultAddress = metadata.FaultAddress,
RoutingKey = metadata.RoutingKey,
RetryCount = metadata.RetryCount,
};
}

Expand All @@ -116,34 +90,77 @@ private static MessageContext<TMessage> BuildTypedMetadata<TMessage>(
ISendEndpointProvider? sendEndpointProvider,
CancellationToken cancellationToken)
{
var props = ea.BasicProperties;
var headers = props.Headers ?? new Dictionary<string, object?>();
var sentTime = GetSentTime(props);
var metadata = ExtractMetadata(ea);
return new MessageContext<TMessage>
{
Message = message,
Publisher = publisher,
SendEndpointProvider = sendEndpointProvider,
CancellationToken = cancellationToken,
MessageId = props.MessageId,
CorrelationId = props.CorrelationId,
RequestId = props.CorrelationId,
RoutingKey = ea.RoutingKey,
Headers = AmqpHeaderValueNormalizer.Normalize(headers),
Redelivered = ea.Redelivered,
RetryCount = RabbitMqConstants.GetRetryCount(headers),
ConversationId = RabbitMqConstants.GetHeaderString(headers, "ConversationId"),
InitiatorId = RabbitMqConstants.GetHeaderString(headers, "InitiatorId"),
SourceAddress = RabbitMqConstants.GetHeaderUri(headers, "SourceAddress"),
DestinationAddress = RabbitMqConstants.GetHeaderUri(headers, "DestinationAddress"),
ResponseAddress = RabbitMqConstants.GetHeaderUri(headers, "ResponseAddress")
?? (string.IsNullOrEmpty(props.ReplyTo) ? null : new Uri($"queue:{props.ReplyTo}")),
FaultAddress = RabbitMqConstants.GetHeaderUri(headers, "FaultAddress"),
SentTime = sentTime,
ExpirationTime = RabbitMqConstants.TryParseExpiration(props.Expiration, sentTime)
MessageId = metadata.MessageId,
CorrelationId = metadata.CorrelationId,
RequestId = metadata.RequestId,
RoutingKey = metadata.RoutingKey,
Headers = metadata.Headers,
Redelivered = metadata.Redelivered,
RetryCount = metadata.RetryCount,
ConversationId = metadata.ConversationId,
InitiatorId = metadata.InitiatorId,
SourceAddress = metadata.SourceAddress,
DestinationAddress = metadata.DestinationAddress,
ResponseAddress = metadata.ResponseAddress,
FaultAddress = metadata.FaultAddress,
SentTime = metadata.SentTime,
ExpirationTime = metadata.ExpirationTime,
};
}

/// <summary>
/// Extracts the transport metadata common to both the fault snapshot and the bare-JSON typed context from a
/// single AMQP delivery: the property/header paths, the response-address fallback to <c>ReplyTo</c>, and the
/// TTL anchored to the delivery's timestamp.
/// </summary>
private static DeliveryMetadata ExtractMetadata(BasicDeliverEventArgs ea)
{
var props = ea.BasicProperties;
var headers = props.Headers ?? new Dictionary<string, object?>();
var sentTime = GetSentTime(props);
return new DeliveryMetadata(
MessageId: props.MessageId,
CorrelationId: props.CorrelationId,
RequestId: props.CorrelationId,
RoutingKey: ea.RoutingKey,
Headers: AmqpHeaderValueNormalizer.Normalize(headers),
Redelivered: ea.Redelivered,
RetryCount: RabbitMqConstants.GetRetryCount(headers),
ConversationId: RabbitMqConstants.GetHeaderString(headers, "ConversationId"),
InitiatorId: RabbitMqConstants.GetHeaderString(headers, "InitiatorId"),
SourceAddress: RabbitMqConstants.GetHeaderUri(headers, "SourceAddress"),
DestinationAddress: RabbitMqConstants.GetHeaderUri(headers, "DestinationAddress"),
ResponseAddress: RabbitMqConstants.GetHeaderUri(headers, "ResponseAddress")
?? (string.IsNullOrEmpty(props.ReplyTo) ? null : new Uri($"queue:{props.ReplyTo}")),
FaultAddress: RabbitMqConstants.GetHeaderUri(headers, "FaultAddress"),
SentTime: sentTime,
ExpirationTime: RabbitMqConstants.TryParseExpiration(props.Expiration, sentTime));
}

private static DateTimeOffset? GetSentTime(IReadOnlyBasicProperties props)
=> props.Timestamp.UnixTime > 0 ? DateTimeOffset.FromUnixTimeSeconds(props.Timestamp.UnixTime) : null;

private readonly record struct DeliveryMetadata(
string? MessageId,
string? CorrelationId,
string? RequestId,
string RoutingKey,
IReadOnlyDictionary<string, object?> Headers,
bool Redelivered,
int RetryCount,
string? ConversationId,
string? InitiatorId,
Uri? SourceAddress,
Uri? DestinationAddress,
Uri? ResponseAddress,
Uri? FaultAddress,
DateTimeOffset? SentTime,
DateTimeOffset? ExpirationTime);
}
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,10 @@ namespace Vulthil.Messaging.RabbitMq.Publishing;

internal interface IInternalPublisher
{
Task InternalPublishAsync<TMessage>(
/// <summary>
/// Publishes the already-serialized message body over the message's fanout/topic exchange.
/// </summary>
Task InternalPublishAsync(
byte[] body,
BasicProperties props,
string routingKey,
Expand Down
53 changes: 13 additions & 40 deletions src/Vulthil.Messaging.RabbitMq/Publishing/RabbitMqPublisher.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
using System.Collections.Concurrent;
using System.Diagnostics;
using System.Text.Json;
using Microsoft.Extensions.Logging;
using RabbitMQ.Client;
using Vulthil.Messaging.Abstractions.Publishers;
Expand Down Expand Up @@ -35,7 +34,7 @@ public RabbitMqPublisher(
/// Channels come from a bounded pool — each leased channel is used non-concurrently and returned for reuse, or
/// discarded if it faults.
/// </remarks>
public async Task InternalPublishAsync<TMessage>(
public async Task InternalPublishAsync(
byte[] body,
BasicProperties props,
string routingKey,
Expand Down Expand Up @@ -102,51 +101,25 @@ public async Task PublishAsync<TMessage>(
?? messageConfiguration.RoutingKeyFormatter?.Invoke(message)
?? string.Empty;

var correlationId = publishContext.CorrelationId
?? messageConfiguration.CorrelationIdFormatter?.Invoke(message)
?? Guid.CreateVersion7().ToString();

var messageId = publishContext.MessageId ?? Guid.CreateVersion7().ToString();
var ids = RabbitMqWireMessageBuilder.ResolveIds(message, publishContext, messageConfiguration);
var exchange = messageConfiguration.Exchange;
var urn = messageConfiguration.Urn;
var urnString = urn.AbsoluteUri;

using var activity = MessagingInstrumentation.ActivitySource.StartActivity(
$"{exchange} publish",
ActivityKind.Producer);
using var activity = RabbitMqWireMessageBuilder.StartProducerActivity(
$"{exchange} publish", "publish", exchange, routingKey, ids.UrnString, ids.MessageId, ids.CorrelationId);

if (activity is not null)
{
activity.SetTag(MessagingInstrumentation.Tags.MessagingSystem, MessagingInstrumentation.SystemValue);
activity.SetTag(MessagingInstrumentation.Tags.MessagingOperation, "publish");
activity.SetTag(MessagingInstrumentation.Tags.MessagingDestination, exchange);
activity.SetTag(MessagingInstrumentation.Tags.MessagingRoutingKey, routingKey);
activity.SetTag(MessagingInstrumentation.Tags.MessageType, urnString);
activity.SetTag(MessagingInstrumentation.Tags.MessagingMessageId, messageId);
activity.SetTag(MessagingInstrumentation.Tags.MessagingCorrelationId, correlationId);
}
var properties = RabbitMqWireMessageBuilder.CreateBaseProperties(ids.UrnString, ids.MessageId, publishContext.Headers);
properties.ReplyTo = RabbitMqAddress.ResolveRoutingKey(publishContext.ResponseAddress);
properties.CorrelationId = ids.CorrelationId;
properties.Persistent = true;

var properties = new BasicProperties()
{
Type = urnString,
MessageId = messageId,
ReplyTo = RabbitMqAddress.ResolveRoutingKey(publishContext.ResponseAddress),
CorrelationId = correlationId,
ContentType = RabbitMqConstants.ContentType,
Headers = new Dictionary<string, object?>(publishContext.Headers),
Persistent = true,
Timestamp = new AmqpTimestamp(DateTimeOffset.UtcNow.ToUnixTimeSeconds()),
};

var envelope = MessageEnvelopeFactory.Create(
message, publishContext, messageId, correlationId, urn, _messageConfigurationProvider.JsonSerializerOptions);
var body = JsonSerializer.SerializeToUtf8Bytes(envelope, _messageConfigurationProvider.JsonSerializerOptions);

MessagingLog.Publishing(_logger, urnString, exchange, routingKey, messageId);
var body = RabbitMqWireMessageBuilder.SerializeEnvelope(
message, publishContext, ids.MessageId, ids.CorrelationId, ids.Urn, _messageConfigurationProvider.JsonSerializerOptions);

MessagingLog.Publishing(_logger, ids.UrnString, exchange, routingKey, ids.MessageId);

try
{
await InternalPublishAsync<TMessage>(body, properties, routingKey, messageConfiguration, cancellationToken);
await InternalPublishAsync(body, properties, routingKey, messageConfiguration, cancellationToken);
activity?.SetStatus(ActivityStatusCode.Ok);
}
catch (Exception ex)
Expand Down
98 changes: 98 additions & 0 deletions src/Vulthil.Messaging.RabbitMq/RabbitMqWireMessageBuilder.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
using System.Diagnostics;
using System.Text.Json;
using RabbitMQ.Client;
using Vulthil.Messaging.RabbitMq.Telemetry;
using Vulthil.Messaging.Transport;

namespace Vulthil.Messaging.RabbitMq;

/// <summary>
/// Shared wire-message construction for the RabbitMQ producer paths (<c>RabbitMqPublisher</c>,
/// <c>RabbitMqSendEndpoint</c>, <c>RabbitMqRequester</c>): resolving the correlation/message identifiers, the common
/// <see cref="BasicProperties"/> fields, the serialized <see cref="MessageEnvelope"/>, and the producer activity
/// tags. Each producer applies its own routing/exchange selection and the <see cref="BasicProperties"/> fields that
/// legitimately differ per operation (reply-to resolution, the wire correlation id, persistence, expiration) on top
/// of the shared result, so the extraction covers only the parts the three sites compute identically.
/// </summary>
internal static class RabbitMqWireMessageBuilder
{
/// <summary>The identifiers resolved for a single outgoing message, shared by every producer path.</summary>
/// <param name="CorrelationId">The resolved business correlation identifier.</param>
/// <param name="MessageId">The resolved message identifier.</param>
/// <param name="Urn">The stable wire URN for the message type.</param>
/// <param name="UrnString">The URN's absolute URI string, used as the AMQP <c>Type</c> and activity tag.</param>
internal readonly record struct ResolvedIds(string CorrelationId, string MessageId, Uri Urn, string UrnString);

/// <summary>
/// Resolves the correlation id, message id, and URN for <paramref name="message"/> the same way across every
/// producer path: an explicit value on <paramref name="context"/> wins, then the type's configured formatter
/// (correlation id only), then a fresh id.
/// </summary>
public static ResolvedIds ResolveIds<TMessage>(TMessage message, PublishContext context, MessageConfiguration messageConfiguration)
where TMessage : notnull
{
var correlationId = context.CorrelationId
?? messageConfiguration.CorrelationIdFormatter?.Invoke(message)
?? Guid.CreateVersion7().ToString();

var messageId = context.MessageId ?? Guid.CreateVersion7().ToString();
var urn = messageConfiguration.Urn;

return new ResolvedIds(correlationId, messageId, urn, urn.AbsoluteUri);
}

/// <summary>
/// Builds and serializes the <see cref="MessageEnvelope"/> for the outgoing message. Callers that convert a
/// publish failure into a typed result (rather than letting it propagate) must call this from within their own
/// try/catch, since serialization can fail for the same reasons the send itself can.
/// </summary>
public static byte[] SerializeEnvelope<TMessage>(
TMessage message,
PublishContext context,
string messageId,
string correlationId,
Uri urn,
JsonSerializerOptions jsonOptions,
string? requestId = null)
where TMessage : notnull
{
var envelope = MessageEnvelopeFactory.Create(message, context, messageId, correlationId, urn, jsonOptions, requestId);
return JsonSerializer.SerializeToUtf8Bytes(envelope, jsonOptions);
}

/// <summary>
/// Creates the <see cref="BasicProperties"/> fields common to every RabbitMQ producer path. Callers set the
/// remaining fields that legitimately differ per operation (<c>ReplyTo</c>, the wire <c>CorrelationId</c>,
/// <c>Persistent</c>, <c>Expiration</c>).
/// </summary>
public static BasicProperties CreateBaseProperties(string urnString, string messageId, IReadOnlyDictionary<string, object?> headers) => new()
{
Type = urnString,
MessageId = messageId,
ContentType = RabbitMqConstants.ContentType,
Headers = new Dictionary<string, object?>(headers),
Timestamp = new AmqpTimestamp(DateTimeOffset.UtcNow.ToUnixTimeSeconds()),
};

/// <summary>
/// Starts a producer <see cref="Activity"/> and applies the standard Vulthil messaging tags. Returns
/// <see langword="null"/> when no listener is recording the transport's activity source.
/// </summary>
public static Activity? StartProducerActivity(
string name, string operation, string destination, string routingKey, string urnString, string messageId, string correlationId)
{
var activity = MessagingInstrumentation.ActivitySource.StartActivity(name, ActivityKind.Producer);
if (activity is not null)
{
activity.SetTag(MessagingInstrumentation.Tags.MessagingSystem, MessagingInstrumentation.SystemValue);
activity.SetTag(MessagingInstrumentation.Tags.MessagingOperation, operation);
activity.SetTag(MessagingInstrumentation.Tags.MessagingDestination, destination);
activity.SetTag(MessagingInstrumentation.Tags.MessagingRoutingKey, routingKey);
activity.SetTag(MessagingInstrumentation.Tags.MessageType, urnString);
activity.SetTag(MessagingInstrumentation.Tags.MessagingMessageId, messageId);
activity.SetTag(MessagingInstrumentation.Tags.MessagingCorrelationId, correlationId);
}

return activity;
}
}
Loading