Skip to content

Commit 07ff32b

Browse files
committed
chore(projections): keep maintained coverage independent of TCP
Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent 895af33 commit 07ff32b

15 files changed

Lines changed: 200 additions & 296 deletions

‎src/EventStore.Projections.Core.Tests/EventStore.Projections.Core.Tests.csproj‎

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,8 +8,6 @@
88
<None Remove="Services\checkpoint_strategy\**" />
99
</ItemGroup>
1010
<ItemGroup>
11-
<Compile Remove="Playground\Launchpad.cs" />
12-
<Compile Remove="Playground\Launchpad2.cs" />
1311
<Compile Remove="Playground\Launchpad3.cs" />
1412
<Compile Remove="Playground\LaunchpadBase.cs" />
1513
</ItemGroup>

‎src/EventStore.Projections.Core.Tests/Playground/Launchpad.cs‎

Lines changed: 0 additions & 86 deletions
This file was deleted.

‎src/EventStore.Projections.Core.Tests/Playground/Launchpad2.cs‎

Lines changed: 0 additions & 68 deletions
This file was deleted.

‎src/EventStore.Projections.Core.Tests/Services/SpecificationWithEmittedStreamsTrackerAndDeleter.cs‎

Lines changed: 120 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,66 @@
1+
using System;
2+
using System.Collections.Generic;
3+
using System.IO;
4+
using System.Linq;
5+
using System.Text;
16
using System.Threading.Tasks;
7+
using EventStore.Client.Streams;
28
using EventStore.Core.Helpers;
39
using EventStore.Core.Messages;
4-
using EventStore.Core.Tests.ClientAPI;
10+
using EventStore.Core.Services.Transport.Grpc;
11+
using EventStore.Core.Tests;
12+
using EventStore.Core.Tests.Helpers;
513
using EventStore.Projections.Core.Services;
614
using EventStore.Projections.Core.Services.Processing;
715
using EventStore.Projections.Core.Services.Processing.Emitting;
16+
using Google.Protobuf;
17+
using Grpc.Core;
18+
using Grpc.Net.Client;
19+
using NUnit.Framework;
20+
using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata;
21+
using ReadEvent = EventStore.Client.Streams.ReadResp.Types.ReadEvent;
22+
using StreamsClient = EventStore.Client.Streams.Streams.StreamsClient;
823

924
namespace EventStore.Projections.Core.Tests.Services;
1025

11-
public abstract class SpecificationWithEmittedStreamsTrackerAndDeleter<TLogFormat, TStreamId> : SpecificationWithMiniNode<TLogFormat, TStreamId>
26+
public abstract class SpecificationWithEmittedStreamsTrackerAndDeleter<TLogFormat, TStreamId> : SpecificationWithDirectoryPerTestFixture
1227
{
28+
private GrpcChannel _channel;
29+
protected MiniNode<TLogFormat, TStreamId> _node;
30+
protected StreamsClient _client;
1331
protected IEmittedStreamsTracker _emittedStreamsTracker;
1432
protected IEmittedStreamsDeleter _emittedStreamsDeleter;
1533
protected ProjectionNamesBuilder _projectionNamesBuilder;
1634
protected ClientMessage.ReadStreamEventsForwardCompleted _readCompleted;
1735
protected IODispatcher _ioDispatcher;
1836
protected bool _trackEmittedStreams = true;
1937
protected string _projectionName = "test_projection";
38+
protected virtual TimeSpan Timeout { get; } = TimeSpan.FromMinutes(1);
2039

21-
protected override Task Given()
40+
protected abstract Task When();
41+
42+
[OneTimeSetUp]
43+
public override async Task TestFixtureSetUp()
44+
{
45+
await base.TestFixtureSetUp();
46+
_node = new MiniNode<TLogFormat, TStreamId>(PathName);
47+
await _node.Start();
48+
_channel = GrpcChannel.ForAddress(new UriBuilder { Scheme = Uri.UriSchemeHttps }.Uri,
49+
new GrpcChannelOptions { HttpClient = _node.HttpClient, DisposeHttpClient = false });
50+
_client = new StreamsClient(_channel);
51+
await Given().WithTimeout(Timeout);
52+
await When().WithTimeout(Timeout);
53+
}
54+
55+
[OneTimeTearDown]
56+
public override async Task TestFixtureTearDown()
57+
{
58+
_channel?.Dispose();
59+
await _node.Shutdown();
60+
await base.TestFixtureTearDown();
61+
}
62+
63+
protected virtual Task Given()
2264
{
2365
_ioDispatcher = new IODispatcher(_node.Node.MainQueue, _node.Node.MainQueue, true);
2466
_node.Node.MainBus.Subscribe<ClientMessage.ReadStreamEventsBackwardCompleted>(_ioDispatcher.BackwardReader);
@@ -38,4 +80,79 @@ protected override Task Given()
3880
_projectionNamesBuilder.GetEmittedStreamsCheckpointName());
3981
return Task.CompletedTask;
4082
}
83+
84+
protected async Task AppendEvent(string stream, string eventType, byte[] data)
85+
{
86+
using var call = _client.Append(AdminCallOptions());
87+
await call.RequestStream.WriteAsync(new AppendReq
88+
{
89+
Options = new()
90+
{
91+
Any = new(),
92+
StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) }
93+
}
94+
});
95+
await call.RequestStream.WriteAsync(new AppendReq
96+
{
97+
ProposedMessage = new()
98+
{
99+
Id = Uuid.NewUuid().ToDto(),
100+
Data = ByteString.CopyFrom(data),
101+
CustomMetadata = ByteString.Empty,
102+
Metadata = {
103+
{ GrpcMetadata.Type, eventType },
104+
{ GrpcMetadata.ContentType, GrpcMetadata.ContentTypes.ApplicationJson }
105+
}
106+
}
107+
});
108+
await call.RequestStream.CompleteAsync();
109+
await call.ResponseAsync;
110+
}
111+
112+
protected async Task<ReadEvent[]> ReadEvents(string stream, int count)
113+
{
114+
using var call = _client.Read(new ReadReq
115+
{
116+
Options = new()
117+
{
118+
Stream = new()
119+
{
120+
StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) },
121+
Start = new()
122+
},
123+
ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Forwards,
124+
Count = (ulong)count,
125+
NoFilter = new(),
126+
UuidOption = new() { Structured = new() }
127+
}
128+
}, AdminCallOptions());
129+
var events = new List<ReadEvent>();
130+
while (await call.ResponseStream.MoveNext(default))
131+
if (call.ResponseStream.Current.Event is { } resolvedEvent)
132+
events.Add(resolvedEvent);
133+
return events.ToArray();
134+
}
135+
136+
protected async Task<ReadEvent[]> WaitForEvents(string stream, int count)
137+
{
138+
var deadline = DateTime.UtcNow + Timeout;
139+
ReadEvent[] events;
140+
do
141+
{
142+
events = await ReadEvents(stream, count);
143+
if (events.Length >= count)
144+
return events;
145+
await Task.Delay(50);
146+
} while (DateTime.UtcNow < deadline);
147+
148+
return events;
149+
}
150+
151+
private static CallOptions AdminCallOptions() => new(
152+
credentials: CallCredentials.FromInterceptor((_, metadata) =>
153+
{
154+
metadata.Add("authorization",
155+
$"Basic {Convert.ToBase64String(Encoding.ASCII.GetBytes("admin:changeit"))}");
156+
return Task.CompletedTask;
157+
}));
41158
}
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,6 @@
11
using System;
22
using System.Threading;
33
using System.Threading.Tasks;
4-
using EventStore.ClientAPI;
54
using EventStore.Core.Tests;
65
using EventStore.Projections.Core.Services.Processing;
76
using EventStore.Projections.Core.Services.Processing.Checkpointing;
@@ -18,38 +17,25 @@ public class with_an_existing_emitted_streams_stream<TLogFormat, TStreamId> : Sp
1817
protected ManualResetEvent _resetEvent = new ManualResetEvent(false);
1918
private string _testStreamName = "test_stream";
2019
private ManualResetEvent _eventAppeared = new ManualResetEvent(false);
21-
private EventStore.ClientAPI.SystemData.UserCredentials _credentials;
2220

2321
protected override async Task Given()
2422
{
25-
_credentials = new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit");
2623
_onDeleteStreamCompleted = () => { _resetEvent.Set(); };
2724

2825
await base.Given();
29-
var sub = await _conn.SubscribeToStreamAsync(_projectionNamesBuilder.GetEmittedStreamsName(), true, (s, evnt) =>
30-
{
31-
_eventAppeared.Set();
32-
return Task.CompletedTask;
33-
}, userCredentials: _credentials);
3426

3527
_emittedStreamsTracker.TrackEmittedStream(new EmittedEvent[] {
3628
new EmittedDataEvent(
3729
_testStreamName, Guid.NewGuid(), "type1", true,
3830
"data", null, CheckpointTag.FromPosition(0, 100, 50), null),
3931
});
4032

41-
if (!_eventAppeared.WaitOne(TimeSpan.FromSeconds(5)))
33+
var events = await WaitForEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 1);
34+
if (events.Length != 1)
4235
{
4336
Assert.Fail("Timed out waiting for emitted stream event");
4437
}
45-
46-
sub.Unsubscribe();
47-
48-
var emittedStreamResult =
49-
await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsName(), 0, 1, false,
50-
_credentials);
51-
Assert.AreEqual(1, emittedStreamResult.Events.Length);
52-
Assert.AreEqual(SliceReadStatus.Success, emittedStreamResult.Status);
38+
_eventAppeared.Set();
5339
}
5440

5541
protected override Task When()
@@ -66,25 +52,22 @@ protected override Task When()
6652
[Test]
6753
public async Task should_have_deleted_the_tracked_emitted_stream()
6854
{
69-
var result = await _conn.ReadStreamEventsForwardAsync(_testStreamName, 0, 1, false,
70-
new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit"));
71-
Assert.AreEqual(SliceReadStatus.StreamNotFound, result.Status);
55+
var events = await ReadEvents(_testStreamName, 1);
56+
Assert.AreEqual(0, events.Length);
7257
}
7358

7459

7560
[Test]
7661
public async Task should_have_deleted_the_checkpoint_stream()
7762
{
78-
var result = await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsCheckpointName(),
79-
0, 1, false, new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit"));
80-
Assert.AreEqual(SliceReadStatus.StreamNotFound, result.Status);
63+
var events = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsCheckpointName(), 1);
64+
Assert.AreEqual(0, events.Length);
8165
}
8266

8367
[Test]
8468
public async Task should_have_deleted_the_emitted_streams_stream()
8569
{
86-
var result = await _conn.ReadStreamEventsForwardAsync(_projectionNamesBuilder.GetEmittedStreamsName(), 0, 1,
87-
false, new EventStore.ClientAPI.SystemData.UserCredentials("admin", "changeit"));
88-
Assert.AreEqual(SliceReadStatus.StreamNotFound, result.Status);
70+
var events = await ReadEvents(_projectionNamesBuilder.GetEmittedStreamsName(), 1);
71+
Assert.AreEqual(0, events.Length);
8972
}
9073
}

0 commit comments

Comments
 (0)