-
Notifications
You must be signed in to change notification settings - Fork 0
/
AkkaStreamsTests.cs
124 lines (103 loc) · 4.79 KB
/
AkkaStreamsTests.cs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
using System;
using System.Reactive.Linq;
using System.Threading.Tasks;
using System.Linq;
using Akka.Actor;
using Akka.Streams;
using Akka.Streams.Dsl;
using Akka.TestKit.Xunit2;
using Akka.Streams.TestKit;
using AkkaExchange.Shared.Events;
using AkkaExchange.Shared.Queries;
using AkkaExchange.Utils;
using Moq;
using Reactive.Streams;
using Xunit;
namespace AkkaExchange.Tests.Akka
{
public class AkkaStreamsTests : TestKit
{
[Fact]
public void AkkaStreams_ActorSourceActorSink_Works()
{
using (var materializer = Sys.Materializer())
{
var probe = CreateTestProbe();
var source = Source.ActorRef<HandlerErrorEvent>(10, OverflowStrategy.DropNew);
var sink = Sink.ActorRef<HandlerErrorEvent>(probe, PoisonPill.Instance);
var graph = source.ToMaterialized(sink, Keep.Left);
var actor = graph.Run(materializer);
var msg = new HandlerErrorEvent("", "", HandlerResult.NotHandled(new object()));
actor.Tell(msg, ActorRefs.Nobody);
var actual = probe.ExpectMsg<HandlerErrorEvent>(TimeSpan.FromSeconds(3));
Assert.Equal("Event type System.Object not supported by handler.", actual.Result.Errors.Single());
}
}
[Fact] // Got a Stack Overflow response. See test below for better way of doing this.
public async Task AkkaStreams_ActorSourcePublisherSink_Works()
{
using (var materializer = Sys.Materializer())
{
var probe = CreateTestProbe();
var source = Source.ActorRef<HandlerErrorEvent>(10, OverflowStrategy.DropNew);
var subscriber = new Mock<ISubscriber<HandlerErrorEvent>>();
// See https://stackoverflow.com/questions/48605870/why-isnt-my-akka-net-stream-subscriber-receiving-messages
subscriber
.Setup(s => s.OnSubscribe(It.IsAny<ISubscription>()))
.Callback((ISubscription sub) => sub.Request(1)); // Subscriptions != Observers. Requires back pressure.
var sink = Sink.FromSubscriber<HandlerErrorEvent>(subscriber.Object);
var graph = source.ToMaterialized(sink, Keep.Both);
var (actor, publisher) = graph.Run(materializer);
await Task.Delay(10);
subscriber.Verify(s => s.OnSubscribe(It.IsAny<ISubscription>()));
var evnt = new HandlerErrorEvent("", "", HandlerResult.NotHandled(new object()));
actor.Tell(evnt, ActorRefs.Nobody);
base.AwaitCondition(() =>
{
try
{
subscriber.Verify(s => s.OnNext(It.IsAny<HandlerErrorEvent>()));
return true;
}
catch(MockException)
{
return false;
}
});
}
}
[Fact]
public void AkkaStreams_ActorSourcePublisherSink_UsingStreamsExtensions_Works()
{
using (var materializer = Sys.Materializer())
{
var probe = this.CreateManualSubscriberProbe<HandlerErrorEvent>();
var source = Source.ActorRef<HandlerErrorEvent>(10, OverflowStrategy.DropNew);
var sink = Sink.FromSubscriber<HandlerErrorEvent>(probe);
var graph = source.ToMaterialized(sink, Keep.Both);
var (actor, publisher) = graph.Run(materializer);
var subscription = probe.ExpectSubscription();
subscription.Request(1);
var evnt = new HandlerErrorEvent("", "", HandlerResult.NotHandled(new object()));
actor.Tell(evnt, ActorRefs.Nobody);
probe.ExpectNext(evnt);
}
}
[Fact]
public async Task AkkaStreams_ActorSourceForeachSink_Works()
{
using (var materializer = Sys.Materializer())
{
var source = Source.ActorRef<HandlerErrorEvent>(10, OverflowStrategy.DropNew);
var observer = new Mock<IObserver<HandlerErrorEvent>>();
var sink = Sink.ForEach<HandlerErrorEvent>(e => observer.Object.OnNext(e));
var graph = source.ToMaterialized(sink, Keep.Both);
var (actor, task) = graph.Run(materializer);
var msg = new HandlerErrorEvent("", "", HandlerResult.NotHandled(new object()));
actor.Tell(msg, ActorRefs.Nobody);
await Task.WhenAny(Task.Delay(1000), task);
observer.Verify(o => o.OnNext(It.IsAny<HandlerErrorEvent>()));
}
}
}
}