Skip to content
Merged
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
41 changes: 29 additions & 12 deletions src/core/Akka.Remote.Tests.MultiNode/RemoteDeliverySpec.cs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading.Tasks;
using Akka.Actor;
using Akka.Remote.TestKit;
using Akka.Configuration;
Expand All @@ -26,6 +27,7 @@ public RemoteDeliveryMultiNetSpec()

CommonConfig = DebugConfig(true)
.WithFallback(ConfigurationFactory.ParseString(@"
# Classic (DotNetty) only; Artery has no batching switch and ignores this key.
akka.remote.dot-netty.tcp.batching.enabled = false # disable batching
"));
}
Expand Down Expand Up @@ -62,7 +64,7 @@ protected override void OnReceive(object message)
public class RemoteDeliverySpec : MultiNodeSpec
{
private readonly RemoteDeliveryMultiNetSpec _config;
private readonly Func<RoleName, string, IActorRef> _identify;
private readonly Func<RoleName, string, Task<IActorRef>> _identify;

public RemoteDeliverySpec() : this(new RemoteDeliveryMultiNetSpec())
{
Expand All @@ -72,11 +74,13 @@ protected RemoteDeliverySpec(RemoteDeliveryMultiNetSpec config) : base(config, t
{
_config = config;

_identify = (role, actorName) => Within(TimeSpan.FromSeconds(10), () =>
_identify = (role, actorName) => WithinAsync(TimeSpan.FromSeconds(10), async () =>
{
Sys.ActorSelection(Node(role)/"user"/actorName)
// NodeAsync, not Node: Node blocks on the controller round trip, which would
// run before the block hands WithinAsync a task and so outside the 10s bound.
Sys.ActorSelection(await NodeAsync(role)/"user"/actorName)
.Tell(new Identify(actorName));
return ExpectMsg<ActorIdentity>()
return (await ExpectMsgAsync<ActorIdentity>())
.Subject;
});
}
Expand All @@ -90,16 +94,16 @@ protected override int InitialParticipantsValueFactory
}

[MultiNodeFact]
public void Remoting_with_TCP_must_not_drop_messages_under_normal_circumstances()
public async Task Remoting_with_TCP_must_not_drop_messages_under_normal_circumstances()
{
Sys.ActorOf<RemoteDeliveryMultiNetSpec.Postman>("postman-" + Myself.Name);
EnterBarrier("actors-started");
await EnterBarrierAsync("actors-started");

RunOn(() =>
await RunOnAsync(async () =>
{
var p1 = _identify(_config.First, "postman-first");
var p2 = _identify(_config.Second, "postman-second");
var p3 = _identify(_config.Third, "postman-third");
var p1 = await _identify(_config.First, "postman-first");
var p2 = await _identify(_config.Second, "postman-second");
var p3 = await _identify(_config.Third, "postman-third");
var route = new List<IActorRef>
{
p2,
Expand All @@ -113,7 +117,7 @@ public void Remoting_with_TCP_must_not_drop_messages_under_normal_circumstances(
{
p1.Tell(new RemoteDeliveryMultiNetSpec.Letter(n, route));
var letterNumber = n;
ExpectMsg<RemoteDeliveryMultiNetSpec.Letter>(
await ExpectMsgAsync<RemoteDeliveryMultiNetSpec.Letter>(
letter => letter.N == letterNumber && letter.Route.Count == 0,
TimeSpan.FromSeconds(5));

Expand All @@ -126,7 +130,20 @@ public void Remoting_with_TCP_must_not_drop_messages_under_normal_circumstances(
},
_config.First);

EnterBarrier("after-1");
// 'second' and 'third' reach this barrier about a second into the run, while 'first'
// is still looping, and BarrierCoordinator arms a barrier's clock on the FIRST
// arrival and only ever shortens it. So this budget, not the 5s per-letter wait, is
// what bounds the whole 500-letter loop, and at the 30s default the loop outlives
// its own barrier on a saturated CI agent. EnterBarrierAsync asks the coordinator
// for RemainingOr(barrier-timeout), so a Within around it sets that budget without
// raising akka.testconductor.barrier-timeout, which also drives the node teardown
// poll in MultiNodeSpecAfterAll (AwaitCondition polls at max/10) and would add
// barrier-timeout/10 of wall clock to every run. 300s over 500 letters permits a
// 600ms mean round trip: 3x the slowest mean measured on such an agent, and far
// under the 5s per-letter deadline, which stays the assertion that reports a real
// drop. A node that fails still aborts the barrier at once, because losing a client
// fails a barrier already in progress.
await WithinAsync(TimeSpan.FromSeconds(300), () => EnterBarrierAsync("after-1"));

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

BIg time budget to allow the other remaining node to finish the loop. If this doesn't work, we'll remove in the future and re-assess how to address this problem.

}
}
}
Loading