1+ using System ;
2+ using System . Collections . Generic ;
3+ using System . IO ;
4+ using System . Linq ;
5+ using System . Text ;
16using System . Threading . Tasks ;
7+ using EventStore . Client . Streams ;
28using EventStore . Core . Helpers ;
39using 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 ;
513using EventStore . Projections . Core . Services ;
614using EventStore . Projections . Core . Services . Processing ;
715using 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
924namespace 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}
0 commit comments