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
3 changes: 3 additions & 0 deletions .github/release.yml
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,9 @@ changelog:
- title: 📚 Documentation
labels:
- documentation
- title: 🪫 Deprecation
labels:
- deprecation
- title: 🔄 Dependency Updates
labels:
- dependencies
Expand Down
6 changes: 5 additions & 1 deletion RabbitMQ.Stream.Client/AddressResolver.cs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,9 @@
// 2.0, and the Mozilla Public License, version 2.0.
// Copyright (c) 2017-2023 Broadcom. All Rights Reserved. The term "Broadcom" refers to Broadcom Inc. and/or its subsidiaries.

using System;
using System.Net;
using System.Threading.Tasks;

namespace RabbitMQ.Stream.Client
{
Expand All @@ -16,6 +18,8 @@ public AddressResolver(EndPoint endPoint)

public EndPoint EndPoint { get; set; }
public bool Enabled { get; set; }
public EndPoint Resolve(string address, int host) => EndPoint;
[Obsolete("Deprecated. Use ResolveAsync instead.")]
public EndPoint Resolve(string address, int port) => EndPoint;
public async Task<EndPoint> ResolveAsync(string address, int port) => EndPoint;
}
}
7 changes: 6 additions & 1 deletion RabbitMQ.Stream.Client/AddressResolverDynamic.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

using System;
using System.Net;
using System.Threading.Tasks;

namespace RabbitMQ.Stream.Client;

Expand All @@ -18,5 +19,9 @@ public AddressResolverDynamic(Func<string, int, EndPoint> resolveFunction)
}

public bool Enabled { get; set; }
public EndPoint Resolve(string address, int host) => _resolveFunction(address, host);
[Obsolete("Deprecated. Use ResolveAsync instead.")]
public EndPoint Resolve(string address, int port) => _resolveFunction(address, port);
#pragma warning disable CS0618
public Task<EndPoint> ResolveAsync(string address, int port) => Task.FromResult(Resolve(address, port));
#pragma warning restore CS0618
}
31 changes: 31 additions & 0 deletions RabbitMQ.Stream.Client/DnsAddressResolver.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
// This source code is dual-licensed under the Apache License, version
// 2.0, and the Mozilla Public License, version 2.0.
// Copyright (c) 2017-2023 Broadcom. All Rights Reserved. The term "Broadcom" refers to Broadcom Inc. and/or its subsidiaries.

using System;
using System.Net;
using System.Threading.Tasks;

namespace RabbitMQ.Stream.Client
{
public class DnsAddressResolver : IAddressResolver
{
public DnsAddressResolver(DnsEndPoint endPoint)
{
EndPoint = endPoint;
Enabled = true;
}

public EndPoint EndPoint { get; set; }
public bool Enabled { get; set; }
[Obsolete("Deprecated. Use ResolveAsync instead.")]
public EndPoint Resolve(string address, int port) => ResolveAsync(address, port).GetAwaiter().GetResult();
public async Task<EndPoint> ResolveAsync(string address, int port)
{
var entries = await Dns.GetHostEntryAsync(((DnsEndPoint)EndPoint).Host).ConfigureAwait(false);
var addressList = entries.AddressList;
var targetIp = addressList[Random.Shared.Next(addressList.Length)];
return new IPEndPoint(targetIp, port);
}
}
}
3 changes: 2 additions & 1 deletion RabbitMQ.Stream.Client/IAddressResolver.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,12 @@
// Copyright (c) 2017-2023 Broadcom. All Rights Reserved. The term "Broadcom" refers to Broadcom Inc. and/or its subsidiaries.

using System.Net;
using System.Threading.Tasks;

namespace RabbitMQ.Stream.Client;

public interface IAddressResolver
{
public bool Enabled { get; }
public EndPoint Resolve(string address, int host);
public Task<EndPoint> ResolveAsync(string address, int port);
}
16 changes: 13 additions & 3 deletions RabbitMQ.Stream.Client/PublicAPI.Unshipped.txt
Original file line number Diff line number Diff line change
@@ -1,3 +1,11 @@
RabbitMQ.Stream.Client.DnsAddressResolver
RabbitMQ.Stream.Client.DnsAddressResolver.DnsAddressResolver(System.Net.DnsEndPoint endPoint) -> void
RabbitMQ.Stream.Client.DnsAddressResolver.Enabled.get -> bool
RabbitMQ.Stream.Client.DnsAddressResolver.Enabled.set -> void
RabbitMQ.Stream.Client.DnsAddressResolver.EndPoint.get -> System.Net.EndPoint
RabbitMQ.Stream.Client.DnsAddressResolver.EndPoint.set -> void
RabbitMQ.Stream.Client.DnsAddressResolver.Resolve(string address, int port) -> System.Net.EndPoint
RabbitMQ.Stream.Client.DnsAddressResolver.ResolveAsync(string address, int port) -> System.Threading.Tasks.Task<System.Net.EndPoint>
abstract RabbitMQ.Stream.Client.AbstractEntity.Close() -> System.Threading.Tasks.Task<RabbitMQ.Stream.Client.ResponseCode>
abstract RabbitMQ.Stream.Client.AbstractEntity.DeleteEntityFromTheServer(bool ignoreIfAlreadyDeleted = false) -> System.Threading.Tasks.Task<RabbitMQ.Stream.Client.ResponseCode>
abstract RabbitMQ.Stream.Client.AbstractEntity.DumpEntityConfiguration() -> string
Expand Down Expand Up @@ -29,12 +37,14 @@ RabbitMQ.Stream.Client.AddressResolver.Enabled.get -> bool
RabbitMQ.Stream.Client.AddressResolver.Enabled.set -> void
RabbitMQ.Stream.Client.AddressResolver.EndPoint.get -> System.Net.EndPoint
RabbitMQ.Stream.Client.AddressResolver.EndPoint.set -> void
RabbitMQ.Stream.Client.AddressResolver.Resolve(string address, int host) -> System.Net.EndPoint
RabbitMQ.Stream.Client.AddressResolver.Resolve(string address, int port) -> System.Net.EndPoint
RabbitMQ.Stream.Client.AddressResolver.ResolveAsync(string address, int port) -> System.Threading.Tasks.Task<System.Net.EndPoint>
RabbitMQ.Stream.Client.AddressResolverDynamic
RabbitMQ.Stream.Client.AddressResolverDynamic.AddressResolverDynamic(System.Func<string, int, System.Net.EndPoint> resolveFunction) -> void
RabbitMQ.Stream.Client.AddressResolverDynamic.Enabled.get -> bool
RabbitMQ.Stream.Client.AddressResolverDynamic.Enabled.set -> void
RabbitMQ.Stream.Client.AddressResolverDynamic.Resolve(string address, int host) -> System.Net.EndPoint
RabbitMQ.Stream.Client.AddressResolverDynamic.Resolve(string address, int port) -> System.Net.EndPoint
RabbitMQ.Stream.Client.AddressResolverDynamic.ResolveAsync(string address, int port) -> System.Threading.Tasks.Task<System.Net.EndPoint>
RabbitMQ.Stream.Client.AlreadyClosedException
RabbitMQ.Stream.Client.AlreadyClosedException.AlreadyClosedException(string s) -> void
RabbitMQ.Stream.Client.AuthMechanism
Expand Down Expand Up @@ -169,7 +179,7 @@ RabbitMQ.Stream.Client.HashRoutingMurmurStrategy.Route(RabbitMQ.Stream.Client.Me
RabbitMQ.Stream.Client.HeartBeatHandler.HeartBeatHandler(System.Func<System.Threading.Tasks.ValueTask<bool>> sendHeartbeatFunc, System.Func<string, string, System.Threading.Tasks.Task<RabbitMQ.Stream.Client.CloseResponse>> close, int heartbeat, Microsoft.Extensions.Logging.ILogger<RabbitMQ.Stream.Client.HeartBeatHandler> logger = null) -> void
RabbitMQ.Stream.Client.IAddressResolver
RabbitMQ.Stream.Client.IAddressResolver.Enabled.get -> bool
RabbitMQ.Stream.Client.IAddressResolver.Resolve(string address, int host) -> System.Net.EndPoint
RabbitMQ.Stream.Client.IAddressResolver.ResolveAsync(string address, int port) -> System.Threading.Tasks.Task<System.Net.EndPoint>
RabbitMQ.Stream.Client.IClient.ClientId.get -> string
RabbitMQ.Stream.Client.IClient.ClientId.init -> void
RabbitMQ.Stream.Client.IClient.Consumers.get -> System.Collections.Generic.IDictionary<byte, (string, RabbitMQ.Stream.Client.ConsumerEvents)>
Expand Down
6 changes: 4 additions & 2 deletions RabbitMQ.Stream.Client/RoutingClient.cs
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,7 @@ clientParameters with
// here it means that there is a AddressResolver configuration
// so there is a load-balancer or proxy we need to get the right connection
// as first we try with the first node given from the LB
var endPoint = clientParameters.AddressResolver.Resolve(broker.Host, (int)broker.Port);
var endPoint = await clientParameters.AddressResolver.ResolveAsync(broker.Host, (int)broker.Port).ConfigureAwait(false);
var client = await routing
.CreateClient(
clientParameters with
Expand All @@ -94,8 +94,10 @@ clientParameters with

var attemptNo = 0;
while (broker.Host != advertisedHost || broker.Port != uint.Parse(advertisedPort))

{
logger?.LogDebug(
endPoint = await clientParameters.AddressResolver.ResolveAsync(broker.Host, (int)broker.Port).ConfigureAwait(false);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why did you add this?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The idea is to have the load balanced implemented on client-side. Here you can see it iterates over the entries. So each call for ResolveAsync would return a different entry.

@lukas8219 lukas8219 May 21, 2026

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

All calls to Dns.GetHostEntryAsync return the same list. The Endpoint being returned on the default approach using Reverse Proxies (where IP doesn't change or doesn't mean anything) it would be constantly hitting the same entry over and over.

logger?.LogInformation(
"advertised_host or advertised_port doesn't match. Expected: {ExpectedHost}:{ExpectedPort}, " +
"Actual: {AdvertisedHost}:{AdvertisedPort}. Attempt number: {AttemptNo}/{MaxAttempts}",
broker.Host, broker.Port, advertisedHost, advertisedPort, attemptNo, maxAttempts);
Expand Down