From 60b9dd85084f76ed80918f4289fcc7eb3e1381ad Mon Sep 17 00:00:00 2001 From: Nadia Aina <113453463+Nadia-Adaptive@users.noreply.github.com> Date: Thu, 20 Aug 2026 14:05:26 +0100 Subject: [PATCH] Port AeronArchive RESOURCE_TEMPORARILY_UNAVAILABLE retry --- .../AeronArchiveTest.cs | 219 ++++++++++++++---- src/Adaptive.Archiver/AeronArchive.cs | 62 +++-- 2 files changed, 215 insertions(+), 66 deletions(-) diff --git a/src/Adaptive.Archiver.Tests/AeronArchiveTest.cs b/src/Adaptive.Archiver.Tests/AeronArchiveTest.cs index 30959fe8..2debfed1 100644 --- a/src/Adaptive.Archiver.Tests/AeronArchiveTest.cs +++ b/src/Adaptive.Archiver.Tests/AeronArchiveTest.cs @@ -33,6 +33,7 @@ public class AeronArchiveTest private const string MtuLengthParamName = "mtu"; private const string TermLengthParamName = "term-length"; private const string SparseParamName = "sparse"; + private const long LongMessageTimeoutNs = 60L * 1000 * 1000 * 1000; private AeronType _aeron; private ControlResponsePoller _controlResponsePoller; @@ -63,15 +64,44 @@ public void AsyncConnectShouldConcludeContext() } [Test] - public void AsyncConnectShouldCloseContext() + public void AsyncConnectShouldCloseResourceInCaseOfExceptionUponStartup() { - const string responseChannel = "aeron:udp?endpoint=localhost:1234"; - const int responseStreamId = 49; + var ctx = A.Fake(); + A.CallTo(() => ctx.AeronClient()).Returns(_aeron); + + var aeronContext = A.Fake(); + A.CallTo(() => _aeron.Ctx).Returns(aeronContext); + var nanoClock = A.Fake(); + A.CallTo(() => aeronContext.NanoClock()).Returns(nanoClock); + var error = new InvalidOperationException("TEST"); + A.CallTo(() => nanoClock.NanoTime()).Throws(error); + + var actualException = Assert.Throws(() => AeronArchive.ConnectAsync(ctx)); + Assert.AreSame(error, actualException); + + A.CallTo(() => nanoClock.NanoTime()) + .MustHaveHappened() + .Then(A.CallTo(() => ctx.Dispose()).MustHaveHappened()); + } + + [Test] + public void ShouldReleasePendingSubscriptionWhenClosedBeforeConnected() + { + const string responseChannel = "aeron:udp?endpoint=localhost:0"; + const int responseStreamId = 42; + const long registrationId = 43L; + var ctx = A.Fake(); A.CallTo(() => ctx.AeronClient()).Returns(_aeron); A.CallTo(() => ctx.ControlResponseChannel()).Returns(responseChannel); A.CallTo(() => ctx.ControlResponseStreamId()).Returns(responseStreamId); - var error = new InvalidOperationException("subscription"); + A.CallTo(() => ctx.MessageTimeoutNs()).Returns(LongMessageTimeoutNs); + A.CallTo(() => ctx.OwnsAeronClient()).Returns(false); + + var aeronContext = A.Fake(); + A.CallTo(() => aeronContext.NanoClock()).Returns(SystemNanoClock.INSTANCE); + A.CallTo(() => _aeron.Ctx).Returns(aeronContext); + A.CallTo(() => _aeron.AsyncAddSubscription( responseChannel, @@ -80,38 +110,29 @@ public void AsyncConnectShouldCloseContext() A._ ) ) - .Throws(error); + .Returns(registrationId); + A.CallTo(() => _aeron.GetSubscription(registrationId)).Returns((Subscription)null); - var actualException = Assert.Throws(() => AeronArchive.ConnectAsync(ctx)); - Assert.AreSame(error, actualException); + using (var asyncConnect = AeronArchive.ConnectAsync(ctx)) + { + Assert.IsNull(asyncConnect.Poll()); + Assert.AreEqual( + AeronArchive.AsyncConnect.AsyncConnectState.AWAIT_SUBSCRIPTION, + asyncConnect.State()); + } - A.CallTo(() => ctx.Conclude()) - .MustHaveHappened() - .Then(A.CallTo(() => ctx.AeronClient()).MustHaveHappened()) - .Then(A.CallTo(() => ctx.ControlResponseChannel()).MustHaveHappened()) - .Then(A.CallTo(() => ctx.ControlResponseStreamId()).MustHaveHappened()) - .Then( - A.CallTo(() => - _aeron.AsyncAddSubscription( - responseChannel, - responseStreamId, - A._, - A._ - ) - ) - .MustHaveHappened() - ) - .Then(A.CallTo(() => ctx.Dispose()).MustHaveHappened()); + A.CallTo(() => _aeron.AsyncRemoveSubscription(registrationId)).MustHaveHappenedOnceExactly(); + A.CallTo(() => ctx.Dispose()).MustHaveHappenedOnceExactly(); } [Test] - public void AsyncConnectShouldCloseResourceInCaseOfExceptionUponStartup() + public void ShouldReleasePendingPublicationWhenClosedBeforeConnected() { const string responseChannel = "aeron:udp?endpoint=localhost:0"; - const int responseStreamId = 49; + const int responseStreamId = 42; const string requestChannel = "aeron:udp?endpoint=localhost:1234"; - const int requestStreamId = -15; - const long subscriptionId = -3275938475934759L; + const int requestStreamId = 43; + const long pubRegistrationId = 45L; var ctx = A.Fake(); A.CallTo(() => ctx.AeronClient()).Returns(_aeron); @@ -119,6 +140,14 @@ public void AsyncConnectShouldCloseResourceInCaseOfExceptionUponStartup() A.CallTo(() => ctx.ControlResponseStreamId()).Returns(responseStreamId); A.CallTo(() => ctx.ControlRequestChannel()).Returns(requestChannel); A.CallTo(() => ctx.ControlRequestStreamId()).Returns(requestStreamId); + A.CallTo(() => ctx.MessageTimeoutNs()).Returns(LongMessageTimeoutNs); + A.CallTo(() => ctx.OwnsAeronClient()).Returns(false); + A.CallTo(() => ctx.ErrorHandler()).Returns(_errorHandler); + + var aeronContext = A.Fake(); + A.CallTo(() => aeronContext.NanoClock()).Returns(SystemNanoClock.INSTANCE); + A.CallTo(() => _aeron.Ctx).Returns(aeronContext); + A.CallTo(() => _aeron.AsyncAddSubscription( responseChannel, @@ -127,31 +156,121 @@ public void AsyncConnectShouldCloseResourceInCaseOfExceptionUponStartup() A._ ) ) - .Returns(subscriptionId); - var error = new IndexOutOfRangeException("exception"); - A.CallTo(() => _aeron.Ctx).Throws(error); + .Returns(44L); + var subscription = A.Fake(); + A.CallTo(() => _aeron.GetSubscription(A._)).Returns(subscription); + A.CallTo(() => _aeron.AsyncAddExclusivePublication(requestChannel, requestStreamId)) + .Returns(pubRegistrationId); + A.CallTo(() => _aeron.GetExclusivePublication(A._)).Returns((ExclusivePublication)null); - var actualException = Assert.Throws(() => AeronArchive.ConnectAsync(ctx)); - Assert.AreSame(error, actualException); + using (var asyncConnect = AeronArchive.ConnectAsync(ctx)) + { + Assert.AreEqual( + AeronArchive.AsyncConnect.AsyncConnectState.AWAIT_SUBSCRIPTION, + asyncConnect.State()); + Assert.IsNull(asyncConnect.Poll()); + Assert.AreEqual( + AeronArchive.AsyncConnect.AsyncConnectState.ADD_PUBLICATION, + asyncConnect.State()); + Assert.IsNull(asyncConnect.Poll()); + Assert.AreEqual( + AeronArchive.AsyncConnect.AsyncConnectState.ADD_PUBLICATION, + asyncConnect.State()); + } - A.CallTo(() => ctx.Conclude()) - .MustHaveHappened() - .Then(A.CallTo(() => ctx.AeronClient()).MustHaveHappened()) - .Then(A.CallTo(() => ctx.ControlResponseChannel()).MustHaveHappened()) - .Then(A.CallTo(() => ctx.ControlResponseStreamId()).MustHaveHappened()) - .Then( - A.CallTo(() => - _aeron.AsyncAddSubscription( - responseChannel, - responseStreamId, - A._, - A._ - ) - ) - .MustHaveHappened() + A.CallTo(() => subscription.Dispose()).MustHaveHappenedOnceExactly(); + A.CallTo(() => _aeron.AsyncRemovePublication(pubRegistrationId)).MustHaveHappenedOnceExactly(); + A.CallTo(() => ctx.Dispose()).MustHaveHappenedOnceExactly(); + } + + [Test] + public void ShouldRetryAddingResourcesWhenResourceTemporarilyUnavailable() + { + const string responseChannel = "aeron:udp?endpoint=localhost:0"; + const int responseStreamId = 42; + const string requestChannel = "aeron:udp?endpoint=localhost:1234"; + const int requestStreamId = 43; + + var ctx = A.Fake(); + A.CallTo(() => ctx.AeronClient()).Returns(_aeron); + A.CallTo(() => ctx.ControlResponseChannel()).Returns(responseChannel); + A.CallTo(() => ctx.ControlResponseStreamId()).Returns(responseStreamId); + A.CallTo(() => ctx.ControlRequestChannel()).Returns(requestChannel); + A.CallTo(() => ctx.ControlRequestStreamId()).Returns(requestStreamId); + A.CallTo(() => ctx.MessageTimeoutNs()).Returns(LongMessageTimeoutNs); + A.CallTo(() => ctx.OwnsAeronClient()).Returns(false); + A.CallTo(() => ctx.ErrorHandler()).Returns(_errorHandler); + + var aeronContext = A.Fake(); + A.CallTo(() => aeronContext.NanoClock()).Returns(SystemNanoClock.INSTANCE); + A.CallTo(() => _aeron.Ctx).Returns(aeronContext); + + var resourceUnavailable = new RegistrationException( + 1, + (int)ErrorCode.RESOURCE_TEMPORARILY_UNAVAILABLE, + ErrorCode.RESOURCE_TEMPORARILY_UNAVAILABLE, + "WARN - resource temporarily unavailable, please retry"); + + A.CallTo(() => + _aeron.AsyncAddSubscription( + responseChannel, + responseStreamId, + A._, + A._ + ) ) - .Then(A.CallTo(() => _aeron.AsyncRemoveSubscription(subscriptionId)).MustHaveHappened()) - .Then(A.CallTo(() => ctx.Dispose()).MustHaveHappened()); + .ReturnsNextFromSequence(1L, 44L); + var subscription = A.Fake(); + A.CallTo(() => _aeron.GetSubscription(A._)) + .Throws(resourceUnavailable).Once() + .Then.Returns(subscription); + A.CallTo(() => _aeron.AsyncAddExclusivePublication(requestChannel, requestStreamId)) + .ReturnsNextFromSequence(3L, 45L); + var publication = A.Fake(); + A.CallTo(() => _aeron.GetExclusivePublication(A._)) + .Throws(resourceUnavailable).Once() + .Then.Returns(publication); + + using (var asyncConnect = AeronArchive.ConnectAsync(ctx)) + { + Assert.AreEqual( + AeronArchive.AsyncConnect.AsyncConnectState.AWAIT_SUBSCRIPTION, + asyncConnect.State()); + Assert.IsNull(asyncConnect.Poll()); + Assert.AreEqual( + AeronArchive.AsyncConnect.AsyncConnectState.AWAIT_SUBSCRIPTION, + asyncConnect.State()); + Assert.IsNull(asyncConnect.Poll()); + Assert.AreEqual( + AeronArchive.AsyncConnect.AsyncConnectState.ADD_PUBLICATION, + asyncConnect.State()); + Assert.IsNull(asyncConnect.Poll()); + Assert.AreEqual( + AeronArchive.AsyncConnect.AsyncConnectState.ADD_PUBLICATION, + asyncConnect.State()); + Assert.IsNull(asyncConnect.Poll()); + Assert.AreEqual( + AeronArchive.AsyncConnect.AsyncConnectState.AWAIT_PUBLICATION_CONNECTED, + asyncConnect.State()); + + A.CallTo(() => + _aeron.AsyncAddSubscription( + responseChannel, + responseStreamId, + A._, + A._ + ) + ) + .MustHaveHappenedTwiceExactly(); + A.CallTo(() => _aeron.GetSubscription(A._)).MustHaveHappenedTwiceExactly(); + A.CallTo(() => _aeron.AsyncAddExclusivePublication(requestChannel, requestStreamId)) + .MustHaveHappenedTwiceExactly(); + A.CallTo(() => _aeron.GetExclusivePublication(A._)).MustHaveHappenedTwiceExactly(); + } + + A.CallTo(() => subscription.Dispose()).MustHaveHappenedOnceExactly(); + A.CallTo(() => publication.Dispose()).MustHaveHappenedOnceExactly(); + A.CallTo(() => ctx.Dispose()).MustHaveHappenedOnceExactly(); } [Test] diff --git a/src/Adaptive.Archiver/AeronArchive.cs b/src/Adaptive.Archiver/AeronArchive.cs index 7fd95b47..e945c14c 100644 --- a/src/Adaptive.Archiver/AeronArchive.cs +++ b/src/Adaptive.Archiver/AeronArchive.cs @@ -3921,20 +3921,6 @@ internal AsyncConnect(Context ctx) Aeron.Aeron aeron = ctx.AeronClient(); - _subscriptionRegistrationId = aeron.AsyncAddSubscription( - ctx.ControlResponseChannel(), - ctx.ControlResponseStreamId(), - null, - (image) => - { - AeronArchive client = _aeronArchive; - if (null != client) - { - client.State(ArchiveState.DISCONNECTED); - } - } - ); - _state = AsyncConnectState.AWAIT_SUBSCRIPTION; _deadlineNs = aeron.Ctx.NanoClock().NanoTime() + ctx.MessageTimeoutNs(); } @@ -4078,7 +4064,20 @@ private void AddPublication() ); } - ExclusivePublication publication = aeron.GetExclusivePublication(_publicationRegistrationId); + ExclusivePublication publication = null; + try + { + publication = aeron.GetExclusivePublication(_publicationRegistrationId); + } + catch (RegistrationException exception) + { + _publicationRegistrationId = Aeron.Aeron.NULL_VALUE; + if (ErrorCode.RESOURCE_TEMPORARILY_UNAVAILABLE != exception.ErrorCode()) + { + throw; + } + } + if (null != publication) { string clientInfo = "name=" + _ctx.ClientName(); @@ -4229,7 +4228,38 @@ private void State(AsyncConnectState newState) private void AwaitSubscription() { - Subscription subscription = _ctx.AeronClient().GetSubscription(_subscriptionRegistrationId); + Aeron.Aeron aeron = _ctx.AeronClient(); + if (Aeron.Aeron.NULL_VALUE == _subscriptionRegistrationId) + { + _subscriptionRegistrationId = aeron.AsyncAddSubscription( + _ctx.ControlResponseChannel(), + _ctx.ControlResponseStreamId(), + null, + (image) => + { + AeronArchive client = _aeronArchive; + if (null != client) + { + client.State(ArchiveState.DISCONNECTED); + } + } + ); + } + + Subscription subscription = null; + try + { + subscription = aeron.GetSubscription(_subscriptionRegistrationId); + } + catch (RegistrationException exception) + { + _subscriptionRegistrationId = Aeron.Aeron.NULL_VALUE; + if (ErrorCode.RESOURCE_TEMPORARILY_UNAVAILABLE != exception.ErrorCode()) + { + throw; + } + } + if (null != subscription) { CheckAndSetupResponseChannel(_ctx, subscription);