diff --git a/playground/DurableJobsJournaling/DurableJobsJournaling.AppHost/DurableJobsJournaling.AppHost.csproj b/playground/DurableJobsJournaling/DurableJobsJournaling.AppHost/DurableJobsJournaling.AppHost.csproj index 3a51b6a864c..68b446a4ca0 100644 --- a/playground/DurableJobsJournaling/DurableJobsJournaling.AppHost/DurableJobsJournaling.AppHost.csproj +++ b/playground/DurableJobsJournaling/DurableJobsJournaling.AppHost/DurableJobsJournaling.AppHost.csproj @@ -18,6 +18,7 @@ + diff --git a/playground/DurableJobsJournaling/DurableJobsJournaling.AppHost/Program.cs b/playground/DurableJobsJournaling/DurableJobsJournaling.AppHost/Program.cs index 2ea376e41ca..166d26f1f94 100644 --- a/playground/DurableJobsJournaling/DurableJobsJournaling.AppHost/Program.cs +++ b/playground/DurableJobsJournaling/DurableJobsJournaling.AppHost/Program.cs @@ -1,7 +1,9 @@ using Aspire.Hosting; +using Aspire.Hosting.ApplicationModel; using Aspire.Hosting.Azure; using Azure.Provisioning; using Azure.Provisioning.Storage; +using DurableJobsJournaling; using DurableJobsJournaling.AppHost.OpenTelemetryCollector; using Microsoft.Extensions.Configuration; @@ -28,57 +30,75 @@ const string OtelMetricExportIntervalMilliseconds = "5000"; -var storageProvider = builder.Configuration.GetValue("Playground:Storage:Provider", "Azurite"); -var useAzurite = storageProvider.Equals("Azurite", StringComparison.OrdinalIgnoreCase); -var useAzure = storageProvider.Equals("Azure", StringComparison.OrdinalIgnoreCase); +var backend = StorageBackendConfiguration.Parse(builder.Configuration.GetValue("Playground:Storage:Provider", "Azurite")); var storage = builder.AddAzureStorage("storage"); -if (useAzurite) +if (backend.IsEmulator()) { storage.RunAsEmulator(); } -else if (useAzure) +else { storage.ConfigureInfrastructure(infrastructure => { var storageAccount = infrastructure.GetProvisionableResources().OfType().Single(); - storageAccount.Kind = StorageKind.BlockBlobStorage; - storageAccount.Sku = new StorageSku { Name = StorageSkuName.PremiumLrs }; - storageAccount.AccessTier.ClearValue(); - RemoveUnsupportedPremiumBlobStorageOutputs(infrastructure); + storageAccount.Kind = backend == StorageBackend.PremiumBlob ? StorageKind.BlockBlobStorage : StorageKind.StorageV2; + storageAccount.Sku = new StorageSku { Name = backend == StorageBackend.PremiumBlob ? StorageSkuName.PremiumLrs : StorageSkuName.StandardLrs }; + if (backend == StorageBackend.PremiumBlob) + { + storageAccount.AccessTier.ClearValue(); + RemoveUnsupportedPremiumBlobStorageOutputs(infrastructure); + } }); } -else +var tableStorage = backend == StorageBackend.PremiumBlob ? builder.AddAzureStorage("clusteringstorage") : storage; +if (backend == StorageBackend.PremiumBlob) { - throw new InvalidOperationException($"Unknown Playground:Storage:Provider value '{storageProvider}'. Use 'Azurite' or 'Azure'."); + tableStorage.ConfigureInfrastructure(infrastructure => + { + var storageAccount = infrastructure.GetProvisionableResources().OfType().Single(); + storageAccount.Sku = new StorageSku { Name = StorageSkuName.StandardLrs }; + }); } -var blobs = storage.AddBlobs("blobs"); -var tableStorage = useAzurite ? storage : builder.AddAzureStorage("clusteringstorage"); -if (useAzure && builder.ExecutionContext.IsPublishMode) +var tables = tableStorage.AddTables("tables"); + +SetStorageRoles(storage, backend.UsesTableJournal() + ? [StorageBuiltInRole.StorageTableDataContributor] + : backend == StorageBackend.PremiumBlob + ? [StorageBuiltInRole.StorageBlobDataContributor] + : [StorageBuiltInRole.StorageBlobDataContributor, StorageBuiltInRole.StorageTableDataContributor]); +if (backend == StorageBackend.PremiumBlob) { - storage.ClearDefaultRoleAssignments(); - tableStorage.ClearDefaultRoleAssignments(); + SetStorageRoles(tableStorage, [StorageBuiltInRole.StorageTableDataContributor]); } -var tables = tableStorage.AddTables("tables"); - var orleans = builder.AddOrleans("cluster") .WithClustering(tables); -var storagePrefix = $"run-{DateTimeOffset.UtcNow:yyyyMMdd-HHmmss}"; +var runId = Guid.NewGuid().ToString("N"); var silo = builder.AddProject("silo") .WithReference(orleans) - .WithReference(blobs) .WithReference(tables) - .WaitFor(blobs) .WaitFor(tables) .WaitFor(otelCollector) .WithReplicas(1) - .WithEnvironment("Playground__Storage__Container", "durablejobs-journaling-playground") - .WithEnvironment("Playground__Storage__Prefix", storagePrefix) + .WithEnvironment("Playground__Storage__Provider", backend.ToString()) + .WithEnvironment("Playground__Storage__Container", $"durablejobs-{runId}") + .WithEnvironment("Playground__Storage__Table", $"durablejobs{runId}") .WithEnvironment("OTEL_METRIC_EXPORT_INTERVAL", OtelMetricExportIntervalMilliseconds); +if (backend.UsesTableJournal()) +{ + var journals = storage.AddTables("journals"); + silo.WithReference(journals).WaitFor(journals); +} +else +{ + var blobs = storage.AddBlobs("blobs"); + silo.WithReference(blobs).WaitFor(blobs); +} + builder.AddProject("web") .WithReference(orleans.AsClient()) .WithReference(tables) @@ -91,6 +111,14 @@ builder.Build().Run(); +static void SetStorageRoles(IResourceBuilder storage, StorageBuiltInRole[] roles) +{ + storage.ClearDefaultRoleAssignments() + .WithAnnotation(new DefaultRoleAssignmentsAnnotation(roles + .Select(role => new RoleDefinition(role.ToString(), StorageBuiltInRole.GetBuiltInRoleName(role))) + .ToHashSet())); +} + static void RemoveUnsupportedPremiumBlobStorageOutputs(AzureResourceInfrastructure infrastructure) { foreach (var output in infrastructure.GetProvisionableResources() diff --git a/playground/DurableJobsJournaling/DurableJobsJournaling.Silo/DurableJobsJournaling.Silo.csproj b/playground/DurableJobsJournaling/DurableJobsJournaling.Silo/DurableJobsJournaling.Silo.csproj index 9e8a9d99aa1..aabcaddb071 100644 --- a/playground/DurableJobsJournaling/DurableJobsJournaling.Silo/DurableJobsJournaling.Silo.csproj +++ b/playground/DurableJobsJournaling/DurableJobsJournaling.Silo/DurableJobsJournaling.Silo.csproj @@ -15,6 +15,7 @@ + diff --git a/playground/DurableJobsJournaling/DurableJobsJournaling.Silo/Program.cs b/playground/DurableJobsJournaling/DurableJobsJournaling.Silo/Program.cs index a1180370e7f..2234279ad01 100644 --- a/playground/DurableJobsJournaling/DurableJobsJournaling.Silo/Program.cs +++ b/playground/DurableJobsJournaling/DurableJobsJournaling.Silo/Program.cs @@ -1,4 +1,6 @@ +using Azure.Data.Tables; using Azure.Storage.Blobs; +using DurableJobsJournaling; using DurableJobsJournaling.Silo; using Orleans.Dashboard; using Orleans.Journaling; @@ -7,27 +9,53 @@ var builder = WebApplication.CreateBuilder(args); builder.AddServiceDefaults(); -builder.AddAzureBlobServiceClient("blobs"); builder.AddKeyedAzureTableServiceClient("tables"); -var storageContainer = builder.Configuration.GetValue("Playground:Storage:Container", "durablejobs-journaling-playground"); -var storagePrefix = builder.Configuration.GetValue("Playground:Storage:Prefix", $"run-{DateTimeOffset.UtcNow:yyyyMMdd-HHmmss}"); +var backend = StorageBackendConfiguration.Parse(builder.Configuration.GetValue("Playground:Storage:Provider", "Azurite")); +if (backend.UsesTableJournal()) +{ + builder.AddAzureTableServiceClient("journals"); + builder.Services.AddOptions() + .Configure((options, client) => + { + options.TableServiceClient = client; + }); +} +else +{ + builder.AddAzureBlobServiceClient("blobs"); + builder.Services.AddOptions() + .Configure((options, client) => + { + options.BlobServiceClient = client; + }); +} builder.UseOrleans(siloBuilder => { #pragma warning disable ORLEANSEXP003 // Type is for evaluation purposes only and is subject to change or removal in future updates. Suppress this diagnostic to proceed. + if (backend.UsesTableJournal()) + { + siloBuilder.UseAzureTableDurableJobs(options => + { + options.TableName = builder.Configuration["Playground:Storage:Table"] + ?? throw new InvalidOperationException("Set Playground:Storage:Table to the shared run-specific table name."); + }); + } + else + { + siloBuilder.UseAzureBlobDurableJobs(options => + { + options.ContainerName = builder.Configuration["Playground:Storage:Container"] + ?? throw new InvalidOperationException("Set Playground:Storage:Container to the shared run-specific container name."); + }); + } + siloBuilder .AddDashboard() .AddActivityPropagation() .AddIncomingGrainCallFilter() .AddDistributedGrainDirectory() - .UseAzureBlobDurableJobs( - options => - { - options.ContainerName = storageContainer; - options.GetWalBlobName = journalId => $"{storagePrefix}/{journalId.Value}/wal"; - options.GetCheckpointBlobName = (journalId, snapshotId) => $"{storagePrefix}/{journalId.Value}/chk.{snapshotId}"; - }) .UseJsonJournalFormat(DurableJobsJournalingJsonContext.Default) .Configure(options => { @@ -46,12 +74,6 @@ #pragma warning restore ORLEANSEXP003 // Type is for evaluation purposes only and is subject to change or removal in future updates. Suppress this diagnostic to proceed. }); -builder.Services.AddOptions() - .Configure((options, blobServiceClient) => - { - options.BlobServiceClient = blobServiceClient; - }); - var app = builder.Build(); app.MapDefaultEndpoints(); app.MapOrleansDashboard(); diff --git a/playground/DurableJobsJournaling/README.md b/playground/DurableJobsJournaling/README.md new file mode 100644 index 00000000000..5cef952a289 --- /dev/null +++ b/playground/DurableJobsJournaling/README.md @@ -0,0 +1,69 @@ +# Durable Jobs journaling playground + +The Aspire app runs durable workflow grains, a web load driver, and +Prometheus/Grafana monitoring. Select journal storage using the shared +`Playground:Storage:Provider` setting: + +| Value | Journal storage | Clustering | +| --- | --- | --- | +| `Azurite` (default) | Azurite Blob | Azurite Table | +| `AzuriteTable` | Azurite Table | Azurite Table | +| `StandardBlob` | Standard Azure Blob account | Table service on the standard account | +| `PremiumBlob` | Premium LRS BlockBlobStorage account | Separate standard account's Table service | +| `Table` | Azure Table | Table service on the same standard account | +| `Azure` | Alias for `PremiumBlob` | Separate standard account's Table service | + +Both Blob tiers use append-blob WALs and block-blob checkpoints. Premium accounts +expose Blob resources; the AppHost removes unsupported queue/table outputs and +supplies clustering through a separate standard account. + +From the repository root, with the .NET SDK and a Docker-compatible container +runtime available: + +```powershell +dotnet run --project playground\DurableJobsJournaling\DurableJobsJournaling.AppHost + +# To select Table journaling on the emulator: +$env:Playground__Storage__Provider = 'AzuriteTable' +dotnet run --project playground\DurableJobsJournaling\DurableJobsJournaling.AppHost +``` + +The AppHost passes the normalized backend to the silo and references `blobs` or +`journals` for journal data, plus `tables` for clustering. The web client uses +the clustering reference. Each app run generates a GUID-based container +(`durablejobs-...`) or journal table (`durablejobs...`) shared by its silo +replicas. Blob layout remains `wal/{journalId}` and +`checkpoints/{journalId}/{snapshotId}`, enabling catalog discovery with the +provider's default layout. + +The run-specific resources preserve journals for inspection. Record the +container/table name from Aspire's environment configuration and delete that +resource when the run's data is no longer needed. For restart/recovery +experiments, keep the shared resource name and workload identity stable across +the processes being restarted. + +Cloud backend selection enables Aspire's Azure resource configuration and can +provision chargeable resources. Use an explicitly approved Azure development +environment and its standard Aspire authentication/deployment procedure. +`PremiumBlob` configures `BlockBlobStorage` with `Premium_LRS`; standard accounts +use `Standard_LRS`. +Storage role defaults match the selected services: Blob and Table Data +Contributor for a shared standard Blob/clustering account, Table Data Contributor +for Table journals, and separate Blob-only journal and Table-only clustering +assignments for premium Blob. Generic manifest publishing retains role modules +with `principalId` and `principalType` inputs for the deployment's intended +identity. An identity-aware Aspire deployment environment can apply these defaults +to its application identities. Configure that environment and identity before +deployment; local Azurite uses its emulator credentials. +Compare backends using the same compute, region, account redundancy, job +distribution, and `Playground:DurableJobs` settings. + +Open the web resource from Aspire's dashboard to configure load concurrency/rate, +start and stop load, drain outstanding work, and inspect workflow metrics. +Existing stage-write and end-to-end latency metrics, durable-job retry policy, +slow start, and scheduling tunings apply to every backend. Use +`Playground__DurableJobs__...` settings to adjust those tunings. + +For a finite provider-only workload with isolated resource cleanup, operation +percentiles, and JSON/CSV reports, use the +[Azure journal benchmarks](../../test/Benchmarks/Journaling/Azure/README.md). diff --git a/playground/DurableJobsJournaling/StorageBackend.cs b/playground/DurableJobsJournaling/StorageBackend.cs new file mode 100644 index 00000000000..0548a0e07de --- /dev/null +++ b/playground/DurableJobsJournaling/StorageBackend.cs @@ -0,0 +1,26 @@ +namespace DurableJobsJournaling; + +internal enum StorageBackend +{ + Azurite, + AzuriteTable, + StandardBlob, + PremiumBlob, + Table +} + +internal static class StorageBackendConfiguration +{ + public static StorageBackend Parse(string value) => value.ToUpperInvariant() switch + { + "AZURITE" => StorageBackend.Azurite, + "AZURITETABLE" => StorageBackend.AzuriteTable, + "STANDARDBLOB" => StorageBackend.StandardBlob, + "PREMIUMBLOB" or "AZURE" => StorageBackend.PremiumBlob, + "TABLE" => StorageBackend.Table, + _ => throw new InvalidOperationException("Playground:Storage:Provider must be Azurite, AzuriteTable, StandardBlob, PremiumBlob, Table, or Azure (PremiumBlob).") + }; + + public static bool IsEmulator(this StorageBackend backend) => backend is StorageBackend.Azurite or StorageBackend.AzuriteTable; + public static bool UsesTableJournal(this StorageBackend backend) => backend is StorageBackend.Table or StorageBackend.AzuriteTable; +} diff --git a/src/Azure/Orleans.DurableJobs.AzureStorage/Hosting/AzureStorageDurableJobsExtensions.cs b/src/Azure/Orleans.DurableJobs.AzureStorage/Hosting/AzureStorageDurableJobsExtensions.cs index 806959deb33..e37d92f3694 100644 --- a/src/Azure/Orleans.DurableJobs.AzureStorage/Hosting/AzureStorageDurableJobsExtensions.cs +++ b/src/Azure/Orleans.DurableJobs.AzureStorage/Hosting/AzureStorageDurableJobsExtensions.cs @@ -10,7 +10,7 @@ namespace Orleans.Hosting; /// -/// Extensions for configuring Azure Blob Storage durable jobs. +/// Extensions for configuring Azure Storage durable jobs. /// public static class AzureStorageDurableJobsExtensions { @@ -65,6 +65,39 @@ public static IServiceCollection UseAzureBlobDurableJobs(this IServiceCollection return services; } + /// + /// Adds durable jobs backed by Azure Table journal storage. + /// + /// The silo builder. + /// The delegate used to configure journal storage. + /// The silo builder, for chaining. + public static ISiloBuilder UseAzureTableDurableJobs(this ISiloBuilder builder, Action configure) + { + ArgumentNullException.ThrowIfNull(builder); + ArgumentNullException.ThrowIfNull(configure); + + builder.AddDurableJobs(); + builder.AddAzureTableJournalStorage(configure); + builder.Configure(options => options.AddTypeInfoResolver(DurableJobsJsonContext.Default)); + builder.Services.UseJournaledDurableJobs(); + return builder; + } + + /// + /// Adds durable jobs backed by Azure Table journal storage. + /// + /// The service collection. + /// The delegate used to configure journal storage. + /// The service collection, for chaining. + public static IServiceCollection UseAzureTableDurableJobs(this IServiceCollection services, Action configure) + { + ArgumentNullException.ThrowIfNull(services); + ArgumentNullException.ThrowIfNull(configure); + + new ServiceCollectionSiloBuilder(services).UseAzureTableDurableJobs(configure); + return services; + } + private static IServiceCollection UseJournaledDurableJobs(this IServiceCollection services) { services.TryAddSingleton(); diff --git a/src/Azure/Orleans.DurableJobs.AzureStorage/Orleans.DurableJobs.AzureStorage.csproj b/src/Azure/Orleans.DurableJobs.AzureStorage/Orleans.DurableJobs.AzureStorage.csproj index 948f046c254..ee584906595 100644 --- a/src/Azure/Orleans.DurableJobs.AzureStorage/Orleans.DurableJobs.AzureStorage.csproj +++ b/src/Azure/Orleans.DurableJobs.AzureStorage/Orleans.DurableJobs.AzureStorage.csproj @@ -4,7 +4,7 @@ README.md Microsoft.Orleans.DurableJobs.AzureStorage Microsoft Orleans Azure Storage Durable Jobs Provider - Microsoft Orleans durable jobs provider backed by Azure Blob Storage + Microsoft Orleans durable jobs provider backed by Azure Blob or Azure Table journal storage $(PackageTags) Azure Storage $(DefaultTargetFrameworks) Orleans.DurableJobs.AzureStorage diff --git a/src/Azure/Orleans.DurableJobs.AzureStorage/README.md b/src/Azure/Orleans.DurableJobs.AzureStorage/README.md index aa23570537d..3e463e4c55d 100644 --- a/src/Azure/Orleans.DurableJobs.AzureStorage/README.md +++ b/src/Azure/Orleans.DurableJobs.AzureStorage/README.md @@ -1,7 +1,7 @@ # Microsoft Orleans Durable Jobs for Azure Storage ## Introduction -Microsoft Orleans Durable Jobs for Azure Storage provides persistent storage for Orleans Durable Jobs using Azure Blob Storage. This allows your Orleans applications to schedule jobs that survive silo restarts, grain deactivation, and cluster reconfigurations. Jobs are stored in append blobs, providing efficient storage and retrieval for time-based job scheduling. +Microsoft Orleans Durable Jobs for Azure Storage persists scheduled jobs through the Azure Blob or Azure Table journal provider. Jobs survive silo restarts, grain deactivation, and cluster reconfigurations. `UseAzureBlobDurableJobs` configures append-blob WALs and block-blob checkpoints; `UseAzureTableDurableJobs` configures Table journal headers and data generations. Both register the journal-backed shard manager and durable-job JSON serialization metadata. ## Getting Started @@ -15,6 +15,20 @@ dotnet add package Microsoft.Orleans.DurableJobs.AzureStorage ### Configuration +Choose `UseAzureBlobDurableJobs` with `AzureBlobJournalStorageOptions`, or +`UseAzureTableDurableJobs` with `AzureTableJournalStorageOptions`. Both support +`ISiloBuilder` and `IServiceCollection` registration. Table configuration accepts +a `TableServiceClient` and `TableName`; Blob configuration accepts a +`BlobServiceClient` and `ContainerName`. Clustering uses an appropriate Table +service, including a separate standard account when journal storage uses a +premium BlockBlobStorage account. + +The [Durable Jobs journaling playground](../../../playground/DurableJobsJournaling/README.md) +provides runnable Blob/Table backend selection. The +[Azure provider benchmarks](../../../test/Benchmarks/Journaling/Azure/README.md) +measure append, checkpoint, recovery, and catalog workloads with bounded work +and explicit resource ownership. + #### Using Connection String ```csharp using Azure.Storage.Blobs; @@ -392,14 +406,15 @@ public enum OrderStatus ## How It Works ### Storage Architecture -1. **Blob Container**: All jobs are stored in a single Azure Blob Storage container -2. **Append Blobs**: Each job shard is stored as an append blob, providing efficient sequential writes -3. **Blob Naming**: Blobs are named with the pattern: `{ShardStartTime:yyyyMMddHHmm}-{SiloAddress}-{Index}` -4. **Metadata**: Blob metadata stores ownership and time range information: - - `Owner`: The silo currently processing this shard - - `Creator`: The silo that created this shard - - `MinDueTime`: Start of the time range for jobs in this shard - - `MaxDueTime`: End of the time range for jobs in this shard +Each durable-job shard has a journal identity and persisted ownership metadata. +The journal-backed shard manager discovers shard identities through the storage +catalog and recovers state through the configured journal format. + +The Blob provider stores the WAL at `wal/{journalId}` and checkpoints at +`checkpoints/{journalId}/{snapshotId}` within the configured container. The +Table provider stores journal headers and data generations in the configured +table. Provider metadata selects the committed checkpoint or generation; caller +metadata carries the durable-job shard's discovery and ownership properties. ### Shard Ownership and High Availability 1. **Optimistic Concurrency**: ETags prevent conflicting updates when multiple silos try to claim a shard @@ -410,12 +425,12 @@ public enum OrderStatus ### Job Lifecycle with Azure Storage ``` ┌─────────────────────┐ -│ Job Scheduled │ ──▶ Written to append blob +│ Job Scheduled │ ──▶ Committed to the shard journal └─────────────────────┘ │ ▼ ┌─────────────────────┐ -│ Waiting in Shard │ ──▶ Persisted in Azure Blob Storage +│ Waiting in Shard │ ──▶ Persisted in Azure journal storage └─────────────────────┘ │ ▼ @@ -428,9 +443,9 @@ public enum OrderStatus │ Job Executed │ ──▶ Handler invoked on target grain └─────────────────────┘ │ - ├──▶ Success ──▶ Job entry removed from blob + ├──▶ Success ──▶ Completion persisted in the journal │ - └──▶ Failure ──▶ Retry: Updated due time in blob + └──▶ Failure ──▶ Retry: Updated due time in the journal No Retry: Job entry removed ``` @@ -446,12 +461,21 @@ services.Configure(options => ``` ### Storage Costs -- **Container**: One container per cluster -- **Blobs**: One blob per active time shard -- **Operations**: - - Schedule job: 1-2 append operations - - Execute job: 1 read + 1 delete operation - - Shard ownership transfer: 1 metadata update +Account for stored WAL/checkpoint bytes or Table generations, append batches, +recovery reads, checkpoint publication and cleanup, ownership metadata updates, +and catalog listing pages. Standard and premium Blob use the same provider +operations with different account pricing and service characteristics. + +Use `orleans-journaling-provider-catalog-pages`, `catalog-items`, and +`catalog-entries` (with the same prefix) to compare catalog traversal and +delivered results. `orleans-journaling-provider-retries` counts explicit +provider retries. The benchmarks pair these counters with independently +measured operation counts, outcomes, payload bytes, and latency. + +Configure request counts, timing, and outcomes through host-owned Azure SDK +diagnostics, Aspire integrations, or other application instrumentation. +Reconcile transaction cost with Azure service metrics and the SDK's retry and +upload behavior. ## Monitoring and Troubleshooting @@ -495,7 +519,7 @@ var blobServiceClient = new BlobServiceClient(storageAccountUri, credential); For more comprehensive documentation, please refer to: - [Microsoft Orleans Documentation](https://dotnet.github.io/orleans/docs/) - [Azure Blob Storage Documentation](https://learn.microsoft.com/azure/storage/blobs/) -- [Orleans Durable Jobs Core Package](../../../Orleans.DurableJobs/README.md) +- [Orleans Durable Jobs Core Package](../../Orleans.DurableJobs/README.md) ## Feedback & Contributing - If you have any issues or would like to provide feedback, please [open an issue on GitHub](https://github.com/dotnet/orleans/issues) diff --git a/src/api/Azure/Orleans.DurableJobs.AzureStorage/Orleans.DurableJobs.AzureStorage.cs b/src/api/Azure/Orleans.DurableJobs.AzureStorage/Orleans.DurableJobs.AzureStorage.cs index a7f71e0c87a..40d9dba65ce 100644 --- a/src/api/Azure/Orleans.DurableJobs.AzureStorage/Orleans.DurableJobs.AzureStorage.cs +++ b/src/api/Azure/Orleans.DurableJobs.AzureStorage/Orleans.DurableJobs.AzureStorage.cs @@ -13,5 +13,9 @@ public static partial class AzureStorageDurableJobsExtensions public static Microsoft.Extensions.DependencyInjection.IServiceCollection UseAzureBlobDurableJobs(this Microsoft.Extensions.DependencyInjection.IServiceCollection services, System.Action configure) { throw null; } public static ISiloBuilder UseAzureBlobDurableJobs(this ISiloBuilder builder, System.Action configure) { throw null; } + + public static Microsoft.Extensions.DependencyInjection.IServiceCollection UseAzureTableDurableJobs(this Microsoft.Extensions.DependencyInjection.IServiceCollection services, System.Action configure) { throw null; } + + public static ISiloBuilder UseAzureTableDurableJobs(this ISiloBuilder builder, System.Action configure) { throw null; } } } \ No newline at end of file diff --git a/test/Benchmarks/Benchmarks.csproj b/test/Benchmarks/Benchmarks.csproj index 22a77b5f72d..fa3e7c28528 100644 --- a/test/Benchmarks/Benchmarks.csproj +++ b/test/Benchmarks/Benchmarks.csproj @@ -41,6 +41,7 @@ + @@ -51,6 +52,7 @@ + diff --git a/test/Benchmarks/Journaling/Azure/AzureJournalBenchmarkConfigurationTests.cs b/test/Benchmarks/Journaling/Azure/AzureJournalBenchmarkConfigurationTests.cs new file mode 100644 index 00000000000..7804493a35c --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/AzureJournalBenchmarkConfigurationTests.cs @@ -0,0 +1,96 @@ +using BenchmarkDotNet.Configs; +using BenchmarkDotNet.ConsoleArguments; +using BenchmarkDotNet.Jobs; +using BenchmarkDotNet.Loggers; +using BenchmarkDotNet.Running; +using BenchmarkDotNet.Validators; +using TestExtensions; +using Xunit; + +namespace Benchmarks.Journaling.Azure; + +[TestCategory("BVT")] +public class AzureJournalBenchmarkConfigurationTests +{ + [Theory] + [InlineData("default")] + [InlineData("dry")] + [InlineData("short")] + [InlineData("medium")] + [InlineData("long")] + public void BdnResolvedPresetsHaveFiniteWorkLimits(string preset) + { + var parsed = ConfigParser.Parse(["--job", preset, "--inProcess"], NullLogger.Instance); + Assert.True(parsed.isSuccess); + using var benchmarks = BenchmarkConverter.TypeToBenchmarks(typeof(AzureJournalBenchmarks), parsed.config); + Assert.NotEmpty(benchmarks.BenchmarksCases); + Assert.Single(benchmarks.Config.GetValidators().OfType()); + Assert.Empty(new AzureJournalBenchmarkValidator().Validate(benchmarks)); + foreach (var benchmark in benchmarks.BenchmarksCases) + { + Assert.InRange(benchmark.Job.Run.IterationCount, 1, 3); + Assert.InRange(benchmark.Job.Run.WarmupCount, 0, 1); + Assert.Equal(1, benchmark.Job.Run.InvocationCount); + Assert.Equal(1, benchmark.Job.Run.UnrollFactor); + Assert.Equal(1, benchmark.Job.Run.LaunchCount); + } + } + + [Theory] + [InlineData("--iterationCount", "100")] + [InlineData("--warmupCount", "100")] + [InlineData("--launchCount", "100")] + [InlineData("--invocationCount", "100")] + [InlineData("--unrollFactor", "100")] + public void BdnCliOverridesAreBoundedOrRejectedBeforeExecution(string name, string value) + { + var parsed = ConfigParser.Parse(["--job", "Dry", name, value], NullLogger.Instance); + Assert.True(parsed.isSuccess); + using var benchmarks = BenchmarkConverter.TypeToBenchmarks(typeof(AzureJournalBenchmarks), parsed.config); + var validator = Assert.Single(benchmarks.Config.GetValidators().OfType()); + var errors = validator.Validate(benchmarks).ToArray(); + if (errors.Length > 0) + { + Assert.All(errors, error => Assert.True(error.IsCritical)); + } + else + { + Assert.All(benchmarks.BenchmarksCases, benchmark => + { + Assert.InRange(benchmark.Job.Run.IterationCount, 1, 3); + Assert.InRange(benchmark.Job.Run.WarmupCount, 0, 1); + Assert.Equal(1, benchmark.Job.Run.InvocationCount); + Assert.Equal(1, benchmark.Job.Run.UnrollFactor); + Assert.Equal(1, benchmark.Job.Run.LaunchCount); + }); + } + } + + [Theory] + [InlineData("iterations")] + [InlineData("warmup")] + [InlineData("launches")] + [InlineData("invocations")] + [InlineData("unroll")] + [InlineData("adaptive")] + public void BdnValidatorRejectsUnsafeFinalJobs(string change) + { + using var benchmarks = BenchmarkConverter.TypeToBenchmarks(typeof(AzureJournalBenchmarks)); + var original = benchmarks.BenchmarksCases[0]; + var job = change switch + { + "iterations" => original.Job.WithIterationCount(100), + "warmup" => original.Job.WithWarmupCount(100), + "launches" => original.Job.WithLaunchCount(100), + "invocations" => original.Job.WithInvocationCount(100), + "unroll" => original.Job.WithUnrollFactor(100), + _ => Job.Default + }; + var candidate = BenchmarkCase.Create(original.Descriptor, job, original.Parameters, original.Config); + var error = Assert.Single(new AzureJournalBenchmarkValidator().Validate( + new ValidationParameters([candidate], original.Config))); + Assert.True(error.IsCritical); + Assert.Same(candidate, error.BenchmarkCase); + Assert.Contains("Azure journal benchmarks require", error.Message); + } +} diff --git a/test/Benchmarks/Journaling/Azure/AzureJournalBenchmarkValidator.cs b/test/Benchmarks/Journaling/Azure/AzureJournalBenchmarkValidator.cs new file mode 100644 index 00000000000..c6b00f95664 --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/AzureJournalBenchmarkValidator.cs @@ -0,0 +1,36 @@ +using BenchmarkDotNet.Characteristics; +using BenchmarkDotNet.Engines; +using BenchmarkDotNet.Environments; +using BenchmarkDotNet.Jobs; +using BenchmarkDotNet.Validators; + +namespace Benchmarks.Journaling.Azure; + +internal sealed class AzureJournalBenchmarkValidator : IValidator +{ + public bool TreatsWarningsAsErrors => true; + + public IEnumerable Validate(ValidationParameters validationParameters) + { + var resolver = new CompositeResolver(EnvironmentResolver.Instance, EngineResolver.Instance); + foreach (var benchmark in validationParameters.Benchmarks) + { + if (benchmark.Descriptor.Type != typeof(AzureJournalBenchmarks)) + { + continue; + } + + var run = benchmark.Job.Run; + if (!run.HasValue(RunMode.IterationCountCharacteristic) || run.IterationCount is < 1 or > 3 + || !run.HasValue(RunMode.WarmupCountCharacteristic) || run.WarmupCount is < 0 or > 1 + || !run.HasValue(RunMode.InvocationCountCharacteristic) || run.InvocationCount != 1 + || run.ResolveValue(RunMode.UnrollFactorCharacteristic, resolver) != 1 + || !run.HasValue(RunMode.LaunchCountCharacteristic) || run.LaunchCount != 1) + { + yield return new ValidationError(true, + "Azure journal benchmarks require 1-3 measured iterations, 0-1 warmup iterations, one invocation per iteration, unroll factor 1, and one launch. Use Journaling.Azure for larger fixed-work runs.", + benchmark); + } + } + } +} diff --git a/test/Benchmarks/Journaling/Azure/AzureJournalBenchmarks.cs b/test/Benchmarks/Journaling/Azure/AzureJournalBenchmarks.cs new file mode 100644 index 00000000000..fa57c571b30 --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/AzureJournalBenchmarks.cs @@ -0,0 +1,151 @@ +using BenchmarkDotNet.Attributes; +using BenchmarkDotNet.Configs; +using BenchmarkDotNet.Engines; +using BenchmarkDotNet.Jobs; + +namespace Benchmarks.Journaling.Azure; + +[BenchmarkCategory("Journaling.Azure")] +[Config(typeof(Configuration))] +[InvocationCount(1, unrollFactor: 1)] +[IterationCount(3)] +[WarmupCount(1)] +public class AzureJournalBenchmarks +{ + public sealed class Configuration : ManualConfig + { + public Configuration() + { + AddJob(Job.Default.WithStrategy(RunStrategy.Monitoring).WithWarmupCount(1).WithIterationCount(3).AsDefault()); + AddJob(Job.Default.DontEnforcePowerPlan().WithLaunchCount(1).AsMutator()); + AddValidator(new AzureJournalBenchmarkValidator()); + } + } + + private AzureJournalScenario _scenario = null!; + private AzureJournalReport _report = null!; + private CancellationTokenSource _iteration = null!; + private int _invocations; + + [ParamsSource(nameof(Backends))] + public AzureJournalBackend Backend { get; set; } + + [ParamsSource(nameof(PayloadSizes))] + public int PayloadBytes { get; set; } + + [ParamsAllValues] + public AzureJournalWorkload Workload { get; set; } + + public IEnumerable Backends => (Environment.GetEnvironmentVariable("JOURNAL_BENCHMARK_BDN_BACKENDS") ?? "AzuriteBlob,AzuriteTable") + .Split(',').Select(AzureJournalOptions.ParseEnum).Distinct().Take(6).ToArray() is { Length: <= 5 } backends + ? backends : throw new ArgumentException("Select at most five backends."); + + public IEnumerable PayloadSizes => (Environment.GetEnvironmentVariable("JOURNAL_BENCHMARK_BDN_PAYLOAD_BYTES") ?? "256") + .Split(',').Select(value => int.Parse(value, System.Globalization.CultureInfo.InvariantCulture)).Distinct().Take(4).ToArray() is { Length: <= 3 } sizes + ? sizes : throw new ArgumentException("Select at most three payload sizes."); + + [GlobalSetup] + public void ReportBuild() => BenchmarkBuildInfo.WriteTo(Console.Out); + + [IterationSetup] + public void Setup() + { + var options = AzureJournalOptions.Parse((Environment.GetEnvironmentVariable("JOURNAL_BENCHMARK_BDN_OPTIONS") ?? "").Split(' ', StringSplitOptions.RemoveEmptyEntries)) with + { + Backend = Backend, + PayloadBytes = PayloadBytes, + Workload = Workload, + Operations = 1, + Concurrency = 1, + AllowAzure = Environment.GetEnvironmentVariable("JOURNAL_BENCHMARK_ALLOW_AZURE") == "true" + }; + options.Validate(); + _report = new AzureJournalReport(options); + _scenario = new AzureJournalScenario(options, _report); + _invocations = 0; + try + { + using var setup = new CancellationTokenSource(TimeSpan.FromSeconds(options.SetupTimeoutSeconds)); + _scenario.PrepareAsync(setup.Token).GetAwaiter().GetResult(); + _iteration = new CancellationTokenSource(TimeSpan.FromSeconds(options.TimeoutSeconds)); + } + catch (Exception exception) + { + try + { + Cleanup(); + } + finally + { + Console.Error.WriteLine($"Azure benchmark setup failed: {exception.GetType().Name}."); + } + + throw new InvalidOperationException("Azure benchmark setup failed. See the sanitized phase and resource output."); + } + } + + [Benchmark] + public async Task Execute() + { + if (Interlocked.Increment(ref _invocations) != 1) + { + throw new InvalidOperationException("Azure benchmarks require exactly one invocation per iteration."); + } + + try + { + await _scenario.ExecuteAsync(0, _iteration.Token); + } + catch (Exception exception) + { + Console.Error.WriteLine($"Azure benchmark measurement failed: {exception.GetType().Name}."); + throw new InvalidOperationException("Azure benchmark measurement failed. See the sanitized phase and resource output."); + } + + return _report.Configuration.ItemsPerOperation; + } + + [IterationCleanup] + public void Cleanup() + { + try + { + if (_invocations == 1) + { + using var verification = new CancellationTokenSource(TimeSpan.FromSeconds(_report.Configuration.SetupTimeoutSeconds)); + try + { + _scenario.VerifyAsync([0], verification.Token).GetAwaiter().GetResult(); + } + catch (Exception exception) + { + Console.Error.WriteLine($"Azure benchmark verification failed: {exception.GetType().Name}."); + throw new InvalidOperationException("Azure benchmark verification failed."); + } + } + } + finally + { + _iteration?.Dispose(); + using var cleanup = new CancellationTokenSource(TimeSpan.FromSeconds(_report.Configuration.CleanupTimeoutSeconds)); + try + { + _scenario.CleanupAsync(cleanup.Token).GetAwaiter().GetResult(); + if (_report.Failures.Count > 0) + { + throw new InvalidOperationException("Azure benchmark lifecycle cleanup failed."); + } + } + catch (Exception exception) + { + _report.Cleanup = "failed"; + Console.Error.WriteLine($"Azure benchmark cleanup failed: {exception.GetType().Name}."); + throw new InvalidOperationException("Azure benchmark cleanup failed. Inspect the owned resource identified below."); + } + finally + { + Console.WriteLine($"Azure benchmark resource={_report.Resource}; ownership={_report.Ownership}; cleanup={_report.Cleanup}; account={_report.AccountVerification}/{_report.AccountKind}/{_report.AccountSku}"); + } + } + } +} diff --git a/test/Benchmarks/Journaling/Azure/AzureJournalOptions.cs b/test/Benchmarks/Journaling/Azure/AzureJournalOptions.cs new file mode 100644 index 00000000000..49536af6275 --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/AzureJournalOptions.cs @@ -0,0 +1,177 @@ +using System.Globalization; +using System.Text.Json.Serialization; + +namespace Benchmarks.Journaling.Azure; + +public enum AzureJournalBackend +{ + AzuriteBlob, + AzuriteTable, + StandardBlob, + PremiumBlob, + Table +} + +public enum AzureJournalWorkload +{ + DurableAppend, + CheckpointReplace, + RecoveryReplay, + CatalogBounded, + CatalogUnbounded +} + +public sealed record AzureJournalOptions +{ + public AzureJournalBackend Backend { get; init; } = AzureJournalBackend.AzuriteBlob; + public AzureJournalWorkload Workload { get; init; } = AzureJournalWorkload.DurableAppend; + public int Operations { get; init; } = 8; + public int Concurrency { get; init; } = 2; + public int PayloadBytes { get; init; } = 256; + public int BatchSize { get; init; } = 4; + public int HistoryBatches { get; init; } = 4; + public int CheckpointBytes { get; init; } = 4096; + public int DueJournals { get; init; } = 16; + public int FutureJournals { get; init; } = 64; + public int Seed { get; init; } = 42; + public bool IncludeMetadata { get; init; } + public bool AllowAzure { get; init; } + public int SetupTimeoutSeconds { get; init; } = 120; + public int TimeoutSeconds { get; init; } = 60; + public int CleanupTimeoutSeconds { get; init; } = 60; + public long MaxWork { get; init; } = 100_000; + public long MaxBytes { get; init; } = 256L * 1024 * 1024; + [JsonIgnore] public string Output { get; init; } = "azure-journal-results"; + + public bool IsEmulator => Backend is AzureJournalBackend.AzuriteBlob or AzureJournalBackend.AzuriteTable; + public bool IsTable => Backend is AzureJournalBackend.Table or AzureJournalBackend.AzuriteTable; + public bool IsCatalog => Workload is AzureJournalWorkload.CatalogBounded or AzureJournalWorkload.CatalogUnbounded; + public int AppendBytes => checked(PayloadBytes * BatchSize); + public long RecoveryBytes => CheckpointBytes + (long)AppendBytes * HistoryBatches; + public long PayloadBytesPerOperation => Workload switch + { + AzureJournalWorkload.DurableAppend => AppendBytes, + AzureJournalWorkload.CheckpointReplace => CheckpointBytes, + AzureJournalWorkload.RecoveryReplay => RecoveryBytes, + _ => 0 + }; + public long ItemsPerOperation => Workload switch + { + AzureJournalWorkload.DurableAppend => BatchSize, + AzureJournalWorkload.CheckpointReplace => 1, + AzureJournalWorkload.RecoveryReplay => 1L + (long)HistoryBatches * BatchSize, + AzureJournalWorkload.CatalogBounded => DueJournals, + AzureJournalWorkload.CatalogUnbounded => (long)DueJournals + FutureJournals, + _ => throw new InvalidOperationException("Unknown workload.") + }; + + // Includes seeding, warmup, measured calls, verification, and catalog traversal work. + public long EstimatedWork => IsCatalog + ? 2L * (DueJournals + FutureJournals + 2) + (Operations + 1L) * (DueJournals + FutureJournals + 1L) + : (Operations + 1L) * (HistoryBatches + 6L); + public long EstimatedPayloadBytes => IsCatalog + ? (DueJournals + FutureJournals + 2L) * 128 + (Operations + 1L) * (DueJournals + FutureJournals) * 128 + : (Operations + 1L) * (2 * RecoveryBytes + 2 * PayloadBytesPerOperation); + + public void Validate() + { + if (!Enum.IsDefined(Backend) || !Enum.IsDefined(Workload)) + { + throw new ArgumentException("Select a named backend and workload."); + } + + CheckRange(Operations, 1, 10_000, nameof(Operations)); + CheckRange(Concurrency, 1, Math.Min(Operations, 256), nameof(Concurrency)); + CheckRange(PayloadBytes, 1, 2 * 1024 * 1024, nameof(PayloadBytes)); + CheckRange(BatchSize, 1, 8192, nameof(BatchSize)); + CheckRange(HistoryBatches, 0, 10_000, nameof(HistoryBatches)); + CheckRange(CheckpointBytes, 1, 32 * 1024 * 1024, nameof(CheckpointBytes)); + CheckRange(DueJournals, 0, 100_000, nameof(DueJournals)); + CheckRange(FutureJournals, 0, 100_000, nameof(FutureJournals)); + CheckRange(SetupTimeoutSeconds, 1, 3600, nameof(SetupTimeoutSeconds)); + CheckRange(TimeoutSeconds, 1, 3600, nameof(TimeoutSeconds)); + CheckRange(CleanupTimeoutSeconds, 1, 600, nameof(CleanupTimeoutSeconds)); + CheckRange(MaxWork, 1, 10_000_000, nameof(MaxWork)); + CheckRange(MaxBytes, 1, 8L * 1024 * 1024 * 1024, nameof(MaxBytes)); + if ((long)PayloadBytes * BatchSize > 2 * 1024 * 1024) + { + throw new ArgumentException("PayloadBytes * BatchSize must fit the common 2 MiB Azure Table append limit."); + } + + if (EstimatedWork > MaxWork || EstimatedPayloadBytes > MaxBytes) + { + throw new ArgumentException("Estimated setup, measurement, and verification work exceeds MaxWork or MaxBytes."); + } + + if (!IsEmulator && !AllowAzure) + { + throw new ArgumentException("Real Azure requires --allow-azure true and incurs storage charges."); + } + + ArgumentException.ThrowIfNullOrWhiteSpace(Output); + } + + internal static AzureJournalOptions Parse(string[] args) + { + var result = new AzureJournalOptions(); + var seen = new HashSet(StringComparer.Ordinal); + for (var index = 0; index < args.Length; index += 2) + { + if (index + 1 == args.Length || !seen.Add(args[index])) + { + throw new ArgumentException("Every option requires one value and may occur once."); + } + + var value = args[index + 1]; + result = args[index] switch + { + "--backend" => result with { Backend = ParseEnum(value) }, + "--workload" => result with { Workload = ParseEnum(value) }, + "--operations" => result with { Operations = ParseInt(value) }, + "--concurrency" => result with { Concurrency = ParseInt(value) }, + "--payload-bytes" => result with { PayloadBytes = ParseInt(value) }, + "--batch-size" => result with { BatchSize = ParseInt(value) }, + "--history-batches" => result with { HistoryBatches = ParseInt(value) }, + "--checkpoint-bytes" => result with { CheckpointBytes = ParseInt(value) }, + "--due-journals" => result with { DueJournals = ParseInt(value) }, + "--future-journals" => result with { FutureJournals = ParseInt(value) }, + "--seed" => result with { Seed = ParseInt(value) }, + "--metadata" => result with { IncludeMetadata = ParseBool(value) }, + "--allow-azure" => result with { AllowAzure = ParseBool(value) }, + "--setup-timeout-seconds" => result with { SetupTimeoutSeconds = ParseInt(value) }, + "--timeout-seconds" => result with { TimeoutSeconds = ParseInt(value) }, + "--cleanup-timeout-seconds" => result with { CleanupTimeoutSeconds = ParseInt(value) }, + "--max-work" => result with { MaxWork = ParseLong(value) }, + "--max-bytes" => result with { MaxBytes = ParseLong(value) }, + "--output" => result with { Output = value }, + _ => throw new ArgumentException("Unknown option. Use Journaling.Azure --help.") + }; + } + + result.Validate(); + return result; + } + + internal static T ParseEnum(string value) where T : struct, Enum + => Enum.GetNames().Contains(value, StringComparer.OrdinalIgnoreCase) && Enum.TryParse(value, true, out var result) + ? result : throw new ArgumentException("Expected a named enum value."); + + private static int ParseInt(string value) + => int.TryParse(value, NumberStyles.Integer, CultureInfo.InvariantCulture, out var result) + ? result : throw new ArgumentException("Expected an integer."); + + private static long ParseLong(string value) + => long.TryParse(value, NumberStyles.Integer, CultureInfo.InvariantCulture, out var result) + ? result : throw new ArgumentException("Expected an integer."); + + private static bool ParseBool(string value) + => bool.TryParse(value, out var result) ? result : throw new ArgumentException("Expected true or false."); + + private static void CheckRange(long value, long min, long max, string name) + { + if (value < min || value > max) + { + throw new ArgumentException($"{name} must be between {min} and {max}."); + } + } +} diff --git a/test/Benchmarks/Journaling/Azure/AzureJournalReport.cs b/test/Benchmarks/Journaling/Azure/AzureJournalReport.cs new file mode 100644 index 00000000000..257473f420e --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/AzureJournalReport.cs @@ -0,0 +1,134 @@ +using System.Globalization; +using System.Text; +using System.Text.Json; +using System.Text.Json.Serialization; +using Azure; + +namespace Benchmarks.Journaling.Azure; + +internal sealed record LatencySummary(int Count, double? MeanMs, double? P50Ms, double? P95Ms, double? P99Ms) +{ + public static LatencySummary Create(IEnumerable values) + { + var sorted = values.Order().ToArray(); + return sorted.Length == 0 ? new(0, null, null, null, null) : new( + sorted.Length, sorted.Average(), Percentile(sorted, .50), Percentile(sorted, .95), Percentile(sorted, .99)); + } + + private static double Percentile(double[] sorted, double percentile) + => sorted[(int)Math.Ceiling(sorted.Length * percentile) - 1]; +} + +internal enum BenchmarkValidationError +{ + JournalCollision, + CatalogUnexpectedIdentity, + CatalogMissingIdentity, + CatalogUnrequestedMetadata, + JournalFormatOrETag, + JournalCallerMetadata, + RecoveryAfterCompletion, + RecoveryFormat, + RecoveryLength, + RecoveryPayload, + RecoveryCompletionOrChecksum +} + +internal sealed class BenchmarkValidationException(BenchmarkValidationError code) : Exception(code.ToString()) +{ + public BenchmarkValidationError Code { get; } = code; +} + +internal sealed record BenchmarkFailure(string Phase, string ExceptionType, int? HttpStatus, int? Operation = null, BenchmarkValidationError? ValidationError = null) +{ + // Exception messages and URIs can contain credentials, SDK request bodies, or connection strings. + public static BenchmarkFailure From(string phase, Exception exception, int? operation = null) + => new(phase, exception.GetType().Name, (exception as RequestFailedException)?.Status, operation, (exception as BenchmarkValidationException)?.Code); +} + +internal sealed record OperationResult(int Index, string Outcome, double LatencyMs, long PayloadBytes, long Items); + +internal sealed class AzureJournalReport(AzureJournalOptions configuration) +{ + public int SchemaVersion => 2; + public string Source => "IJournalStorage / IJournalStorageCatalog"; + public string LoadModel => "closed-loop fixed work; one independent journal per mutation"; + public string LatencyUnit => "ms; Stopwatch around each awaited operation, including recovery/catalog validation"; + public string PayloadByteUnit => "caller payload bytes committed or replayed; catalog = 0"; + public string ItemUnit => Configuration.Workload switch + { + AzureJournalWorkload.DurableAppend => "payload records", + AzureJournalWorkload.CheckpointReplace => "checkpoints", + AzureJournalWorkload.RecoveryReplay => "checkpoint plus payload records", + _ => "catalog entries" + }; + public string TelemetryUnit => "counts of catalog pages, candidate items, delivered entries, and explicit provider retries"; + public AzureJournalOptions Configuration { get; } = configuration; + public BenchmarkBuildInfo Build { get; } = BenchmarkBuildInfo.Current; + public string Runtime { get; } = System.Runtime.InteropServices.RuntimeInformation.FrameworkDescription; + public string OperatingSystem { get; } = System.Runtime.InteropServices.RuntimeInformation.OSDescription; + public string? Resource { get; set; } + public string Ownership { get; set; } = "not-created"; + public string Cleanup { get; set; } = "not-required"; + public string AccountVerification { get; set; } = "not-checked"; + public string? AccountKind { get; set; } + public string? AccountSku { get; set; } + public string Phase { get; set; } = "setup"; + public double ElapsedSeconds { get; set; } + public List Operations { get; } = []; + public List Failures { get; } = []; + public IReadOnlyList ProviderMetrics { get; set; } = []; + public int Completed => Operations.Count(operation => operation.Outcome == "completed"); + public int Failed => Operations.Count(operation => operation.Outcome == "failed"); + public int Cancelled => Operations.Count(operation => operation.Outcome == "cancelled"); + public int NotStarted => Configuration.Operations - Operations.Count; + public long CompletedPayloadBytes => Operations.Where(operation => operation.Outcome == "completed").Sum(operation => operation.PayloadBytes); + public long CompletedItems => Operations.Where(operation => operation.Outcome == "completed").Sum(operation => operation.Items); + public double CompletedOperationsPerSecond => ElapsedSeconds > 0 ? Completed / ElapsedSeconds : 0; + public double PayloadBytesPerSecond => ElapsedSeconds > 0 ? CompletedPayloadBytes / ElapsedSeconds : 0; + public LatencySummary SuccessfulLatency => LatencySummary.Create(Operations.Where(operation => operation.Outcome == "completed").Select(operation => operation.LatencyMs)); + public LatencySummary FailedLatency => LatencySummary.Create(Operations.Where(operation => operation.Outcome == "failed").Select(operation => operation.LatencyMs)); + public LatencySummary CancelledLatency => LatencySummary.Create(Operations.Where(operation => operation.Outcome == "cancelled").Select(operation => operation.LatencyMs)); + public bool Success => Phase == "complete" && Failures.Count == 0 && Completed == Configuration.Operations && Cleanup == "deleted"; + + internal static JsonSerializerOptions JsonOptions { get; } = new() + { + WriteIndented = true, + Converters = { new JsonStringEnumConverter() } + }; + + public async Task ExportAsync(string prefix) + { + var json = JsonSerializer.Serialize(this, JsonOptions); + // JSON is the complete report; CSV carries the same nested configuration, failures and metrics as quoted JSON fields. + string[] headers = + [ + "schema_version", "backend", "workload", "phase", "success", "resource", "ownership", "cleanup", + "account_verification", "account_kind", "account_sku", "completed", "failed", "cancelled", "not_started", + "elapsed_seconds", "operations_per_second", "payload_bytes_per_second", "completed_payload_bytes", "completed_items", + "p50_ms", "p95_ms", "p99_ms", "item_unit", "configuration_json", "operations_json", "failures_json", "provider_metrics_json", "build_json" + ]; + object?[] values = + [ + SchemaVersion, Configuration.Backend, Configuration.Workload, Phase, Success, Resource, Ownership, Cleanup, + AccountVerification, AccountKind, AccountSku, Completed, Failed, Cancelled, NotStarted, + ElapsedSeconds, CompletedOperationsPerSecond, PayloadBytesPerSecond, CompletedPayloadBytes, CompletedItems, + SuccessfulLatency.P50Ms, SuccessfulLatency.P95Ms, SuccessfulLatency.P99Ms, ItemUnit, + JsonSerializer.Serialize(Configuration, JsonOptions), JsonSerializer.Serialize(Operations, JsonOptions), + JsonSerializer.Serialize(Failures, JsonOptions), JsonSerializer.Serialize(ProviderMetrics, JsonOptions), + JsonSerializer.Serialize(Build, JsonOptions) + ]; + var csv = string.Join(',', headers) + "\n" + string.Join(',', values.Select(CsvField)) + "\n"; + await WriteNewAsync(prefix + ".json", json); + await WriteNewAsync(prefix + ".csv", csv); + } + + private static string CsvField(object? value) + => "\"" + (Convert.ToString(value, CultureInfo.InvariantCulture) ?? "").Replace("\"", "\"\"", StringComparison.Ordinal) + "\""; + + private static async Task WriteNewAsync(string path, string content) + { + await using var stream = new FileStream(path, FileMode.CreateNew, FileAccess.Write, FileShare.None, 4096, useAsync: true); + await stream.WriteAsync(Encoding.UTF8.GetBytes(content)); + } +} diff --git a/test/Benchmarks/Journaling/Azure/AzureJournalRunner.cs b/test/Benchmarks/Journaling/Azure/AzureJournalRunner.cs new file mode 100644 index 00000000000..74b07012990 --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/AzureJournalRunner.cs @@ -0,0 +1,182 @@ +using System.Diagnostics; +using System.Text.Json; + +namespace Benchmarks.Journaling.Azure; + +internal static class AzureJournalRunner +{ + public static async Task RunCommandAsync(string[] args) + { + if (args is ["--help"]) + { + Console.WriteLine(""" + Journaling.Azure --backend AzuriteBlob|AzuriteTable|StandardBlob|PremiumBlob|Table + --workload DurableAppend|CheckpointReplace|RecoveryReplay|CatalogBounded|CatalogUnbounded + --operations 8 --concurrency 2 --payload-bytes 256 --batch-size 4 --history-batches 4 + --checkpoint-bytes 4096 --due-journals 16 --future-journals 64 --metadata false --seed 42 + --setup-timeout-seconds 120 --timeout-seconds 60 --cleanup-timeout-seconds 60 + --max-work 100000 --max-bytes 268435456 --allow-azure false --output azure-journal-results + Azure endpoints and emulator connection settings are environment-only. See Journaling\Azure\README.md. + Each run creates an isolated resource and new .json/.csv files. Existing output files are preserved. + """); + return 0; + } + + AzureJournalOptions options; + try + { + options = AzureJournalOptions.Parse(args); + } + catch (ArgumentException exception) + { + Console.Error.WriteLine($"Invalid benchmark configuration: {exception.Message} Use Journaling.Azure --help."); + return 2; + } + + using var cancellation = new CancellationTokenSource(); + ConsoleCancelEventHandler handler = (_, eventArgs) => + { + eventArgs.Cancel = true; + cancellation.Cancel(); + }; + Console.CancelKeyPress += handler; + try + { + var report = await RunAsync(options, cancellation.Token); + try + { + await report.ExportAsync(options.Output); + } + catch (Exception exception) + { + report.Failures.Add(BenchmarkFailure.From("export", exception)); + Console.WriteLine(JsonSerializer.Serialize(report, AzureJournalReport.JsonOptions)); + return 1; + } + + Console.WriteLine(JsonSerializer.Serialize(report, AzureJournalReport.JsonOptions)); + return report.Success ? 0 : 1; + } + finally + { + Console.CancelKeyPress -= handler; + } + } + + internal static async Task RunAsync( + AzureJournalOptions options, + CancellationToken cancellationToken = default, + Func? createScenario = null) + { + options.Validate(); + var report = new AzureJournalReport(options); + var scenario = createScenario is null ? new AzureJournalScenario(options, report) : createScenario(options, report); + try + { + using (var setup = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken)) + { + setup.CancelAfter(TimeSpan.FromSeconds(options.SetupTimeoutSeconds)); + await scenario.PrepareAsync(setup.Token); + } + + report.Phase = "measurement"; + using (var measurement = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken)) + { + measurement.CancelAfter(TimeSpan.FromSeconds(options.TimeoutSeconds)); + await MeasureAsync(scenario, options, report, measurement); + } + + if (report.Failures.Count == 0 && report.Completed == options.Operations) + { + report.Phase = "verification"; + using var verification = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + verification.CancelAfter(TimeSpan.FromSeconds(options.SetupTimeoutSeconds)); + await scenario.VerifyAsync(report.Operations.Where(result => result.Outcome == "completed").Select(result => result.Index), verification.Token); + report.Phase = "complete"; + } + } + catch (Exception exception) + { + report.Failures.Add(BenchmarkFailure.From(report.Phase, exception)); + } + finally + { + using var cleanup = new CancellationTokenSource(TimeSpan.FromSeconds(options.CleanupTimeoutSeconds)); + try + { + // Every worker has joined before ownership cleanup, including after caller cancellation. + await scenario.CleanupAsync(cleanup.Token); + } + catch (Exception exception) + { + report.Cleanup = "failed"; + report.Failures.Add(BenchmarkFailure.From("cleanup", exception)); + } + } + + return report; + } + + private static async Task MeasureAsync( + IAzureJournalScenario scenario, AzureJournalOptions options, AzureJournalReport report, CancellationTokenSource cancellation) + { + var results = new OperationResult?[options.Operations]; + var failures = new BenchmarkFailure?[options.Operations]; + var next = -1; + scenario.StartMetrics(); + var started = Stopwatch.GetTimestamp(); + try + { + // Allocate one task per worker, rather than one task per requested operation. + var workers = new Task[options.Concurrency]; + for (var worker = 0; worker < workers.Length; worker++) + { + workers[worker] = WorkAsync(); + } + + await Task.WhenAll(workers); + } + finally + { + report.ElapsedSeconds = Stopwatch.GetElapsedTime(started).TotalSeconds; + report.ProviderMetrics = scenario.StopMetrics(); + report.Operations.AddRange(results.OfType()); + report.Failures.AddRange(failures.OfType()); + if (cancellation.IsCancellationRequested && report.Failures.Count == 0) + { + report.Failures.Add(BenchmarkFailure.From("measurement", new OperationCanceledException(cancellation.Token))); + } + } + + async Task WorkAsync() + { + while (!cancellation.IsCancellationRequested) + { + var index = Interlocked.Increment(ref next); + if (index >= options.Operations) + { + return; + } + + var timestamp = Stopwatch.GetTimestamp(); + try + { + await scenario.ExecuteAsync(index, cancellation.Token); + results[index] = new(index, "completed", Stopwatch.GetElapsedTime(timestamp).TotalMilliseconds, + options.PayloadBytesPerOperation, options.ItemsPerOperation); + } + catch (OperationCanceledException exception) when (cancellation.IsCancellationRequested) + { + results[index] = new(index, "cancelled", Stopwatch.GetElapsedTime(timestamp).TotalMilliseconds, 0, 0); + failures[index] = BenchmarkFailure.From("measurement", exception, index); + } + catch (Exception exception) + { + results[index] = new(index, "failed", Stopwatch.GetElapsedTime(timestamp).TotalMilliseconds, 0, 0); + failures[index] = BenchmarkFailure.From("measurement", exception, index); + await cancellation.CancelAsync(); + } + } + } + } +} diff --git a/test/Benchmarks/Journaling/Azure/AzureJournalRunnerTests.cs b/test/Benchmarks/Journaling/Azure/AzureJournalRunnerTests.cs new file mode 100644 index 00000000000..346cd244e4b --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/AzureJournalRunnerTests.cs @@ -0,0 +1,826 @@ +using System.Diagnostics.Metrics; +using System.Reflection; +using System.Reflection.Emit; +using System.Text.Json; +using Azure; +using Azure.Core; +using Azure.Data.Tables; +using Azure.Data.Tables.Models; +using Azure.Storage.Blobs; +using Azure.Storage.Blobs.Models; +using DurableJobsJournaling; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging.Abstractions; +using Microsoft.Extensions.Options; +using Orleans; +using Orleans.Journaling; +using Orleans.Runtime; +using Orleans.Serialization.Buffers; +using TestExtensions; +using Xunit; + +namespace Benchmarks.Journaling.Azure; + +[TestCategory("BVT")] +public class AzureJournalRunnerTests +{ + [Theory] + [InlineData("--concurrency", "0")] + [InlineData("--operations", "10001")] + [InlineData("--payload-bytes", "2097153")] + [InlineData("--max-work", "1")] + [InlineData("--max-bytes", "1")] + [InlineData("--timeout-seconds", "0")] + [InlineData("--backend", "PremiumBlob")] + [InlineData("--backend", "2")] + [InlineData("--unknown", "value")] + public void ParserRejectsUnsafeWork(string name, string value) + => Assert.Throws(() => AzureJournalOptions.Parse([name, value])); + + [Fact] + public void ParserEnforcesCommonAppendLimitAndTotalWork() + { + var valid = AzureJournalOptions.Parse(["--payload-bytes", "1048576", "--batch-size", "2", "--history-batches", "0"]); + Assert.Equal(2 * 1024 * 1024, valid.AppendBytes); + Assert.Throws(() => AzureJournalOptions.Parse(["--payload-bytes", "1048577", "--batch-size", "2"])); + Assert.Throws(() => AzureJournalOptions.Parse(["--operations", "1", "--operations", "2"])); + Assert.Throws(() => AzureJournalOptions.Parse(["--operations"])); + Assert.Throws(() => AzureJournalOptions.Parse(["--workload", "CatalogUnbounded", "--future-journals", "100000"])); + } + + [Fact] + public void AccountVerificationDistinguishesEmulatorStandardAndPremium() + { + AzureJournalScenario.VerifyAccount(AzureJournalBackend.StandardBlob, AccountKind.StorageV2, SkuName.StandardLrs); + AzureJournalScenario.VerifyAccount(AzureJournalBackend.PremiumBlob, AccountKind.BlockBlobStorage, SkuName.PremiumLrs); + Assert.Throws(() => AzureJournalScenario.VerifyAccount(AzureJournalBackend.AzuriteBlob, AccountKind.StorageV2, SkuName.StandardLrs)); + Assert.Throws(() => AzureJournalScenario.VerifyAccount(AzureJournalBackend.AzuriteTable, AccountKind.StorageV2, SkuName.StandardLrs)); + Assert.Throws(() => AzureJournalScenario.VerifyAccount(AzureJournalBackend.StandardBlob, (AccountKind)int.MaxValue, SkuName.StandardLrs)); + Assert.Throws(() => AzureJournalScenario.VerifyAccount(AzureJournalBackend.PremiumBlob, AccountKind.BlockBlobStorage, (SkuName)int.MaxValue)); + Assert.Throws(() => AzureJournalScenario.VerifyAccount(AzureJournalBackend.StandardBlob, AccountKind.StorageV2, (SkuName)int.MaxValue)); + Assert.True(AzureJournalOptions.Parse(["--backend", "Table", "--allow-azure", "true"]).AllowAzure); + } + + [Theory] + [InlineData(AzureJournalBackend.StandardBlob, AccountKind.Storage, SkuName.StandardLrs)] + [InlineData(AzureJournalBackend.StandardBlob, AccountKind.BlobStorage, SkuName.StandardLrs)] + [InlineData(AzureJournalBackend.StandardBlob, AccountKind.StorageV2, SkuName.StandardLrs)] + [InlineData(AzureJournalBackend.StandardBlob, AccountKind.StorageV2, SkuName.StandardGrs)] + [InlineData(AzureJournalBackend.StandardBlob, AccountKind.StorageV2, SkuName.StandardRagrs)] + [InlineData(AzureJournalBackend.StandardBlob, AccountKind.StorageV2, SkuName.StandardZrs)] + [InlineData(AzureJournalBackend.PremiumBlob, AccountKind.BlockBlobStorage, SkuName.PremiumLrs)] + public void AccountVerificationAcceptsSupportedSdkSkus(AzureJournalBackend backend, AccountKind kind, SkuName sku) + => AzureJournalScenario.VerifyAccount(backend, kind, sku); + + [Theory] + [InlineData(AzureJournalBackend.PremiumBlob, AccountKind.StorageV2, SkuName.PremiumLrs)] + [InlineData(AzureJournalBackend.PremiumBlob, AccountKind.BlockBlobStorage, SkuName.StandardLrs)] + [InlineData(AzureJournalBackend.StandardBlob, AccountKind.StorageV2, SkuName.PremiumLrs)] + [InlineData(AzureJournalBackend.StandardBlob, AccountKind.FileStorage, SkuName.StandardLrs)] + [InlineData(AzureJournalBackend.StandardBlob, AccountKind.BlockBlobStorage, SkuName.StandardLrs)] + public void AccountVerificationRejectsMismatchedSdkKindAndTier(AzureJournalBackend backend, AccountKind kind, SkuName sku) + => Assert.Throws(() => AzureJournalScenario.VerifyAccount(backend, kind, sku)); + + [Theory] + [InlineData("https://account.blob.core.windows.net/?sig=secret", false)] + [InlineData("https://user:secret@account.blob.core.windows.net/", false)] + [InlineData("http://account.blob.core.windows.net/", false)] + [InlineData("https://account.blob.core.windows.net/", true)] + [InlineData("http://127.0.0.1:10000/devstoreaccount1", false)] + public void EndpointsEnforceAuthenticationAndEmulatorIsolation(string endpoint, bool emulator) + => Assert.Throws(() => AzureJournalScenario.ValidateEndpoint(new Uri(endpoint), emulator)); + + [Fact] + public void PercentilesUseNearestRankAndExcludeFailures() + { + var report = new AzureJournalReport(new() { Operations = 100, Concurrency = 1 }) { ElapsedSeconds = 2 }; + report.Operations.AddRange(Enumerable.Range(1, 99).Select(value => new OperationResult(value, "completed", value, 10, 1))); + report.Operations.Add(new(100, "failed", 10000, 0, 0)); + Assert.Equal(new LatencySummary(99, 50, 50, 95, 99), report.SuccessfulLatency); + Assert.Equal(49.5, report.CompletedOperationsPerSecond); + Assert.Equal(495, report.PayloadBytesPerSecond); + Assert.Equal(1, report.Failed); + Assert.Equal(0, report.NotStarted); + Assert.Null(LatencySummary.Create([]).P99Ms); + Assert.False(report.Success); + } + + [Fact] + public void MetricsUseExactMeterAndMeasurementWindow() + { + using var services = new ServiceCollection().AddMetrics().BuildServiceProvider(); + var instruments = new OrleansInstruments(services.GetRequiredService()); + using var foreign = new Meter("Microsoft.Orleans"); + var foreignCount = foreign.CreateCounter("orleans-journaling-provider-catalog-pages"); + using var collector = new ProviderMetrics(instruments.Meter); + var telemetry = new JournalStorageTelemetry(instruments); + telemetry.OnCatalogPage("orders-primary", 100); + collector.Start(); + foreignCount.Add(1000); + telemetry.OnCatalogPage("orders-primary", 0); + telemetry.OnCatalogPage("orders-primary", 3); + telemetry.OnCatalogEntry("orders-primary"); + telemetry.OnCatalogEntry("orders-primary"); + telemetry.OnRetry("orders-primary", "metadata_conflict"); + telemetry.OnRetry("orders-primary", "metadata_conflict"); + telemetry.OnRetry("orders-secondary", "metadata_conflict"); + var metrics = collector.Stop(); + telemetry.OnCatalogPage("orders-primary", 100); + Assert.Equal(new ProviderMetric[] + { + new("orleans-journaling-provider-catalog-entries", "count", "orders-primary", "", 2, 2), + new("orleans-journaling-provider-catalog-items", "count", "orders-primary", "", 1, 3), + new("orleans-journaling-provider-catalog-pages", "count", "orders-primary", "", 2, 2), + new("orleans-journaling-provider-retries", "count", "orders-primary", "metadata_conflict", 2, 2), + new("orleans-journaling-provider-retries", "count", "orders-secondary", "metadata_conflict", 1, 1) + }, metrics); + collector.Start(); + Assert.Empty(collector.Stop()); + } + + [Fact] + public void MetricsEnableCountersOnlyDuringExplicitCollectionWindows() + { + using var meter = new Meter("Microsoft.Orleans"); + var pages = meter.CreateCounter("orleans-journaling-provider-catalog-pages"); + using var collector = new ProviderMetrics(meter); + var entries = meter.CreateCounter("orleans-journaling-provider-catalog-entries"); + Assert.False(pages.Enabled); + Assert.False(entries.Enabled); + pages.Add(100); + collector.Start(); + Assert.True(pages.Enabled); + Assert.True(entries.Enabled); + Assert.Throws(collector.Start); + var retries = meter.CreateCounter("orleans-journaling-provider-retries"); + Assert.True(retries.Enabled); + pages.Add(1, new KeyValuePair("provider", "orders")); + entries.Add(2, new KeyValuePair("provider", "orders")); + Assert.Equal(new ProviderMetric[] + { + new("orleans-journaling-provider-catalog-entries", "count", "orders", "", 1, 2), + new("orleans-journaling-provider-catalog-pages", "count", "orders", "", 1, 1) + }, collector.Stop()); + Assert.False(pages.Enabled); + Assert.False(entries.Enabled); + Assert.False(retries.Enabled); + entries.Add(100); + collector.Start(); + Assert.True(pages.Enabled); + Assert.True(entries.Enabled); + Assert.True(retries.Enabled); + entries.Add(3, new KeyValuePair("provider", "orders")); + Assert.Equal(new ProviderMetric("orleans-journaling-provider-catalog-entries", "count", "orders", "", 1, 3), + Assert.Single(collector.Stop())); + collector.Start(); + Assert.True(pages.Enabled); + collector.Dispose(); + Assert.False(pages.Enabled); + Assert.False(entries.Enabled); + Assert.False(retries.Enabled); + } + + [Fact] + public void MetricsCollectOnlyCatalogCountersAndPreserveIntegerPrecision() + { + using var meter = new Meter("Microsoft.Orleans"); + using var collector = new ProviderMetrics(meter); + var items = meter.CreateCounter("orleans-journaling-provider-catalog-items"); + var oldCalls = meter.CreateCounter("orleans-journaling-provider-api-calls"); + var oldOperations = meter.CreateCounter("orleans-journaling-provider-operations"); + var duration = meter.CreateHistogram("orleans-journaling-provider-operation-duration", "ms"); + var wrongType = meter.CreateHistogram("orleans-journaling-provider-catalog-pages"); + var futureMetric = meter.CreateCounter("orleans-journaling-provider-unrelated"); + collector.Start(); + const long Count = 9_007_199_254_740_993; + items.Add(Count, new KeyValuePair("provider", "orders")); + items.Add(2, new KeyValuePair("provider", "orders")); + oldCalls.Add(100); + oldOperations.Add(100); + duration.Record(5); + wrongType.Record(10); + futureMetric.Add(100); + Assert.Equal(new ProviderMetric("orleans-journaling-provider-catalog-items", "count", "orders", "", 2, Count + 2), + Assert.Single(collector.Stop())); + } + + [Fact] + public async Task MetricsCaptureActualVolatileCatalogEntries() + { + using var services = new ServiceCollection().AddMetrics().BuildServiceProvider(); + var instruments = new OrleansInstruments(services.GetRequiredService()); + using var collector = new ProviderMetrics(instruments.Meter); + var provider = new VolatileJournalStorageProvider(Options.Create(new JournaledStateManagerOptions()), instruments); + var token = TestContext.Current.CancellationToken; + Assert.True(await provider.CreateStorage(new("due/a")).CreateIfNotExistsAsync(cancellationToken: token)); + Assert.True(await provider.CreateStorage(new("due/b")).CreateIfNotExistsAsync(cancellationToken: token)); + Assert.True(await provider.CreateStorage(new("future/c")).CreateIfNotExistsAsync(cancellationToken: token)); + collector.Start(); + var ids = new List(); + await foreach (var entry in provider.ListAsync(new() { Prefix = new("due/") }, token)) + { + ids.Add(entry.Id); + } + + Assert.Equal(new JournalId[] { new("due/a"), new("due/b") }, ids); + Assert.Equal(new ProviderMetric("orleans-journaling-provider-catalog-entries", "count", "volatile", "", 2, 2), + Assert.Single(collector.Stop())); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task MetricsCaptureActualAzureCatalogPagesCandidatesAndEntries(bool table) + { + var builder = new MetricsSiloBuilder(); + builder.Services.AddLogging().AddMetrics().AddSingleton(); + if (table) + { + builder.AddAzureTableJournalStorage(options => options.TableServiceClient = new CatalogTableService()); + } + else + { + builder.AddAzureBlobJournalStorage(options => options.BlobServiceClient = new CatalogBlobService()); + } + + await using var services = builder.Services.BuildServiceProvider(); + var meter = services.GetRequiredService().Meter; + var counters = new List(); + using var observer = new MeterListener + { + InstrumentPublished = (instrument, _) => + { + if (ReferenceEquals(instrument.Meter, meter) + && instrument.Name.StartsWith("orleans-journaling-provider-", StringComparison.Ordinal)) + { + counters.Add(instrument); + } + } + }; + observer.Start(); + using var collector = new ProviderMetrics(meter); + var catalog = services.GetRequiredService(); + Assert.Equal(4, counters.Count); + Assert.All(counters, counter => Assert.False(counter.Enabled)); + var lifecycle = new SiloLifecycleSubject(NullLogger.Instance); + foreach (var participant in services.GetServices>()) + { + participant.Participate(lifecycle); + } + + var token = TestContext.Current.CancellationToken; + await lifecycle.OnStart(token); + try + { + var options = new ListOptions { Prefix = new("catalog/"), MaxId = new("catalog/b") }; + Assert.Equal(new JournalId[] { new("catalog/a"), new("catalog/b") }, await ReadIdsAsync()); + Assert.All(counters, counter => Assert.False(counter.Enabled)); + collector.Start(); + Assert.All(counters, counter => Assert.True(counter.Enabled)); + Assert.Equal(new JournalId[] { new("catalog/a"), new("catalog/b") }, await ReadIdsAsync()); + var provider = table ? "azure_table" : "azure_blob"; + Assert.Equal(new ProviderMetric[] + { + new("orleans-journaling-provider-catalog-entries", "count", provider, "", 2, 2), + new("orleans-journaling-provider-catalog-items", "count", provider, "", 1, 3), + new("orleans-journaling-provider-catalog-pages", "count", provider, "", 2, 2) + }, collector.Stop()); + Assert.All(counters, counter => Assert.False(counter.Enabled)); + Assert.Equal(new JournalId[] { new("catalog/a"), new("catalog/b") }, await ReadIdsAsync()); + Assert.All(counters, counter => Assert.False(counter.Enabled)); + + async Task> ReadIdsAsync() + { + var ids = new List(); + await foreach (var entry in catalog.ListAsync(options, token)) + { + ids.Add(entry.Id); + } + + return ids; + } + } + finally + { + await lifecycle.OnStop(token); + } + } + + [Fact] + public async Task RunnerJoinsBoundedWorkersBeforeCleanupAndPreservesCancellation() + { + using var cancellation = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken); + var entered = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var active = 0; + var cleanupCalled = false; + var task = AzureJournalRunner.RunAsync(new() { Operations = 6, Concurrency = 2 }, cancellation.Token, (_, report) => + new FakeScenario(report) + { + Execute = async (_, token) => + { + if (Interlocked.Increment(ref active) == 2) + { + entered.SetResult(); + } + + try + { + await Task.Delay(Timeout.InfiniteTimeSpan, token); + } + finally + { + Interlocked.Decrement(ref active); + } + }, + CleanupAction = token => + { + Assert.False(token.IsCancellationRequested); + Assert.Equal(0, active); + cleanupCalled = true; + return Task.CompletedTask; + } + }); + await entered.Task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + await cancellation.CancelAsync(); + var result = await task.WaitAsync(TimeSpan.FromSeconds(10), TestContext.Current.CancellationToken); + Assert.True(cleanupCalled); + Assert.Equal(0, result.Completed); + Assert.Equal(0, result.Failed); + Assert.Equal(2, result.Cancelled); + Assert.Equal(4, result.NotStarted); + Assert.Equal("deleted", result.Cleanup); + Assert.False(result.Success); + Assert.All(result.Failures, failure => Assert.Equal("measurement", failure.Phase)); + } + + [Fact] + public async Task RunnerReportsPartialFailureAndIndependentCleanupFailure() + { + var result = await AzureJournalRunner.RunAsync(new() { Operations = 4, Concurrency = 1 }, TestContext.Current.CancellationToken, createScenario: (_, report) => + new FakeScenario(report) + { + Execute = (index, _) => index == 1 ? throw new InvalidDataException("sensitive message") : Task.CompletedTask, + CleanupAction = _ => throw new TimeoutException("secret endpoint") + }); + Assert.Equal(1, result.Completed); + Assert.Equal(1, result.Failed); + Assert.Equal(2, result.NotStarted); + Assert.Equal("measurement", result.Phase); + Assert.Equal("failed", result.Cleanup); + Assert.Collection(result.Failures, + failure => Assert.Equal(new BenchmarkFailure("measurement", "InvalidDataException", null, 1), failure), + failure => Assert.Equal(new BenchmarkFailure("cleanup", "TimeoutException", null), failure)); + Assert.False(result.Success); + } + + [Fact] + public async Task SetupFailureStillCleansOwnedResourceAndMeasuresNothing() + { + var result = await AzureJournalRunner.RunAsync(new(), TestContext.Current.CancellationToken, createScenario: (_, report) => + new FakeScenario(report) { Prepare = _ => throw new InvalidDataException() }); + Assert.Equal("setup", Assert.Single(result.Failures).Phase); + Assert.Empty(result.Operations); + Assert.Empty(result.ProviderMetrics); + Assert.Equal(0, result.ElapsedSeconds); + Assert.Equal("deleted", result.Cleanup); + } + + [Theory] + [InlineData(false)] + [InlineData(true)] + public async Task PartialProviderInitializationStopsAttemptedStagesBeforeDisposal(bool cancel) + { + var events = new List(); + using var cancellation = CancellationTokenSource.CreateLinkedTokenSource(TestContext.Current.CancellationToken); + var builder = new MetricsSiloBuilder(); + builder.Services.AddLogging().AddMetrics().AddSingleton(); + builder.AddAzureBlobJournalStorage(options => options.BlobServiceClient = new CatalogBlobService()); + var startupFailure = new InvalidOperationException("startup failed"); + builder.Services.AddSingleton>(_ => new FailingLifecycleParticipant( + events, () => + { + if (cancel) + { + cancellation.Cancel(); + cancellation.Token.ThrowIfCancellationRequested(); + } + + throw startupFailure; + })); + var options = new AzureJournalOptions(); + var report = new AzureJournalReport(options); + var scenario = new AzureJournalScenario(options, report); + try + { + var error = await Record.ExceptionAsync(() => scenario.InitializeProviderAsync(builder.Services, cancellation.Token)); + if (cancel) + { + Assert.IsAssignableFrom(error); + } + else + { + Assert.Same(startupFailure, error); + } + + report.Failures.Add(BenchmarkFailure.From(report.Phase, error!)); + Assert.Equal(new[] { "start:early", "start:failed" }, events); + await scenario.CleanupAsync(TestContext.Current.CancellationToken); + Assert.Equal(new[] { "start:early", "start:failed", "stop:failed", "stop:early", "dispose" }, events); + Assert.Equal("initialize", Assert.Single(report.Failures).Phase); + Assert.False(report.Success); + await scenario.CleanupAsync(TestContext.Current.CancellationToken); + Assert.Equal(5, events.Count); + } + finally + { + await scenario.CleanupAsync(TestContext.Current.CancellationToken); + } + } + + [Fact] + public async Task VerificationFailureRetainsCompletedOperationsAndFailsTheRun() + { + var result = await AzureJournalRunner.RunAsync(new() { Operations = 2 }, TestContext.Current.CancellationToken, createScenario: (_, report) => + new FakeScenario(report) { Verify = _ => throw new BenchmarkValidationException(BenchmarkValidationError.RecoveryCompletionOrChecksum) }); + Assert.Equal(2, result.Completed); + Assert.Equal(0, result.NotStarted); + Assert.Equal("verification", result.Phase); + Assert.Equal("verification", Assert.Single(result.Failures).Phase); + Assert.Equal(BenchmarkValidationError.RecoveryCompletionOrChecksum, result.Failures[0].ValidationError); + Assert.Equal("deleted", result.Cleanup); + Assert.False(result.Success); + } + + [Fact] + public async Task ResourceCollisionAndUnconfirmedCreationNeverDelete() + { + foreach (var exception in new Exception[] { new RequestFailedException(409, "collision"), new OperationCanceledException() }) + { + var report = new AzureJournalReport(new()); + var deleted = false; + var resource = new OwnedBenchmarkResource(report, _ => Task.FromException(exception), _ => + { + deleted = true; + return Task.CompletedTask; + }); + Assert.Same(exception, await Record.ExceptionAsync(() => resource.CreateAsync(CancellationToken.None))); + await resource.CleanupAsync(CancellationToken.None); + Assert.False(deleted); + Assert.Equal("creation-unconfirmed", report.Ownership); + Assert.Equal("manual-check-required", report.Cleanup); + } + } + + [Fact] + public async Task LostCreateResponseFollowedByRetryConflictRequiresManualInspection() + { + var report = new AzureJournalReport(new()) { Resource = "journalbenchownedbythisattempt" }; + var existsOnServer = false; + var deletes = 0; + var resource = new OwnedBenchmarkResource(report, _ => + { + existsOnServer = true; + // The SDK observes a conflict on retry after losing the successful create response. + return Task.FromException(new RequestFailedException(409, "ResourceAlreadyExists")); + }, _ => + { + deletes++; + return Task.CompletedTask; + }); + var failure = await Assert.ThrowsAsync(() => resource.CreateAsync(TestContext.Current.CancellationToken)); + report.Failures.Add(BenchmarkFailure.From("setup", failure)); + await resource.CleanupAsync(TestContext.Current.CancellationToken); + Assert.True(existsOnServer); + Assert.Equal(0, deletes); + Assert.Equal("creation-unconfirmed", report.Ownership); + Assert.Equal("manual-check-required", report.Cleanup); + Assert.Equal("journalbenchownedbythisattempt", report.Resource); + Assert.Equal(409, Assert.Single(report.Failures).HttpStatus); + Assert.False(report.Success); + } + + [Fact] + public async Task AcknowledgedResourceIsDeletedExactlyOnce() + { + var report = new AzureJournalReport(new()); + var deletes = 0; + var resource = new OwnedBenchmarkResource(report, _ => Task.CompletedTask, _ => + { + deletes++; + return Task.CompletedTask; + }); + await resource.CreateAsync(CancellationToken.None); + Assert.Equal("owned", report.Ownership); + await resource.CleanupAsync(CancellationToken.None); + await resource.CleanupAsync(CancellationToken.None); + Assert.Equal(1, deletes); + Assert.Equal("deleted", report.Cleanup); + } + + [Fact] + public async Task ReportsExportConfigurationUnitsFailuresAndProviderMetricsWithoutSecrets() + { + var prefix = Path.Combine(Path.GetTempPath(), "journal-report-" + Guid.NewGuid().ToString("N")); + try + { + const long Items = 9_007_199_254_740_993; + var report = new AzureJournalReport(new() { Output = "secret-path", Operations = 1, Concurrency = 1 }) + { + ProviderMetrics = + [ + new("orleans-journaling-provider-catalog-items", "count", "orders", "", 1, Items), + new("orleans-journaling-provider-retries", "count", "orders", "metadata_conflict", 1, 1) + ] + }; + report.Failures.Add(BenchmarkFailure.From("setup", new RequestFailedException(403, "sig=secret"))); + await report.ExportAsync(prefix); + var json = await File.ReadAllTextAsync(prefix + ".json", TestContext.Current.CancellationToken); + var csv = await File.ReadAllTextAsync(prefix + ".csv", TestContext.Current.CancellationToken); + using var document = JsonDocument.Parse(json); + Assert.Equal(2, document.RootElement.GetProperty("SchemaVersion").GetInt32()); + Assert.Equal("counts of catalog pages, candidate items, delivered entries, and explicit provider retries", + document.RootElement.GetProperty("TelemetryUnit").GetString()); + var metric = document.RootElement.GetProperty("ProviderMetrics")[0]; + Assert.Equal(new[] { "Instrument", "Unit", "Provider", "Reason", "Observations", "Sum" }, + metric.EnumerateObject().Select(property => property.Name)); + Assert.Equal(Items, metric.GetProperty("Sum").GetInt64()); + Assert.Equal("metadata_conflict", document.RootElement.GetProperty("ProviderMetrics")[1].GetProperty("Reason").GetString()); + Assert.Equal("AzuriteBlob", document.RootElement.GetProperty("Configuration").GetProperty("Backend").GetString()); + Assert.Equal(403, document.RootElement.GetProperty("Failures")[0].GetProperty("HttpStatus").GetInt32()); + Assert.Contains("provider_metrics_json", csv); + Assert.Contains("payload_bytes_per_second", csv); + Assert.Contains(JsonSerializer.Serialize(report.ProviderMetrics, AzureJournalReport.JsonOptions).Replace("\"", "\"\"", StringComparison.Ordinal), csv); + Assert.Contains("build_json", csv); + Assert.Equal(report.Build, document.RootElement.GetProperty("Build").Deserialize()); + Assert.Contains(JsonSerializer.Serialize(report.Build, AzureJournalReport.JsonOptions).Replace("\"", "\"\"", StringComparison.Ordinal), csv); + Assert.Equal("IJournalStorage / IJournalStorageCatalog", document.RootElement.GetProperty("Source").GetString()); + Assert.DoesNotContain("secret", json); + Assert.DoesNotContain("secret", csv); + await Assert.ThrowsAsync(() => report.ExportAsync(prefix)); + } + finally + { + File.Delete(prefix + ".json"); + File.Delete(prefix + ".csv"); + } + } + + [Theory] + [InlineData(null, null)] + [InlineData("10.0.0. Commit Hash: ", null)] + [InlineData("10.0.0+short", null)] + [InlineData("10.0.0+https://example.invalid/?sig=sensitive", null)] + [InlineData("10.0.0+012345678901234567890123456789012345678G", null)] + [InlineData("10.0.0+0123456789012345678901234567890123456789-dirty", null)] + [InlineData("10.0.0. Commit Hash: +ABCDEF0123456789ABCDEF0123456789ABCDEF01", "abcdef0123456789abcdef0123456789abcdef01")] + [InlineData("10.0.0+0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef", "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef")] + public void BuildProvenanceExportsOnlySanitizedSourceRevision(string? version, string? expected) + => Assert.Equal(expected, BenchmarkAssemblyBuild.ParseSourceRevision(version)); + + [Fact] + public void BuildProvenanceUsesLoadedAssemblyMetadata() + { + const string Revision = "0123456789012345678901234567890123456789"; + var assembly = AssemblyBuilder.DefineDynamicAssembly(new AssemblyName("BenchmarkProvenanceFixture"), AssemblyBuilderAccess.Run); + var module = assembly.DefineDynamicModule("Fixture"); + assembly.SetCustomAttribute(new CustomAttributeBuilder( + typeof(AssemblyInformationalVersionAttribute).GetConstructor([typeof(string)])!, ["10.0.0+" + Revision])); + Assert.Equal(new BenchmarkAssemblyBuild("BenchmarkProvenanceFixture", Revision, module.ModuleVersionId), + BenchmarkAssemblyBuild.FromAssembly(assembly)); + } + + [Fact] + public void BuildProvenanceIdentifiesBenchmarkAndProviderModulesInBdnLog() + { + var build = new AzureJournalReport(new()).Build; + Assert.Equal("Benchmarks", build.Benchmark.Name); + Assert.Equal(typeof(AzureJournalReport).Module.ModuleVersionId, build.Benchmark.ModuleVersionId); + Assert.Equal("Orleans.Journaling", build.Journaling.Name); + Assert.Equal(typeof(IJournalStorage).Module.ModuleVersionId, build.Journaling.ModuleVersionId); + Assert.Equal("Orleans.Journaling.AzureStorage", build.AzureStorage.Name); + Assert.Equal(typeof(AzureBlobJournalStorageOptions).Module.ModuleVersionId, build.AzureStorage.ModuleVersionId); + Assert.NotEqual(Guid.Empty, build.Benchmark.ModuleVersionId); + Assert.NotEqual(Guid.Empty, build.Journaling.ModuleVersionId); + Assert.NotEqual(Guid.Empty, build.AzureStorage.ModuleVersionId); + using var output = new StringWriter(); + BenchmarkBuildInfo.WriteTo(output); + var line = Assert.Single(output.ToString().Split(Environment.NewLine, StringSplitOptions.RemoveEmptyEntries)); + const string Prefix = "Azure benchmark build: "; + Assert.StartsWith(Prefix, line); + Assert.Equal(build, JsonSerializer.Deserialize(line[Prefix.Length..])); + Assert.NotNull(typeof(AzureJournalBenchmarks).GetMethod(nameof(AzureJournalBenchmarks.ReportBuild))! + .GetCustomAttribute()); + } + + [Fact] + public void RecoveryConsumesEveryChunkAndRequiresCompletionCountAndChecksum() + { + byte[] checkpoint = [1, 2, 3]; + byte[] append = [4, 5]; + var metadata = new JournalMetadata(AzureJournalScenario.FormatKey, "etag", AzureJournalScenario.CallerMetadata); + using var consumer = new AzureJournalScenario.PayloadConsumer(checkpoint, append, 2); + using var buffer = new ArcBufferWriter(); + buffer.Write(new byte[] { 1, 2, 3, 4 }); + consumer.Read(new JournalBufferReader(buffer.Reader, false), metadata); + Assert.Equal(0, buffer.Reader.Length); + Assert.Throws(consumer.AssertComplete); + buffer.Write(new byte[] { 5, 4, 5 }); + consumer.Read(new JournalBufferReader(buffer.Reader, true), metadata); + Assert.Equal(0, buffer.Reader.Length); + consumer.AssertComplete(); + Assert.Throws(() => consumer.Read(new JournalBufferReader(buffer.Reader, true), metadata)); + + using var corrupt = new AzureJournalScenario.PayloadConsumer(checkpoint, append, 0); + buffer.Write(new byte[] { 1, 2, 0 }); + Assert.Throws(() => corrupt.Read(new JournalBufferReader(buffer.Reader, true), metadata)); + using var shortRead = new AzureJournalScenario.PayloadConsumer(checkpoint, append, 1); + using var shortBuffer = new ArcBufferWriter(); + shortBuffer.Write(checkpoint); + shortRead.Read(new JournalBufferReader(shortBuffer.Reader, true), metadata); + Assert.Throws(shortRead.AssertComplete); + } + + [Fact] + public void CatalogBoundsSelectDueIdentitiesAndPreserveMetadataProjection() + { + var options = new AzureJournalOptions { Workload = AzureJournalWorkload.CatalogBounded, DueJournals = 3, FutureJournals = 9, IncludeMetadata = true }; + var scenario = new AzureJournalScenario(options, new(options)); + var listing = scenario.CatalogOptions(); + Assert.Equal(new JournalId("catalog/"), listing.Prefix); + Assert.Equal(AzureJournalScenario.CatalogId(0, false), listing.MinId); + Assert.Equal(AzureJournalScenario.CatalogId(2, false), listing.MaxId); + Assert.True(listing.IncludeMetadata); + Assert.Equal(3, options.ItemsPerOperation); + var unbounded = options with { Workload = AzureJournalWorkload.CatalogUnbounded }; + Assert.True(new AzureJournalScenario(unbounded, new(unbounded)).CatalogOptions().MaxId.IsDefault); + Assert.Equal(12, unbounded.ItemsPerOperation); + Assert.Throws(() => AzureJournalScenario.ValidateMetadata(new JournalMetadata(AzureJournalScenario.FormatKey, "etag"), true)); + Assert.Throws(() => AzureJournalScenario.ValidateMetadata(new JournalMetadata("json", "etag", AzureJournalScenario.CallerMetadata), true)); + } + + [Fact] + public async Task CatalogConsumesAndValidatesExactMembershipAndMetadata() + { + var first = AzureJournalScenario.CatalogId(0, false); + var second = AzureJournalScenario.CatalogId(1, false); + var future = AzureJournalScenario.CatalogId(0, true); + HashSet expected = [first, second]; + var metadata = new JournalMetadata(AzureJournalScenario.FormatKey, "etag", AzureJournalScenario.CallerMetadata); + var token = TestContext.Current.CancellationToken; + await AzureJournalScenario.ValidateCatalogAsync(Entries([]), new HashSet(), false, token); + await AzureJournalScenario.ValidateCatalogAsync(Entries([new(second, metadata), new(first, metadata)]), expected, true, token); + var missing = await Assert.ThrowsAsync(() => + AzureJournalScenario.ValidateCatalogAsync(Entries([new(first, metadata)]), expected, true, token)); + Assert.Equal(BenchmarkValidationError.CatalogMissingIdentity, missing.Code); + foreach (var extra in new[] { first, future }) + { + var invalid = await Assert.ThrowsAsync(() => + AzureJournalScenario.ValidateCatalogAsync(Entries([new(first, metadata), new(extra, metadata)]), expected, true, token)); + Assert.Equal(BenchmarkValidationError.CatalogUnexpectedIdentity, invalid.Code); + } + + await AzureJournalScenario.ValidateCatalogAsync(Entries([new(first, null), new(second, null)]), expected, false, token); + var unexpectedMetadata = await Assert.ThrowsAsync(() => + AzureJournalScenario.ValidateCatalogAsync(Entries([new(first, metadata), new(second, metadata)]), expected, false, token)); + Assert.Equal(BenchmarkValidationError.CatalogUnrequestedMetadata, unexpectedMetadata.Code); + var missingMetadata = await Assert.ThrowsAsync(() => + AzureJournalScenario.ValidateCatalogAsync(Entries([new(first, null), new(second, metadata)]), expected, true, token)); + Assert.Equal(BenchmarkValidationError.JournalFormatOrETag, missingMetadata.Code); + + static async IAsyncEnumerable Entries(JournalCatalogEntry[] values) + { + await Task.Yield(); + foreach (var value in values) + { + yield return value; + } + } + } + + [Theory] + [InlineData("Azurite", "Azurite", true, false)] + [InlineData("AzuriteTable", "AzuriteTable", true, true)] + [InlineData("StandardBlob", "StandardBlob", false, false)] + [InlineData("PremiumBlob", "PremiumBlob", false, false)] + [InlineData("Azure", "PremiumBlob", false, false)] + [InlineData("Table", "Table", false, true)] + public void PlaygroundRoutesBackendAndKeepsAzureAlias(string name, string expected, bool emulator, bool table) + { + var backend = StorageBackendConfiguration.Parse(name); + Assert.Equal(expected, backend.ToString()); + Assert.Equal(emulator, backend.IsEmulator()); + Assert.Equal(table, backend.UsesTableJournal()); + } + + private sealed class MetricsSiloBuilder : ISiloBuilder + { + public IServiceCollection Services { get; } = new ServiceCollection(); + public IConfiguration Configuration { get; } = new ConfigurationBuilder().Build(); + } + + private sealed class FailingLifecycleParticipant(List events, Action fail) : ILifecycleParticipant, IDisposable + { + public void Participate(ISiloLifecycle lifecycle) + { + lifecycle.Subscribe("early", ServiceLifecycleStage.RuntimeInitialize + 1, + onStart: _ => { events.Add("start:early"); return Task.CompletedTask; }, + onStop: token => { token.ThrowIfCancellationRequested(); events.Add("stop:early"); return Task.CompletedTask; }); + lifecycle.Subscribe("failed", ServiceLifecycleStage.RuntimeInitialize + 2, + onStart: _ => { events.Add("start:failed"); fail(); return Task.CompletedTask; }, + onStop: token => { token.ThrowIfCancellationRequested(); events.Add("stop:failed"); return Task.CompletedTask; }); + lifecycle.Subscribe("unreached", ServiceLifecycleStage.RuntimeInitialize + 3, + onStart: _ => { events.Add("start:unreached"); return Task.CompletedTask; }, + onStop: _ => { events.Add("stop:unreached"); return Task.CompletedTask; }); + } + + public void Dispose() => events.Add("dispose"); + } + + private sealed class CatalogBlobService : BlobServiceClient + { + private readonly CatalogContainer _container = new(); + public override BlobContainerClient GetBlobContainerClient(string blobContainerName) => _container; + } + + private sealed class CatalogContainer : BlobContainerClient + { + public override Task> CreateIfNotExistsAsync( + PublicAccessType publicAccessType = PublicAccessType.None, + IDictionary? metadata = null, + BlobContainerEncryptionScopeOptions? encryptionScopeOptions = null, + CancellationToken cancellationToken = default) + => Task.FromResult(Response.FromValue( + BlobsModelFactory.BlobContainerInfo(new ETag("created"), DateTimeOffset.UnixEpoch), new CatalogResponse())); + + public override AsyncPageable GetBlobsAsync(GetBlobsOptions options, CancellationToken cancellationToken = default) + { + Assert.Equal("wal/catalog/", options.Prefix); + var items = new[] { "catalog/a", "catalog/b", "catalog/c" }.Select(id => BlobsModelFactory.BlobItem( + name: $"wal/{id}", deleted: false, + properties: BlobsModelFactory.BlobItemProperties(accessTierInferred: false, blobType: BlobType.Append))).ToArray(); + return AsyncPageable.FromPages( + [ + Page.FromValues([], "next", new CatalogResponse()), + Page.FromValues(items, null, new CatalogResponse()) + ]); + } + } + + private sealed class CatalogTableService : TableServiceClient + { + private readonly CatalogTable _table = new(); + public override TableClient GetTableClient(string tableName) => _table; + } + + private sealed class CatalogTable : TableClient + { + public override Task> CreateIfNotExistsAsync(CancellationToken cancellationToken = default) + => Task.FromResult(Response.FromValue(new TableItem("journal"), new CatalogResponse())); + + public override AsyncPageable QueryAsync( + string? filter = null, int? maxPerPage = null, IEnumerable? select = null, CancellationToken cancellationToken = default) + { + Assert.Contains("PartitionKey le", filter); + Assert.Equal(1000, maxPerPage); + var items = new[] { "catalog/a", "catalog/b", "catalog/c" } + .Select(id => (T)(ITableEntity)new TableEntity { ["JournalId"] = id }).ToArray(); + return AsyncPageable.FromPages( + [ + Page.FromValues([], "next", new CatalogResponse()), + Page.FromValues(items, null, new CatalogResponse()) + ]); + } + } + + private sealed class CatalogResponse : Response + { + public override int Status => 200; + public override string ReasonPhrase => "OK"; + public override Stream? ContentStream { get; set; } + public override string ClientRequestId { get; set; } = ""; + public override void Dispose() { } + protected override bool ContainsHeader(string name) => false; + protected override IEnumerable EnumerateHeaders() => []; + protected override bool TryGetHeader(string name, out string value) { value = ""; return false; } + protected override bool TryGetHeaderValues(string name, out IEnumerable values) { values = []; return false; } + } + + private sealed class FakeScenario(AzureJournalReport report) : IAzureJournalScenario + { + public Func Prepare { get; init; } = _ => Task.CompletedTask; + public Func Execute { get; init; } = (_, _) => Task.CompletedTask; + public Func CleanupAction { get; init; } = _ => Task.CompletedTask; + public Func Verify { get; init; } = _ => Task.CompletedTask; + public Task PrepareAsync(CancellationToken cancellationToken) => Prepare(cancellationToken); + public Task ExecuteAsync(int index, CancellationToken cancellationToken) => Execute(index, cancellationToken); + public Task VerifyAsync(IEnumerable completed, CancellationToken cancellationToken) => Verify(cancellationToken); + public void StartMetrics() { } + public IReadOnlyList StopMetrics() => []; + public async Task CleanupAsync(CancellationToken cancellationToken) + { + await CleanupAction(cancellationToken); + report.Cleanup = "deleted"; + } + } +} diff --git a/test/Benchmarks/Journaling/Azure/AzureJournalScenario.cs b/test/Benchmarks/Journaling/Azure/AzureJournalScenario.cs new file mode 100644 index 00000000000..a32d14c6714 --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/AzureJournalScenario.cs @@ -0,0 +1,484 @@ +using System.Buffers; +using System.Security.Cryptography; +using Azure.Core; +using Azure.Data.Tables; +using Azure.Identity; +using Azure.Storage.Blobs; +using Azure.Storage.Blobs.Models; +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging.Abstractions; +using Orleans; +using Orleans.Journaling; +using Orleans.Runtime; + +namespace Benchmarks.Journaling.Azure; + +internal interface IAzureJournalScenario +{ + Task PrepareAsync(CancellationToken cancellationToken); + Task ExecuteAsync(int index, CancellationToken cancellationToken); + Task VerifyAsync(IEnumerable completed, CancellationToken cancellationToken); + void StartMetrics(); + IReadOnlyList StopMetrics(); + Task CleanupAsync(CancellationToken cancellationToken); +} + +internal sealed class AzureJournalScenario(AzureJournalOptions options, AzureJournalReport report) : IAzureJournalScenario +{ + internal const string FormatKey = "benchmark-bytes-v1"; + private readonly byte[] _append = CreatePayload(options.AppendBytes, options.Seed); + private readonly byte[] _checkpoint = CreatePayload(options.CheckpointBytes, unchecked(options.Seed + 1)); + private readonly string _resourceName = "journalbench" + Guid.NewGuid().ToString("N"); + private ServiceProvider? _services; + private SiloLifecycleSubject? _lifecycle; + private IJournalStorageProvider _provider = null!; + private IJournalStorageCatalog _catalog = null!; + private ProviderMetrics? _metrics; + private BlobContainerClient? _container; + private TableServiceClient? _tableService; + private readonly List _journals = []; + private OwnedBenchmarkResource? _resource; + private byte[] _expectedHash = null!; + private HashSet _expectedCatalog = []; + + public async Task PrepareAsync(CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + report.Resource = _resourceName; + Console.Error.WriteLine($"Azure journal benchmark resource={_resourceName}; backend={options.Backend}"); + var expectedBatches = options.Workload switch + { + AzureJournalWorkload.DurableAppend => options.HistoryBatches + 1, + AzureJournalWorkload.CheckpointReplace => 0, + _ => options.HistoryBatches + }; + _expectedHash = PayloadConsumer.ComputeHash(_checkpoint, _append, expectedBatches); + _expectedCatalog = Enumerable.Range(0, options.DueJournals).Select(index => CatalogId(index, future: false)).ToHashSet(); + if (options.Workload == AzureJournalWorkload.CatalogUnbounded) + { + _expectedCatalog.UnionWith(Enumerable.Range(0, options.FutureJournals).Select(index => CatalogId(index, future: true))); + } + + var builder = new ProviderSiloBuilder(); + builder.Services.AddLogging().AddMetrics().AddSingleton(); + if (options.IsTable) + { + var clientOptions = new TableClientOptions(); + ConfigureClient(clientOptions); + _tableService = options.IsEmulator + ? new TableServiceClient(EmulatorConnectionString(), clientOptions) + : new TableServiceClient(ReadAzureEndpoint("JOURNAL_BENCHMARK_TABLE_ENDPOINT"), new DefaultAzureCredential(), clientOptions); + ValidateEndpoint(_tableService.Uri, options.IsEmulator); + report.AccountVerification = options.IsEmulator ? "emulator" : "requested-unverified"; + report.AccountKind = options.IsEmulator ? "Azurite" : null; + _resource = new OwnedBenchmarkResource(report, + async token => { await _tableService.CreateTableAsync(_resourceName, token); }, + async token => { await _tableService.DeleteTableAsync(_resourceName, token); }); + await _resource.CreateAsync(cancellationToken); + builder.AddAzureTableJournalStorage(value => + { + value.TableName = _resourceName; + value.TableServiceClient = _tableService; + }); + } + else + { + var clientOptions = new BlobClientOptions(); + ConfigureClient(clientOptions); + var client = options.IsEmulator + ? new BlobServiceClient(EmulatorConnectionString(), clientOptions) + : new BlobServiceClient(ReadAzureEndpoint("JOURNAL_BENCHMARK_BLOB_ENDPOINT"), new DefaultAzureCredential(), clientOptions); + ValidateEndpoint(client.Uri, options.IsEmulator); + if (options.IsEmulator) + { + report.AccountVerification = "emulator"; + report.AccountKind = "Azurite"; + } + else + { + var account = (await client.GetAccountInfoAsync(cancellationToken)).Value; + report.AccountKind = account.AccountKind.ToString(); + report.AccountSku = account.SkuName.ToString(); + VerifyAccount(options.Backend, account.AccountKind, account.SkuName); + report.AccountVerification = "data-plane-verified"; + } + + _container = client.GetBlobContainerClient(_resourceName); + _resource = new OwnedBenchmarkResource(report, + async token => { await _container.CreateAsync(cancellationToken: token); }, + async token => { await _container.DeleteAsync(cancellationToken: token); }); + await _resource.CreateAsync(cancellationToken); + builder.AddAzureBlobJournalStorage(value => + { + value.ContainerName = _resourceName; + value.BlobServiceClient = client; + }); + } + + builder.Services.AddKeyedSingleton(FormatKey, new BenchmarkByteFormat()); + builder.Services.Configure(value => value.JournalFormatKey = FormatKey); + await InitializeProviderAsync(builder.Services, cancellationToken); + report.Phase = "seed"; + if (options.IsCatalog) + { + await SeedCatalogAsync(cancellationToken); + } + else + { + // The last journal is reserved for warmup; measured operations each get an identical baseline. + for (var index = 0; index <= options.Operations; index++) + { + var storage = _provider.CreateStorage(WorkId(index)); + await CreateJournalAsync(storage, cancellationToken); + await storage.ReplaceAsync(new ReadOnlySequence(_checkpoint), cancellationToken); + for (var batch = 0; batch < options.HistoryBatches; batch++) + { + await storage.AppendAsync(new ReadOnlySequence(_append), cancellationToken); + } + + _journals.Add(storage); + } + } + + report.Phase = "warmup"; + await ExecuteAsync(options.Operations, cancellationToken); + await VerifyAsync([options.Operations], cancellationToken); + } + + internal async Task InitializeProviderAsync(IServiceCollection services, CancellationToken cancellationToken) + { + _services = services.BuildServiceProvider(); + _metrics = new ProviderMetrics(_services.GetRequiredService().Meter); + _provider = _services.GetRequiredService(); + _catalog = _services.GetRequiredService(); + _lifecycle = new SiloLifecycleSubject(NullLogger.Instance); + foreach (var participant in _services.GetServices>()) + { + participant.Participate(_lifecycle); + } + + report.Phase = "initialize"; + await _lifecycle.OnStart(cancellationToken); + } + + public async Task ExecuteAsync(int index, CancellationToken cancellationToken) + { + cancellationToken.ThrowIfCancellationRequested(); + switch (options.Workload) + { + case AzureJournalWorkload.DurableAppend: + await _journals[index].AppendAsync(new ReadOnlySequence(_append), cancellationToken); + break; + case AzureJournalWorkload.CheckpointReplace: + await _journals[index].ReplaceAsync(new ReadOnlySequence(_checkpoint), cancellationToken); + break; + case AzureJournalWorkload.RecoveryReplay: + await VerifyJournalAsync(index, options.HistoryBatches, cancellationToken); + break; + case AzureJournalWorkload.CatalogBounded: + case AzureJournalWorkload.CatalogUnbounded: + await VerifyCatalogAsync(cancellationToken); + break; + default: + throw new InvalidOperationException("Unknown workload."); + } + } + + public async Task VerifyAsync(IEnumerable completed, CancellationToken cancellationToken) + { + if (options.IsCatalog) + { + return; + } + + var batches = options.Workload == AzureJournalWorkload.DurableAppend ? options.HistoryBatches + 1 : 0; + foreach (var index in completed) + { + if (options.Workload != AzureJournalWorkload.RecoveryReplay) + { + await VerifyJournalAsync(index, batches, cancellationToken); + } + + ValidateMetadata(await _provider.CreateStorage(WorkId(index)).GetMetadataAsync(cancellationToken), options.IncludeMetadata); + } + } + + private async Task VerifyJournalAsync(int index, int batches, CancellationToken cancellationToken) + { + using var consumer = new PayloadConsumer(_checkpoint, _append, batches, _expectedHash); + // Recovery uses a fresh storage object, exercising persisted state rather than a writer's cache. + await _provider.CreateStorage(WorkId(index)).ReadAsync(consumer, cancellationToken); + consumer.AssertComplete(); + } + + private async Task SeedCatalogAsync(CancellationToken cancellationToken) + { + for (var index = 0; index < options.DueJournals; index++) + { + await CreateJournalAsync(_provider.CreateStorage(CatalogId(index, future: false)), cancellationToken); + } + + for (var index = 0; index < options.FutureJournals; index++) + { + await CreateJournalAsync(_provider.CreateStorage(CatalogId(index, future: true)), cancellationToken); + } + + await CreateJournalAsync(_provider.CreateStorage(new JournalId("catalog/00000000/expired")), cancellationToken); + await CreateJournalAsync(_provider.CreateStorage(new JournalId("unrelated/00000000")), cancellationToken); + } + + private async Task CreateJournalAsync(IJournalStorage storage, CancellationToken cancellationToken) + { + if (!await storage.CreateIfNotExistsAsync(options.IncludeMetadata ? CallerMetadata : null, cancellationToken)) + { + throw new BenchmarkValidationException(BenchmarkValidationError.JournalCollision); + } + } + + internal static IReadOnlyDictionary CallerMetadata { get; } = new Dictionary + { + ["benchmark"] = "azure-journaling", + ["version"] = "1" + }; + + internal static JournalId CatalogId(int index, bool future) => new($"catalog/{(future ? "21000101" : "20000101")}/{index:D8}"); + private static JournalId WorkId(int index) => new($"work/{index:D8}"); + + internal ListOptions CatalogOptions() => new() + { + Prefix = new JournalId("catalog/"), + MinId = CatalogId(0, future: false), + MaxId = options.Workload == AzureJournalWorkload.CatalogBounded ? CatalogId(Math.Max(0, options.DueJournals - 1), future: false) : default, + IncludeMetadata = options.IncludeMetadata + }; + + private Task VerifyCatalogAsync(CancellationToken cancellationToken) + => ValidateCatalogAsync(_catalog.ListAsync(CatalogOptions(), cancellationToken), _expectedCatalog, options.IncludeMetadata, cancellationToken); + + internal static async Task ValidateCatalogAsync( + IAsyncEnumerable entries, IReadOnlySet expectedIds, bool includeMetadata, CancellationToken cancellationToken) + { + var expected = new HashSet(expectedIds); + await foreach (var entry in entries.WithCancellation(cancellationToken)) + { + if (!expected.Remove(entry.Id)) + { + throw new BenchmarkValidationException(BenchmarkValidationError.CatalogUnexpectedIdentity); + } + + if (includeMetadata) + { + ValidateMetadata(entry.Metadata, includeMetadata: true); + } + else if (entry.Metadata is not null) + { + throw new BenchmarkValidationException(BenchmarkValidationError.CatalogUnrequestedMetadata); + } + } + + if (expected.Count != 0) + { + throw new BenchmarkValidationException(BenchmarkValidationError.CatalogMissingIdentity); + } + } + + internal static void ValidateMetadata(IJournalMetadata? metadata, bool includeMetadata) + { + if (metadata?.Format != FormatKey || string.IsNullOrEmpty(metadata.ETag)) + { + throw new BenchmarkValidationException(BenchmarkValidationError.JournalFormatOrETag); + } + + var properties = metadata.Properties.Where(pair => !pair.Key.StartsWith('$')).ToDictionary(pair => pair.Key, pair => pair.Value); + if (properties.Count != (includeMetadata ? CallerMetadata.Count : 0) + || (includeMetadata && CallerMetadata.Any(pair => !properties.TryGetValue(pair.Key, out var value) || value != pair.Value))) + { + throw new BenchmarkValidationException(BenchmarkValidationError.JournalCallerMetadata); + } + } + + public void StartMetrics() => _metrics!.Start(); + public IReadOnlyList StopMetrics() => _metrics!.Stop(); + + public async Task CleanupAsync(CancellationToken cancellationToken) + { + try + { + if (_lifecycle is { } lifecycle) + { + _lifecycle = null; + try + { + await lifecycle.OnStop(cancellationToken); + } + catch (Exception exception) + { + report.Failures.Add(BenchmarkFailure.From("cleanup-lifecycle", exception)); + } + } + + if (_resource is not null) + { + await _resource.CleanupAsync(cancellationToken); + } + } + finally + { + _metrics?.Dispose(); + if (_services is not null) + { + await _services.DisposeAsync(); + _services = null; + } + } + } + + internal static void VerifyAccount(AzureJournalBackend backend, AccountKind kind, SkuName sku) + { + var matches = backend switch + { + AzureJournalBackend.PremiumBlob => kind == AccountKind.BlockBlobStorage && sku == SkuName.PremiumLrs, + AzureJournalBackend.StandardBlob => kind is AccountKind.Storage or AccountKind.StorageV2 or AccountKind.BlobStorage + && sku is SkuName.StandardLrs or SkuName.StandardGrs or SkuName.StandardRagrs or SkuName.StandardZrs, + _ => false + }; + if (!matches) + { + throw new InvalidOperationException("The Azure account kind/SKU does not match the requested backend."); + } + } + + private static string EmulatorConnectionString() + => Environment.GetEnvironmentVariable("JOURNAL_BENCHMARK_AZURITE_CONNECTION_STRING") ?? "UseDevelopmentStorage=true"; + + private static Uri ReadAzureEndpoint(string name) + { + if (!Uri.TryCreate(Environment.GetEnvironmentVariable(name), UriKind.Absolute, out var uri)) + { + throw new ArgumentException($"Set {name} to the Azure service endpoint."); + } + + ValidateEndpoint(uri, emulator: false); + return uri; + } + + internal static void ValidateEndpoint(Uri uri, bool emulator) + { + if (!string.IsNullOrEmpty(uri.Query) || !string.IsNullOrEmpty(uri.UserInfo) || !string.IsNullOrEmpty(uri.Fragment) + || (emulator ? !uri.IsLoopback || uri.Scheme is not ("http" or "https") : uri.IsLoopback || uri.Scheme != "https")) + { + throw new ArgumentException("Use loopback emulator endpoints or HTTPS Azure service endpoints authenticated with Entra ID."); + } + } + + private static void ConfigureClient(ClientOptions options) + { + options.Retry.MaxRetries = 2; + options.Retry.NetworkTimeout = TimeSpan.FromSeconds(10); + options.Diagnostics.IsLoggingEnabled = false; + options.Diagnostics.IsLoggingContentEnabled = false; + } + + internal static byte[] CreatePayload(int length, int seed) + { + var bytes = new byte[length]; + new Random(seed).NextBytes(bytes); + return bytes; + } + + private sealed class ProviderSiloBuilder : ISiloBuilder + { + public IServiceCollection Services { get; } = new ServiceCollection(); + public IConfiguration Configuration { get; } = new ConfigurationBuilder().Build(); + } + + private sealed class BenchmarkByteFormat : IJournalFormat + { + public string FormatKey => AzureJournalScenario.FormatKey; + public string MimeType => "application/octet-stream"; + public JournalBufferWriter CreateWriter() => throw new NotSupportedException("The provider benchmark writes deterministic bytes directly."); + public void Replay(JournalBufferReader input, JournalReplayContext context) => throw new NotSupportedException("The provider benchmark replays bytes through PayloadConsumer."); + } + + internal sealed class PayloadConsumer : IJournalStorageConsumer, IDisposable + { + private readonly byte[] _checkpoint; + private readonly byte[] _append; + private readonly long _expectedLength; + private readonly byte[] _expectedHash; + private readonly IncrementalHash _hash = IncrementalHash.CreateHash(HashAlgorithmName.SHA256); + private long _length; + private bool _completed; + + public PayloadConsumer(byte[] checkpoint, byte[] append, int batches, byte[]? expectedHash = null) + { + _checkpoint = checkpoint; + _append = append; + _expectedLength = checkpoint.Length + (long)append.Length * batches; + _expectedHash = expectedHash ?? ComputeHash(checkpoint, append, batches); + } + + internal static byte[] ComputeHash(byte[] checkpoint, byte[] append, int batches) + { + using var hash = IncrementalHash.CreateHash(HashAlgorithmName.SHA256); + hash.AppendData(checkpoint); + for (var index = 0; index < batches; index++) + { + hash.AppendData(append); + } + + return hash.GetHashAndReset(); + } + + public void Read(JournalBufferReader buffer, IJournalMetadata? metadata) + { + if (_completed) + { + throw new BenchmarkValidationException(BenchmarkValidationError.RecoveryAfterCompletion); + } + + if (metadata?.Format != FormatKey) + { + throw new BenchmarkValidationException(BenchmarkValidationError.RecoveryFormat); + } + + Span scratch = stackalloc byte[4096]; + while (buffer.Length > 0) + { + var count = Math.Min(buffer.Length, scratch.Length); + var bytes = scratch[..count]; + buffer.Read(bytes); + if (_length + count > _expectedLength) + { + throw new BenchmarkValidationException(BenchmarkValidationError.RecoveryLength); + } + + for (var index = 0; index < count; index++) + { + var position = _length + index; + var expected = position < _checkpoint.Length ? _checkpoint[position] : _append[(position - _checkpoint.Length) % _append.Length]; + if (bytes[index] != expected) + { + throw new BenchmarkValidationException(BenchmarkValidationError.RecoveryPayload); + } + } + + _hash.AppendData(bytes); + _length += count; + } + + _completed = buffer.IsCompleted; + } + + public void AssertComplete() + { + if (!_completed || _length != _expectedLength || !_hash.GetHashAndReset().AsSpan().SequenceEqual(_expectedHash)) + { + throw new BenchmarkValidationException(BenchmarkValidationError.RecoveryCompletionOrChecksum); + } + } + + public void Dispose() => _hash.Dispose(); + } +} diff --git a/test/Benchmarks/Journaling/Azure/BenchmarkBuildInfo.cs b/test/Benchmarks/Journaling/Azure/BenchmarkBuildInfo.cs new file mode 100644 index 00000000000..b75782c490a --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/BenchmarkBuildInfo.cs @@ -0,0 +1,45 @@ +using System.Reflection; +using System.Text.Json; +using Orleans.Journaling; + +namespace Benchmarks.Journaling.Azure; + +internal sealed record BenchmarkAssemblyBuild(string Name, string? SourceRevision, Guid ModuleVersionId) +{ + public static BenchmarkAssemblyBuild FromAssembly(Assembly assembly) + => new( + assembly.GetName().Name!, + ParseSourceRevision(assembly.GetCustomAttribute()?.InformationalVersion), + assembly.ManifestModule.ModuleVersionId); + + internal static string? ParseSourceRevision(string? informationalVersion) + { + if (informationalVersion is null) + { + return null; + } + + var separator = informationalVersion.LastIndexOf('+'); + if (separator < 0) + { + return null; + } + + var revision = informationalVersion[(separator + 1)..]; + return revision.Length is 40 or 64 && revision.All(Uri.IsHexDigit) ? revision.ToLowerInvariant() : null; + } +} + +internal sealed record BenchmarkBuildInfo( + BenchmarkAssemblyBuild Benchmark, + BenchmarkAssemblyBuild Journaling, + BenchmarkAssemblyBuild AzureStorage) +{ + public static BenchmarkBuildInfo Current { get; } = new( + BenchmarkAssemblyBuild.FromAssembly(typeof(BenchmarkBuildInfo).Assembly), + BenchmarkAssemblyBuild.FromAssembly(typeof(IJournalStorage).Assembly), + BenchmarkAssemblyBuild.FromAssembly(typeof(AzureBlobJournalStorageOptions).Assembly)); + + public static void WriteTo(TextWriter writer) + => writer.WriteLine($"Azure benchmark build: {JsonSerializer.Serialize(Current)}"); +} diff --git a/test/Benchmarks/Journaling/Azure/OwnedBenchmarkResource.cs b/test/Benchmarks/Journaling/Azure/OwnedBenchmarkResource.cs new file mode 100644 index 00000000000..97a2bd2e865 --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/OwnedBenchmarkResource.cs @@ -0,0 +1,33 @@ +namespace Benchmarks.Journaling.Azure; + +internal sealed class OwnedBenchmarkResource( + AzureJournalReport report, + Func create, + Func delete) +{ + private bool _owned; + + public async Task CreateAsync(CancellationToken cancellationToken) + { + report.Ownership = "creation-unconfirmed"; + report.Cleanup = "manual-check-required"; + // A retried create can return 409 after the original success response was lost. + // Only an acknowledged success establishes ownership; every failure requires inspection. + await create(cancellationToken); + + _owned = true; + report.Ownership = "owned"; + report.Cleanup = "pending"; + } + + public async Task CleanupAsync(CancellationToken cancellationToken) + { + if (_owned) + { + report.Cleanup = "deleting"; + await delete(cancellationToken); + _owned = false; + report.Cleanup = "deleted"; + } + } +} diff --git a/test/Benchmarks/Journaling/Azure/ProviderMetrics.cs b/test/Benchmarks/Journaling/Azure/ProviderMetrics.cs new file mode 100644 index 00000000000..8310ac6dbfb --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/ProviderMetrics.cs @@ -0,0 +1,122 @@ +using System.Diagnostics.Metrics; + +namespace Benchmarks.Journaling.Azure; + +internal sealed record ProviderMetric( + string Instrument, + string Unit, + string Provider, + string Reason, + long Observations, + long Sum); + +internal sealed class ProviderMetrics(Meter meter) : IDisposable +{ + private MeterListener? _listener; + private readonly Dictionary _values = []; + private readonly object _lock = new(); + private bool _recording; + + public void Start() + { + lock (_lock) + { + if (_listener is not null) + { + throw new InvalidOperationException("Provider metric collection is already active."); + } + + _values.Clear(); + _recording = true; + } + + var listener = new MeterListener(); + listener.InstrumentPublished = (instrument, subscriber) => + { + if (ReferenceEquals(instrument.Meter, meter) + && instrument is Counter + && instrument.Name is "orleans-journaling-provider-catalog-pages" + or "orleans-journaling-provider-catalog-items" + or "orleans-journaling-provider-catalog-entries" + or "orleans-journaling-provider-retries") + { + subscriber.EnableMeasurementEvents(instrument); + } + }; + listener.SetMeasurementEventCallback((instrument, value, tags, _) => Record(instrument, value, tags)); + _listener = listener; + listener.Start(); + } + + public IReadOnlyList Stop() + { + Dispose(); + lock (_lock) + { + return _values.OrderBy(pair => pair.Key.Instrument, StringComparer.Ordinal) + .ThenBy(pair => pair.Key.Provider, StringComparer.Ordinal) + .ThenBy(pair => pair.Key.Reason, StringComparer.Ordinal) + .Select(pair => pair.Value.Snapshot(pair.Key)).ToArray(); + } + } + + private void Record(Instrument instrument, long value, ReadOnlySpan> tags) + { + var key = new MetricKey(instrument.Name, instrument.Unit ?? "count", "", ""); + foreach (var tag in tags) + { + var text = tag.Value as string ?? ""; + key = tag.Key switch + { + "provider" => key with { Provider = text }, + "reason" => key with { Reason = text }, + _ => key + }; + } + + lock (_lock) + { + if (!_recording) + { + return; + } + + if (!_values.TryGetValue(key, out var aggregate)) + { + _values.Add(key, aggregate = new Aggregate()); + } + + aggregate.Record(value); + } + } + + public void Dispose() + { + MeterListener? listener; + lock (_lock) + { + _recording = false; + listener = _listener; + _listener = null; + } + + listener?.Dispose(); + } + + private sealed record MetricKey(string Instrument, string Unit, string Provider, string Reason); + + private sealed class Aggregate + { + private long _count; + private long _sum; + + public void Record(long value) + { + _count++; + _sum += value; + } + + public ProviderMetric Snapshot(MetricKey key) + => new(key.Instrument, key.Unit, key.Provider, key.Reason, _count, _sum); + } +} diff --git a/test/Benchmarks/Journaling/Azure/README.md b/test/Benchmarks/Journaling/Azure/README.md new file mode 100644 index 00000000000..6702c1ffda3 --- /dev/null +++ b/test/Benchmarks/Journaling/Azure/README.md @@ -0,0 +1,316 @@ +# Azure journal provider benchmarks + +`Journaling.Azure` compares durable append, checkpoint replacement, byte recovery, +and catalog traversal through `IJournalStorage` and `IJournalStorageCatalog`. +`Journaling.Azure.Bdn` runs the same operations under BenchmarkDotNet. The +[Durable Jobs playground](../../../../playground/DurableJobsJournaling/README.md) +exercises grain scheduling and workflow latency with the same backend choices. + +## Backends and authentication + +| Backend | Storage implementation | Account verification | +| --- | --- | --- | +| `AzuriteBlob` | Append-blob WAL and block-blob checkpoints | Emulator, loopback only | +| `AzuriteTable` | Table journal headers and data generations | Emulator, loopback only | +| `StandardBlob` | Append-blob WAL and block-blob checkpoints | `GetAccountInfoAsync`: standard SKU and compatible account kind | +| `PremiumBlob` | The same Blob provider on a premium account | `GetAccountInfoAsync`: `BlockBlobStorage` and premium SKU | +| `Table` | Azure Table journal storage | Requested backend; account tier is unverified | + +Premium **BlockBlobStorage accounts support append blobs**. The account tier is +the comparison axis; both Blob cases execute the same WAL and checkpoint code. +The standalone runner accepts the supported standard SDK SKUs for `StandardBlob` +and records the actual SKU in `AccountSku`; `data-plane-verified` describes the +account kind/SKU check. Operators select redundancy on their existing accounts. +The playground provisions matched `Standard_LRS` and `Premium_LRS` presets. + +For account-tier comparisons, place the benchmark host and accounts in the same +Azure region, match redundancy settings, and keep compute, payload, history, +concurrency, and catalog distribution constant between runs. Different +redundancy settings introduce an additional comparison dimension. Record account +configuration and service-side metrics alongside each report. + +**Real Azure runs incur storage, transaction, and possible network charges.** +They require `--allow-azure true`. Use existing test accounts and an Entra +identity with Storage Blob Data Contributor or Storage Table Data Contributor, +including permission to create and delete the isolated container/table. +`DefaultAzureCredential` supports local Azure CLI sign-in and managed identity. +Configure service endpoints through environment variables: + +```powershell +# Sign in using your organization's normal Azure CLI procedure. +az login +$env:JOURNAL_BENCHMARK_BLOB_ENDPOINT = 'https://yourstandardaccount.blob.core.windows.net' +$env:JOURNAL_BENCHMARK_TABLE_ENDPOINT = 'https://yourtableaccount.table.core.windows.net' +``` + +Azure endpoints use HTTPS with Entra authentication. The runner accepts endpoint +settings from the environment and exports only backend, verified kind/SKU, and +generated resource identity. Exception reporting preserves phase, exception +type, HTTP status, operation index, and benchmark validation code. Credential +values, connection strings, URI queries, and SDK exception messages stay out of +reports. Account-info permission or tier mismatches fail before resource creation. + +## Build and emulator smoke + +Run commands from the repository root with the SDK selected by `global.json`. + +```powershell +dotnet build test\Benchmarks\Benchmarks.csproj -c Release -f net10.0 + +# In a separate terminal, start an installed Azurite on its default loopback ports. +azurite --skipApiVersionCheck --location .\.azurite-journal-bench + +# In the benchmark terminal: +$run = Join-Path $env:TEMP ('journal-bench-' + [guid]::NewGuid().ToString('N')) +dotnet run --project test\Benchmarks\Benchmarks.csproj -c Release -f net10.0 --no-build -- ` + Journaling.Azure --backend AzuriteBlob --workload DurableAppend --output $run +``` + +`UseDevelopmentStorage=true` supplies the default emulator connection. +For a separately configured loopback emulator, set +`JOURNAL_BENCHMARK_AZURITE_CONNECTION_STRING` in the environment. All resolved +emulator service endpoints must be loopback addresses. + +Exercise the two implementations and five operations with metadata enabled: + +```powershell +foreach ($backend in 'AzuriteBlob', 'AzuriteTable') { + foreach ($workload in 'DurableAppend', 'CheckpointReplace', 'RecoveryReplay', 'CatalogBounded', 'CatalogUnbounded') { + $run = Join-Path $env:TEMP ("$backend-$workload-" + [guid]::NewGuid().ToString('N')) + dotnet run --project test\Benchmarks\Benchmarks.csproj -c Release -f net10.0 --no-build -- ` + Journaling.Azure --backend $backend --workload $workload --operations 3 --concurrency 2 ` + --due-journals 3 --future-journals 8 --metadata true --output $run + if ($LASTEXITCODE -ne 0) { throw "Benchmark failed; inspect $run.json and console output." } + } +} +``` + +Repeat with `--metadata false` to compare projection and stored-metadata overhead. +Emulator results establish correctness and a local baseline. Measure Azure +account performance on the intended Azure compute and storage environment. + +## Fixed-work comparisons and scale + +Each invocation selects one backend and workload. For example: + +```powershell +$run = Join-Path $env:TEMP ('standard-append-' + [guid]::NewGuid().ToString('N')) +dotnet run --project test\Benchmarks\Benchmarks.csproj -c Release -f net10.0 --no-build -- ` + Journaling.Azure --backend StandardBlob --allow-azure true --workload DurableAppend ` + --operations 100 --concurrency 8 --payload-bytes 1024 --batch-size 8 ` + --history-batches 16 --checkpoint-bytes 65536 --metadata true --output $run +``` + +For the premium comparison, point `JOURNAL_BENCHMARK_BLOB_ENDPOINT` at the premium +BlockBlobStorage account and select `PremiumBlob`. For Table, set +`JOURNAL_BENCHMARK_TABLE_ENDPOINT` and select `Table`. Sweep concurrency, batch +size, and history in separate runs, retaining the full configuration with each +result. Increasing operations and concurrency supplies a bounded, closed-loop +saturation workload: each worker starts its next operation after the previous +one finishes. + +| Option | Default | Meaning / hard limit | +| --- | --- | --- | +| `--operations` | 8 | Total measured operations, 1-10,000 | +| `--concurrency` | 2 | Fixed workers, at most operations and 256 | +| `--payload-bytes` | 256 | Bytes per fixed-length payload record | +| `--batch-size` | 4 | Records per append; product with payload bytes at most 2 MiB | +| `--history-batches` | 4 | Append batches following each seeded checkpoint, 0-10,000 | +| `--checkpoint-bytes` | 4096 | Complete replacement payload, 1 byte-32 MiB | +| `--due-journals` | 16 | Catalog due range, 0-100,000 identities | +| `--future-journals` | 64 | Catalog future tail, 0-100,000 identities | +| `--metadata` | false | Store two caller metadata properties and project catalog metadata | +| `--seed` | 42 | Deterministic payload seed | +| `--max-work` | 100,000 | Estimated logical/setup/traversal work cap, at most 10,000,000 | +| `--max-bytes` | 268,435,456 | Estimated payload/metadata processing cap, at most 8 GiB | +| `--setup-timeout-seconds` | 120 | Setup + seed + warmup deadline; also the separate verification deadline | +| `--timeout-seconds` | 60 | Whole measurement deadline | +| `--cleanup-timeout-seconds` | 60 | Independent cleanup deadline | +| `--output` | azure-journal-results | Prefix for new JSON and CSV files; parent directory must exist | + +Work estimates include warmup and verification. They bound requested logical +work and payload processing; Azure SDK retries and storage wire overhead add +service-side work. The clients use at most two SDK retries and a 10-second +network timeout within the enclosing phase deadline. + +## Operation contracts + +Setup creates deterministic buffers, a new resource, and the baseline before +measurement. Each mutation gets an independent journal and one writer. All +journals begin with one checkpoint and the requested append history. One +additional journal supplies warmup. This gives every measured append or +replacement the same initial state and bounds WAL growth. + +| Workload | One measured operation | Verification | +| --- | --- | --- | +| `DurableAppend` | Await one atomic `AppendAsync` batch commit | Fresh-handle read after measurement checks checkpoint + history + appended batch | +| `CheckpointReplace` | Await one `ReplaceAsync`, including default obsolete-generation cleanup | Fresh-handle read checks only the replacement | +| `RecoveryReplay` | Create a fresh storage handle and consume the complete checkpoint/WAL byte stream | Every chunk is consumed; ordered bytes, byte count, SHA-256, format and `IsCompleted` are checked inside the operation | +| `CatalogBounded` | Fully enumerate the prefix and inclusive due-id lower/upper bounds | Exact due membership, count, uniqueness and requested metadata | +| `CatalogUnbounded` | Fully enumerate the same prefix/lower bound with an open upper bound | Exact due + future membership, count, uniqueness and requested metadata | + +The storage format is registered as `benchmark-bytes-v1` with +`application/octet-stream`. Its input consists of deterministic fixed-length +byte records. This measures provider I/O plus validation; the existing +`Journaling` dispatch supplies the in-memory durable-state codec benchmarks. +Recovery read callbacks guarantee a consistent format; complete caller metadata +and ETags are checked using the storage metadata API outside measurement. +Catalog metadata is checked in the listing operation. + +Catalog setup also creates an expired identity below the lower bound and an +unrelated identity outside the prefix. Raising the future tail above 5,000 +Blob entries or 1,000 Table entries exercises service paging. Full enumeration +and membership checks ensure the benchmark measures retrieval and traversal. +Set `--due-journals 0` for an empty due range with a future tail, or set both +catalog counts to zero for an empty selected namespace. + +## Reports and units + +Schema version 2 JSON contains runner-owned operation measurements and +catalog/retry counters: + +| Field | Meaning | +| --- | --- | +| `Configuration` | Backend, workload, seed, sizes, work caps, requested concurrency/count, deadlines, and opt-in | +| `Build` | Loaded benchmark, Journaling, and Azure Storage assembly names, source revisions, and module version IDs | +| `Resource`, `Ownership`, `Cleanup` | Generated resource identity and lifecycle outcome | +| `AccountVerification`, `AccountKind`, `AccountSku` | Emulator, data-plane-verified Blob tier, or requested-unverified Table | +| `Completed`, `Failed`, `Cancelled`, `NotStarted` | Mutually exclusive outcomes summing to the requested operation count | +| `Operations` | Index, outcome, per-operation `Stopwatch` latency in milliseconds, successful payload bytes and items | +| `SuccessfulLatency`, `FailedLatency`, `CancelledLatency` | Count, arithmetic mean, nearest-rank p50/p95/p99, separately by outcome | +| `ElapsedSeconds`, `CompletedOperationsPerSecond` | Measurement window and completed count divided by that window | +| `CompletedPayloadBytes`, `PayloadBytesPerSecond` | Successfully committed or replayed caller bytes and bytes/window-second; catalog contributes zero | +| `CompletedItems`, `ItemUnit` | Append records, replacements, checkpoint + replayed records, or catalog entries | +| `Failures`, `Phase`, `Success` | Partial-result failure diagnostics and whole-run outcome, including verification/cleanup | +| `ProviderMetrics` | Catalog pages, candidate items, delivered entries, and explicit retries from the provider's exact DI `Microsoft.Orleans` meter | + +An append with eight records is **one** timed operation. Its p99 is the p99 of +eight-record batch commits. Recovery and catalog latency include their complete +stream consumption and correctness checks. Setup, warmup, post-write +verification, cleanup, and export are outside latency, throughput, and provider +metric windows. Small sample counts produce coarse percentiles; increase +operations for tail-latency analysis. +The fixed-work runner activates its provider metric listener at the start of +measurement and detaches it when the workers finish, before verification and +cleanup. Metric callbacks and aggregation are part of the measured fixed-work +configuration. + +CSV has one summary row, invariant-culture numeric values, quoted fields, and +the same configuration, raw operations, failures, and provider metrics in JSON +columns. Use `Import-Csv` to read it and `ConvertFrom-Json` for nested columns. +The JSON includes runtime/OS and explicit source/unit descriptions. `Build` +(also exported as CSV `build_json`) identifies each of the three loaded +assemblies using its compiler-generated `ModuleVersionId` and the hexadecimal +source revision appended by the .NET SDK to its informational version. A +`SourceRevision` of `null` explicitly records missing or unrecognized revision +metadata. Module IDs distinguish compiled code, including local changes built +on the same revision; retain the binaries with reports for exact reproduction. +This provenance comes from the loaded assemblies, so running an older build +after changing the checkout still reports that build's identity. Existing +output files are preserved by create-new writes; export failure produces a +nonzero exit and console diagnostics. + +The four `orleans-journaling-provider-` counters describe provider semantics: + +| Suffix | Count | +| --- | --- | +| `catalog-pages` | Successful pages received during catalog listing, including empty pages | +| `catalog-items` | Candidate items in those whole pages, before local filtering | +| `catalog-entries` | Entries actually yielded to the catalog consumer | +| `retries` | Explicit provider retries, grouped by reason | + +Every metric row contains `Instrument`, `Unit` (`count`), `Provider`, `Reason`, +`Observations`, and `Sum`. `Sum` is an integer counter total; `Observations` is +the number of measurements received. `Provider` preserves the provider-supplied +label, and `Reason` is populated for retry counters. A completed catalog scan +can receive more candidates than it yields. Append, replacement, and recovery +runs can have an empty `ProviderMetrics` array when they perform no explicit +provider retry. Their operation counts, outcomes, payload bytes, and latency +percentiles come from the runner's own samples. + +Use host-owned Azure SDK diagnostics, Aspire integrations, or other application +instrumentation for request counts, duration, outcomes, and transport attempts. +Reconcile cost accounting with Azure service metrics and the configured SDK's +retry and upload behavior. Catalog page/item counters measure traversal work. + +Schema version 2 identifies the catalog-only metric row shape and semantics. +Historical schema version 1 reports retain the instruments and measured build +from their original run; compare reports using their recorded schema and units. + +## Bounded BenchmarkDotNet adapter + +Default parameters select the two emulator backends, 256-byte payloads, and +five workloads: ten cases. Every iteration sets up a fresh resource and executes +one guarded invocation, then verifies and cleans up with synchronous lifecycle +bridges. Job mutators enforce one warmup and three measured iterations, +including when CLI toolchain/job selection replaces the default monitoring job. +The `Dry` smoke preset also retains the three-iteration cap. BDN's reported +operation is one provider call or complete traversal. +Before execution, a benchmark validator checks the final resolved job: 1-3 +measured iterations, 0-1 warmup iterations, one invocation, unroll factor 1, +and one process launch. It rejects configurations outside those limits before +resource setup, including adaptive iteration counts. Use the fixed-work runner +for larger explicitly sized comparisons. +Its mean describes iteration samples; use the fixed-work runner's independently +timed operation distribution for request-tail analysis. +The BDN path leaves provider metric collection inactive, so its timings include +the workload and correctness checks without the fixed-work runner's metric +listener callbacks or aggregation. +Each benchmark process emits the same build information in an `Azure benchmark +build:` JSON line during global setup. Retain the BDN log under the configured +`--artifacts` directory alongside its summary reports, especially for +multi-runtime or isolated-process comparisons. + +```powershell +# Correctness-only smoke uses the already built assembly in process. +dotnet run --project test\Benchmarks\Benchmarks.csproj -c Release -f net10.0 --no-build -- ` + Journaling.Azure.Bdn --filter '*AzureJournalBenchmarks*' --job Dry --inProcess ` + --noOverwrite --artifacts "$env:TEMP\journal-bdn-smoke" > "$env:TEMP\journal-bdn-smoke.log" 2>&1 +``` + +Use the normal isolated-process toolchain for comparisons (omit `--inProcess`); +allow sufficient `--buildTimeout` seconds for this project's dependency graph: + +```powershell +# Run on the intended benchmark host with an existing standard Azure account. +$env:JOURNAL_BENCHMARK_BLOB_ENDPOINT = 'https://yourstandardaccount.blob.core.windows.net' +$env:JOURNAL_BENCHMARK_ALLOW_AZURE = 'true' +$env:JOURNAL_BENCHMARK_BDN_BACKENDS = 'StandardBlob' +$env:JOURNAL_BENCHMARK_BDN_PAYLOAD_BYTES = '256,4096' +dotnet run --project test\Benchmarks\Benchmarks.csproj -c Release -f net10.0 --no-build -- ` + Journaling.Azure.Bdn --filter '*AzureJournalBenchmarks*DurableAppend*' --buildTimeout 600 ` + --noOverwrite --artifacts "$env:TEMP\journal-bdn-standard" > "$env:TEMP\journal-bdn-standard.log" 2>&1 +``` + +Select backends using `JOURNAL_BENCHMARK_BDN_BACKENDS` (comma-separated, at most +five) and sizes using `JOURNAL_BENCHMARK_BDN_PAYLOAD_BYTES` (at most three). +`JOURNAL_BENCHMARK_BDN_OPTIONS` accepts space-separated fixed-runner options +for history, checkpoint, catalog, metadata, and deadlines. The adapter fixes +operations/concurrency to one. `JOURNAL_BENCHMARK_ALLOW_AZURE=true` supplies +explicit cloud opt-in for BDN; endpoints remain in the same environment settings. +Use a filter selecting one workload and only the relevant backend/size values +to keep case count and service work intentional. + +## Resource ownership and failures + +Each run creates a GUID-named `journalbench...` container/table using the service +**create** operation before starting the provider lifecycle. A collision fails +setup and retains the existing resource. Cleanup deletes only a resource whose +creation was acknowledged to this run, after all workers join. Ctrl+C cancels +setup/measurement/verification while cleanup retains its own finite deadline. +An operation failure stops further scheduling, joins active workers, and exports +partial results with a nonzero exit. +Cleanup stops every attempted lifecycle stage in reverse order, including a +stage whose startup failed or was cancelled, before deleting the owned resource +and disposing the service provider. + +If creation's response is lost, `creation-unconfirmed` and +`manual-check-required` identify the resource to inspect. A final HTTP 409 also +requires inspection: it can mean a pre-existing resource or an SDK retry after +this run's successful create response was lost. Both outcomes fail setup and +preserve the resource until ownership is established. Cleanup errors retain +the resource name and `failed` outcome for operator follow-up. Abrupt process +termination can also leave a run resource; use the generated identity to inspect +and remove that individual test resource after confirming ownership. +The resource identity is printed to standard error before creation so it is +available even when the process stops before report export. diff --git a/test/Benchmarks/Program.cs b/test/Benchmarks/Program.cs index c23de2d63ff..419fe3be633 100644 --- a/test/Benchmarks/Program.cs +++ b/test/Benchmarks/Program.cs @@ -1,6 +1,7 @@ using System.Diagnostics; using BenchmarkDotNet.Running; using Benchmarks.Journaling; +using Benchmarks.Journaling.Azure; using Benchmarks.MapReduce; using Benchmarks.Ping; using Benchmarks.Placement; @@ -341,6 +342,18 @@ internal class Program typeof(DurableCommandReaderBenchmarks) ]).Run(args); }, + ["Journaling.Azure"] = args => + { + Environment.ExitCode = AzureJournalRunner.RunCommandAsync(args).GetAwaiter().GetResult(); + }, + ["Journaling.Azure.Bdn"] = args => + { + var summaries = BenchmarkSwitcher.FromTypes([typeof(AzureJournalBenchmarks)]).Run(args).ToArray(); + if (summaries.Length == 0 || summaries.Any(summary => summary.HasCriticalValidationErrors || summary.Reports.Any(report => !report.Success))) + { + Environment.ExitCode = 1; + } + }, ["suite"] = args => { _ = BenchmarkSwitcher.FromAssembly(typeof(Program).Assembly).Run(args); diff --git a/test/Extensions/Orleans.DurableJobs.AzureStorage.Tests/DurableJobs/AzureStorageDurableJobsConfigurationTests.cs b/test/Extensions/Orleans.DurableJobs.AzureStorage.Tests/DurableJobs/AzureStorageDurableJobsConfigurationTests.cs index 65918ba6d50..fd8e1c12079 100644 --- a/test/Extensions/Orleans.DurableJobs.AzureStorage.Tests/DurableJobs/AzureStorageDurableJobsConfigurationTests.cs +++ b/test/Extensions/Orleans.DurableJobs.AzureStorage.Tests/DurableJobs/AzureStorageDurableJobsConfigurationTests.cs @@ -4,6 +4,7 @@ using Orleans.DurableJobs; using Orleans.Journaling.Json; using Orleans.Hosting; +using Orleans.Journaling; using TestExtensions; using Xunit; @@ -48,6 +49,36 @@ public void UseAzureBlobDurableJobs_ConfiguresDurableJobsJsonMetadata() Assert.Contains(options.SerializerOptions.TypeInfoResolverChain, resolver => durableJobsJsonContextType.IsInstanceOfType(resolver)); } + [Theory] + [InlineData(false)] + [InlineData(true)] + public void UseAzureTableDurableJobs_ConfiguresJournalProviderAndDurableJobsJsonMetadata(bool useServiceCollection) + { + var builder = new TestSiloBuilder(); + static void Configure(AzureTableJournalStorageOptions options) + { + options.ConfigureTableServiceClient("UseDevelopmentStorage=true"); + options.TableName = "durablejobstest"; + } + + if (useServiceCollection) + { + Assert.Same(builder.Services, builder.Services.UseAzureTableDurableJobs(Configure)); + } + else + { + Assert.Same(builder, builder.UseAzureTableDurableJobs(Configure)); + } + + Assert.Contains(builder.Services, descriptor => descriptor.ServiceType == typeof(IJournalStorageProvider)); + Assert.Contains(builder.Services, descriptor => descriptor.ServiceType == typeof(IJournalStorageCatalog)); + using var serviceProvider = builder.Services.BuildServiceProvider(); + Assert.Equal("durablejobstest", serviceProvider.GetRequiredService>().Value.TableName); + var options = serviceProvider.GetRequiredService>().Value; + var contextType = typeof(DurableJob).Assembly.GetType("Orleans.DurableJobs.DurableJobsJsonContext", throwOnError: true)!; + Assert.Contains(options.SerializerOptions.TypeInfoResolverChain, resolver => contextType.IsInstanceOfType(resolver)); + } + private sealed class TestSiloBuilder : ISiloBuilder { public IServiceCollection Services { get; } = new ServiceCollection();