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
10 changes: 10 additions & 0 deletions src/AWS/Orleans.Journaling.S3/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,3 +7,13 @@ The provider uses S3 Express append writes (`WriteOffsetBytes`) for WAL appends.
Buckets should be created ahead of time for AWS S3 Express One Zone. `CreateBucketIfNotExists` is intended for local emulators.

Metadata updates rewrite the current WAL using a conditional single-object upload. Publish a checkpoint to compact the WAL before updating metadata when the replacement object would exceed S3's 5 GB (5,000,000,000 byte) single-upload limit. Checkpoint snapshots use the same upload limit.

## Catalog enumeration

`IJournalStorageCatalog.ListAsync` returns journal identities incrementally in S3 traversal order, including unordered directory-bucket listings. Set `ListOptions.Prefix` to select an exact journal id and its descendants.

The provider handles `ListObjectsV2` continuations internally, requests up to 1000 objects per page, and yields canonical WAL identities from that page before fetching more objects. The bucket traversal supports `GetObjectKey` and `TryParseJournalId` mappings; checkpoints, aliases, and unrelated objects consume space in the native page before filtering.

Client traversal memory is proportional to the current native page. An enumerator advance can cross multiple filtered or empty pages, and the storage service determines scan work, latency, and retries. Enumeration observes the live bucket; concurrent changes follow S3 listing semantics. Use subsequent enumerations to discover later changes and tolerate repeated identities during changes.

Dispose the enumerator when stopping early and use a cancellation token covering its lifetime. Cancellation and service failures propagate through enumeration.
46 changes: 25 additions & 21 deletions src/AWS/Orleans.Journaling.S3/S3JournalStorageProvider.cs
Original file line number Diff line number Diff line change
Expand Up @@ -46,12 +46,13 @@ public IJournalStorage CreateStorage(JournalId journalId)
}

public async IAsyncEnumerable<JournalId> ListAsync(
JournalId prefix = default,
ListOptions? options = null,
[EnumeratorCancellation] CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();
var prefix = options?.Prefix ?? default;
var client = GetClient();
var bucketName = GetBucketName();
var journalIds = new HashSet<JournalId>();
string? continuationToken = null;

do
Expand All @@ -60,42 +61,45 @@ public async IAsyncEnumerable<JournalId> ListAsync(
new ListObjectsV2Request
{
BucketName = bucketName,
MaxKeys = 1000,
ContinuationToken = continuationToken,
},
cancellationToken).ConfigureAwait(false);

cancellationToken.ThrowIfCancellationRequested();
foreach (var item in response.S3Objects)
{
cancellationToken.ThrowIfCancellationRequested();
if (!item.Key.EndsWith("/wal", StringComparison.Ordinal))
{
continue;
}

var storageIdValue = item.Key[..^"/wal".Length];
var journalId = _options.TryParseJournalId(storageIdValue);
if (journalId is not { IsDefault: false } id || !prefix.IsPrefixOf(id))
{
continue;
}

var journalObjectKey = _options.GetObjectKeyForJournal(id);
var canonicalWalObjectKey = S3JournalStorageOptions.GetWalObjectKeyForJournal(id, journalObjectKey);
if (string.Equals(item.Key, canonicalWalObjectKey, StringComparison.Ordinal))
if (TryGetJournalId(item.Key, prefix, out var id))
{
journalIds.Add(id);
yield return id;
}
}

continuationToken = response.IsTruncated == true ? response.NextContinuationToken : null;
}
while (continuationToken is not null);

foreach (var journalId in journalIds.OrderBy(static journalId => journalId.Value, StringComparer.Ordinal))
cancellationToken.ThrowIfCancellationRequested();
}

private bool TryGetJournalId(string objectKey, JournalId prefix, out JournalId journalId)
{
if (objectKey.EndsWith("/wal", StringComparison.Ordinal)
&& _options.TryParseJournalId(objectKey[..^"/wal".Length]) is { IsDefault: false } id
&& prefix.IsPrefixOf(id))
{
cancellationToken.ThrowIfCancellationRequested();
yield return journalId;
var journalObjectKey = _options.GetObjectKeyForJournal(id);
var canonicalWalObjectKey = S3JournalStorageOptions.GetWalObjectKeyForJournal(id, journalObjectKey);
if (string.Equals(objectKey, canonicalWalObjectKey, StringComparison.Ordinal))
{
journalId = id;
return true;
}
}

journalId = default;
return false;
}

public void Participate(ISiloLifecycle observer)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,9 +40,10 @@ public AzureBlobJournalStorageProvider(
private async Task Initialize(CancellationToken cancellationToken)
{
var client = await _options.CreateClient!(cancellationToken);
_defaultContainer = client.GetBlobContainerClient(_options.ContainerName);
await _defaultContainer.CreateIfNotExistsAsync(cancellationToken: cancellationToken).ConfigureAwait(false);
var container = client.GetBlobContainerClient(_options.ContainerName);
await container.CreateIfNotExistsAsync(cancellationToken: cancellationToken).ConfigureAwait(false);
await _containerFactory.InitializeAsync(client, cancellationToken).ConfigureAwait(false);
_defaultContainer = container;
}

public IJournalStorage CreateStorage(JournalId journalId)
Expand All @@ -56,40 +57,37 @@ public IJournalStorage CreateStorage(JournalId journalId)
}

public async IAsyncEnumerable<JournalId> ListAsync(
JournalId prefix = default,
ListOptions? options = null,
[EnumeratorCancellation] CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();
var prefix = options?.Prefix ?? default;
var container = GetDefaultContainerClient();
var blobPrefix = prefix.IsDefault ? null : prefix.Value;
var journalIds = new List<JournalId>();
await foreach (var item in container.GetBlobsAsync(

await foreach (var page in container.GetBlobsAsync(
traits: BlobTraits.None,
states: BlobStates.None,
prefix: blobPrefix,
cancellationToken: cancellationToken))
prefix: prefix.IsDefault ? null : prefix.Value,
cancellationToken: cancellationToken).AsPages(pageSizeHint: 5000))
{
if (item.Properties.BlobType is { } blobType && blobType != BlobType.Append)
{
continue;
}

if (!item.Name.EndsWith("/wal", StringComparison.Ordinal))
{
continue;
}

var storageIdValue = item.Name[..^"/wal".Length];
if (TryParseJournalId(storageIdValue, out var journalId) && prefix.IsPrefixOf(journalId))
cancellationToken.ThrowIfCancellationRequested();
foreach (var item in page.Values)
{
journalIds.Add(journalId);
cancellationToken.ThrowIfCancellationRequested();
if (item.Properties.BlobType is { } blobType && blobType != BlobType.Append
|| !item.Name.EndsWith("/wal", StringComparison.Ordinal))
{
continue;
}

if (TryParseJournalId(item.Name[..^"/wal".Length], out var journalId) && prefix.IsPrefixOf(journalId))
{
yield return journalId;
}
}
}

foreach (var journalId in journalIds.OrderBy(static journalId => journalId.Value, StringComparer.Ordinal))
{
cancellationToken.ThrowIfCancellationRequested();
yield return journalId;
}
cancellationToken.ThrowIfCancellationRequested();
}

public void Participate(ISiloLifecycle observer)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,25 +57,31 @@ public IJournalStorage CreateStorage(JournalId journalId)
}

public async IAsyncEnumerable<JournalId> ListAsync(
JournalId prefix = default,
ListOptions? options = null,
[EnumeratorCancellation] CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();
var prefix = options?.Prefix ?? default;
var table = _tableClientProvider.GetTableClient();
var filter = TableClient.CreateQueryFilter($"RowKey eq {AzureTableJournalStorage.HeaderRowKey}");
var journalIds = new List<JournalId>();
await foreach (var entity in table.QueryAsync<TableEntity>(filter, select: JournalIdSelect, cancellationToken: cancellationToken))
await foreach (var page in table.QueryAsync<TableEntity>(
filter,
maxPerPage: 1000,
select: JournalIdSelect,
cancellationToken: cancellationToken).AsPages(pageSizeHint: 1000))
{
if (TryGetJournalId(entity, out var journalId) && prefix.IsPrefixOf(journalId))
cancellationToken.ThrowIfCancellationRequested();
foreach (var entity in page.Values)
{
journalIds.Add(journalId);
cancellationToken.ThrowIfCancellationRequested();
if (TryGetJournalId(entity, out var journalId) && prefix.IsPrefixOf(journalId))
{
yield return journalId;
}
}
}

foreach (var journalId in journalIds.OrderBy(static journalId => journalId.Value, StringComparer.Ordinal))
{
cancellationToken.ThrowIfCancellationRequested();
yield return journalId;
}
cancellationToken.ThrowIfCancellationRequested();
}

private static bool TryGetJournalId(TableEntity entity, out JournalId journalId)
Expand Down
10 changes: 10 additions & 0 deletions src/Azure/Orleans.Journaling.AzureStorage/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,16 @@ siloBuilder.AddAzureTableJournalStorage(options =>

`AzureTableJournalStorageOptions.GetPartitionKey` can be used to apply a custom partition layout. The returned key must be unique per journal, satisfy Azure Table partition-key restrictions, and contain at most 1,024 characters. The canonical journal id is stored in the header so catalog listing remains accurate with custom mappings.

## Catalog enumeration

Both Azure providers implement `IJournalStorageCatalog.ListAsync`, returning identities incrementally in service traversal order. Set `ListOptions.Prefix` to enumerate an exact journal id and its descendants. The provider handles service continuations internally and yields identities from the current page before fetching the next page.

The Blob catalog scans the configured `ContainerName` and interprets append blobs named `<journalId>/wal` as journal identities. This traversal applies equally when a custom naming delegate or container factory produces the same entries. Each internal page requests up to 5000 blobs, including checkpoints and other entries before filtering.

Table enumeration supports custom partition mappings through the canonical journal id stored in each header. It requests up to 1000 header rows per service page; Table Storage determines the internal scan work required by that query. Prefix filtering occurs on the returned headers, so 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.

## Getting Started
To use this package, install it via NuGet:

Expand Down
2 changes: 1 addition & 1 deletion src/Orleans.DurableJobs/JournaledJobShardManager.cs
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ public override async Task<List<IJobShard>> AssignJobShardsAsync(DateTimeOffset
var newClaimCount = 0;
var membershipSnapshot = _membershipService.CurrentSnapshot;

await foreach (var storageId in _catalog.ListAsync(JobShardId.StoragePrefix, cancellationToken))
await foreach (var storageId in _catalog.ListAsync(new() { Prefix = JobShardId.StoragePrefix }, cancellationToken))
{
var descriptor = await GetDescriptorAsync(storageId, cancellationToken);
if (descriptor is null || descriptor.Poisoned || descriptor.StartTime > maxDueTime)
Expand Down
46 changes: 46 additions & 0 deletions src/Orleans.Journaling/CompatibilitySuppressions.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
<?xml version="1.0" encoding="utf-8"?>
<!-- https://learn.microsoft.com/dotnet/fundamentals/package-validation/diagnostic-ids -->
<Suppressions xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xmlns:xsd="http://www.w3.org/2001/XMLSchema">
<Suppression>
<DiagnosticId>CP0002</DiagnosticId>
<Target>M:Orleans.Journaling.IJournalStorageCatalog.ListAsync(Orleans.Journaling.JournalId,System.Threading.CancellationToken)</Target>
<Left>lib/net10.0/Orleans.Journaling.dll</Left>
<Right>lib/net10.0/Orleans.Journaling.dll</Right>
<IsBaselineSuppression>true</IsBaselineSuppression>
</Suppression>
<Suppression>
<DiagnosticId>CP0002</DiagnosticId>
<Target>M:Orleans.Journaling.VolatileJournalStorageProvider.ListAsync(Orleans.Journaling.JournalId,System.Threading.CancellationToken)</Target>
<Left>lib/net10.0/Orleans.Journaling.dll</Left>
<Right>lib/net10.0/Orleans.Journaling.dll</Right>
<IsBaselineSuppression>true</IsBaselineSuppression>
</Suppression>
<Suppression>
<DiagnosticId>CP0002</DiagnosticId>
<Target>M:Orleans.Journaling.IJournalStorageCatalog.ListAsync(Orleans.Journaling.JournalId,System.Threading.CancellationToken)</Target>
<Left>lib/net8.0/Orleans.Journaling.dll</Left>
<Right>lib/net8.0/Orleans.Journaling.dll</Right>
<IsBaselineSuppression>true</IsBaselineSuppression>
</Suppression>
<Suppression>
<DiagnosticId>CP0002</DiagnosticId>
<Target>M:Orleans.Journaling.VolatileJournalStorageProvider.ListAsync(Orleans.Journaling.JournalId,System.Threading.CancellationToken)</Target>
<Left>lib/net8.0/Orleans.Journaling.dll</Left>
<Right>lib/net8.0/Orleans.Journaling.dll</Right>
<IsBaselineSuppression>true</IsBaselineSuppression>
</Suppression>
<Suppression>
<DiagnosticId>CP0006</DiagnosticId>
<Target>M:Orleans.Journaling.IJournalStorageCatalog.ListAsync(Orleans.Journaling.ListOptions,System.Threading.CancellationToken)</Target>
<Left>lib/net10.0/Orleans.Journaling.dll</Left>
<Right>lib/net10.0/Orleans.Journaling.dll</Right>
<IsBaselineSuppression>true</IsBaselineSuppression>
</Suppression>
<Suppression>
<DiagnosticId>CP0006</DiagnosticId>
<Target>M:Orleans.Journaling.IJournalStorageCatalog.ListAsync(Orleans.Journaling.ListOptions,System.Threading.CancellationToken)</Target>
<Left>lib/net8.0/Orleans.Journaling.dll</Left>
<Right>lib/net8.0/Orleans.Journaling.dll</Right>
<IsBaselineSuppression>true</IsBaselineSuppression>
</Suppression>
</Suppressions>
24 changes: 18 additions & 6 deletions src/Orleans.Journaling/IJournalStorageCatalog.cs
Original file line number Diff line number Diff line change
Expand Up @@ -4,16 +4,28 @@ namespace Orleans.Journaling;
/// Provides catalog operations for journal storage instances.
/// </summary>
/// <remarks>
/// A catalog only discovers storage identities. Storage lifecycle, metadata, and data mutation
/// operations remain on <see cref="IJournalStorage"/>.
/// A catalog discovers storage identities. <see cref="IJournalStorage"/> provides storage lifecycle,
/// metadata, and data mutation operations.
/// </remarks>
public interface IJournalStorageCatalog
{
/// <summary>
/// Lists journal ids which match <paramref name="prefix"/>.
/// Enumerates journal ids matching the supplied options.
/// </summary>
/// <param name="prefix">The journal id prefix, or the default value to list all ids.</param>
/// <param name="options">The listing options, or <see langword="null"/> to list all ids.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>Matching ids in lexicographic <see cref="JournalId.Value"/> order.</returns>
IAsyncEnumerable<JournalId> ListAsync(JournalId prefix = default, CancellationToken cancellationToken = default);
/// <returns>Matching ids in provider traversal order.</returns>
/// <remarks>
/// Options are read when enumeration begins. Providers fetch storage pages internally and yield matching ids
/// as they are discovered. Advancing the enumerator can traverse multiple empty or filtered storage pages.
/// Storage services determine request latency, retries, and internal scan work.
/// Enumeration observes live storage; concurrent changes follow the provider's listing semantics.
/// Callers should tolerate repeated identities during concurrent changes and start a new enumeration to
/// discover later changes. Dispose the enumerator when stopping early.
/// Storage and cancellation errors propagate through enumeration. Start a new enumeration after a listing error.
/// </remarks>
/// <exception cref="OperationCanceledException"><paramref name="cancellationToken"/> is canceled.</exception>
IAsyncEnumerable<JournalId> ListAsync(
ListOptions? options = null,
CancellationToken cancellationToken = default);
}
15 changes: 15 additions & 0 deletions src/Orleans.Journaling/ListOptions.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
namespace Orleans.Journaling;

/// <summary>
/// Options for enumerating journal storage identities.
/// </summary>
public sealed class ListOptions
{
/// <summary>
/// Gets or sets the journal id prefix. The default value matches all ids.
/// </summary>
/// <remarks>
/// A prefix matches the exact journal id and its descendant segments.
/// </remarks>
public JournalId Prefix { get; set; }
}
Loading
Loading