From 8ea009c53e0a258319ddd0352e600348d4a09add Mon Sep 17 00:00:00 2001 From: Vlad Raikov Date: Thu, 16 Jul 2026 12:43:01 -0700 Subject: [PATCH 1/3] Fix Cosmos reminder pagination to not stop on empty pages For cross-partition Cosmos queries, empty pages are valid while HasMoreResults is true. Drain FeedIterator until HasMoreResults is false in CosmosReminderTable ReadRows methods. Fixes #10273 --- .../CosmosReminderTable.cs | 16 ++++------------ 1 file changed, 4 insertions(+), 12 deletions(-) diff --git a/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs b/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs index efcd8d694c8..f4511bbcdc4 100644 --- a/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs +++ b/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs @@ -79,18 +79,14 @@ public async Task ReadRows(GrainId grainId) var query = self._container.GetItemLinqQueryable(requestOptions: requestOptions).ToFeedIterator(); var reminders = new List(); - do + while (query.HasMoreResults) { var queryResponse = await query.ReadNextAsync().ConfigureAwait(false); if (queryResponse != null && queryResponse.Count > 0) { reminders.AddRange(queryResponse); } - else - { - break; - } - } while (query.HasMoreResults); + } return reminders; }, @@ -122,18 +118,14 @@ public async Task ReadRows(uint begin, uint end) var iterator = query.ToFeedIterator(); var reminders = new List(); - do + while (iterator.HasMoreResults) { var queryResponse = await iterator.ReadNextAsync().ConfigureAwait(false); if (queryResponse != null && queryResponse.Count > 0) { reminders.AddRange(queryResponse); } - else - { - break; - } - } while (iterator.HasMoreResults); + } return reminders; }, From be8311daeef3c4e60cf0fc7ef5dd11fd470c9e1a Mon Sep 17 00:00:00 2001 From: Vlad Raikov Date: Thu, 16 Jul 2026 14:51:55 -0700 Subject: [PATCH 2/3] Extracting iterator, adding tests --- .../CosmosReminderTable.cs | 52 +------ .../FeedIteratorExtensions.cs | 32 ++++ .../FeedIteratorExtensionsTests.cs | 146 ++++++++++++++++++ 3 files changed, 186 insertions(+), 44 deletions(-) create mode 100644 src/Azure/Orleans.Reminders.Cosmos/FeedIteratorExtensions.cs create mode 100644 test/Extensions/Orleans.Cosmos.Tests/FeedIteratorExtensionsTests.cs diff --git a/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs b/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs index f4511bbcdc4..ffa650e9753 100644 --- a/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs +++ b/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs @@ -73,22 +73,11 @@ public async Task ReadRows(GrainId grainId) { var pk = new PartitionKey(ReminderEntity.ConstructPartitionKey(_clusterOptions.ServiceId, grainId)); var requestOptions = new QueryRequestOptions { PartitionKey = pk }; - var response = await _executor.ExecuteOperation(static async args => + var response = await _executor.ExecuteOperation(static args => { var (self, grainId, requestOptions) = args; - var query = self._container.GetItemLinqQueryable(requestOptions: requestOptions).ToFeedIterator(); - - var reminders = new List(); - while (query.HasMoreResults) - { - var queryResponse = await query.ReadNextAsync().ConfigureAwait(false); - if (queryResponse != null && queryResponse.Count > 0) - { - reminders.AddRange(queryResponse); - } - } - - return reminders; + var iterator = self._container.GetItemLinqQueryable(requestOptions: requestOptions).ToFeedIterator(); + return iterator.DrainAsync(); }, (this, grainId, requestOptions)).ConfigureAwait(false); @@ -106,7 +95,7 @@ public async Task ReadRows(uint begin, uint end) { try { - var response = await _executor.ExecuteOperation(static async args => + var response = await _executor.ExecuteOperation(static args => { var (self, begin, end) = args; var query = self._container.GetItemLinqQueryable() @@ -116,18 +105,7 @@ public async Task ReadRows(uint begin, uint end) ? query.Where(r => r.GrainHash > begin && r.GrainHash <= end) : query.Where(r => r.GrainHash > begin || r.GrainHash <= end); - var iterator = query.ToFeedIterator(); - var reminders = new List(); - while (iterator.HasMoreResults) - { - var queryResponse = await iterator.ReadNextAsync().ConfigureAwait(false); - if (queryResponse != null && queryResponse.Count > 0) - { - reminders.AddRange(queryResponse); - } - } - - return reminders; + return query.ToFeedIterator().DrainAsync(); }, (this, begin, end)).ConfigureAwait(false); @@ -230,26 +208,12 @@ public async Task TestOnlyClearTable() { try { - var entities = await _executor.ExecuteOperation(static async self => + var entities = await _executor.ExecuteOperation(static self => { - var query = self._container.GetItemLinqQueryable() + var iterator = self._container.GetItemLinqQueryable() .Where(entity => entity.ServiceId == self._clusterOptions.ServiceId) .ToFeedIterator(); - var reminders = new List(); - do - { - var queryResponse = await query.ReadNextAsync().ConfigureAwait(false); - if (queryResponse != null && queryResponse.Count > 0) - { - reminders.AddRange(queryResponse); - } - else - { - break; - } - } while (query.HasMoreResults); - - return reminders; + return iterator.DrainAsync(); }, this).ConfigureAwait(false); var deleteTasks = new List(); diff --git a/src/Azure/Orleans.Reminders.Cosmos/FeedIteratorExtensions.cs b/src/Azure/Orleans.Reminders.Cosmos/FeedIteratorExtensions.cs new file mode 100644 index 00000000000..84cb2a42027 --- /dev/null +++ b/src/Azure/Orleans.Reminders.Cosmos/FeedIteratorExtensions.cs @@ -0,0 +1,32 @@ +using System.Threading; + +namespace Orleans.Reminders.Cosmos; + +internal static class FeedIteratorExtensions +{ + /// + /// Fully drains a Cosmos DB , collecting every item across + /// all pages. Empty pages are skipped but do not terminate iteration: may remain true after an empty + /// + /// result (for example when the previous page consumed the RU budget while scanning a + /// partition with no matching items), so iteration must continue until + /// HasMoreResults is false. + /// + public static async Task> DrainAsync( + this FeedIterator iterator, + CancellationToken cancellationToken = default) + { + var items = new List(); + while (iterator.HasMoreResults) + { + var page = await iterator.ReadNextAsync(cancellationToken).ConfigureAwait(false); + if (page is { Count: > 0 }) + { + items.AddRange(page); + } + } + + return items; + } +} diff --git a/test/Extensions/Orleans.Cosmos.Tests/FeedIteratorExtensionsTests.cs b/test/Extensions/Orleans.Cosmos.Tests/FeedIteratorExtensionsTests.cs new file mode 100644 index 00000000000..0bb7fc18943 --- /dev/null +++ b/test/Extensions/Orleans.Cosmos.Tests/FeedIteratorExtensionsTests.cs @@ -0,0 +1,146 @@ +using System; +using System.Collections; +using System.Collections.Generic; +using System.Linq; +using System.Net; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Azure.Cosmos; +using Orleans.Reminders.Cosmos; +using Xunit; + +namespace Tester.Cosmos.Reminders; + +/// +/// Unit tests for . These validate the +/// specific pagination invariant that once tripped a real production bug in +/// CosmosReminderTable.ReadRows: a can return an +/// empty page while remains true +/// (e.g. when the previous page exhausted the RU budget while scanning a partition +/// with no matches). The drain helper must keep iterating past empty pages. +/// +/// A live Cosmos DB (or emulator) cannot reliably reproduce this pattern in a +/// deterministic way, so these tests drive the helper with an in-memory +/// subclass that plays back a scripted sequence of +/// pages including empty ones. +/// +public class FeedIteratorExtensionsTests +{ + [Fact] + public async Task DrainAsync_EmptyPageInMiddle_ContinuesIterating() + { + // Simulates the pathological page layout the fix guards against: results, + // then an empty page while HasMoreResults is still true, then more results. + // A "break on first empty page" drain would drop the trailing rows. + var iterator = new FakeFeedIterator( + new[] { 1, 2, 3 }, + Array.Empty(), + new[] { 4, 5 }); + + var drained = await iterator.DrainAsync(); + + Assert.Equal(new[] { 1, 2, 3, 4, 5 }, drained); + } + + [Fact] + public async Task DrainAsync_LeadingEmptyPage_ContinuesIterating() + { + // First page empty while HasMoreResults is still true. A break-on-empty + // implementation would return zero rows even though matches exist further on. + var iterator = new FakeFeedIterator( + Array.Empty(), + new[] { 1, 2, 3 }); + + var drained = await iterator.DrainAsync(); + + Assert.Equal(new[] { 1, 2, 3 }, drained); + } + + [Fact] + public async Task DrainAsync_AllEmptyPages_ReturnsEmpty() + { + var iterator = new FakeFeedIterator( + Array.Empty(), + Array.Empty(), + Array.Empty()); + + var drained = await iterator.DrainAsync(); + + Assert.Empty(drained); + } + + [Fact] + public async Task DrainAsync_SinglePage_ReturnsAllItems() + { + var iterator = new FakeFeedIterator(new[] { 1, 2, 3, 4, 5 }); + + var drained = await iterator.DrainAsync(); + + Assert.Equal(new[] { 1, 2, 3, 4, 5 }, drained); + } + + [Fact] + public async Task DrainAsync_TrailingEmptyPage_ReturnsAllItems() + { + var iterator = new FakeFeedIterator( + new[] { 1, 2 }, + Array.Empty()); + + var drained = await iterator.DrainAsync(); + + Assert.Equal(new[] { 1, 2 }, drained); + } + + /// + /// A that plays back a scripted list of pages. + /// stays true until every scripted page has + /// been consumed via , so an empty + /// page never terminates iteration on its own. + /// + private sealed class FakeFeedIterator : FeedIterator + { + private readonly Queue> _pages; + + public FakeFeedIterator(params IReadOnlyList[] pages) + { + _pages = new Queue>(pages); + } + + public override bool HasMoreResults => _pages.Count > 0; + + public override Task> ReadNextAsync(CancellationToken cancellationToken = default) + { + if (_pages.Count == 0) + { + throw new InvalidOperationException("ReadNextAsync called after all pages consumed."); + } + + return Task.FromResult>(new FakeFeedResponse(_pages.Dequeue())); + } + } + + /// + /// Minimal stub exposing the members the drain + /// helper actually reads ( and the enumerator). Everything + /// else returns a safe default so nothing throws on unrelated access. + /// + private sealed class FakeFeedResponse : FeedResponse + { + private readonly IReadOnlyList _items; + + public FakeFeedResponse(IReadOnlyList items) => _items = items; + + public override int Count => _items.Count; + public override string ContinuationToken => null!; + public override string IndexMetrics => null!; + public override Headers Headers => null!; + public override IEnumerable Resource => _items; + public override HttpStatusCode StatusCode => HttpStatusCode.OK; + public override double RequestCharge => 0; + public override string ActivityId => string.Empty; + public override string ETag => null!; + public override CosmosDiagnostics Diagnostics => null!; + + public override IEnumerator GetEnumerator() => _items.GetEnumerator(); + } +} From 99eab9c7f0307c278ec6391d9b7c14a9f39316be Mon Sep 17 00:00:00 2001 From: Reuben Bond Date: Mon, 10 Aug 2026 14:21:06 -0700 Subject: [PATCH 3/3] refactor(cosmos): rename feed iterator helper Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 2a7080e2-0260-404d-82db-13416b40a339 --- .../CosmosReminderTable.cs | 6 ++--- .../FeedIteratorExtensions.cs | 2 +- .../FeedIteratorExtensionsTests.cs | 22 +++++++++---------- 3 files changed, 15 insertions(+), 15 deletions(-) diff --git a/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs b/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs index ffa650e9753..0092f9c9302 100644 --- a/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs +++ b/src/Azure/Orleans.Reminders.Cosmos/CosmosReminderTable.cs @@ -77,7 +77,7 @@ public async Task ReadRows(GrainId grainId) { var (self, grainId, requestOptions) = args; var iterator = self._container.GetItemLinqQueryable(requestOptions: requestOptions).ToFeedIterator(); - return iterator.DrainAsync(); + return iterator.ToListAsync(); }, (this, grainId, requestOptions)).ConfigureAwait(false); @@ -105,7 +105,7 @@ public async Task ReadRows(uint begin, uint end) ? query.Where(r => r.GrainHash > begin && r.GrainHash <= end) : query.Where(r => r.GrainHash > begin || r.GrainHash <= end); - return query.ToFeedIterator().DrainAsync(); + return query.ToFeedIterator().ToListAsync(); }, (this, begin, end)).ConfigureAwait(false); @@ -213,7 +213,7 @@ public async Task TestOnlyClearTable() var iterator = self._container.GetItemLinqQueryable() .Where(entity => entity.ServiceId == self._clusterOptions.ServiceId) .ToFeedIterator(); - return iterator.DrainAsync(); + return iterator.ToListAsync(); }, this).ConfigureAwait(false); var deleteTasks = new List(); diff --git a/src/Azure/Orleans.Reminders.Cosmos/FeedIteratorExtensions.cs b/src/Azure/Orleans.Reminders.Cosmos/FeedIteratorExtensions.cs index 84cb2a42027..566b76441e4 100644 --- a/src/Azure/Orleans.Reminders.Cosmos/FeedIteratorExtensions.cs +++ b/src/Azure/Orleans.Reminders.Cosmos/FeedIteratorExtensions.cs @@ -13,7 +13,7 @@ internal static class FeedIteratorExtensions /// partition with no matching items), so iteration must continue until /// HasMoreResults is false. /// - public static async Task> DrainAsync( + public static async Task> ToListAsync( this FeedIterator iterator, CancellationToken cancellationToken = default) { diff --git a/test/Extensions/Orleans.Cosmos.Tests/FeedIteratorExtensionsTests.cs b/test/Extensions/Orleans.Cosmos.Tests/FeedIteratorExtensionsTests.cs index 0bb7fc18943..c901b021b7e 100644 --- a/test/Extensions/Orleans.Cosmos.Tests/FeedIteratorExtensionsTests.cs +++ b/test/Extensions/Orleans.Cosmos.Tests/FeedIteratorExtensionsTests.cs @@ -12,7 +12,7 @@ namespace Tester.Cosmos.Reminders; /// -/// Unit tests for . These validate the +/// Unit tests for . These validate the /// specific pagination invariant that once tripped a real production bug in /// CosmosReminderTable.ReadRows: a can return an /// empty page while remains true @@ -27,7 +27,7 @@ namespace Tester.Cosmos.Reminders; public class FeedIteratorExtensionsTests { [Fact] - public async Task DrainAsync_EmptyPageInMiddle_ContinuesIterating() + public async Task ToListAsync_EmptyPageInMiddle_ContinuesIterating() { // Simulates the pathological page layout the fix guards against: results, // then an empty page while HasMoreResults is still true, then more results. @@ -37,13 +37,13 @@ public async Task DrainAsync_EmptyPageInMiddle_ContinuesIterating() Array.Empty(), new[] { 4, 5 }); - var drained = await iterator.DrainAsync(); + var drained = await iterator.ToListAsync(); Assert.Equal(new[] { 1, 2, 3, 4, 5 }, drained); } [Fact] - public async Task DrainAsync_LeadingEmptyPage_ContinuesIterating() + public async Task ToListAsync_LeadingEmptyPage_ContinuesIterating() { // First page empty while HasMoreResults is still true. A break-on-empty // implementation would return zero rows even though matches exist further on. @@ -51,42 +51,42 @@ public async Task DrainAsync_LeadingEmptyPage_ContinuesIterating() Array.Empty(), new[] { 1, 2, 3 }); - var drained = await iterator.DrainAsync(); + var drained = await iterator.ToListAsync(); Assert.Equal(new[] { 1, 2, 3 }, drained); } [Fact] - public async Task DrainAsync_AllEmptyPages_ReturnsEmpty() + public async Task ToListAsync_AllEmptyPages_ReturnsEmpty() { var iterator = new FakeFeedIterator( Array.Empty(), Array.Empty(), Array.Empty()); - var drained = await iterator.DrainAsync(); + var drained = await iterator.ToListAsync(); Assert.Empty(drained); } [Fact] - public async Task DrainAsync_SinglePage_ReturnsAllItems() + public async Task ToListAsync_SinglePage_ReturnsAllItems() { var iterator = new FakeFeedIterator(new[] { 1, 2, 3, 4, 5 }); - var drained = await iterator.DrainAsync(); + var drained = await iterator.ToListAsync(); Assert.Equal(new[] { 1, 2, 3, 4, 5 }, drained); } [Fact] - public async Task DrainAsync_TrailingEmptyPage_ReturnsAllItems() + public async Task ToListAsync_TrailingEmptyPage_ReturnsAllItems() { var iterator = new FakeFeedIterator( new[] { 1, 2 }, Array.Empty()); - var drained = await iterator.DrainAsync(); + var drained = await iterator.ToListAsync(); Assert.Equal(new[] { 1, 2 }, drained); }