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
219 changes: 169 additions & 50 deletions src/Adaptive.Archiver.Tests/AeronArchiveTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Context>();
A.CallTo(() => ctx.AeronClient()).Returns(_aeron);

var aeronContext = A.Fake<AeronType.Context>();
A.CallTo(() => _aeron.Ctx).Returns(aeronContext);
var nanoClock = A.Fake<INanoClock>();
A.CallTo(() => aeronContext.NanoClock()).Returns(nanoClock);
var error = new InvalidOperationException("TEST");
A.CallTo(() => nanoClock.NanoTime()).Throws(error);

var actualException = Assert.Throws<InvalidOperationException>(() => 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<Context>();
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<AeronType.Context>();
A.CallTo(() => aeronContext.NanoClock()).Returns(SystemNanoClock.INSTANCE);
A.CallTo(() => _aeron.Ctx).Returns(aeronContext);

A.CallTo(() =>
_aeron.AsyncAddSubscription(
responseChannel,
Expand All @@ -80,45 +110,44 @@ public void AsyncConnectShouldCloseContext()
A<UnavailableImageHandler>._
)
)
.Throws(error);
.Returns(registrationId);
A.CallTo(() => _aeron.GetSubscription(registrationId)).Returns((Subscription)null);

var actualException = Assert.Throws<InvalidOperationException>(() => 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<AvailableImageHandler>._,
A<UnavailableImageHandler>._
)
)
.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<Context>();
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<AeronType.Context>();
A.CallTo(() => aeronContext.NanoClock()).Returns(SystemNanoClock.INSTANCE);
A.CallTo(() => _aeron.Ctx).Returns(aeronContext);

A.CallTo(() =>
_aeron.AsyncAddSubscription(
responseChannel,
Expand All @@ -127,31 +156,121 @@ public void AsyncConnectShouldCloseResourceInCaseOfExceptionUponStartup()
A<UnavailableImageHandler>._
)
)
.Returns(subscriptionId);
var error = new IndexOutOfRangeException("exception");
A.CallTo(() => _aeron.Ctx).Throws(error);
.Returns(44L);
var subscription = A.Fake<Subscription>();
A.CallTo(() => _aeron.GetSubscription(A<long>._)).Returns(subscription);
A.CallTo(() => _aeron.AsyncAddExclusivePublication(requestChannel, requestStreamId))
.Returns(pubRegistrationId);
A.CallTo(() => _aeron.GetExclusivePublication(A<long>._)).Returns((ExclusivePublication)null);

var actualException = Assert.Throws<IndexOutOfRangeException>(() => 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<AvailableImageHandler>._,
A<UnavailableImageHandler>._
)
)
.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<Context>();
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<AeronType.Context>();
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<AvailableImageHandler>._,
A<UnavailableImageHandler>._
)
)
.Then(A.CallTo(() => _aeron.AsyncRemoveSubscription(subscriptionId)).MustHaveHappened())
.Then(A.CallTo(() => ctx.Dispose()).MustHaveHappened());
.ReturnsNextFromSequence(1L, 44L);
var subscription = A.Fake<Subscription>();
A.CallTo(() => _aeron.GetSubscription(A<long>._))
.Throws(resourceUnavailable).Once()
.Then.Returns(subscription);
A.CallTo(() => _aeron.AsyncAddExclusivePublication(requestChannel, requestStreamId))
.ReturnsNextFromSequence(3L, 45L);
var publication = A.Fake<ExclusivePublication>();
A.CallTo(() => _aeron.GetExclusivePublication(A<long>._))
.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<AvailableImageHandler>._,
A<UnavailableImageHandler>._
)
)
.MustHaveHappenedTwiceExactly();
A.CallTo(() => _aeron.GetSubscription(A<long>._)).MustHaveHappenedTwiceExactly();
A.CallTo(() => _aeron.AsyncAddExclusivePublication(requestChannel, requestStreamId))
.MustHaveHappenedTwiceExactly();
A.CallTo(() => _aeron.GetExclusivePublication(A<long>._)).MustHaveHappenedTwiceExactly();
}

A.CallTo(() => subscription.Dispose()).MustHaveHappenedOnceExactly();
A.CallTo(() => publication.Dispose()).MustHaveHappenedOnceExactly();
A.CallTo(() => ctx.Dispose()).MustHaveHappenedOnceExactly();
}

[Test]
Expand Down
62 changes: 46 additions & 16 deletions src/Adaptive.Archiver/AeronArchive.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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);
Expand Down
Loading