diff --git a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/implementation/ServiceBusReceiveLinkProcessor.java b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/implementation/ServiceBusReceiveLinkProcessor.java index 47aa31ebcd4b..4d7f18b4902d 100644 --- a/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/implementation/ServiceBusReceiveLinkProcessor.java +++ b/sdk/servicebus/azure-messaging-servicebus/src/main/java/com/azure/messaging/servicebus/implementation/ServiceBusReceiveLinkProcessor.java @@ -5,6 +5,8 @@ import com.azure.core.amqp.AmqpEndpointState; import com.azure.core.amqp.AmqpRetryPolicy; +import com.azure.core.amqp.exception.AmqpErrorCondition; +import com.azure.core.amqp.exception.AmqpException; import com.azure.core.amqp.implementation.AmqpReceiveLink; import com.azure.core.util.AsyncCloseable; import com.azure.core.util.logging.ClientLogger; @@ -115,7 +117,15 @@ public Mono updateDisposition(String lockToken, DeliveryState deliveryStat "lockToken[%s]. state[%s]. Cannot update disposition with no link.", lockToken, deliveryState))); } - return link.updateDisposition(lockToken, deliveryState); + return link.updateDisposition(lockToken, deliveryState).onErrorResume(error -> { + if (error instanceof AmqpException) { + AmqpException amqpException = (AmqpException) error; + if (AmqpErrorCondition.TIMEOUT_ERROR.equals(amqpException.getErrorCondition())) { + return link.closeAsync().then(Mono.error(error)); + } + } + return Mono.error(error); + }); } /**