Skip to content

Commit 847be02

Browse files
committed
refactor(streaming): hide cache start capabilities
1 parent 6c6de74 commit 847be02

8 files changed

Lines changed: 115 additions & 22 deletions

File tree

src/Azure/Orleans.Streaming.EventHubs/Providers/Streams/EventHub/EventHubQueueCache.cs

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -137,8 +137,7 @@ public object GetCursor(StreamId streamId, StreamSequenceToken? sequenceToken)
137137
return cache.GetCursor(streamId, sequenceToken);
138138
}
139139

140-
/// <inheritdoc />
141-
public object GetCursorAtPosition(StreamId streamId, StreamSubscriptionStartPosition startPosition)
140+
object IEventHubQueueCache.GetCursorAtPosition(StreamId streamId, StreamSubscriptionStartPosition startPosition)
142141
{
143142
return cache.GetCursorAtPosition(streamId, startPosition);
144143
}

src/Orleans.Streaming/Common/SimpleCache/SimpleQueueCache.cs

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -149,8 +149,7 @@ public virtual IQueueCacheCursor GetCacheCursor(StreamId streamId, StreamSequenc
149149
return cursor;
150150
}
151151

152-
/// <inheritdoc />
153-
public virtual IQueueCacheCursor GetCacheCursorAtPosition(StreamId streamId, StreamSubscriptionStartPosition startPosition)
152+
IQueueCacheCursor IQueueCache.GetCacheCursorAtPosition(StreamId streamId, StreamSubscriptionStartPosition startPosition)
154153
{
155154
if (startPosition == StreamSubscriptionStartPosition.Latest)
156155
{

src/Orleans.Streaming/Generator/GeneratorPooledCache.cs

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -186,8 +186,7 @@ public IQueueCacheCursor GetCacheCursor(StreamId streamId, StreamSequenceToken?
186186
return new Cursor(cache, streamId, token);
187187
}
188188

189-
/// <inheritdoc />
190-
public IQueueCacheCursor GetCacheCursorAtPosition(StreamId streamId, StreamSubscriptionStartPosition startPosition)
189+
IQueueCacheCursor IQueueCache.GetCacheCursorAtPosition(StreamId streamId, StreamSubscriptionStartPosition startPosition)
191190
{
192191
return new Cursor(cache, cache.GetCursorAtPosition(streamId, startPosition));
193192
}

src/Orleans.Streaming/MemoryStreams/MemoryPooledCache.cs

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -180,8 +180,7 @@ public IQueueCacheCursor GetCacheCursor(StreamId streamId, StreamSequenceToken?
180180
return new Cursor(cache, streamId, token);
181181
}
182182

183-
/// <inheritdoc/>
184-
public IQueueCacheCursor GetCacheCursorAtPosition(StreamId streamId, StreamSubscriptionStartPosition startPosition)
183+
IQueueCacheCursor IQueueCache.GetCacheCursorAtPosition(StreamId streamId, StreamSubscriptionStartPosition startPosition)
185184
{
186185
return new Cursor(cache, cache.GetCursorAtPosition(streamId, startPosition));
187186
}

src/api/Azure/Orleans.Streaming.EventHubs/Orleans.Streaming.EventHubs.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -484,10 +484,10 @@ public void Dispose() { }
484484

485485
public object GetCursor(Runtime.StreamId streamId, Streams.StreamSequenceToken? sequenceToken) { throw null; }
486486

487-
public object GetCursorAtPosition(Runtime.StreamId streamId, Streams.StreamSubscriptionStartPosition startPosition) { throw null; }
488-
489487
public int GetMaxAddCount() { throw null; }
490488

489+
object IEventHubQueueCache.GetCursorAtPosition(Runtime.StreamId streamId, Streams.StreamSubscriptionStartPosition startPosition) { throw null; }
490+
491491
public void Refresh(object cursor, Streams.StreamSequenceToken? sequenceToken) { }
492492

493493
public void SignalPurge() { }

src/api/Orleans.Streaming/Orleans.Streaming.cs

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -449,14 +449,14 @@ public void AddToCache(System.Collections.Generic.IList<Orleans.Streams.IBatchCo
449449

450450
public Orleans.Streams.IQueueCacheCursor GetCacheCursor(Runtime.StreamId streamId, Orleans.Streams.StreamSequenceToken? token) { throw null; }
451451

452-
public Orleans.Streams.IQueueCacheCursor GetCacheCursorAtPosition(Runtime.StreamId streamId, Orleans.Streams.StreamSubscriptionStartPosition startPosition) { throw null; }
453-
454452
public int GetMaxAddCount() { throw null; }
455453

456454
public Orleans.Streams.StreamSequenceToken GetSequenceToken(ref Streams.Common.CachedMessage cachedMessage) { throw null; }
457455

458456
public bool IsUnderPressure() { throw null; }
459457

458+
Orleans.Streams.IQueueCacheCursor Orleans.Streams.IQueueCache.GetCacheCursorAtPosition(Runtime.StreamId streamId, Orleans.Streams.StreamSubscriptionStartPosition startPosition) { throw null; }
459+
460460
public bool TryPurgeFromCache(out System.Collections.Generic.IList<Orleans.Streams.IBatchContainer> purgedItems) { throw null; }
461461
}
462462

@@ -887,12 +887,12 @@ public virtual void AddToCache(System.Collections.Generic.IList<Orleans.Streams.
887887

888888
public virtual Orleans.Streams.IQueueCacheCursor GetCacheCursor(Runtime.StreamId streamId, Orleans.Streams.StreamSequenceToken? token) { throw null; }
889889

890-
public virtual Orleans.Streams.IQueueCacheCursor GetCacheCursorAtPosition(Runtime.StreamId streamId, Orleans.Streams.StreamSubscriptionStartPosition startPosition) { throw null; }
891-
892890
public int GetMaxAddCount() { throw null; }
893891

894892
public virtual bool IsUnderPressure() { throw null; }
895893

894+
Orleans.Streams.IQueueCacheCursor Orleans.Streams.IQueueCache.GetCacheCursorAtPosition(Runtime.StreamId streamId, Orleans.Streams.StreamSubscriptionStartPosition startPosition) { throw null; }
895+
896896
public virtual bool TryPurgeFromCache(out System.Collections.Generic.IList<Orleans.Streams.IBatchContainer> purgedItems) { throw null; }
897897
}
898898

@@ -1010,14 +1010,14 @@ public void AddToCache(System.Collections.Generic.IList<Orleans.Streams.IBatchCo
10101010

10111011
public Orleans.Streams.IQueueCacheCursor GetCacheCursor(Runtime.StreamId streamId, Orleans.Streams.StreamSequenceToken? token) { throw null; }
10121012

1013-
public Orleans.Streams.IQueueCacheCursor GetCacheCursorAtPosition(Runtime.StreamId streamId, Orleans.Streams.StreamSubscriptionStartPosition startPosition) { throw null; }
1014-
10151013
public int GetMaxAddCount() { throw null; }
10161014

10171015
public Orleans.Streams.StreamSequenceToken GetSequenceToken(ref Common.CachedMessage cachedMessage) { throw null; }
10181016

10191017
public bool IsUnderPressure() { throw null; }
10201018

1019+
Orleans.Streams.IQueueCacheCursor Orleans.Streams.IQueueCache.GetCacheCursorAtPosition(Runtime.StreamId streamId, Orleans.Streams.StreamSubscriptionStartPosition startPosition) { throw null; }
1020+
10211021
public bool TryPurgeFromCache(out System.Collections.Generic.IList<Orleans.Streams.IBatchContainer> purgedItems) { throw null; }
10221022
}
10231023

test/Orleans.Streaming.Tests/OrleansRuntime/Streams/SimpleQueueCacheTests.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@ public void EarliestAvailableStartsAtOldestRetainedMessageForStream()
2424
new TestBatchContainer(targetStream, 3),
2525
]);
2626

27-
var cursor = cache.GetCacheCursorAtPosition(targetStream, StreamSubscriptionStartPosition.EarliestAvailable);
27+
var cursor = ((IQueueCache)cache).GetCacheCursorAtPosition(targetStream, StreamSubscriptionStartPosition.EarliestAvailable);
2828

2929
Assert.True(cursor.MoveNext());
3030
Assert.Equal(1, cursor.GetCurrent(out _)!.SequenceToken.SequenceNumber);
@@ -40,7 +40,7 @@ public void EarliestAvailableWaitsForFirstFutureMatchingMessage()
4040
var targetStream = StreamId.Create("namespace", Guid.NewGuid());
4141
var otherStream = StreamId.Create("namespace", Guid.NewGuid());
4242
cache.AddToCache([new TestBatchContainer(otherStream, 100)]);
43-
var cursor = cache.GetCacheCursorAtPosition(targetStream, StreamSubscriptionStartPosition.EarliestAvailable);
43+
var cursor = ((IQueueCache)cache).GetCacheCursorAtPosition(targetStream, StreamSubscriptionStartPosition.EarliestAvailable);
4444

4545
Assert.False(cursor.MoveNext());
4646

test/Orleans.Streaming.Tests/StreamingTests/StreamSubscriptionHandleImplTests.cs

Lines changed: 101 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,36 @@ await Assert.ThrowsAsync<NotSupportedException>(
111111
StreamSubscriptionStartPosition.EarliestAvailable));
112112
}
113113

114+
[Fact]
115+
public async Task ItemSubscribeOverloadForwardsToStreamConsumer()
116+
{
117+
var consumer = new RecordingInternalObservable();
118+
IAsyncObservable<int> stream = CreateStream(isRewindable: true, consumer);
119+
var observer = Substitute.For<IAsyncObserver<int>>();
120+
121+
await stream.SubscribeAsync(
122+
observer,
123+
StreamSubscriptionStartPosition.EarliestAvailable,
124+
filterData: "filter");
125+
126+
Assert.Same(observer, consumer.ItemObserver);
127+
Assert.Equal(StreamSubscriptionStartPosition.EarliestAvailable, consumer.StartPosition);
128+
Assert.Equal("filter", consumer.FilterData);
129+
}
130+
131+
[Fact]
132+
public async Task BatchSubscribeOverloadForwardsToStreamConsumer()
133+
{
134+
var consumer = new RecordingInternalObservable();
135+
IAsyncBatchObservable<int> stream = CreateStream(isRewindable: true, consumer);
136+
var observer = Substitute.For<IAsyncBatchObserver<int>>();
137+
138+
await stream.SubscribeAsync(observer, StreamSubscriptionStartPosition.EarliestAvailable);
139+
140+
Assert.Same(observer, consumer.BatchObserver);
141+
Assert.Equal(StreamSubscriptionStartPosition.EarliestAvailable, consumer.StartPosition);
142+
}
143+
114144
[Fact]
115145
public void ActiveImplicitSubscriptionRejectsOlderAcknowledgedToken()
116146
{
@@ -179,20 +209,87 @@ private static GuidId CreateSubscriptionId(bool implicitSubscription)
179209
return GuidId.GetGuidId(subscriptionGuid);
180210
}
181211

182-
private static StreamImpl<int> CreateStream(bool isRewindable)
212+
private static StreamImpl<int> CreateStream(
213+
bool isRewindable,
214+
IInternalAsyncObservable<int>? consumer = null)
183215
{
184216
return new StreamImpl<int>(
185217
new QualifiedStreamId("provider", StreamId.Create("namespace", Guid.NewGuid())),
186-
new TestStreamProvider(),
218+
new TestStreamProvider(consumer),
187219
isRewindable,
188220
Substitute.For<IRuntimeClient>());
189221
}
190222

191-
private sealed class TestStreamProvider : IInternalStreamProvider
223+
private sealed class TestStreamProvider(IInternalAsyncObservable<int>? consumer) : IInternalStreamProvider
192224
{
193225
public IInternalAsyncBatchObserver<T> GetProducerInterface<T>(IAsyncStream<T> streamId) => null!;
194226

195-
public IInternalAsyncObservable<T> GetConsumerInterface<T>(IAsyncStream<T> streamId) => null!;
227+
public IInternalAsyncObservable<T> GetConsumerInterface<T>(IAsyncStream<T> streamId)
228+
=> (IInternalAsyncObservable<T>)(object)consumer!;
229+
}
230+
231+
private sealed class RecordingInternalObservable : IInternalAsyncObservable<int>
232+
{
233+
public IAsyncObserver<int>? ItemObserver { get; private set; }
234+
public IAsyncBatchObserver<int>? BatchObserver { get; private set; }
235+
public StreamSubscriptionStartPosition StartPosition { get; private set; }
236+
public string? FilterData { get; private set; }
237+
238+
public Task<StreamSubscriptionHandle<int>> SubscribeAsync(IAsyncObserver<int> observer)
239+
=> Task.FromResult<StreamSubscriptionHandle<int>>(null!);
240+
241+
public Task<StreamSubscriptionHandle<int>> SubscribeAsync(
242+
IAsyncObserver<int> observer,
243+
StreamSequenceToken? token,
244+
string? filterData = null)
245+
=> Task.FromResult<StreamSubscriptionHandle<int>>(null!);
246+
247+
public Task<StreamSubscriptionHandle<int>> SubscribeAsync(
248+
IAsyncObserver<int> observer,
249+
StreamSubscriptionStartPosition startPosition,
250+
string? filterData = null)
251+
{
252+
ItemObserver = observer;
253+
StartPosition = startPosition;
254+
FilterData = filterData;
255+
return Task.FromResult<StreamSubscriptionHandle<int>>(null!);
256+
}
257+
258+
public Task<StreamSubscriptionHandle<int>> SubscribeAsync(IAsyncBatchObserver<int> observer)
259+
=> Task.FromResult<StreamSubscriptionHandle<int>>(null!);
260+
261+
public Task<StreamSubscriptionHandle<int>> SubscribeAsync(
262+
IAsyncBatchObserver<int> observer,
263+
StreamSequenceToken? token)
264+
=> Task.FromResult<StreamSubscriptionHandle<int>>(null!);
265+
266+
public Task<StreamSubscriptionHandle<int>> SubscribeAsync(
267+
IAsyncBatchObserver<int> observer,
268+
StreamSubscriptionStartPosition startPosition)
269+
{
270+
BatchObserver = observer;
271+
StartPosition = startPosition;
272+
return Task.FromResult<StreamSubscriptionHandle<int>>(null!);
273+
}
274+
275+
public Task<StreamSubscriptionHandle<int>> ResumeAsync(
276+
StreamSubscriptionHandle<int> handle,
277+
IAsyncObserver<int> observer,
278+
StreamSequenceToken? token = null)
279+
=> Task.FromResult<StreamSubscriptionHandle<int>>(null!);
280+
281+
public Task<StreamSubscriptionHandle<int>> ResumeAsync(
282+
StreamSubscriptionHandle<int> handle,
283+
IAsyncBatchObserver<int> observer,
284+
StreamSequenceToken? token = null)
285+
=> Task.FromResult<StreamSubscriptionHandle<int>>(null!);
286+
287+
public Task UnsubscribeAsync(StreamSubscriptionHandle<int> handle) => Task.CompletedTask;
288+
289+
public Task<IList<StreamSubscriptionHandle<int>>> GetAllSubscriptions()
290+
=> Task.FromResult<IList<StreamSubscriptionHandle<int>>>([]);
291+
292+
public Task Cleanup() => Task.CompletedTask;
196293
}
197294

198295
private sealed class LegacyObservable : IAsyncObservable<int>

0 commit comments

Comments
 (0)