Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
5ff78e4
feat: Add opt-in non-exclusive session locking for session receivers
EldertGrootenboerMS Jun 18, 2026
3a1bc02
Merge branch 'main' into feature/servicebus-nonexclusive-session-3771…
EldertGrootenboer Jun 18, 2026
a7778cc
test: Fix EventSourceTests mock callback arity for new session receiv…
EldertGrootenboerMS Jun 19, 2026
6d2fc8b
Merge branch 'main' into feature/servicebus-nonexclusive-session-3771…
EldertGrootenboer Jun 19, 2026
9078bbd
Switch Track 2 non-exclusive session to composite AmqpNonExclusiveSes…
Jun 29, 2026
0da1d79
Allow accept-any (next available) session in non-exclusive mode
Jun 29, 2026
53d6ca0
test: Add codec round-trip and accept-any tests for non-exclusive ses…
EldertGrootenboerMS Jun 29, 2026
6581ba3
Add ToString() to non-exclusive session filter codec
EldertGrootenboerMS Jul 8, 2026
11a3839
test: Add accept-any non-exclusive session takeover live test
EldertGrootenboerMS Jul 8, 2026
57b1cd0
Merge branch 'main' into feature/servicebus-nonexclusive-session-3771…
EldertGrootenboer Jul 8, 2026
126d5a2
Clarify non-exclusive session filter comments and error message
EldertGrootenboerMS Jul 8, 2026
863df64
Merge branch 'main' into feature/servicebus-nonexclusive-session-3771…
EldertGrootenboer Jul 8, 2026
6495224
Clarify session filter comment in AmqpConnectionScope
EldertGrootenboerMS Jul 8, 2026
ac76210
Merge remote-tracking branch 'origin/main' into feature/servicebus-no…
EldertGrootenboerMS Jul 8, 2026
09db390
Refine non-exclusive session test comment and enable exclusive live test
EldertGrootenboerMS Jul 8, 2026
c25436e
Merge branch 'main' into feature/servicebus-nonexclusive-session-3771…
EldertGrootenboer Jul 8, 2026
710a5e2
Dispose session receivers in live tests and rename filter tests
EldertGrootenboerMS Jul 8, 2026
b500df7
Fix SessionLockedUntil and SessionId doc comments
EldertGrootenboerMS Jul 8, 2026
f85a327
Match any values for new receiver params in Moq setups
EldertGrootenboerMS Jul 8, 2026
efbfd9e
Merge branch 'main' into feature/servicebus-nonexclusive-session-3771…
EldertGrootenboer Jul 8, 2026
48ff0d5
Enable non-exclusive session live tests
EldertGrootenboerMS Jul 8, 2026
1d47580
Merge branch 'main' into feature/servicebus-nonexclusive-session-3771…
EldertGrootenboer Jul 8, 2026
7b7be08
Merge branch 'main' into feature/servicebus-nonexclusive-session-3771…
EldertGrootenboer Jul 8, 2026
e72e53b
Merge remote-tracking branch 'origin/main' into feature/servicebus-no…
EldertGrootenboer Jul 29, 2026
ab60191
test: Cover new CreateTransportReceiver parameters in GetMessageSessi…
EldertGrootenboer Jul 29, 2026
16da25d
docs: correct non-exclusive session comments to match observed servic…
EldertGrootenboer Jul 31, 2026
3c0a9d7
Address pre-PR review board findings for non-exclusive session locking
EldertGrootenboer Aug 10, 2026
0719e16
Remove EditorBrowsable(Advanced) tags per review feedback
EldertGrootenboer Aug 10, 2026
8c2af85
Report unsupported non-exclusive sessions consistently and de-race th…
EldertGrootenboer Aug 13, 2026
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
1 change: 1 addition & 0 deletions sdk/servicebus/Azure.Messaging.ServiceBus/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
### Features Added

- Added `GetMessageSessionsAsync` overloads on `ServiceBusClient` for queues and subscriptions. The no-filter overload returns the IDs of sessions that have active messages or session state, and the `sessionStateUpdatedAfter` overload returns session IDs whose session state was updated after the specified timestamp. Implements the `com.microsoft:get-message-sessions` AMQP management operation. ([#58761](https://github.com/Azure/azure-sdk-for-net/pull/58761))
- Added opt-in support for non-exclusive session locking on `ServiceBusSessionReceiver`, allowing a session to be cooperatively taken over by another receiver. Set `ServiceBusSessionReceiverOptions.EnableNonExclusiveSession` to accept a session non-exclusively, then read the token from `ServiceBusSessionReceiver.SessionLockToken` and pass it as `ServiceBusSessionReceiverOptions.SessionLockToken = Guid.Parse(token)` to take that session over. `ServiceBusSessionReceiver.IsSessionExclusive` reports the mode the session was established under. Dispositions for a non-exclusive session are routed over the management link so that settlement keeps working across a takeover, which lowers settlement throughput compared to an exclusive session. This applies to `ServiceBusSessionReceiver` only; `ServiceBusSessionProcessor` continues to lock sessions exclusively. Accepting a session with `EnableNonExclusiveSession` set throws `NotSupportedException` when the endpoint declines it, either by refusing the request outright or by accepting it without assigning a lock token, which is how a caller detects whether the feature is available for a namespace. An endpoint that declines in some other way surfaces the exception its own error maps to. ([#60060](https://github.com/Azure/azure-sdk-for-net/pull/60060))

Comment thread
EldertGrootenboer marked this conversation as resolved.
### Breaking Changes

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -499,18 +499,22 @@ public partial class ServiceBusSessionReceiver : Azure.Messaging.ServiceBus.Serv
{
protected ServiceBusSessionReceiver() { }
public override bool IsClosed { get { throw null; } }
public virtual bool IsSessionExclusive { get { throw null; } }
public virtual string SessionId { get { throw null; } }
public virtual System.DateTimeOffset SessionLockedUntil { get { throw null; } }
public virtual string SessionLockToken { get { throw null; } }
public virtual System.Threading.Tasks.Task<System.BinaryData> GetSessionStateAsync(System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
public virtual System.Threading.Tasks.Task RenewSessionLockAsync(System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
public virtual System.Threading.Tasks.Task SetSessionStateAsync(System.BinaryData sessionState, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
}
public partial class ServiceBusSessionReceiverOptions
{
public ServiceBusSessionReceiverOptions() { }
public bool EnableNonExclusiveSession { get { throw null; } set { } }
public string Identifier { get { throw null; } set { } }
public int PrefetchCount { get { throw null; } set { } }
public Azure.Messaging.ServiceBus.ServiceBusReceiveMode ReceiveMode { get { throw null; } set { } }
public System.Guid? SessionLockToken { get { throw null; } set { } }
Comment thread
EldertGrootenboer marked this conversation as resolved.
public override bool Equals(object obj) { throw null; }
public override int GetHashCode() { throw null; }
public override string ToString() { throw null; }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -499,18 +499,22 @@ public partial class ServiceBusSessionReceiver : Azure.Messaging.ServiceBus.Serv
{
protected ServiceBusSessionReceiver() { }
public override bool IsClosed { get { throw null; } }
public virtual bool IsSessionExclusive { get { throw null; } }
public virtual string SessionId { get { throw null; } }
public virtual System.DateTimeOffset SessionLockedUntil { get { throw null; } }
public virtual string SessionLockToken { get { throw null; } }
public virtual System.Threading.Tasks.Task<System.BinaryData> GetSessionStateAsync(System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
public virtual System.Threading.Tasks.Task RenewSessionLockAsync(System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
public virtual System.Threading.Tasks.Task SetSessionStateAsync(System.BinaryData sessionState, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
}
public partial class ServiceBusSessionReceiverOptions
{
public ServiceBusSessionReceiverOptions() { }
public bool EnableNonExclusiveSession { get { throw null; } set { } }
public string Identifier { get { throw null; } set { } }
public int PrefetchCount { get { throw null; } set { } }
public Azure.Messaging.ServiceBus.ServiceBusReceiveMode ReceiveMode { get { throw null; } set { } }
public System.Guid? SessionLockToken { get { throw null; } set { } }
public override bool Equals(object obj) { throw null; }
public override int GetHashCode() { throw null; }
public override string ToString() { throw null; }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -499,18 +499,22 @@ public partial class ServiceBusSessionReceiver : Azure.Messaging.ServiceBus.Serv
{
protected ServiceBusSessionReceiver() { }
public override bool IsClosed { get { throw null; } }
public virtual bool IsSessionExclusive { get { throw null; } }
public virtual string SessionId { get { throw null; } }
public virtual System.DateTimeOffset SessionLockedUntil { get { throw null; } }
public virtual string SessionLockToken { get { throw null; } }
public virtual System.Threading.Tasks.Task<System.BinaryData> GetSessionStateAsync(System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
public virtual System.Threading.Tasks.Task RenewSessionLockAsync(System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
public virtual System.Threading.Tasks.Task SetSessionStateAsync(System.BinaryData sessionState, System.Threading.CancellationToken cancellationToken = default(System.Threading.CancellationToken)) { throw null; }
}
public partial class ServiceBusSessionReceiverOptions
{
public ServiceBusSessionReceiverOptions() { }
public bool EnableNonExclusiveSession { get { throw null; } set { } }
public string Identifier { get { throw null; } set { } }
public int PrefetchCount { get { throw null; } set { } }
public Azure.Messaging.ServiceBus.ServiceBusReceiveMode ReceiveMode { get { throw null; } set { } }
public System.Guid? SessionLockToken { get { throw null; } set { } }
public override bool Equals(object obj) { throw null; }
public override int GetHashCode() { throw null; }
public override string ToString() { throw null; }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -168,6 +168,8 @@ public override TransportSender CreateSender(
/// <param name="sessionId">The session ID to receive messages for.</param>
/// <param name="isSessionReceiver">Whether or not this is a sessionful receiver link.</param>
/// <param name="isProcessor">Whether or not the receiver is being created for a processor.</param>
/// <param name="isSessionExclusive">Whether or not the session is locked exclusively. Only applicable for session receivers.</param>
/// <param name="sessionLockToken">The session lock token to present when cooperatively taking over a non-exclusive session. Only applicable for session receivers.</param>
/// <param name="cancellationToken">An optional <see cref="CancellationToken"/> instance to signal the request to cancel the
/// open link operation. Only applicable for session receivers.</param>
/// <returns>A <see cref="TransportReceiver" /> configured in the requested manner.</returns>
Expand All @@ -180,7 +182,9 @@ public override TransportReceiver CreateReceiver(
string sessionId,
bool isSessionReceiver,
bool isProcessor,
CancellationToken cancellationToken)
bool isSessionExclusive = true,
Guid? sessionLockToken = null,
CancellationToken cancellationToken = default)
{
Argument.AssertNotDisposed(_closed, nameof(AmqpClient));

Expand All @@ -196,6 +200,8 @@ public override TransportReceiver CreateReceiver(
isSessionReceiver,
isProcessor,
_messageConverter,
isSessionExclusive,
sessionLockToken,
cancellationToken
);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,17 @@ internal class AmqpClientConstants
public const string FilterReceivedAt = FilterReceivedAtPartNameV2 + " > ";
public const string FilterReceivedAtFormatString = FilterReceivedAt + "{0}";
public static readonly AmqpSymbol SessionFilterName = AmqpConstants.Vendor + ":session-filter";
public static readonly AmqpSymbol NonExclusiveSessionFilterName = AmqpConstants.Vendor + ":non-exclusive-session-filter";

/// <summary>
/// The portion of an attach rejection that names the filter as one the endpoint does not recognize. An
/// endpoint without non-exclusive session locking refuses <see cref="NonExclusiveSessionFilterName" /> with
/// a description carrying this. The service authors the wording under no contract, so a rewording returns
/// callers to the exception the refusal's condition maps to. Observed in the form "The link '&lt;name&gt;'
/// contains invalid filter type. System only support Jms or Apache selector filter type."
/// </summary>
public const string UnrecognizedFilterErrorFragment = "invalid filter type";

public static readonly AmqpSymbol MessageReceiptsFilterName = AmqpConstants.Vendor + ":message-receipts-filter";
Comment thread
EldertGrootenboer marked this conversation as resolved.
public static readonly AmqpSymbol ClientSideCursorFilterName = AmqpConstants.Vendor + ":client-side-filter";
public static readonly TimeSpan ClientMinimumTokenRefreshInterval = TimeSpan.FromMinutes(4);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
using System.Threading.Tasks;
using Azure.Core;
using Azure.Core.Diagnostics;
using Azure.Messaging.ServiceBus.Amqp.Framing;
using Azure.Messaging.ServiceBus.Authorization;
using Azure.Messaging.ServiceBus.Core;
using Azure.Messaging.ServiceBus.Diagnostics;
Expand Down Expand Up @@ -362,6 +363,8 @@ public virtual async Task<RequestResponseAmqpLink> OpenManagementLinkAsync(
/// <param name="receiveMode">The <see cref="ServiceBusReceiveMode"/> used to specify how messages are received. Defaults to PeekLock mode.</param>
/// <param name="sessionId">The session to connect to.</param>
/// <param name="isSessionReceiver">Whether or not this is a sessionful receiver.</param>
/// <param name="isSessionExclusive">Whether or not the session is locked exclusively. Only applicable for session receivers.</param>
/// <param name="sessionLockToken">The session lock token to present when cooperatively taking over a non-exclusive session. Only applicable for session receivers.</param>
/// <param name="cancellationToken">An optional <see cref="CancellationToken"/> instance to signal the request to cancel the operation.</param>
/// <returns>A link for use with consumer operations.</returns>
public virtual async Task<ReceivingAmqpLink> OpenReceiverLinkAsync(
Expand All @@ -372,7 +375,9 @@ public virtual async Task<ReceivingAmqpLink> OpenReceiverLinkAsync(
ServiceBusReceiveMode receiveMode,
string sessionId,
bool isSessionReceiver,
CancellationToken cancellationToken)
bool isSessionExclusive = true,
Guid? sessionLockToken = null,
CancellationToken cancellationToken = default)
{
Argument.AssertNotDisposed(_disposed, nameof(AmqpConnectionScope));
cancellationToken.ThrowIfCancellationRequested<TaskCanceledException>();
Expand All @@ -392,6 +397,8 @@ public virtual async Task<ReceivingAmqpLink> OpenReceiverLinkAsync(
receiveMode: receiveMode,
sessionId: sessionId,
isSessionReceiver: isSessionReceiver,
isSessionExclusive: isSessionExclusive,
sessionLockToken: sessionLockToken,
cancellationToken: cancellationToken
).ConfigureAwait(false);

Expand Down Expand Up @@ -637,6 +644,8 @@ protected virtual async Task<RequestResponseAmqpLink> CreateManagementLinkAsync(
/// <param name="receiveMode">The <see cref="ServiceBusReceiveMode"/> used to specify how messages are received. Defaults to PeekLock mode.</param>
/// <param name="sessionId">The session to receive from.</param>
/// <param name="isSessionReceiver">Whether or not this is a sessionful receiver.</param>
/// <param name="isSessionExclusive">Whether or not the session is locked exclusively. Only applicable for session receivers.</param>
/// <param name="sessionLockToken">The session lock token to present when cooperatively taking over a non-exclusive session. Only applicable for session receivers.</param>
/// <param name="cancellationToken">An optional <see cref="CancellationToken"/> instance to signal the request to cancel the operation.</param>
/// <returns>A link for use for operations related to receiving events.</returns>
protected virtual async Task<ReceivingAmqpLink> CreateReceivingLinkAsync(
Expand All @@ -649,7 +658,9 @@ protected virtual async Task<ReceivingAmqpLink> CreateReceivingLinkAsync(
ServiceBusReceiveMode receiveMode,
string sessionId,
bool isSessionReceiver,
CancellationToken cancellationToken)
bool isSessionExclusive = true,
Guid? sessionLockToken = null,
CancellationToken cancellationToken = default)
{
Argument.AssertNotDisposed(IsDisposed, nameof(AmqpConnectionScope));
cancellationToken.ThrowIfCancellationRequested<TaskCanceledException>();
Expand Down Expand Up @@ -683,10 +694,25 @@ protected virtual async Task<ReceivingAmqpLink> CreateReceivingLinkAsync(

var filters = new FilterSet();

// even if supplied sessionId is null, we need to add the Session filter if it is a session receiver
// even if the supplied sessionId is null, a session receiver needs a session filter on the link:
// the plain session filter for an exclusive session, or the composite non-exclusive session filter otherwise.
if (isSessionReceiver)
{
filters.Add(AmqpClientConstants.SessionFilterName, sessionId);
if (isSessionExclusive)
{
filters.Add(AmqpClientConstants.SessionFilterName, sessionId);
}
else
{
// Non-exclusive locking: a single composite filter carries the session id and (for takeover)
// the lock token. Presence of this filter implies non-exclusive mode. The plain session
// filter is omitted. A service that predates this change does not silently ignore the unknown
// filter - it rejects the attach with an "invalid filter type" error - so this path depends on
// the service-side change being deployed in the target region.
filters.Add(
AmqpClientConstants.NonExclusiveSessionFilterName,
new AmqpNonExclusiveSessionFilterCodec { SessionId = sessionId, LockToken = sessionLockToken });
}
}

var linkSettings = new AmqpLinkSettings
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,9 @@ public static Exception ToMessagingContractException(string condition, string me
return new NotSupportedException(EnrichMessage(message));
}

// A receiver opening a non-exclusive session link reports one refusal carrying this condition as a
// NotSupportedException instead, because that refusal is the endpoint declining the feature rather than
// the operation. See AmqpReceiver.IsUnrecognizedFilterRejection; every other refusal maps here.
if (string.Equals(condition, AmqpErrorCode.NotAllowed.Value, StringComparison.InvariantCultureIgnoreCase))
{
return new InvalidOperationException(EnrichMessage(message));
Expand Down
Loading
Loading