Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ namespace Orleans.Journaling;

internal sealed class AzureTableJournalStorageProvider : ILifecycleParticipant<ISiloLifecycle>, IJournalStorageProvider, IJournalStorageCatalog
{
private const int MaximumDefaultJournalIdLength = 512;
private static readonly string[] JournalIdSelect = [AzureTableJournalStorage.JournalIdPropertyName];
private static readonly string[] JournalMetadataSelect =
[
Expand Down Expand Up @@ -74,7 +75,8 @@ public async IAsyncEnumerable<JournalCatalogEntry> ListAsync(
if (range.IsEmpty
|| (_options.UsesDefaultPartitionKey
&& range.Prefix is { } prefix
&& prefix.AsSpan().IndexOfAnyExceptInRange(' ', '~') >= 0))
&& (prefix.Length > MaximumDefaultJournalIdLength
|| prefix.AsSpan().IndexOfAnyExceptInRange(' ', '~') >= 0)))
{
yield break;
}
Expand Down Expand Up @@ -114,7 +116,10 @@ private string GetCatalogFilter(JournalCatalogRange range)
if (range.LowerBound is { } lowerBound)
{
var lowerKey = GetPartitionKeyBound(lowerBound);
filter += TableClient.CreateQueryFilter($" and PartitionKey ge {lowerKey}");
// A longer bound sorts after its longest possible stored prefix.
filter += lowerBound.Length > MaximumDefaultJournalIdLength
? TableClient.CreateQueryFilter($" and PartitionKey gt {lowerKey}")
: TableClient.CreateQueryFilter($" and PartitionKey ge {lowerKey}");
}

if (range.MaxId is { } maxId)
Expand All @@ -125,9 +130,11 @@ private string GetCatalogFilter(JournalCatalogRange range)

if (range.Prefix is { } prefix)
{
// Encoded keys contain only 0..F, so G bounds every suffix of the encoded prefix.
var prefixEnd = AzureTableJournalStorageOptions.EncodePartitionKey(prefix) + "G";
filter += TableClient.CreateQueryFilter($" and PartitionKey lt {prefixEnd}");
var prefixKey = AzureTableJournalStorageOptions.EncodePartitionKey(prefix);
// At the length limit only the exact id can match; shorter prefixes can have descendants.
filter += prefix.Length == MaximumDefaultJournalIdLength
? TableClient.CreateQueryFilter($" and PartitionKey le {prefixKey}")
: TableClient.CreateQueryFilter($" and PartitionKey lt {prefixKey + "G"}");
}
}
else
Expand All @@ -151,6 +158,7 @@ private string GetCatalogFilter(JournalCatalogRange range)

private static string GetPartitionKeyBound(string value)
{
value = value[..Math.Min(value.Length, MaximumDefaultJournalIdLength)];
var index = value.AsSpan().IndexOfAnyExceptInRange(' ', '~');
if (index < 0)
{
Expand Down
2 changes: 1 addition & 1 deletion src/Azure/Orleans.Journaling.AzureStorage/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ ASCII bounds narrow traversal using `GetBlobsOptions.StartFrom` and an upper cut

Shared prefixes keep timestamp range scans narrow. Ranges crossing punctuation or directory boundaries can transfer additional candidates. Upper-bound termination happens while consuming results: the provider reads the page crossing the conservative upper bound and then completes the traversal. Empty intersections issue no request. The same range planning applies to both account types using the configured Blob client and ordinary listing requests.

Table enumeration requests up to 1000 header rows per service page and reads identities from the canonical `JournalId` header property. The default partition mapping accepts printable ASCII journal ids (`0x20` through `0x7E`) and encodes each byte as two uppercase hexadecimal digits, preserving ordinal ordering and raw prefixes, including partial segments. The 1,024-character partition-key limit therefore permits at most 512 characters per journal id. Validation occurs when creating storage. This restriction belongs to the default Table mapping; custom mappings and other providers retain their journal-id contracts. Default-mapping queries combine the header row key with indexed partition-key constraints for the raw prefix, inclusive `MinId`, and inclusive `MaxId`. Query bounds retain their original ordinal meaning, including bounds outside the stored ASCII alphabet. An empty intersection issues no query. Custom mappings apply well-formed Unicode bounds to the canonical `JournalId` property; bounds containing unpaired surrogates are enforced locally. This preserves ordinal range semantics and can require a full table scan. Unbounded enumeration also scans headers across the table. Constraints are checked again before yielding, and one enumerator advance can cross multiple empty or filtered pages.
Table enumeration requests up to 1000 header rows per service page and reads identities from the canonical `JournalId` header property. The default partition mapping accepts printable ASCII journal ids (`0x20` through `0x7E`) and encodes each byte as two uppercase hexadecimal digits, preserving ordinal ordering and raw prefixes, including partial segments. The 1,024-character partition-key limit therefore permits at most 512 characters per journal id. Validation occurs when creating storage. This restriction belongs to the default Table mapping; custom mappings and other providers retain their journal-id contracts. Default-mapping queries combine the header row key with indexed partition-key constraints for the raw prefix, inclusive `MinId`, and inclusive `MaxId`. Query bounds retain their original ordinal meaning, including bounds outside the stored ASCII alphabet or beyond the stored-id length limit. Partition-key filter literals stay within 1,024 characters: a 512-character prefix selects its exact id, longer prefixes produce an empty result, and longer bounds compare against their first 512 characters with a strict lower comparison and an inclusive upper comparison. An empty intersection issues no query. Custom mappings apply well-formed Unicode bounds to the canonical `JournalId` property; bounds containing unpaired surrogates are enforced locally. This preserves ordinal range semantics and can require a full table scan. Unbounded enumeration also scans headers across the table. Constraints are checked again before yielding, and one enumerator advance can cross multiple empty or filtered pages.

Use `await foreach` or dispose a retained enumerator when stopping early. Pass a cancellation token covering the traversal lifetime; cancellation and service errors propagate to the caller.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -605,6 +605,126 @@ public async Task ListAsync_BoundsBeyondStoredIdLengthLimit_DoNotThrowOrExcludeS
Assert.Equal(1, Assert.Single(table.QueryCalls).ReturnedCount);
}

[Theory]
[InlineData(511, false)]
[InlineData(511, true)]
[InlineData(512, false)]
[InlineData(512, true)]
public async Task ListAsync_MaximumLengthPrefix_UsesServiceSizedExactBounds(int length, bool bounded)
{
var prefix = new string('a', length);
var ids = new[] { new string('a', 510), new string('a', 511), new string('a', 512), new string('a', 511) + "b", "b" };
var table = new FakeTableClient();
foreach (var id in ids)
{
table.AddHeader(new(id));
}

using var context = await CreateStartedProviderAsync(table, TestContext.Current.CancellationToken);
var maximum = prefix + "a";
var result = await ToListAsync(
context.Provider.ListAsync(new()
{
Prefix = new(prefix),
MinId = bounded ? new(prefix) : default,
MaxId = bounded ? new(maximum) : default,
}, TestContext.Current.CancellationToken),
TestContext.Current.CancellationToken);

var expected = ids.Where(id => id.StartsWith(prefix, StringComparison.Ordinal)
&& (!bounded || string.CompareOrdinal(id, maximum) <= 0)).ToArray();
Assert.Equal(expected, result.Select(entry => entry.Id.Value));
var query = Assert.Single(table.QueryCalls);
Assert.Equal(expected.Length, query.ReturnedCount);
var prefixKey = AzureTableJournalStorageOptions.EncodePartitionKey(prefix);
Assert.Contains(
length == 512
? TableClient.CreateQueryFilter($"PartitionKey le {prefixKey}")
: TableClient.CreateQueryFilter($"PartitionKey lt {prefixKey + "G"}"),
query.Filter);
}

[Theory]
[InlineData(513)]
[InlineData(2048)]
public async Task ListAsync_OversizedDefaultPrefix_ReturnsEmptyBeforeClientAccess(int length)
{
using var context = CreateProvider();

var result = await ToListAsync(
context.Provider.ListAsync(new() { Prefix = new(new string('a', length)) }, TestContext.Current.CancellationToken),
TestContext.Current.CancellationToken);

Assert.Empty(result);
}

[Theory]
[InlineData(511, -1)]
[InlineData(512, -1)]
[InlineData(513, -1)]
[InlineData(2048, -1)]
[InlineData(510, 0x00)]
[InlineData(511, 0x00)]
[InlineData(512, 0x00)]
[InlineData(2048, 0x00)]
[InlineData(510, 0x7F)]
[InlineData(511, 0x7F)]
[InlineData(512, 0x7F)]
[InlineData(2048, 0x7F)]
[InlineData(510, 0xD800)]
[InlineData(511, 0xD800)]
[InlineData(512, 0xD800)]
[InlineData(2048, 0xD800)]
public async Task ListAsync_LongBounds_UseServiceSizedExactOrdinalFilters(int length, int unsupportedCharacter)
{
var bound = new string('a', length)
+ (unsupportedCharacter < 0 ? string.Empty : (char)unsupportedCharacter + "/suffix");
string[] ids = ["a", new string('a', 511), new string('a', 512), new string('a', 511) + "b", "b"];
var table = new FakeTableClient();
foreach (var id in ids)
{
table.AddHeader(new(id));
}

using var context = await CreateStartedProviderAsync(table, TestContext.Current.CancellationToken);
var lower = await ToListAsync(
context.Provider.ListAsync(new() { MinId = new(bound) }, TestContext.Current.CancellationToken),
TestContext.Current.CancellationToken);
var upper = await ToListAsync(
context.Provider.ListAsync(new() { MaxId = new(bound) }, TestContext.Current.CancellationToken),
TestContext.Current.CancellationToken);

var expectedLower = ids.Where(id => string.CompareOrdinal(id, bound) >= 0).ToArray();
var expectedUpper = ids.Where(id => string.CompareOrdinal(id, bound) <= 0).ToArray();
Assert.Equal(expectedLower, lower.Select(entry => entry.Id.Value));
Assert.Equal(expectedUpper, upper.Select(entry => entry.Id.Value));
Assert.Equal([expectedLower.Length, expectedUpper.Length], table.QueryCalls.Select(query => query.ReturnedCount));
Assert.All(table.QueryCalls, query => Assert.DoesNotContain("JournalId", query.Filter));
}

[Fact]
public async Task ListAsync_CustomMapping_LongIdRetainsCanonicalPropertyFilters()
{
var id = new JournalId(new string('a', 513));
var table = new FakeTableClient();
table.AddHeader(id, "custom");
var options = CreateOptions(table);
options.GetPartitionKey = static _ => "custom";
using var context = CreateProvider(options);
await StartAsync(context.Provider, TestContext.Current.CancellationToken);

var result = await ToListAsync(
context.Provider.ListAsync(new() { Prefix = id, MinId = id, MaxId = id }, TestContext.Current.CancellationToken),
TestContext.Current.CancellationToken);

Assert.Equal([id], result.Select(entry => entry.Id));
var query = Assert.Single(table.QueryCalls);
Assert.Equal(1, query.ReturnedCount);
Assert.DoesNotContain("PartitionKey", query.Filter);
Assert.Contains(TableClient.CreateQueryFilter($"JournalId ge {id.Value}"), query.Filter);
Assert.Contains(TableClient.CreateQueryFilter($"JournalId le {id.Value}"), query.Filter);
}

[Fact]
public async Task ListAsync_DefaultMapping_RetainsExactLocalRangeFilter()
{
Expand Down Expand Up @@ -1169,8 +1289,11 @@ public override AsyncPageable<T> QueryAsync<T>(
cancellationToken.ThrowIfCancellationRequested();
// OData filter text must survive transport encoding without surrogate replacement.
_ = StrictUtf8.GetByteCount(filter!);
var clauses = Regex.Matches(filter!, @"(RowKey|PartitionKey|JournalId) (eq|ge|le|lt) '((?:[^']|'')*)'");
var clauses = Regex.Matches(filter!, @"(RowKey|PartitionKey|JournalId) (eq|ge|gt|le|lt) '((?:[^']|'')*)'");
Assert.Equal(filter, string.Join(" and ", clauses.Select(clause => clause.Value)));
Assert.All(
clauses.Where(clause => clause.Groups[1].Value == "PartitionKey"),
clause => Assert.InRange(clause.Groups[3].Value.Replace("''", "'").Length, 0, 1024));
var values = _entities
.Where(entity => clauses.All(clause =>
{
Expand All @@ -1191,6 +1314,7 @@ public override AsyncPageable<T> QueryAsync<T>(
{
"eq" => comparison == 0,
"ge" => comparison >= 0,
"gt" => comparison > 0,
"le" => comparison <= 0,
"lt" => comparison < 0,
_ => throw new InvalidOperationException("Unexpected table filter comparison."),
Expand Down
Loading