Skip to content
Open
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
1 change: 1 addition & 0 deletions sdk/servicebus/azure-servicebus/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
### Other Changes

- When using the async `AmqpOverWebsocket` transport on Python 3.10 or later, `aiohttp>=3.14.0` is now recommended. Earlier `aiohttp` versions have a WebSocket heartbeat bug ([aio-libs/aiohttp#12030](https://github.com/aio-libs/aiohttp/pull/12030)) that can cause the connection to be dropped during long message processing, surfacing as a `SocketError` ("Cannot write to closing transport"). Python 3.9 users must upgrade Python to install an `aiohttp` release containing this fix. ([#44028](https://github.com/Azure/azure-sdk-for-python/issues/44028))
- Management operations (peek, deferred receive, message settlement over the management link, lock renewal, session state, session listing, schedule/cancel) now send `com.microsoft:server-timeout`: the caller's remaining time less one second, or 60 seconds when none was given. Previously no bound was sent, so a stalled service held the call until the AMQP link failed; it now raises a retryable `OperationTimeoutError`, so a persistently stalled service surfaces after roughly four minutes at default retry settings. Matches the .NET, Java and Go SDKs.

## 7.14.3 (2025-11-11)

Expand Down
2 changes: 1 addition & 1 deletion sdk/servicebus/azure-servicebus/api.metadata.yml
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
apiMdSha256: 68c237729c78165bea132f330a2e89519b2e6ba19373d40c48c55f09e8a994de
parserVersion: 0.3.30
parserVersion: 0.3.31
pythonVersion: 3.12.10
Original file line number Diff line number Diff line change
Expand Up @@ -25,13 +25,19 @@
OperationTimeoutError,
SessionLockLostError,
)
from ._common.utils import create_properties, strip_protocol_from_uri, parse_sas_credential
from ._common.utils import (
create_properties,
strip_protocol_from_uri,
parse_sas_credential,
get_server_timeout_ms,
)
from ._common.constants import (
CONTAINER_PREFIX,
MANAGEMENT_PATH_SUFFIX,
TOKEN_TYPE_SASTOKEN,
MGMT_REQUEST_OP_TYPE_ENTITY_MGMT,
ASSOCIATEDLINKPROPERTYNAME,
REQUEST_RESPONSE_TIMEOUT,
)

if TYPE_CHECKING:
Expand Down Expand Up @@ -499,6 +505,10 @@ def _mgmt_request_response(
except AttributeError:
pass

application_properties[REQUEST_RESPONSE_TIMEOUT] = self._amqp_transport.AMQP_UINT_VALUE(
get_server_timeout_ms(timeout)
)

mgmt_msg = self._amqp_transport.create_mgmt_msg(
message=message,
application_properties=application_properties,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,14 @@
PYAMQP_LIBRARY = "pyamqp"
OPERATION_TIMEOUT = VENDOR + b":timeout"

# Bounds a management operation on the service side when the caller gave no timeout.
DEFAULT_SERVER_TIMEOUT_MS = 60000
# Subtracted from the caller's remaining time so the service answers before the client
# gives up.
SERVER_TIMEOUT_BUFFER_MS = 1000
# The property is encoded as an AMQP uint, so cap at its maximum (about 49.7 days).
MAX_SERVER_TIMEOUT_MS = 2**32 - 1

MANAGEMENT_PATH_SUFFIX = "/$management"

MGMT_RESPONSE_SESSION_STATE = b"session-state"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,9 @@
DEAD_LETTER_QUEUE_SUFFIX,
TRANSFER_DEAD_LETTER_QUEUE_SUFFIX,
USER_AGENT_PREFIX,
DEFAULT_SERVER_TIMEOUT_MS,
SERVER_TIMEOUT_BUFFER_MS,
MAX_SERVER_TIMEOUT_MS,
)
from ..amqp import AmqpAnnotatedMessage

Expand Down Expand Up @@ -89,6 +92,23 @@ def utc_now():
return datetime.datetime.now(timezone.utc)


def get_server_timeout_ms(timeout: Optional[float]) -> int:
"""Return the server-timeout for a management operation, in milliseconds.

This is a service-side bound, not a client-side one. It clamps at zero, since under
a second of remaining time there is no room for the service to answer first, and at
the AMQP uint maximum, since the value is encoded as one.

:param float or None timeout: The caller's remaining timeout in seconds, or None.
:rtype: int
:returns: The remaining time less the buffer, or the default if no timeout was given.
"""
if timeout is None:
return DEFAULT_SERVER_TIMEOUT_MS
capped = min(timeout, MAX_SERVER_TIMEOUT_MS / 1000)
return max(int(capped * 1000) - SERVER_TIMEOUT_BUFFER_MS, 0)


def build_uri(address, entity):
parsed = urlparse(address)
if parsed.path:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,13 +13,19 @@
from ._transport._pyamqp_transport_async import PyamqpTransportAsync
from .._base_handler import _generate_sas_token, BaseHandler as BaseHandlerSync, _get_backoff_time
from .._common._configuration import Configuration
from .._common.utils import create_properties, strip_protocol_from_uri, parse_sas_credential
from .._common.utils import (
create_properties,
strip_protocol_from_uri,
parse_sas_credential,
get_server_timeout_ms,
)
from .._common.constants import (
TOKEN_TYPE_SASTOKEN,
MGMT_REQUEST_OP_TYPE_ENTITY_MGMT,
ASSOCIATEDLINKPROPERTYNAME,
CONTAINER_PREFIX,
MANAGEMENT_PATH_SUFFIX,
REQUEST_RESPONSE_TIMEOUT,
)
from ..exceptions import (
ServiceBusConnectionError,
Expand Down Expand Up @@ -340,6 +346,10 @@ async def _mgmt_request_response(
except AttributeError:
pass

application_properties[REQUEST_RESPONSE_TIMEOUT] = self._amqp_transport.AMQP_UINT_VALUE(
get_server_timeout_ms(timeout)
)

mgmt_msg = self._amqp_transport.create_mgmt_msg( # type: ignore # TODO: fix mypy
message=message,
application_properties=application_properties,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,248 @@
# Copyright (c) Microsoft Corporation. All rights reserved.

"""Unit tests for the management server-timeout.

Management operations send the `com.microsoft:server-timeout` property, asking the service to
bound the operation on its side. The value is the caller's remaining time less a one second
buffer, so the service answers before the client gives up, or 60 seconds when the caller
supplied no timeout.

Previously no bound was sent at all, so a stalled service could hold a management call until
the AMQP link itself failed. Matches the .NET, Java and Go SDKs.
"""

from unittest.mock import MagicMock

import struct

import pytest

from azure.servicebus._common import mgmt_handlers
from azure.servicebus._common.constants import (
ERROR_CODE_TIMEOUT,
MAX_SERVER_TIMEOUT_MS,
MGMT_RESPONSE_MESSAGE_ERROR_CONDITION,
REQUEST_RESPONSE_TIMEOUT,
)
from azure.servicebus._common.utils import get_server_timeout_ms
from azure.servicebus._transport._pyamqp_transport import PyamqpTransport
from azure.servicebus.exceptions import OperationTimeoutError


class TestServerTimeoutMillis:
"""`get_server_timeout_ms` converts the caller's remaining time into the
value advertised to the service."""

@pytest.mark.parametrize(
"remaining_seconds,expected_ms",
[
(None, 60000), # no caller timeout: the default, not "no bound at all"
(120, 119000),
(60, 59000), # equals the default, but must still take the buffer path
(10, 9000),
(1.5, 500),
(1, 0), # at the buffer, nothing left to give the service
(0.5, 0),
(-5, 0), # deadline already passed: clamped, never a negative uint
],
)
def test_remaining_time_less_buffer(self, remaining_seconds, expected_ms):
assert get_server_timeout_ms(remaining_seconds) == expected_ms

def test_wire_contract(self):
# The key .NET and Java send, encoded as an unsigned int of milliseconds.
assert REQUEST_RESPONSE_TIMEOUT == b"com.microsoft:server-timeout"
encoded = PyamqpTransport.AMQP_UINT_VALUE(get_server_timeout_ms(None))
assert encoded == {"TYPE": "UINT", "VALUE": 60000}

@pytest.mark.parametrize("remaining_seconds", [4294968, 5_000_000, 1e12, float("inf")])
def test_capped_at_the_amqp_uint_maximum(self, remaining_seconds):
# `timeout` is only bounded at zero, so a large value would overflow the uint encoder.
result = get_server_timeout_ms(remaining_seconds)
assert result <= MAX_SERVER_TIMEOUT_MS
struct.pack(">I", result) # raises if it does not fit

def test_just_below_the_cap_is_not_clamped(self):
assert get_server_timeout_ms(4294967) == 4294966000


class TestManagementRequestSetsServerTimeout:
"""The property must actually reach the outgoing management message, on both
the associated-link and no-associated-link paths."""

def _make_handler(self):
from azure.servicebus._base_handler import BaseHandler

captured = {}

def fake_create_mgmt_msg(message, application_properties, config, reply_to, **kwargs):
captured.clear()
captured.update(application_properties)
return MagicMock()

handler = BaseHandler.__new__(BaseHandler)
handler._amqp_transport = MagicMock()
handler._amqp_transport.create_mgmt_msg = fake_create_mgmt_msg
handler._amqp_transport.AMQP_UINT_VALUE = PyamqpTransport.AMQP_UINT_VALUE
handler._amqp_transport.get_handler_link_name = lambda h: "link-1"
handler._amqp_transport.mgmt_client_request = lambda *args, **kwargs: "response"
handler._amqp_transport.TIMEOUT_ERROR = TimeoutError
handler._open = lambda: None
handler._handler = MagicMock()
handler._config = MagicMock(encoding="UTF-8")
handler._mgmt_target = "queue/$management"
return handler, captured

def test_default_sent_when_caller_gave_no_timeout(self):
# The gap this closes: previously no bound was sent at all.
handler, captured = self._make_handler()
handler._mgmt_request_response(b"op", {}, lambda *a: None, timeout=None)
assert captured[REQUEST_RESPONSE_TIMEOUT] == {"TYPE": "UINT", "VALUE": 60000}

def test_remaining_time_less_buffer_sent(self):
handler, captured = self._make_handler()
handler._mgmt_request_response(b"op", {}, lambda *a: None, timeout=10)
assert captured[REQUEST_RESPONSE_TIMEOUT] == {"TYPE": "UINT", "VALUE": 9000}

def test_clamped_below_buffer(self):
handler, captured = self._make_handler()
handler._mgmt_request_response(b"op", {}, lambda *a: None, timeout=0.4)
assert captured[REQUEST_RESPONSE_TIMEOUT] == {"TYPE": "UINT", "VALUE": 0}

def test_associated_link_name_preserved(self):
handler, captured = self._make_handler()
handler._mgmt_request_response(b"op", {}, lambda *a: None, timeout=None)
assert b"associated-link-name" in captured
assert REQUEST_RESPONSE_TIMEOUT in captured

def test_sent_on_calls_without_an_associated_link(self):
# list_sessions passes keep_alive_associated_link=False, starting from an empty map.
handler, captured = self._make_handler()
handler._mgmt_request_response(
b"op",
{},
lambda *a: None,
keep_alive_associated_link=False,
timeout=None,
)
assert captured == {REQUEST_RESPONSE_TIMEOUT: {"TYPE": "UINT", "VALUE": 60000}}


class TestServiceTimeoutResponse:
"""The response half: the service's answer must surface as a retryable
`OperationTimeoutError`, which rests on `com.microsoft:timeout` in `errorCondition`."""

@staticmethod
def _response(condition):
message = MagicMock()
message.application_properties = {MGMT_RESPONSE_MESSAGE_ERROR_CONDITION: condition}
return message

def test_timeout_condition_raises_retryable_operation_timeout_error(self):
with pytest.raises(OperationTimeoutError) as exc_info:
mgmt_handlers.default(408, self._response(ERROR_CODE_TIMEOUT), "The operation timed out.", PyamqpTransport)

assert exc_info.value._retryable is True

def test_success_returns_the_value_untouched(self):
message = self._response(None)
message.value = {"ok": True}
assert mgmt_handlers.default(200, message, None, PyamqpTransport) == {"ok": True}

def test_uamqp_maps_the_same_condition(self):
# Skip on uamqp itself: the transport module imports without it, but UamqpTransport is not defined.
pytest.importorskip("uamqp", reason="uamqp not installed")
from azure.servicebus._transport._uamqp_transport import UamqpTransport

with pytest.raises(OperationTimeoutError):
mgmt_handlers.default(
408,
self._response(ERROR_CODE_TIMEOUT),
"The operation timed out.",
UamqpTransport,
)


class TestUamqpValueType:
"""uamqp is the one place the value type changes: pyamqp uses a plain dict, uamqp an `AMQPuInt`."""

def test_transport_supplied_type_is_what_reaches_the_message(self):

from azure.servicebus._base_handler import BaseHandler

sentinel = object()
captured = {}

def fake_create_mgmt_msg(message, application_properties, config, reply_to, **kwargs):
captured.update(application_properties)
return MagicMock()

handler = BaseHandler.__new__(BaseHandler)
handler._amqp_transport = MagicMock()
handler._amqp_transport.create_mgmt_msg = fake_create_mgmt_msg
handler._amqp_transport.AMQP_UINT_VALUE = lambda ms: sentinel
handler._amqp_transport.get_handler_link_name = lambda h: "link-1"
handler._amqp_transport.mgmt_client_request = lambda *args, **kwargs: "response"
handler._amqp_transport.TIMEOUT_ERROR = TimeoutError
handler._open = lambda: None
handler._handler = MagicMock()
handler._config = MagicMock(encoding="UTF-8")
handler._mgmt_target = "queue/$management"

handler._mgmt_request_response(b"op", {}, lambda *a: None, timeout=None)
assert captured[REQUEST_RESPONSE_TIMEOUT] is sentinel

def test_real_uamqp_type_encodes_in_application_properties(self):
uamqp = pytest.importorskip("uamqp", reason="uamqp not installed")
from azure.servicebus._transport._uamqp_transport import UamqpTransport

value = UamqpTransport.AMQP_UINT_VALUE(get_server_timeout_ms(None))
assert isinstance(value, uamqp.types.AMQPuInt)

message = UamqpTransport.create_mgmt_msg(
message={"operation": "peek"},
application_properties={REQUEST_RESPONSE_TIMEOUT: value},
config=MagicMock(encoding="UTF-8"),
reply_to="queue/$management",
)

assert message.encode_message()


class TestAsyncParity:
"""The async management path is a separate implementation and can drift from the sync one."""

@pytest.mark.asyncio
async def test_async_management_sets_server_timeout(self):
from azure.servicebus.aio._base_handler_async import BaseHandler as AsyncBaseHandler

captured = {}

def fake_create_mgmt_msg(message, application_properties, config, reply_to, **kwargs):
captured.clear()
captured.update(application_properties)
return MagicMock()

async def fake_request(*args, **kwargs):
return "response"

async def fake_open():
return None

handler = AsyncBaseHandler.__new__(AsyncBaseHandler)
handler._amqp_transport = MagicMock()
handler._amqp_transport.create_mgmt_msg = fake_create_mgmt_msg
handler._amqp_transport.AMQP_UINT_VALUE = PyamqpTransport.AMQP_UINT_VALUE
handler._amqp_transport.get_handler_link_name = lambda h: "link-1"
handler._amqp_transport.mgmt_client_request_async = fake_request
handler._amqp_transport.TIMEOUT_ERROR = TimeoutError
handler._open = fake_open
handler._handler = MagicMock()
handler._config = MagicMock(encoding="UTF-8")
handler._mgmt_target = "queue/$management"

await handler._mgmt_request_response(b"op", {}, lambda *a: None, timeout=None)
assert captured[REQUEST_RESPONSE_TIMEOUT] == {"TYPE": "UINT", "VALUE": 60000}

await handler._mgmt_request_response(b"op", {}, lambda *a: None, timeout=10)
assert captured[REQUEST_RESPONSE_TIMEOUT] == {"TYPE": "UINT", "VALUE": 9000}
Loading