From d1413a3903657f65cd6773ec1475619b9da619ae Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 22 Sep 2026 12:44:25 -0700 Subject: [PATCH 01/10] Fix external client partition routing with worker metadata Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../ClientPartitionTests.cs | 616 ++++++++++++++++++ docs/providers/azure-storage.md | 11 + .../AzureStorageOrchestrationService.cs | 107 ++- src/DurableTask.AzureStorage/Storage/Queue.cs | 11 + 4 files changed, 733 insertions(+), 12 deletions(-) create mode 100644 Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs diff --git a/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs b/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs new file mode 100644 index 000000000..b3cb56b07 --- /dev/null +++ b/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs @@ -0,0 +1,616 @@ +// Copyright Microsoft Corporation +// Licensed under the Apache License, Version 2.0. + +namespace DurableTask.AzureStorage.Tests +{ + using System; + using System.Collections.Generic; + using System.Linq; + using System.Threading; + using System.Threading.Tasks; + using Azure.Core; + using Azure.Core.Pipeline; + using Azure.Data.Tables; + using Azure.Storage.Blobs; + using Azure.Storage.Blobs.Specialized; + using Azure.Storage.Queues; + using DurableTask.AzureStorage.Storage; + using DurableTask.Core; + using Microsoft.VisualStudio.TestTools.UnitTesting; + using Newtonsoft.Json.Linq; + + [TestClass] + public class ClientPartitionTests + { + readonly string connection = TestHelpers.GetTestStorageAccountConnectionString(); + + [DataTestMethod] + [DataRow(16, 4, "Table")] + [DataRow(4, 16, "Table")] + [DataRow(4, 4, "Table")] + [DataRow(16, 4, "Safe")] + [DataRow(4, 16, "Safe")] + [DataRow(4, 4, "Safe")] + [DataRow(16, 4, "Legacy")] + [DataRow(4, 16, "Legacy")] + [DataRow(4, 4, "Legacy")] + public async Task ExistingHubClientPreservesTopologyAndRouting(int clientPartitions, int targetPartitions, string manager) + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, targetPartitions, manager)); + var clientSettings = this.Settings(hub, clientPartitions, manager); + using var service = new AzureStorageOrchestrationService(clientSettings); + try + { + await target.CreateIfNotExistsAsync(); + string[] before = await this.QueueNamesAsync(hub); + Assert.AreEqual(targetPartitions, before.Length); + Assert.AreEqual(targetPartitions, await this.PartitionCountAsync(hub, manager)); + + var client = new TaskHubClient(service); + await client.CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); + await client.RaiseEventAsync(new OrchestrationInstance { InstanceId = "partition-probe" }, "Probe", "payload"); + + CollectionAssert.AreEqual(before, await this.QueueNamesAsync(hub)); + Assert.AreEqual(targetPartitions, await this.PartitionCountAsync(hub, manager)); + Assert.AreEqual(clientPartitions, clientSettings.PartitionCount, "Client discovery must not mutate caller settings."); + + string expectedQueue = AzureStorageOrchestrationService.GetControlQueueName( + hub, (int)(Fnv1aHashHelper.ComputeHash("partition-probe") % targetPartitions)); + var queues = new QueueServiceClient(this.connection); + foreach (string queueName in before) + { + int messages = (await queues.GetQueueClient(queueName).PeekMessagesAsync(2)).Value.Length; + Assert.AreEqual(queueName == expectedQueue ? 2 : 0, messages, queueName); + } + } + finally + { + await this.CleanupAsync(hub, manager); + } + } + + [DataTestMethod] + [DataRow("Table")] + [DataRow("Safe")] + [DataRow("Legacy")] + public async Task ConcurrentClientQueriesDoNotReconfigureExistingHub(string manager) + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, manager)); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, 16, manager)); + try + { + await target.CreateIfNotExistsAsync(); + var client = new TaskHubClient(service); + await Task.WhenAll(Enumerable.Range(0, 8).Select(i => client.GetOrchestrationStateAsync("missing" + i))); + Assert.AreEqual(4, (await this.QueueNamesAsync(hub)).Length); + Assert.AreEqual(4, await this.PartitionCountAsync(hub, manager)); + } + finally + { + await this.CleanupAsync(hub, manager); + } + } + + [DataTestMethod] + [DataRow(4, "Table")] + [DataRow(16, "Table")] + [DataRow(4, "Safe")] + [DataRow(16, "Safe")] + [DataRow(4, "Legacy")] + [DataRow(16, "Legacy")] + public async Task ClientCreatesMissingHubAndCanRecreateAfterDelete(int partitions, string manager) + { + string hub = NewHub(); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, partitions, manager)); + var client = new TaskHubClient(service); + try + { + Assert.AreEqual(0, (await this.QueueNamesAsync(hub)).Length); + for (int attempt = 0; attempt < 2; attempt++) + { + await client.CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); + Assert.AreEqual(partitions, (await this.QueueNamesAsync(hub)).Length); + Assert.AreEqual(partitions, await this.PartitionCountAsync(hub, manager)); + await service.DeleteAsync(); + } + } + finally + { + await this.CleanupAsync(hub, manager); + } + } + + [DataTestMethod] + [DataRow("Table")] + [DataRow("Safe")] + [DataRow("Legacy")] + public async Task ExplicitCreationRetainsConfiguredWorkerPartitions(string manager) + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, manager)); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, 16, manager)); + try + { + await target.CreateIfNotExistsAsync(); + var client = new TaskHubClient(service); + await client.GetOrchestrationStateAsync("missing"); + await service.CreateIfNotExistsAsync(); + Assert.AreEqual(16, (await this.QueueNamesAsync(hub)).Length); + Assert.AreEqual(16, await this.PartitionCountAsync(hub, manager)); + await client.CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); + var queue = new QueueClient(this.connection, AzureStorageOrchestrationService.GetControlQueueName(hub, 14)); + Assert.AreEqual(1, (await queue.PeekMessagesAsync(1)).Value.Length); + } + finally + { + await this.CleanupAsync(hub, manager); + } + } + + [DataTestMethod] + [DataRow("Table")] + [DataRow("Safe")] + [DataRow("Legacy")] + public async Task ConcurrentClientsCreateMissingHub(string manager) + { + string hub = NewHub(); + var services = Enumerable.Range(0, 8) + .Select(_ => new AzureStorageOrchestrationService(this.Settings(hub, 16, manager))) + .ToArray(); + try + { + await Task.WhenAll(services.Select((service, i) => new TaskHubClient(service) + .CreateOrchestrationInstanceAsync("Probe", string.Empty, "probe" + i, "hello"))); + Assert.AreEqual(16, (await this.QueueNamesAsync(hub)).Length); + Assert.AreEqual(16, await this.PartitionCountAsync(hub, manager)); + var expected = Enumerable.Range(0, 8) + .GroupBy(i => Fnv1aHashHelper.ComputeHash("probe" + i) % 16) + .ToDictionary(g => (int)g.Key, g => g.Count()); + var queues = new QueueServiceClient(this.connection); + for (int i = 0; i < 16; i++) + { + var queue = queues.GetQueueClient(AzureStorageOrchestrationService.GetControlQueueName(hub, i)); + Assert.AreEqual(expected.TryGetValue(i, out int count) ? count : 0, + (await queue.PeekMessagesAsync(8)).Value.Length, $"partition {i}"); + } + } + finally + { + foreach (var service in services) + { + service.Dispose(); + } + await this.CleanupAsync(hub, manager); + } + } + + [DataTestMethod] + [DataRow("null")] + [DataRow("{broken")] + [DataRow("{\"TaskHubName\":\"wrong\",\"PartitionCount\":4}")] + [DataRow("{\"PartitionCount\":0}")] + [DataRow("{\"PartitionCount\":17}")] + public async Task InvalidWorkerMetadataIsNotReplacedAndCanBeRetried(string metadata) + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, 16, "Table")); + try + { + await target.CreateIfNotExistsAsync(); + await this.WriteWorkerMetadataAsync(hub, metadata); + Exception failure = null; + try + { + await new TaskHubClient(service).GetOrchestrationStateAsync("missing"); + } + catch (Exception e) + { + failure = e; + } + Assert.IsNotNull(failure, "Malformed metadata must not be interpreted as an absent hub."); + Assert.AreEqual(metadata, await this.ReadWorkerMetadataAsync(hub)); + Assert.AreEqual(4, (await this.QueueNamesAsync(hub)).Length); + + await this.WriteWorkerMetadataAsync(hub, $"{{\"TaskHubName\":\"{hub}\",\"PartitionCount\":4}}"); + await new TaskHubClient(service).GetOrchestrationStateAsync("missing"); + Assert.AreEqual(4, (await this.QueueNamesAsync(hub)).Length); + } + finally + { + await this.CleanupAsync(hub, "Table"); + } + } + + [DataTestMethod] + [DataRow("Table")] + [DataRow("Safe")] + [DataRow("Legacy")] + public async Task UnmarkedExistingHubRetainsLegacyInitialization(string manager) + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, manager)); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, 16, manager)); + try + { + await target.CreateIfNotExistsAsync(); + await this.WriteWorkerMetadataAsync(hub, null); + await new TaskHubClient(service).GetOrchestrationStateAsync("missing"); + Assert.AreEqual(16, (await this.QueueNamesAsync(hub)).Length); + Assert.AreEqual(16, await this.PartitionCountAsync(hub, manager)); + Assert.IsNull(await this.ReadWorkerMetadataAsync(hub), "Clients must not publish authoritative worker settings."); + } + finally + { + await this.CleanupAsync(hub, manager); + } + } + + [TestMethod] + public async Task ImplicitCreationDoesNotPublishWorkerMetadata() + { + string hub = NewHub(); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); + try + { + await new TaskHubClient(service).GetOrchestrationStateAsync("missing"); + Assert.IsNull(await this.ReadWorkerMetadataAsync(hub)); + await service.CreateIfNotExistsAsync(); + Assert.AreEqual(4, JObject.Parse(await this.ReadWorkerMetadataAsync(hub)).Value("PartitionCount")); + await service.DeleteAsync(); + Assert.IsNull(await this.ReadWorkerMetadataAsync(hub)); + } + finally + { + await this.CleanupAsync(hub, "Table"); + } + } + + [DataTestMethod] + [DataRow("Table", "Safe")] + [DataRow("Safe", "Table")] + [DataRow("Legacy", "Table")] + public async Task MarkedHubClientDoesNotCreateItsOwnPartitionManager(string targetManager, string clientManager) + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, targetManager)); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, 16, clientManager)); + try + { + await target.CreateIfNotExistsAsync(); + await new TaskHubClient(service).GetOrchestrationStateAsync("missing"); + Assert.AreEqual(4, (await this.QueueNamesAsync(hub)).Length); + Assert.AreEqual(4, await this.PartitionCountAsync(hub, targetManager)); + if (clientManager == "Table") + { + var tables = new TableServiceClient(this.connection); + await foreach (var table in tables.QueryAsync(filter: $"TableName eq '{hub}Partitions'")) + { + Assert.Fail("A client must not introduce a different partition-manager store."); + } + } + } + finally + { + await this.CleanupAsync(hub, targetManager); + if (clientManager != targetManager) + { + await this.CleanupAsync(hub, clientManager); + } + } + } + + [TestMethod] + public async Task ClientDuringWorkerCreationDoesNotCachePartialPartitionTable() + { + string hub = NewHub(); + var policy = new PausePartitionCreationPolicy(hub + "Partitions"); + var tableOptions = new TableClientOptions(); + tableOptions.AddPolicy(policy, HttpPipelinePosition.PerCall); + var targetSettings = this.Settings(hub, 16, "Table"); + targetSettings.StorageAccountClientProvider = new StorageAccountClientProvider( + StorageServiceClientProvider.ForBlob(this.connection), + StorageServiceClientProvider.ForQueue(this.connection), + StorageServiceClientProvider.ForTable(this.connection, tableOptions)); + using var target = new AzureStorageOrchestrationService(targetSettings); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, 16, "Table")); + Task creation = target.CreateIfNotExistsAsync(); + try + { + Assert.AreSame(policy.FirstFourCreated, await Task.WhenAny(policy.FirstFourCreated, Task.Delay(TimeSpan.FromSeconds(10)))); + Assert.AreEqual(16, JObject.Parse(await this.ReadWorkerMetadataAsync(hub)).Value("PartitionCount")); + await new TaskHubClient(service).CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); + policy.Release(); + await creation; + var queue = new QueueClient(this.connection, AzureStorageOrchestrationService.GetControlQueueName(hub, 14)); + Assert.AreEqual(1, (await queue.PeekMessagesAsync(1)).Value.Length, + "A matching 16-partition client must not route using a partially populated four-row lease table."); + } + finally + { + policy.Release(); + await creation; + await this.CleanupAsync(hub, "Table"); + } + } + + [TestMethod] + public async Task InterruptedWorkerCreationCanBeRetried() + { + string hub = NewHub(); + var policy = new PausePartitionCreationPolicy(hub + "Partitions") { FailAfterRelease = true }; + var tableOptions = new TableClientOptions(); + tableOptions.AddPolicy(policy, HttpPipelinePosition.PerCall); + var settings = this.Settings(hub, 16, "Table"); + settings.StorageAccountClientProvider = new StorageAccountClientProvider( + StorageServiceClientProvider.ForBlob(this.connection), + StorageServiceClientProvider.ForQueue(this.connection), + StorageServiceClientProvider.ForTable(this.connection, tableOptions)); + using var target = new AzureStorageOrchestrationService(settings); + using var clientService = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); + Task creation = target.CreateIfNotExistsAsync(); + try + { + Assert.AreSame(policy.FirstFourCreated, await Task.WhenAny(policy.FirstFourCreated, Task.Delay(TimeSpan.FromSeconds(10)))); + policy.Release(); + await Assert.ThrowsExceptionAsync(() => creation); + await new TaskHubClient(clientService).CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); + Assert.AreEqual(4, await this.PartitionCountAsync(hub, "Table"), "The client must not take over lease provisioning."); + Assert.AreEqual(1, (await new QueueClient(this.connection, + AzureStorageOrchestrationService.GetControlQueueName(hub, 14)).PeekMessagesAsync(1)).Value.Length); + + policy.FailAfterRelease = false; + await target.CreateIfNotExistsAsync(); + Assert.AreEqual(16, await this.PartitionCountAsync(hub, "Table")); + } + finally + { + policy.Release(); + try + { + await creation; + } + catch (InvalidOperationException) + { + // The injected provisioning failure is expected; cleanup still owns this hub. + } + await this.CleanupAsync(hub, "Table"); + } + } + + [TestMethod] + public async Task UnreadableWorkerMetadataDoesNotFallBackAndCanBeRetried() + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); + var policy = new FailMetadataReadPolicy(); + var options = new QueueClientOptions(); + options.AddPolicy(policy, HttpPipelinePosition.PerCall); + var settings = this.Settings(hub, 16, "Table"); + settings.StorageAccountClientProvider = new StorageAccountClientProvider( + StorageServiceClientProvider.ForBlob(this.connection), + StorageServiceClientProvider.ForQueue(this.connection, options), + StorageServiceClientProvider.ForTable(this.connection)); + using var service = new AzureStorageOrchestrationService(settings); + try + { + await target.CreateIfNotExistsAsync(); + var client = new TaskHubClient(service); + await Assert.ThrowsExceptionAsync(() => client.GetOrchestrationStateAsync("missing")); + Assert.AreEqual(4, (await this.QueueNamesAsync(hub)).Length); + policy.Fail = false; + await client.GetOrchestrationStateAsync("missing"); + Assert.AreEqual(4, (await this.QueueNamesAsync(hub)).Length); + } + finally + { + await this.CleanupAsync(hub, "Table"); + } + } + + [TestMethod] + public async Task DiscoveredClientQueuesDoNotBecomeRecreatedWorkerLeases() + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 16, "Table")); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); + try + { + await target.CreateIfNotExistsAsync(); + await new TaskHubClient(service).CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); + await service.CreateAsync(); + var actual = new List(); + await foreach (TableEntity row in new TableClient(this.connection, hub + "Partitions").QueryAsync()) + { + actual.Add(row.RowKey); + } + CollectionAssert.AreEquivalent( + Enumerable.Range(0, 4).Select(i => AzureStorageOrchestrationService.GetControlQueueName(hub, i)).ToArray(), + actual.ToArray()); + Assert.AreEqual(4, service.AllControlQueues.Count()); + } + finally + { + await this.CleanupAsync(hub, "Table"); + } + } + + [TestMethod] + public async Task HubDeletionRemovesMetadataWhenAppLeaseContainerCannotBeDeleted() + { + string hub = NewHub(); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); + var container = new BlobContainerClient(this.connection, hub.ToLowerInvariant() + "-applease"); + var lease = container.GetBlobLeaseClient(); + try + { + await service.CreateIfNotExistsAsync(); + await lease.AcquireAsync(TimeSpan.FromSeconds(60)); + await service.DeleteAsync(); + Assert.IsTrue(await container.ExistsAsync(), "Legacy app-lease cleanup is best-effort while a lease is active."); + Assert.IsNull(await this.ReadWorkerMetadataAsync(hub), "Deleted hubs must not leave authoritative topology behind."); + } + finally + { + if (await container.ExistsAsync()) + { + await lease.ReleaseAsync(); + } + await this.CleanupAsync(hub, "Table"); + } + } + + [TestMethod] + public async Task WorkerPublicationPreservesUnrelatedQueueMetadata() + { + string hub = NewHub(); + var queue = new QueueClient(this.connection, hub.ToLowerInvariant() + "-workitems"); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); + using var clientService = new AzureStorageOrchestrationService(this.Settings(hub, 16, "Table")); + try + { + await queue.CreateIfNotExistsAsync(new Dictionary { ["application"] = "preserved" }); + await target.CreateIfNotExistsAsync(); + string published = await this.ReadWorkerMetadataAsync(hub); + await new TaskHubClient(clientService).GetOrchestrationStateAsync("missing"); + Assert.AreEqual(published, await this.ReadWorkerMetadataAsync(hub)); + Assert.AreEqual("preserved", (await queue.GetPropertiesAsync()).Value.Metadata["application"]); + } + finally + { + await this.CleanupAsync(hub, "Table"); + } + } + + static string NewHub() => "ClientPartitions" + Guid.NewGuid().ToString("N"); + + async Task ReadWorkerMetadataAsync(string hub) + { + var queue = new QueueClient(this.connection, hub.ToLowerInvariant() + "-workitems"); + if (!await queue.ExistsAsync()) + { + return null; + } + var metadata = (await queue.GetPropertiesAsync()).Value.Metadata; + return metadata.TryGetValue("durabletask_taskhub", out string value) ? value : null; + } + + async Task WriteWorkerMetadataAsync(string hub, string value) + { + var queue = new QueueClient(this.connection, hub.ToLowerInvariant() + "-workitems"); + var metadata = (await queue.GetPropertiesAsync()).Value.Metadata; + if (value == null) + { + metadata.Remove("durabletask_taskhub"); + } + else + { + metadata["durabletask_taskhub"] = value; + } + await queue.SetMetadataAsync(metadata); + } + + AzureStorageOrchestrationServiceSettings Settings(string hub, int partitions, string manager) => + new AzureStorageOrchestrationServiceSettings + { + StorageAccountClientProvider = new StorageAccountClientProvider(this.connection), + TaskHubName = hub, + PartitionCount = partitions, + UseTablePartitionManagement = manager == "Table", + UseLegacyPartitionManagement = manager == "Legacy", + }; + + async Task CleanupAsync(string hub, string manager) + { + using var cleanup = new AzureStorageOrchestrationService(this.Settings(hub, 16, manager)); + await cleanup.DeleteAsync(); + } + + async Task QueueNamesAsync(string hub) + { + var queues = new QueueServiceClient(this.connection); + var names = new List(); + await foreach (var queue in queues.GetQueuesAsync(prefix: hub.ToLowerInvariant() + "-control-")) + { + names.Add(queue.Name); + } + return names.OrderBy(n => n, StringComparer.Ordinal).ToArray(); + } + + async Task PartitionCountAsync(string hub, string manager) + { + if (manager != "Table") + { + var blob = new BlobClient(this.connection, hub.ToLowerInvariant() + "-leases", "taskhub.json"); + return JObject.Parse((await blob.DownloadContentAsync()).Value.Content.ToString()).Value("PartitionCount"); + } + + var table = new TableClient(this.connection, hub + "Partitions"); + int count = 0; + await foreach (TableEntity row in table.QueryAsync()) + { + count++; + } + return count; + } + + class PausePartitionCreationPolicy : HttpPipelinePolicy + { + readonly string tableName; + readonly TaskCompletionSource firstFourCreated = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + readonly TaskCompletionSource released = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + int started; + int completed; + + public PausePartitionCreationPolicy(string tableName) => this.tableName = tableName; + + public Task FirstFourCreated => this.firstFourCreated.Task; + + public bool FailAfterRelease { get; set; } + + public void Release() => this.released.TrySetResult(true); + + public override void Process(HttpMessage message, ReadOnlyMemory pipeline) => + throw new NotSupportedException("This test uses asynchronous table operations."); + + public override async ValueTask ProcessAsync(HttpMessage message, ReadOnlyMemory pipeline) + { + bool partitionInsert = message.Request.Method == RequestMethod.Post + && message.Request.Uri.ToUri().AbsolutePath.EndsWith("/" + this.tableName, StringComparison.Ordinal); + if (partitionInsert && Interlocked.Increment(ref this.started) > 4) + { + await this.released.Task; + if (this.FailAfterRelease) + { + throw new InvalidOperationException("Injected partition provisioning failure."); + } + } + await ProcessNextAsync(message, pipeline); + if (partitionInsert && Interlocked.Increment(ref this.completed) == 4) + { + this.firstFourCreated.TrySetResult(true); + } + } + + } + + class FailMetadataReadPolicy : HttpPipelinePolicy + { + public bool Fail { get; set; } = true; + + public override void Process(HttpMessage message, ReadOnlyMemory pipeline) => + throw new NotSupportedException("This test uses asynchronous queue operations."); + + public override ValueTask ProcessAsync(HttpMessage message, ReadOnlyMemory pipeline) + { + if (this.Fail && message.Request.Uri.ToUri().AbsolutePath.EndsWith("-workitems", StringComparison.Ordinal)) + { + throw new Azure.RequestFailedException(403, "Injected metadata authorization failure."); + } + return ProcessNextAsync(message, pipeline); + } + } + } +} diff --git a/docs/providers/azure-storage.md b/docs/providers/azure-storage.md index cc9960838..1590826ac 100644 --- a/docs/providers/azure-storage.md +++ b/docs/providers/azure-storage.md @@ -151,6 +151,7 @@ The Azure Storage provider creates these resources: | **Instances Table** | `{taskhub}Instances` | Instance metadata | | **Partitions Table** | `{taskhub}Partitions` | Partition leases (table manager) | | **Lease Blobs** | `{taskhub}-leases/` | Partition leases (blob manager) | +| **Worker Partition Metadata** | `durabletask_taskhub` metadata on `{taskhub}-workitems` | Complete partition configuration published during explicit hub initialization | ### Partitioning @@ -228,6 +229,16 @@ The work item queue is a simple, non-partitioned queue for activity function mes > [!IMPORTANT] > Partition count **cannot be changed** after task hub creation. Set it high enough to accommodate future scale-out needs. The maximum number of workers that can process orchestrations concurrently equals the partition count. Note that higher partition counts increase Azure Storage costs due to more queue and table operations. +#### Clients targeting another worker's hub + +Explicit hub initialization (`CreateIfNotExistsAsync`, `CreateAsync`, or worker startup) publishes the worker's configured partition count before creating partition leases. Client operations use that published count for queue initialization and message routing, rather than applying the caller's `PartitionCount` to an existing target hub. Clients do not publish this metadata or modify partition leases. + +Upgrade and initialize the target worker before creating clients that rely on partition discovery. Both the worker and client must use a version that supports this metadata. Clients cache initialization, so recreate clients that were already used before the target worker was upgraded. + +For compatibility, an absent metadata entry retains the previous behavior: client operations use their configured partition count and automatically create missing hub resources. Hubs initialized only by older workers therefore still require matching client and worker partition counts. The client does not infer a count from partition-table rows, which can be incomplete while a worker is starting. Invalid or unreadable metadata causes an error rather than falling back to a potentially incorrect count. + +The hub creation APIs remain administrative operations that apply their configured worker settings; they should not be used to discover another worker's configuration. Hub deletion removes the metadata along with the existing work-item queue, including when deletion is performed by an older provider. Unlike best-effort app-lease cleanup, work-item queue deletion failures are propagated. + ### Lease Management Workers compete for partition ownership using one of two partition managers: diff --git a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs index 8a283d22c..b7bb8ff69 100644 --- a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs +++ b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs @@ -48,6 +48,8 @@ public sealed class AzureStorageOrchestrationService : IOrchestrationServicePurgeClient, IEntityOrchestrationService { + const string WorkerTaskHubInfoMetadataKey = "durabletask_taskhub"; + static readonly HistoryEvent[] EmptyHistoryEventList = new HistoryEvent[0]; static readonly OrchestrationInstance EmptySourceInstance = new OrchestrationInstance @@ -67,6 +69,7 @@ public sealed class AzureStorageOrchestrationService : readonly ITrackingStore trackingStore; readonly ResettableLazy taskHubCreator; + readonly ResettableLazy> clientTaskHubInitializer; readonly BlobPartitionLeaseManager leaseManager; readonly AppLeaseManager appLeaseManager; readonly OrchestrationSessionManager orchestrationSessionManager; @@ -120,12 +123,7 @@ public AzureStorageOrchestrationService(AzureStorageOrchestrationServiceSettings this.messageManager = new MessageManager(this.settings, this.azureStorageClient, compressedMessageBlobContainerName); this.allControlQueues = new ConcurrentDictionary(); - for (int index = 0; index < this.settings.PartitionCount; index++) - { - var controlQueueName = GetControlQueueName(this.settings.TaskHubName, index); - ControlQueue controlQueue = new ControlQueue(this.azureStorageClient, controlQueueName, this.messageManager); - this.allControlQueues.TryAdd(controlQueue.Name, controlQueue); - } + this.InitializeControlQueues(); var workItemQueueName = GetWorkItemQueueName(this.settings.TaskHubName); this.workItemQueue = new WorkItemQueue(this.azureStorageClient, workItemQueueName, this.messageManager); @@ -145,6 +143,9 @@ public AzureStorageOrchestrationService(AzureStorageOrchestrationServiceSettings this.taskHubCreator = new ResettableLazy( this.GetTaskHubCreatorTask, LazyThreadSafetyMode.ExecutionAndPublication); + this.clientTaskHubInitializer = new ResettableLazy>( + this.InitializeClientTaskHubAsync, + LazyThreadSafetyMode.ExecutionAndPublication); this.leaseManager = GetBlobLeaseManager( this.azureStorageClient, @@ -192,6 +193,16 @@ public AzureStorageOrchestrationService(AzureStorageOrchestrationServiceSettings internal IEnumerable AllControlQueues => this.allControlQueues.Values; + void InitializeControlQueues() + { + this.allControlQueues.Clear(); + for (int index = 0; index < this.settings.PartitionCount; index++) + { + string name = GetControlQueueName(this.settings.TaskHubName, index); + this.allControlQueues.TryAdd(name, new ControlQueue(this.azureStorageClient, name, this.messageManager)); + } + } + internal IEnumerable OwnedControlQueues => this.orchestrationSessionManager.Queues; internal WorkItemQueue WorkItemQueue => this.workItemQueue; @@ -334,18 +345,84 @@ Task IEntityOrchestrationService.LockNextEntityWorkIt public async Task CreateAsync() { await this.DeleteAsync(); - await this.EnsureTaskHubAsync(); + await this.CreateIfNotExistsAsync(); } /// /// Creates the necessary Azure Storage resources for the orchestration service if they don't already exist. /// - public Task CreateIfNotExistsAsync() + public async Task CreateIfNotExistsAsync() { - return this.EnsureTaskHubAsync(); + // Publish the worker's complete topology before leases are created individually. + // Queue deletion removes this metadata, including when performed by older versions. + Queue queue = GetWorkItemQueue(this.azureStorageClient); + await queue.CreateIfNotExistsAsync(); + IDictionary metadata = await queue.GetMetadataAsync(); + metadata[WorkerTaskHubInfoMetadataKey] = + Utils.SerializeToJson(GetTaskHubInfo(this.settings.TaskHubName, this.settings.PartitionCount)); + await queue.SetMetadataAsync(metadata); + await this.EnsureTaskHubCreatedAsync(); + this.clientTaskHubInitializer.Reset(); } async Task EnsureTaskHubAsync() + { + await this.GetClientPartitionCountAsync(); + } + + async Task GetClientPartitionCountAsync() + { + try + { + return await this.clientTaskHubInitializer.Value; + } + catch (Exception e) + { + this.settings.Logger.GeneralError( + this.azureStorageClient.QueueAccountName, + this.settings.TaskHubName, + $"Failed to initialize the task hub client: {e}"); + this.clientTaskHubInitializer.Reset(); + throw; + } + } + + async Task InitializeClientTaskHubAsync() + { + Queue queue = GetWorkItemQueue(this.azureStorageClient); + IDictionary metadata = await queue.ExistsAsync() ? await queue.GetMetadataAsync() : null; + if (metadata == null || !metadata.TryGetValue(WorkerTaskHubInfoMetadataKey, out string serializedHubInfo)) + { + // Older workers do not publish topology. Preserve their initialization contract + // rather than treating a partially populated partition table as authoritative. + await this.EnsureTaskHubCreatedAsync(); + return this.settings.PartitionCount; + } + + TaskHubInfo hubInfo = Utils.DeserializeFromJson(serializedHubInfo); + if (hubInfo == null || + !string.Equals(hubInfo.TaskHubName, this.settings.TaskHubName, StringComparison.OrdinalIgnoreCase) || + hubInfo.PartitionCount < 1 || hubInfo.PartitionCount > 16) + { + throw new InvalidOperationException($"Task hub '{this.settings.TaskHubName}' has invalid worker partition metadata."); + } + + // Clients may finish creating resources after an interrupted worker initialization, + // but must not create leases or overwrite the worker's partition configuration. + var tasks = new List + { + this.trackingStore.CreateAsync(), + this.workItemQueue.CreateIfNotExistsAsync(), + }; + for (int i = 0; i < hubInfo.PartitionCount; i++) + { + tasks.Add(this.azureStorageClient.GetQueueReference(GetControlQueueName(this.settings.TaskHubName, i)).CreateIfNotExistsAsync()); + } + await Task.WhenAll(tasks); + return hubInfo.PartitionCount; + } + + async Task EnsureTaskHubCreatedAsync() { try { @@ -378,8 +455,11 @@ async Task GetTaskHubCreatorTask() tasks.Add(this.workItemQueue.CreateIfNotExistsAsync()); - foreach (ControlQueue controlQueue in this.allControlQueues.Values) + for (int index = 0; index < this.settings.PartitionCount; index++) { + string name = GetControlQueueName(this.settings.TaskHubName, index); + ControlQueue controlQueue = this.allControlQueues.GetOrAdd( + name, queueName => new ControlQueue(this.azureStorageClient, queueName, this.messageManager)); tasks.Add(controlQueue.CreateIfNotExistsAsync()); tasks.Add(this.partitionManager.CreateLease(controlQueue.Name)); } @@ -405,7 +485,7 @@ public async Task CreateAsync(bool recreateInstanceStore) this.taskHubCreator.Reset(); } - await this.taskHubCreator.Value; + await this.CreateIfNotExistsAsync(); } /// @@ -436,6 +516,8 @@ public async Task DeleteAsync(bool deleteInstanceStore) await Task.WhenAll(tasks.ToArray()); this.taskHubCreator.Reset(); + this.clientTaskHubInitializer.Reset(); + this.InitializeControlQueues(); } private Task DeleteTrackingStore() @@ -2170,7 +2252,8 @@ public Task DownloadBlobAsync(string blobUri) // be supported: https://github.com/Azure/azure-functions-durable-extension/issues/1 async Task GetControlQueueAsync(string instanceId) { - uint partitionIndex = Fnv1aHashHelper.ComputeHash(instanceId) % (uint)this.settings.PartitionCount; + int partitionCount = await this.GetClientPartitionCountAsync(); + uint partitionIndex = Fnv1aHashHelper.ComputeHash(instanceId) % (uint)partitionCount; string queueName = GetControlQueueName(this.settings.TaskHubName, (int)partitionIndex); ControlQueue cachedQueue; diff --git a/src/DurableTask.AzureStorage/Storage/Queue.cs b/src/DurableTask.AzureStorage/Storage/Queue.cs index a31913306..910ca04f3 100644 --- a/src/DurableTask.AzureStorage/Storage/Queue.cs +++ b/src/DurableTask.AzureStorage/Storage/Queue.cs @@ -39,6 +39,17 @@ public Queue(AzureStorageClient azureStorageClient, QueueServiceClient queueServ public Uri Uri => this.queueClient.Uri; + public async Task> GetMetadataAsync(CancellationToken cancellationToken = default) + { + QueueProperties properties = await this.queueClient.GetPropertiesAsync(cancellationToken).DecorateFailure(); + return properties.Metadata; + } + + public async Task SetMetadataAsync(IDictionary metadata, CancellationToken cancellationToken = default) + { + await this.queueClient.SetMetadataAsync(metadata, cancellationToken).DecorateFailure(); + } + public async Task GetApproximateMessagesCountAsync(CancellationToken cancellationToken = default) { QueueProperties properties = await this.queueClient.GetPropertiesAsync(cancellationToken).DecorateFailure(); From 0552c85bbfefab6c1457fbd334f2527194c69006 Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 22 Sep 2026 14:00:52 -0700 Subject: [PATCH 02/10] Preserve initialized worker activity admission ordering Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../ClientPartitionTests.cs | 25 +++++++++++++++++++ .../AzureStorageOrchestrationService.cs | 8 +++++- 2 files changed, 32 insertions(+), 1 deletion(-) diff --git a/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs b/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs index b3cb56b07..63780a733 100644 --- a/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs +++ b/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs @@ -410,6 +410,31 @@ public async Task UnreadableWorkerMetadataDoesNotFallBackAndCanBeRetried() } } + [TestMethod] + public async Task ExplicitInitializationDoesNotRediscoverWorkerMetadata() + { + string hub = NewHub(); + var policy = new FailMetadataReadPolicy { Fail = false }; + var options = new QueueClientOptions(); + options.AddPolicy(policy, HttpPipelinePosition.PerCall); + var settings = this.Settings(hub, 4, "Table"); + settings.StorageAccountClientProvider = new StorageAccountClientProvider( + StorageServiceClientProvider.ForBlob(this.connection), + StorageServiceClientProvider.ForQueue(this.connection, options), + StorageServiceClientProvider.ForTable(this.connection)); + using var service = new AzureStorageOrchestrationService(settings); + try + { + await service.CreateIfNotExistsAsync(); + policy.Fail = true; + Assert.IsNull(await new TaskHubClient(service).GetOrchestrationStateAsync("missing")); + } + finally + { + await this.CleanupAsync(hub, "Table"); + } + } + [TestMethod] public async Task DiscoveredClientQueuesDoNotBecomeRecreatedWorkerLeases() { diff --git a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs index b7bb8ff69..5b2171ce6 100644 --- a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs +++ b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs @@ -362,7 +362,8 @@ public async Task CreateIfNotExistsAsync() Utils.SerializeToJson(GetTaskHubInfo(this.settings.TaskHubName, this.settings.PartitionCount)); await queue.SetMetadataAsync(metadata); await this.EnsureTaskHubCreatedAsync(); - this.clientTaskHubInitializer.Reset(); + // Worker admission must not be delayed by rediscovery after explicit initialization. + this.clientTaskHubInitializer.Reset(Task.FromResult(this.settings.PartitionCount)); } async Task EnsureTaskHubAsync() @@ -2365,6 +2366,11 @@ public void Reset() { this.lazy = new Lazy(this.valueFactory, this.threadSafetyMode); } + + public void Reset(T value) + { + this.lazy = new Lazy(() => value, this.threadSafetyMode); + } } struct TaskHubQueueMessage From 8919ab5817f6b6f3f9bc9fd95f00a7b347756b5c Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 22 Sep 2026 17:26:23 -0700 Subject: [PATCH 03/10] Require worker metadata before client initialization Reject missing partition metadata without client-side hub creation and validate rewind before persistent changes. This intentionally removes unmarked-hub compatibility and client-only auto-provisioning. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../ClientPartitionTests.cs | 225 ++++++++++++++++-- docs/providers/azure-storage.md | 9 +- .../AzureStorageOrchestrationService.cs | 14 +- 3 files changed, 222 insertions(+), 26 deletions(-) diff --git a/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs b/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs index 63780a733..c8e3813a8 100644 --- a/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs +++ b/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs @@ -100,7 +100,7 @@ public async Task ConcurrentClientQueriesDoNotReconfigureExistingHub(string mana [DataRow(16, "Safe")] [DataRow(4, "Legacy")] [DataRow(16, "Legacy")] - public async Task ClientCreatesMissingHubAndCanRecreateAfterDelete(int partitions, string manager) + public async Task ClientRequiresExplicitCreationAndRecreationAfterDelete(int partitions, string manager) { string hub = NewHub(); using var service = new AzureStorageOrchestrationService(this.Settings(hub, partitions, manager)); @@ -110,6 +110,8 @@ public async Task ClientCreatesMissingHubAndCanRecreateAfterDelete(int partition Assert.AreEqual(0, (await this.QueueNamesAsync(hub)).Length); for (int attempt = 0; attempt < 2; attempt++) { + await this.AssertMissingMetadataRejectedAsync(client, hub); + await service.CreateIfNotExistsAsync(); await client.CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); Assert.AreEqual(partitions, (await this.QueueNamesAsync(hub)).Length); Assert.AreEqual(partitions, await this.PartitionCountAsync(hub, manager)); @@ -153,7 +155,7 @@ public async Task ExplicitCreationRetainsConfiguredWorkerPartitions(string manag [DataRow("Table")] [DataRow("Safe")] [DataRow("Legacy")] - public async Task ConcurrentClientsCreateMissingHub(string manager) + public async Task ConcurrentClientsRequireExplicitHubCreation(string manager) { string hub = NewHub(); var services = Enumerable.Range(0, 8) @@ -161,6 +163,8 @@ public async Task ConcurrentClientsCreateMissingHub(string manager) .ToArray(); try { + await Task.WhenAll(services.Select(service => this.AssertMissingMetadataRejectedAsync(new TaskHubClient(service), hub))); + await Task.WhenAll(services.Select(service => service.CreateIfNotExistsAsync())); await Task.WhenAll(services.Select((service, i) => new TaskHubClient(service) .CreateOrchestrationInstanceAsync("Probe", string.Empty, "probe" + i, "hello"))); Assert.AreEqual(16, (await this.QueueNamesAsync(hub)).Length); @@ -225,22 +229,36 @@ public async Task InvalidWorkerMetadataIsNotReplacedAndCanBeRetried(string metad } [DataTestMethod] - [DataRow("Table")] - [DataRow("Safe")] - [DataRow("Legacy")] - public async Task UnmarkedExistingHubRetainsLegacyInitialization(string manager) + [DataRow(16, 4, "Table")] + [DataRow(4, 16, "Table")] + [DataRow(4, 4, "Table")] + [DataRow(16, 4, "Safe")] + [DataRow(4, 16, "Safe")] + [DataRow(4, 4, "Safe")] + [DataRow(16, 4, "Legacy")] + [DataRow(4, 16, "Legacy")] + [DataRow(4, 4, "Legacy")] + public async Task UnmarkedExistingHubIsRejectedWithoutReconfiguration(int clientPartitions, int targetPartitions, string manager) { string hub = NewHub(); - using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, manager)); - using var service = new AzureStorageOrchestrationService(this.Settings(hub, 16, manager)); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, targetPartitions, manager)); + var writes = new RecordingWritePolicy(); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, clientPartitions, manager, writes)); try { await target.CreateIfNotExistsAsync(); await this.WriteWorkerMetadataAsync(hub, null); - await new TaskHubClient(service).GetOrchestrationStateAsync("missing"); - Assert.AreEqual(16, (await this.QueueNamesAsync(hub)).Length); - Assert.AreEqual(16, await this.PartitionCountAsync(hub, manager)); + var client = new TaskHubClient(service); + await this.AssertMissingMetadataRejectedAsync(client, hub); + Assert.AreEqual(0, writes.Count, "Rejected clients must not issue storage writes."); + Assert.AreEqual(targetPartitions, (await this.QueueNamesAsync(hub)).Length); + Assert.AreEqual(targetPartitions, await this.PartitionCountAsync(hub, manager)); Assert.IsNull(await this.ReadWorkerMetadataAsync(hub), "Clients must not publish authoritative worker settings."); + await target.CreateIfNotExistsAsync(); + await client.CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); + int partition = (int)(Fnv1aHashHelper.ComputeHash("partition-probe") % targetPartitions); + Assert.AreEqual(1, (await new QueueClient(this.connection, + AzureStorageOrchestrationService.GetControlQueueName(hub, partition)).PeekMessagesAsync(1)).Value.Length); } finally { @@ -249,15 +267,17 @@ public async Task UnmarkedExistingHubRetainsLegacyInitialization(string manager) } [TestMethod] - public async Task ImplicitCreationDoesNotPublishWorkerMetadata() + public async Task MissingHubQueriesDoNotCreateResourcesAndCanRetryAfterWorkerInitialization() { string hub = NewHub(); using var service = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); try { - await new TaskHubClient(service).GetOrchestrationStateAsync("missing"); + var client = new TaskHubClient(service); + await this.AssertMissingMetadataRejectedAsync(client, hub); Assert.IsNull(await this.ReadWorkerMetadataAsync(hub)); await service.CreateIfNotExistsAsync(); + Assert.IsNull(await client.GetOrchestrationStateAsync("missing")); Assert.AreEqual(4, JObject.Parse(await this.ReadWorkerMetadataAsync(hub)).Value("PartitionCount")); await service.DeleteAsync(); Assert.IsNull(await this.ReadWorkerMetadataAsync(hub)); @@ -268,6 +288,48 @@ public async Task ImplicitCreationDoesNotPublishWorkerMetadata() } } + [DataTestMethod] + [DataRow(4, 16)] + [DataRow(16, 4)] + public async Task ClientBeforeWorkerMetadataPublicationFailsWithoutWritesAndRetries(int clientPartitions, int targetPartitions) + { + string hub = NewHub(); + var pause = new PauseMetadataPublicationPolicy(); + var queueOptions = new QueueClientOptions(); + queueOptions.AddPolicy(pause, HttpPipelinePosition.PerCall); + var settings = this.Settings(hub, targetPartitions, "Table"); + settings.StorageAccountClientProvider = new StorageAccountClientProvider( + StorageServiceClientProvider.ForBlob(this.connection), + StorageServiceClientProvider.ForQueue(this.connection, queueOptions), + StorageServiceClientProvider.ForTable(this.connection)); + using var target = new AzureStorageOrchestrationService(settings); + var writes = new RecordingWritePolicy(); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, clientPartitions, "Table", writes)); + Task creation = target.CreateIfNotExistsAsync(); + try + { + Assert.AreSame(pause.Reached, await Task.WhenAny(pause.Reached, Task.Delay(TimeSpan.FromSeconds(10)))); + var client = new TaskHubClient(service); + await this.AssertMissingMetadataRejectedAsync(client, hub); + Assert.AreEqual(0, writes.Count); + Assert.AreEqual(0, (await this.QueueNamesAsync(hub)).Length); + pause.Release(); + await creation; + await client.CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); + Assert.AreEqual(targetPartitions, (await this.QueueNamesAsync(hub)).Length); + Assert.AreEqual(targetPartitions, await this.PartitionCountAsync(hub, "Table")); + int partition = (int)(Fnv1aHashHelper.ComputeHash("partition-probe") % targetPartitions); + Assert.AreEqual(1, (await new QueueClient(this.connection, + AzureStorageOrchestrationService.GetControlQueueName(hub, partition)).PeekMessagesAsync(1)).Value.Length); + } + finally + { + pause.Release(); + await creation; + await this.CleanupAsync(hub, "Table"); + } + } + [DataTestMethod] [DataRow("Table", "Safe")] [DataRow("Safe", "Table")] @@ -435,6 +497,47 @@ public async Task ExplicitInitializationDoesNotRediscoverWorkerMetadata() } } + [TestMethod] + public async Task UnmarkedHubRewindDoesNotModifyHistoryBeforeRejection() + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); + using var worker = new TaskHubWorker(target); + worker.AddTaskOrchestrations(typeof(FailedOrchestration)); + await worker.StartAsync(); + try + { + var targetClient = new TaskHubClient(target); + var instance = await targetClient.CreateOrchestrationInstanceAsync(typeof(FailedOrchestration), "hello"); + var state = await targetClient.WaitForOrchestrationAsync(instance, TimeSpan.FromSeconds(30)); + Assert.AreEqual(OrchestrationStatus.Failed, state.OrchestrationStatus); + await worker.StopAsync(); + await this.WriteWorkerMetadataAsync(hub, null); + string history = await target.GetOrchestrationHistoryAsync(instance.InstanceId, instance.ExecutionId); + string[] inventory = await this.ResourceInventoryAsync(hub); + var writes = new RecordingWritePolicy(); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, 16, "Table", writes)); + InvalidOperationException error = await Assert.ThrowsExceptionAsync( + () => service.RewindTaskOrchestrationAsync(instance.InstanceId, "retry")); + StringAssert.Contains(error.Message, "CreateIfNotExistsAsync"); + Assert.AreEqual(0, writes.Count, "Rejected rewind must not mutate history or instance status."); + Assert.AreEqual(history, await target.GetOrchestrationHistoryAsync(instance.InstanceId, instance.ExecutionId)); + Assert.AreEqual(OrchestrationStatus.Failed, (await targetClient.GetOrchestrationStateAsync(instance.InstanceId)).OrchestrationStatus); + CollectionAssert.AreEqual(inventory, await this.ResourceInventoryAsync(hub)); + + await target.CreateIfNotExistsAsync(); + await service.RewindTaskOrchestrationAsync(instance.InstanceId, "retry"); + Assert.AreEqual(OrchestrationStatus.Pending, (await targetClient.GetOrchestrationStateAsync(instance.InstanceId)).OrchestrationStatus); + string queue = AzureStorageOrchestrationService.GetControlQueueName(hub, (int)(Fnv1aHashHelper.ComputeHash(instance.InstanceId) % 4)); + Assert.AreEqual(1, (await new QueueClient(this.connection, queue).PeekMessagesAsync(2)).Value.Length); + } + finally + { + await worker.StopAsync(isForced: true); + await this.CleanupAsync(hub, "Table"); + } + } + [TestMethod] public async Task DiscoveredClientQueuesDoNotBecomeRecreatedWorkerLeases() { @@ -537,16 +640,66 @@ async Task WriteWorkerMetadataAsync(string hub, string value) await queue.SetMetadataAsync(metadata); } - AzureStorageOrchestrationServiceSettings Settings(string hub, int partitions, string manager) => + async Task AssertMissingMetadataRejectedAsync(TaskHubClient client, string hub) + { + string[] before = await this.ResourceInventoryAsync(hub); + InvalidOperationException error = await Assert.ThrowsExceptionAsync( + () => client.CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello")); + StringAssert.Contains(error.Message, hub); + StringAssert.Contains(error.Message, "CreateIfNotExistsAsync"); + await Assert.ThrowsExceptionAsync(() => client.GetOrchestrationStateAsync("missing")); + await Assert.ThrowsExceptionAsync(() => client.RaiseEventAsync( + new OrchestrationInstance { InstanceId = "partition-probe" }, "event", "hello")); + CollectionAssert.AreEqual(before, await this.ResourceInventoryAsync(hub)); + } + + async Task ResourceInventoryAsync(string hub) + { + var resources = new List(); + await foreach (var queue in new QueueServiceClient(this.connection).GetQueuesAsync(prefix: hub.ToLowerInvariant())) + { + resources.Add("queue:" + queue.Name); + Assert.AreEqual(0, (await new QueueClient(this.connection, queue.Name).PeekMessagesAsync(1)).Value.Length); + } + await foreach (var container in new BlobServiceClient(this.connection).GetBlobContainersAsync(prefix: hub.ToLowerInvariant())) + { + resources.Add("container:" + container.Name); + } + await foreach (var table in new TableServiceClient(this.connection).QueryAsync(filter: $"TableName ge '{hub}' and TableName lt '{hub}z'")) + { + resources.Add("table:" + table.Name); + } + return resources.OrderBy(x => x, StringComparer.Ordinal).ToArray(); + } + + AzureStorageOrchestrationServiceSettings Settings(string hub, int partitions, string manager, RecordingWritePolicy writes = null) => new AzureStorageOrchestrationServiceSettings { - StorageAccountClientProvider = new StorageAccountClientProvider(this.connection), + StorageAccountClientProvider = this.CreateStorageProvider(writes), TaskHubName = hub, PartitionCount = partitions, UseTablePartitionManagement = manager == "Table", UseLegacyPartitionManagement = manager == "Legacy", }; + StorageAccountClientProvider CreateStorageProvider(RecordingWritePolicy writes) + { + if (writes == null) + { + return new StorageAccountClientProvider(this.connection); + } + var blobs = new BlobClientOptions(); + var queues = new QueueClientOptions(); + var tables = new TableClientOptions(); + blobs.AddPolicy(writes, HttpPipelinePosition.PerCall); + queues.AddPolicy(writes, HttpPipelinePosition.PerCall); + tables.AddPolicy(writes, HttpPipelinePosition.PerCall); + return new StorageAccountClientProvider( + StorageServiceClientProvider.ForBlob(this.connection, blobs), + StorageServiceClientProvider.ForQueue(this.connection, queues), + StorageServiceClientProvider.ForTable(this.connection, tables)); + } + async Task CleanupAsync(string hub, string manager) { using var cleanup = new AzureStorageOrchestrationService(this.Settings(hub, 16, manager)); @@ -581,6 +734,48 @@ async Task PartitionCountAsync(string hub, string manager) return count; } + public class FailedOrchestration : TaskOrchestration + { + public override Task RunTask(OrchestrationContext context, string input) => + throw new InvalidOperationException("Expected test orchestration failure."); + } + + class RecordingWritePolicy : HttpPipelinePolicy + { + int count; + public int Count => this.count; + public override void Process(HttpMessage message, ReadOnlyMemory pipeline) => + throw new NotSupportedException("Tests use asynchronous storage operations."); + public override ValueTask ProcessAsync(HttpMessage message, ReadOnlyMemory pipeline) + { + if (message.Request.Method != RequestMethod.Get && message.Request.Method != RequestMethod.Head) + { + Interlocked.Increment(ref this.count); + } + return ProcessNextAsync(message, pipeline); + } + } + + class PauseMetadataPublicationPolicy : HttpPipelinePolicy + { + readonly TaskCompletionSource reached = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + readonly TaskCompletionSource released = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + public Task Reached => this.reached.Task; + public void Release() => this.released.TrySetResult(true); + public override void Process(HttpMessage message, ReadOnlyMemory pipeline) => + throw new NotSupportedException("Tests use asynchronous queue operations."); + public override async ValueTask ProcessAsync(HttpMessage message, ReadOnlyMemory pipeline) + { + if (message.Request.Method == RequestMethod.Put && + message.Request.Uri.ToUri().Query.Contains("comp=metadata")) + { + this.reached.TrySetResult(true); + await this.released.Task; + } + await ProcessNextAsync(message, pipeline); + } + } + class PausePartitionCreationPolicy : HttpPipelinePolicy { readonly string tableName; diff --git a/docs/providers/azure-storage.md b/docs/providers/azure-storage.md index 1590826ac..bac8b50dc 100644 --- a/docs/providers/azure-storage.md +++ b/docs/providers/azure-storage.md @@ -233,9 +233,14 @@ The work item queue is a simple, non-partitioned queue for activity function mes Explicit hub initialization (`CreateIfNotExistsAsync`, `CreateAsync`, or worker startup) publishes the worker's configured partition count before creating partition leases. Client operations use that published count for queue initialization and message routing, rather than applying the caller's `PartitionCount` to an existing target hub. Clients do not publish this metadata or modify partition leases. -Upgrade and initialize the target worker before creating clients that rely on partition discovery. Both the worker and client must use a version that supports this metadata. Clients cache initialization, so recreate clients that were already used before the target worker was upgraded. +> [!WARNING] +> **Breaking change:** Client operations now require published worker partition metadata. Client-only automatic hub creation is no longer supported. Missing metadata is rejected even if the client's partition count matches the target's. Existing hubs initialized only by older workers must be explicitly initialized with an upgraded provider before these clients can use them. + +Deploy the upgraded target worker first and let it publish metadata, or explicitly call `CreateIfNotExistsAsync` using the target's existing partition configuration, before issuing client operations. `CreateAsync` and worker startup also remain able to initialize a hub. Do not change an existing hub's partition count during this upgrade. + +If the work-item queue or its `durabletask_taskhub` metadata entry is absent, the client throws an actionable exception without creating queues, tables, blobs, or leases or enqueuing the message. This includes calls racing with worker startup before metadata publication. Failed initialization is not cached: retrying the same client after publication succeeds. There is no fallback to the caller's count or a partially populated partition table. Malformed metadata and storage access failures remain explicit errors. -For compatibility, an absent metadata entry retains the previous behavior: client operations use their configured partition count and automatically create missing hub resources. Hubs initialized only by older workers therefore still require matching client and worker partition counts. The client does not infer a count from partition-table rows, which can be incomplete while a worker is starting. Invalid or unreadable metadata causes an error rather than falling back to a potentially incorrect count. +Both worker and client binaries must be upgraded to benefit from this behavior. Older client binaries ignore the metadata and are not fixed by upgrading the worker alone. Successfully initialized clients cache topology; recreate them after deleting and recreating a hub externally. Explicit worker initialization preserves its already-completed initialization path for local clients and activity admission. The hub creation APIs remain administrative operations that apply their configured worker settings; they should not be used to discover another worker's configuration. Hub deletion removes the metadata along with the existing work-item queue, including when deletion is performed by an older provider. Unlike best-effort app-lease cleanup, work-item queue deletion failures are propagated. diff --git a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs index 5b2171ce6..9864edee7 100644 --- a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs +++ b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs @@ -394,10 +394,10 @@ async Task InitializeClientTaskHubAsync() IDictionary metadata = await queue.ExistsAsync() ? await queue.GetMetadataAsync() : null; if (metadata == null || !metadata.TryGetValue(WorkerTaskHubInfoMetadataKey, out string serializedHubInfo)) { - // Older workers do not publish topology. Preserve their initialization contract - // rather than treating a partially populated partition table as authoritative. - await this.EnsureTaskHubCreatedAsync(); - return this.settings.PartitionCount; + throw new InvalidOperationException( + $"Task hub '{this.settings.TaskHubName}' is missing worker partition metadata. " + + "Initialize the target with an upgraded worker or explicitly call CreateIfNotExistsAsync " + + "using the target's partition configuration before retrying the client operation."); } TaskHubInfo hubInfo = Utils.DeserializeFromJson(serializedHubInfo); @@ -1847,7 +1847,6 @@ public async Task CreateTaskOrchestrationAsync(TaskMessage creationMessage, Orch Utils.ConvertDateTimeInHistoryEventsToUTC(creationMessage.Event); - // Client operations will auto-create the task hub if it doesn't already exist. await this.EnsureTaskHubAsync(); InstanceStatus existingInstance = await this.trackingStore.FetchInstanceStatusAsync( @@ -1916,7 +1915,6 @@ public Task SendTaskOrchestrationMessageBatchAsync(params TaskMessage[] messages /// The message to send. public async Task SendTaskOrchestrationMessageAsync(TaskMessage message) { - // Client operations will auto-create the task hub if it doesn't already exist. await this.EnsureTaskHubAsync(); ControlQueue controlQueue = await this.GetControlQueueAsync(message.OrchestrationInstance.InstanceId); await this.SendTaskOrchestrationMessageInternalAsync(EmptySourceInstance, controlQueue, message); @@ -1938,7 +1936,6 @@ internal Task SendTaskOrchestrationMessageInternalAsync( /// List of objects that represent the list of orchestrations. public async Task> GetOrchestrationStateAsync(string instanceId, bool allExecutions) { - // Client operations will auto-create the task hub if it doesn't already exist. await this.EnsureTaskHubAsync(); return new OrchestrationState[] { @@ -1954,7 +1951,6 @@ await this.trackingStore.GetStateAsync(instanceId, allExecutions, fetchInput: tr /// The object that represents the orchestration. public async Task GetOrchestrationStateAsync(string instanceId, string executionId) { - // Client operations will auto-create the task hub if it doesn't already exist. await this.EnsureTaskHubAsync(); return await this.trackingStore.GetStateAsync(instanceId, executionId, fetchInput: true); } @@ -1969,7 +1965,6 @@ public async Task GetOrchestrationStateAsync(string instance /// List of objects that represent the list of orchestrations. public async Task> GetOrchestrationStateAsync(string instanceId, bool allExecutions, bool fetchInput = true) { - // Client operations will auto-create the task hub if it doesn't already exist. await this.EnsureTaskHubAsync(); return await this.trackingStore.GetStateAsync(instanceId, allExecutions, fetchInput).ToListAsync(); } @@ -2065,6 +2060,7 @@ public async Task ForceTerminateTaskOrchestrationAsync(string instanceId, string /// The reason for rewinding. public async Task RewindTaskOrchestrationAsync(string instanceId, string reason) { + await this.EnsureTaskHubAsync(); List queueIds = await this.trackingStore.RewindHistoryAsync(instanceId).ToListAsync(); foreach (string id in queueIds) From d59751b6b950cd1773e3658821511339ef38e283 Mon Sep 17 00:00:00 2001 From: wangbill Date: Tue, 22 Sep 2026 18:18:59 -0700 Subject: [PATCH 04/10] Initialize partition-manager test hubs before client use Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../TestTablePartitionManager.cs | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/test/DurableTask.AzureStorage.Tests/TestTablePartitionManager.cs b/test/DurableTask.AzureStorage.Tests/TestTablePartitionManager.cs index 7fdb42e04..f89ca4177 100644 --- a/test/DurableTask.AzureStorage.Tests/TestTablePartitionManager.cs +++ b/test/DurableTask.AzureStorage.Tests/TestTablePartitionManager.cs @@ -529,6 +529,8 @@ public async Task TestUnhealthyWorker() taskHubWorkers[i].AddTaskOrchestrations(typeof(LongRunningOrchestrator)); } + await services[0].CreateIfNotExistsAsync(); + // Create 50 orchestration instances. var client = new TaskHubClient(services[0]); var createInstanceTasks = new Task[InstanceCount]; @@ -611,6 +613,8 @@ public async Task EnsureOwnedQueueExclusive() taskHubWorkers[i].AddTaskActivities(typeof(Hello)); } + await services[0].CreateIfNotExistsAsync(); + // Create 100 orchestration instances. var client = new TaskHubClient(services[0]); var createInstanceTasks = new Task[InstanceCount]; From 957b7c4a51c07ebbe21e1497feda33c8ab1dc94c Mon Sep 17 00:00:00 2001 From: wangbill Date: Wed, 23 Sep 2026 09:36:18 -0700 Subject: [PATCH 05/10] Retain discovered queues for complete client hub deletion Cache validated topology for cleanup, assert precise metadata errors, and place regression tests in the canonical lowercase project tree. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 53306d22-2f23-4522-b14f-fc52789a4292 --- docs/providers/azure-storage.md | 2 + .../AzureStorageOrchestrationService.cs | 5 +- .../ClientPartitionTests.cs | 131 ++++++++++++++++-- 3 files changed, 124 insertions(+), 14 deletions(-) rename {Test => test}/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs (87%) diff --git a/docs/providers/azure-storage.md b/docs/providers/azure-storage.md index bac8b50dc..4dafa111c 100644 --- a/docs/providers/azure-storage.md +++ b/docs/providers/azure-storage.md @@ -244,6 +244,8 @@ Both worker and client binaries must be upgraded to benefit from this behavior. The hub creation APIs remain administrative operations that apply their configured worker settings; they should not be used to discover another worker's configuration. Hub deletion removes the metadata along with the existing work-item queue, including when deletion is performed by an older provider. Unlike best-effort app-lease cleanup, work-item queue deletion failures are propagated. +Client initialization retains references to every discovered control queue. Deleting the hub through that service therefore removes the full discovered topology, including queues the client has never sent to, even when the caller's configured partition count is smaller. + ### Lease Management Workers compete for partition ownership using one of two partition managers: diff --git a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs index 9864edee7..90cf42d36 100644 --- a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs +++ b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs @@ -417,7 +417,10 @@ async Task InitializeClientTaskHubAsync() }; for (int i = 0; i < hubInfo.PartitionCount; i++) { - tasks.Add(this.azureStorageClient.GetQueueReference(GetControlQueueName(this.settings.TaskHubName, i)).CreateIfNotExistsAsync()); + string name = GetControlQueueName(this.settings.TaskHubName, i); + ControlQueue controlQueue = this.allControlQueues.GetOrAdd( + name, queueName => new ControlQueue(this.azureStorageClient, queueName, this.messageManager)); + tasks.Add(controlQueue.CreateIfNotExistsAsync()); } await Task.WhenAll(tasks); return hubInfo.PartitionCount; diff --git a/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs b/test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs similarity index 87% rename from Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs rename to test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs index c8e3813a8..661cf1e12 100644 --- a/Test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs +++ b/test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs @@ -17,6 +17,7 @@ namespace DurableTask.AzureStorage.Tests using DurableTask.AzureStorage.Storage; using DurableTask.Core; using Microsoft.VisualStudio.TestTools.UnitTesting; + using Newtonsoft.Json; using Newtonsoft.Json.Linq; [TestClass] @@ -191,12 +192,12 @@ public async Task ConcurrentClientsRequireExplicitHubCreation(string manager) } [DataTestMethod] - [DataRow("null")] - [DataRow("{broken")] - [DataRow("{\"TaskHubName\":\"wrong\",\"PartitionCount\":4}")] - [DataRow("{\"PartitionCount\":0}")] - [DataRow("{\"PartitionCount\":17}")] - public async Task InvalidWorkerMetadataIsNotReplacedAndCanBeRetried(string metadata) + [DataRow("null", false)] + [DataRow("{broken", true)] + [DataRow("{\"TaskHubName\":\"wrong\",\"PartitionCount\":4}", false)] + [DataRow("{\"PartitionCount\":0}", false)] + [DataRow("{\"PartitionCount\":17}", false)] + public async Task InvalidWorkerMetadataIsNotReplacedAndCanBeRetried(string metadata, bool invalidJson) { string hub = NewHub(); using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); @@ -205,21 +206,20 @@ public async Task InvalidWorkerMetadataIsNotReplacedAndCanBeRetried(string metad { await target.CreateIfNotExistsAsync(); await this.WriteWorkerMetadataAsync(hub, metadata); - Exception failure = null; - try + var client = new TaskHubClient(service); + if (invalidJson) { - await new TaskHubClient(service).GetOrchestrationStateAsync("missing"); + await Assert.ThrowsExceptionAsync(() => client.GetOrchestrationStateAsync("missing")); } - catch (Exception e) + else { - failure = e; + await Assert.ThrowsExceptionAsync(() => client.GetOrchestrationStateAsync("missing")); } - Assert.IsNotNull(failure, "Malformed metadata must not be interpreted as an absent hub."); Assert.AreEqual(metadata, await this.ReadWorkerMetadataAsync(hub)); Assert.AreEqual(4, (await this.QueueNamesAsync(hub)).Length); await this.WriteWorkerMetadataAsync(hub, $"{{\"TaskHubName\":\"{hub}\",\"PartitionCount\":4}}"); - await new TaskHubClient(service).GetOrchestrationStateAsync("missing"); + await client.GetOrchestrationStateAsync("missing"); Assert.AreEqual(4, (await this.QueueNamesAsync(hub)).Length); } finally @@ -497,6 +497,93 @@ public async Task ExplicitInitializationDoesNotRediscoverWorkerMetadata() } } + [DataTestMethod] + [DataRow(4, 16, "Table", false)] + [DataRow(4, 16, "Table", true)] + [DataRow(16, 4, "Table", false)] + [DataRow(16, 4, "Table", true)] + [DataRow(4, 4, "Table", false)] + [DataRow(4, 4, "Table", true)] + [DataRow(4, 16, "Safe", false)] + [DataRow(4, 16, "Safe", true)] + [DataRow(16, 4, "Safe", false)] + [DataRow(16, 4, "Safe", true)] + [DataRow(4, 4, "Safe", false)] + [DataRow(4, 4, "Safe", true)] + [DataRow(4, 16, "Legacy", false)] + [DataRow(4, 16, "Legacy", true)] + [DataRow(16, 4, "Legacy", false)] + [DataRow(16, 4, "Legacy", true)] + [DataRow(4, 4, "Legacy", false)] + [DataRow(4, 4, "Legacy", true)] + public async Task ClientDeletionRemovesEntireDiscoveredTopology( + int clientPartitions, int targetPartitions, string manager, bool send) + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, targetPartitions, manager)); + using var service = new AzureStorageOrchestrationService(this.Settings(hub, clientPartitions, manager)); + var client = new TaskHubClient(service); + try + { + await target.CreateIfNotExistsAsync(); + Assert.AreEqual(targetPartitions, (await this.QueueNamesAsync(hub)).Length); + Assert.IsNull(await client.GetOrchestrationStateAsync("missing")); + if (send) + { + await client.CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); + } + + await service.DeleteAsync(); + Assert.AreEqual(0, (await this.QueueNamesAsync(hub)).Length, + "Deletion must include every discovered control queue, not only configured or routed queues."); + await this.AssertMissingMetadataRejectedAsync(client, hub); + await service.CreateIfNotExistsAsync(); + CollectionAssert.AreEqual( + Enumerable.Range(0, clientPartitions).Select(i => AzureStorageOrchestrationService.GetControlQueueName(hub, i)).ToArray(), + await this.QueueNamesAsync(hub)); + Assert.AreEqual(clientPartitions, await this.PartitionCountAsync(hub, manager)); + Assert.AreEqual(clientPartitions, service.AllControlQueues.Count()); + await client.CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); + } + finally + { + await this.CleanupAsync(hub, manager); + } + } + + [TestMethod] + public async Task FailedClientInitializationRetainsDiscoveredQueuesForRetryAndDeletion() + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 16, "Table")); + var failure = new FailControlQueueCreationPolicy(); + var queueOptions = new QueueClientOptions(); + queueOptions.AddPolicy(failure, HttpPipelinePosition.PerCall); + var settings = this.Settings(hub, 4, "Table"); + settings.StorageAccountClientProvider = new StorageAccountClientProvider( + StorageServiceClientProvider.ForBlob(this.connection), + StorageServiceClientProvider.ForQueue(this.connection, queueOptions), + StorageServiceClientProvider.ForTable(this.connection)); + using var service = new AzureStorageOrchestrationService(settings); + var client = new TaskHubClient(service); + try + { + await target.CreateIfNotExistsAsync(); + await Assert.ThrowsExceptionAsync(() => client.GetOrchestrationStateAsync("missing")); + Assert.AreEqual(16, service.AllControlQueues.Count(), + "References created before an initialization failure must remain available for cleanup."); + failure.Fail = false; + Assert.IsNull(await client.GetOrchestrationStateAsync("missing")); + await service.DeleteAsync(); + Assert.AreEqual(0, (await this.QueueNamesAsync(hub)).Length); + Assert.AreEqual(4, service.AllControlQueues.Count()); + } + finally + { + await this.CleanupAsync(hub, "Table"); + } + } + [TestMethod] public async Task UnmarkedHubRewindDoesNotModifyHistoryBeforeRejection() { @@ -832,5 +919,23 @@ public override ValueTask ProcessAsync(HttpMessage message, ReadOnlyMemory pipeline) => + throw new NotSupportedException("Tests use asynchronous queue operations."); + + public override ValueTask ProcessAsync(HttpMessage message, ReadOnlyMemory pipeline) + { + if (this.Fail && message.Request.Method == RequestMethod.Put && + message.Request.Uri.ToUri().AbsolutePath.EndsWith("-control-15", StringComparison.Ordinal)) + { + throw new Azure.RequestFailedException(500, "Injected control-queue initialization failure."); + } + return ProcessNextAsync(message, pipeline); + } + } } } From 678602c4ba84ff8cf0f24caa48f16650ccc11bd1 Mon Sep 17 00:00:00 2001 From: wangbill Date: Wed, 23 Sep 2026 10:41:05 -0700 Subject: [PATCH 06/10] Delete the shared Safe partition lease container once Intent and ownership leases share a container. Remove the competing duplicate delete while preserving idempotent absence and propagation of other storage errors. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 53306d22-2f23-4522-b14f-fc52789a4292 --- docs/providers/azure-storage.md | 2 + .../Partitioning/SafePartitionManager.cs | 21 +-- .../LeaseContainerDeletionTests.cs | 155 ++++++++++++++++++ 3 files changed, 159 insertions(+), 19 deletions(-) create mode 100644 test/DurableTask.AzureStorage.Tests/LeaseContainerDeletionTests.cs diff --git a/docs/providers/azure-storage.md b/docs/providers/azure-storage.md index 4dafa111c..bbfa53937 100644 --- a/docs/providers/azure-storage.md +++ b/docs/providers/azure-storage.md @@ -266,6 +266,8 @@ When `UseTablePartitionManagement = false`: - Uses Azure Blob leases for concurrency control - Available in "safe" (`UseLegacyPartitionManagement = false`) and "legacy" (`UseLegacyPartitionManagement = true`) variants +Safe-mode intent and ownership leases use different prefixes in the same container. Hub deletion deletes that shared container once, not once per prefix. An already-absent container is handled idempotently; other storage deletion failures are propagated. + #### Partition lifecycle 1. Workers acquire leases to claim partition ownership diff --git a/src/DurableTask.AzureStorage/Partitioning/SafePartitionManager.cs b/src/DurableTask.AzureStorage/Partitioning/SafePartitionManager.cs index 7a95f4c84..a5801265a 100644 --- a/src/DurableTask.AzureStorage/Partitioning/SafePartitionManager.cs +++ b/src/DurableTask.AzureStorage/Partitioning/SafePartitionManager.cs @@ -15,10 +15,8 @@ namespace DurableTask.AzureStorage.Partitioning { using System; using System.Collections.Generic; - using System.Runtime.ExceptionServices; using System.Threading; using System.Threading.Tasks; - using Azure; using DurableTask.AzureStorage.Storage; class SafePartitionManager : IPartitionManager @@ -104,23 +102,8 @@ Task IPartitionManager.CreateLeaseStore() Task IPartitionManager.DeleteLeases() { - return Task.WhenAll( - this.intentLeaseManager.DeleteAllAsync(), - this.ownershipLeaseManager.DeleteAllAsync() - ).ContinueWith(t => - { - if (t.Exception?.InnerExceptions?.Count > 0) - { - foreach (Exception e in t.Exception.InnerExceptions) - { - RequestFailedException storageException = e as RequestFailedException; - if (storageException == null || storageException.Status != 404) - { - ExceptionDispatchInfo.Capture(e).Throw(); - } - } - } - }); + // Intent and ownership leases share one container; deleting it removes both. + return this.intentLeaseManager.DeleteAllAsync(); } async Task IPartitionManager.StartAsync() diff --git a/test/DurableTask.AzureStorage.Tests/LeaseContainerDeletionTests.cs b/test/DurableTask.AzureStorage.Tests/LeaseContainerDeletionTests.cs new file mode 100644 index 000000000..71749809f --- /dev/null +++ b/test/DurableTask.AzureStorage.Tests/LeaseContainerDeletionTests.cs @@ -0,0 +1,155 @@ +// Copyright Microsoft Corporation +// Licensed under the Apache License, Version 2.0. + +namespace DurableTask.AzureStorage.Tests +{ + using System; + using System.Collections.Concurrent; + using System.Linq; + using System.Threading; + using System.Threading.Tasks; + using Azure; + using Azure.Core; + using Azure.Core.Pipeline; + using Azure.Storage.Blobs; + using DurableTask.AzureStorage.Storage; + using Microsoft.VisualStudio.TestTools.UnitTesting; + + [TestClass] + public class LeaseContainerDeletionTests + { + readonly string connection = TestHelpers.GetTestStorageAccountConnectionString(); + + [DataTestMethod] + [DataRow(false)] + [DataRow(true)] + public async Task BlobPartitionManagerDeletesSharedContainerOnce(bool legacy) + { + string hub = "SingleDelete" + Guid.NewGuid().ToString("N"); + var observer = new ContainerDeletePolicy(hub); + using var service = this.CreateService(hub, legacy, observer); + await service.CreateIfNotExistsAsync(); + var container = new BlobContainerClient(this.connection, hub.ToLowerInvariant() + "-leases"); + var names = new ConcurrentBag(); + await foreach (var blob in container.GetBlobsAsync()) + { + names.Add(blob.Name); + } + if (!legacy) + { + Assert.IsTrue(names.Any(n => n.StartsWith("intent/", StringComparison.Ordinal))); + Assert.IsTrue(names.Any(n => n.StartsWith("ownership/", StringComparison.Ordinal))); + } + + Task deletion = service.DeleteAsync(); + try + { + Assert.AreSame(observer.FirstDelete, await Task.WhenAny(observer.FirstDelete, Task.Delay(TimeSpan.FromSeconds(10)))); + observer.Release(); + await deletion; + Assert.AreEqual(1, observer.Requests.Count, "One container must not receive concurrent duplicate DELETE requests."); + Assert.AreEqual(1, observer.Requests.Distinct().Count()); + Assert.IsFalse(await container.ExistsAsync()); + await service.DeleteAsync(); + Assert.AreEqual(2, observer.Requests.Count, "A later idempotent delete should issue only one request of its own."); + } + finally + { + observer.Release(); + Console.WriteLine($"Lease container DELETE requests: {observer.Requests.Count}; distinct container paths: {observer.Requests.Distinct().Count()}."); + await container.DeleteIfExistsAsync(); + } + } + + [DataTestMethod] + [DataRow(403, "AuthorizationPermissionMismatch")] + [DataRow(409, "ConcurrentContainerOperationInProgress")] + public async Task SafeDeletePropagatesStorageFailureAndCanBeRetried(int status, string errorCode) + { + string hub = "DeleteFailure" + Guid.NewGuid().ToString("N"); + var observer = new ContainerDeletePolicy(hub) { FailureStatus = status, FailureCode = errorCode }; + observer.Release(); + using var service = this.CreateService(hub, legacy: false, observer); + try + { + await service.CreateIfNotExistsAsync(); + DurableTaskStorageException error = await Assert.ThrowsExceptionAsync(() => service.DeleteAsync()); + Assert.AreEqual(status, error.HttpStatusCode); + Assert.AreEqual(errorCode, error.ErrorCode); + Assert.AreEqual(1, observer.Requests.Count); + observer.FailureStatus = null; + await service.DeleteAsync(); + Assert.AreEqual(2, observer.Requests.Count); + Assert.IsFalse(await new BlobContainerClient(this.connection, hub.ToLowerInvariant() + "-leases").ExistsAsync()); + } + finally + { + observer.Release(); + await new BlobContainerClient(this.connection, hub.ToLowerInvariant() + "-leases").DeleteIfExistsAsync(); + } + } + + AzureStorageOrchestrationService CreateService(string hub, bool legacy, ContainerDeletePolicy observer) + { + var options = new BlobClientOptions(); + options.AddPolicy(observer, HttpPipelinePosition.PerCall); + return new AzureStorageOrchestrationService(new AzureStorageOrchestrationServiceSettings + { + TaskHubName = hub, + PartitionCount = 4, + UseTablePartitionManagement = false, + UseLegacyPartitionManagement = legacy, + StorageAccountClientProvider = new StorageAccountClientProvider( + StorageServiceClientProvider.ForBlob(this.connection, options), + StorageServiceClientProvider.ForQueue(this.connection), + StorageServiceClientProvider.ForTable(this.connection)), + }); + } + + class ContainerDeletePolicy : HttpPipelinePolicy + { + readonly string suffix; + readonly TaskCompletionSource firstDelete = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + readonly TaskCompletionSource released = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + int active; + + public ContainerDeletePolicy(string hub) => this.suffix = "/" + hub.ToLowerInvariant() + "-leases"; + public ConcurrentQueue Requests { get; } = new ConcurrentQueue(); + public Task FirstDelete => this.firstDelete.Task; + public int? FailureStatus { get; set; } + public string FailureCode { get; set; } + public void Release() => this.released.TrySetResult(true); + public override void Process(HttpMessage message, ReadOnlyMemory pipeline) => + throw new NotSupportedException("Tests use asynchronous storage requests."); + + public override async ValueTask ProcessAsync(HttpMessage message, ReadOnlyMemory pipeline) + { + if (message.Request.Method != RequestMethod.Delete || + !message.Request.Uri.ToUri().AbsolutePath.EndsWith(this.suffix, StringComparison.Ordinal)) + { + await ProcessNextAsync(message, pipeline); + return; + } + this.Requests.Enqueue(message.Request.Uri.ToUri().AbsolutePath); + if (Interlocked.CompareExchange(ref this.active, 1, 0) != 0) + { + throw new RequestFailedException(409, "Duplicate concurrent container deletion.", "ConcurrentContainerOperationInProgress", null); + } + try + { + this.firstDelete.TrySetResult(true); + await this.released.Task; + if (this.FailureStatus.HasValue) + { + throw new RequestFailedException(this.FailureStatus.Value, "Injected container deletion failure.", this.FailureCode, null); + } + await ProcessNextAsync(message, pipeline); + } + finally + { + Interlocked.Exchange(ref this.active, 0); + } + } + } + } +} From 3ae35fb87b9d9c5c60d523a6b107f5813e68efbb Mon Sep 17 00:00:00 2001 From: wangbill Date: Wed, 23 Sep 2026 14:29:25 -0700 Subject: [PATCH 07/10] Clarify client task hub initialization helper name Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 87bb5e78-73ce-4a5c-9374-09c192f3a938 --- .../AzureStorageOrchestrationService.cs | 28 +++++++++---------- 1 file changed, 14 insertions(+), 14 deletions(-) diff --git a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs index 90cf42d36..4b1235862 100644 --- a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs +++ b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs @@ -310,7 +310,7 @@ EntityBackendQueries IEntityOrchestrationService.EntityBackendQueries => new EntityTrackingStoreQueries( this.messageManager, this.trackingStore, - this.EnsureTaskHubAsync, + this.EnsureClientTaskHubInitializedAsync, ((IEntityOrchestrationService)this).EntityBackendProperties, this.SendTaskOrchestrationMessageAsync); @@ -366,7 +366,7 @@ public async Task CreateIfNotExistsAsync() this.clientTaskHubInitializer.Reset(Task.FromResult(this.settings.PartitionCount)); } - async Task EnsureTaskHubAsync() + async Task EnsureClientTaskHubInitializedAsync() { await this.GetClientPartitionCountAsync(); } @@ -793,7 +793,7 @@ async Task LockNextTaskOrchestrationWorkItemAsync(boo { Guid traceActivityId = StartNewLogicalTraceScope(useExisting: true); - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); using (var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, this.shutdownSource.Token)) { @@ -1634,7 +1634,7 @@ public async Task LockNextTaskActivityWorkItem( TimeSpan receiveTimeout, CancellationToken cancellationToken) { - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); using (var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, this.shutdownSource.Token)) { @@ -1850,7 +1850,7 @@ public async Task CreateTaskOrchestrationAsync(TaskMessage creationMessage, Orch Utils.ConvertDateTimeInHistoryEventsToUTC(creationMessage.Event); - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); InstanceStatus existingInstance = await this.trackingStore.FetchInstanceStatusAsync( creationMessage.OrchestrationInstance.InstanceId); @@ -1918,7 +1918,7 @@ public Task SendTaskOrchestrationMessageBatchAsync(params TaskMessage[] messages /// The message to send. public async Task SendTaskOrchestrationMessageAsync(TaskMessage message) { - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); ControlQueue controlQueue = await this.GetControlQueueAsync(message.OrchestrationInstance.InstanceId); await this.SendTaskOrchestrationMessageInternalAsync(EmptySourceInstance, controlQueue, message); } @@ -1939,7 +1939,7 @@ internal Task SendTaskOrchestrationMessageInternalAsync( /// List of objects that represent the list of orchestrations. public async Task> GetOrchestrationStateAsync(string instanceId, bool allExecutions) { - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); return new OrchestrationState[] { await this.trackingStore.GetStateAsync(instanceId, allExecutions, fetchInput: true).FirstOrDefaultAsync(), @@ -1954,7 +1954,7 @@ await this.trackingStore.GetStateAsync(instanceId, allExecutions, fetchInput: tr /// The object that represents the orchestration. public async Task GetOrchestrationStateAsync(string instanceId, string executionId) { - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); return await this.trackingStore.GetStateAsync(instanceId, executionId, fetchInput: true); } @@ -1968,7 +1968,7 @@ public async Task GetOrchestrationStateAsync(string instance /// List of objects that represent the list of orchestrations. public async Task> GetOrchestrationStateAsync(string instanceId, bool allExecutions, bool fetchInput = true) { - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); return await this.trackingStore.GetStateAsync(instanceId, allExecutions, fetchInput).ToListAsync(); } @@ -1978,7 +1978,7 @@ public async Task> GetOrchestrationStateAsync(string i /// List of public async Task> GetOrchestrationStateAsync(CancellationToken cancellationToken = default(CancellationToken)) { - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); return await this.trackingStore.GetStateAsync(cancellationToken).ToListAsync(); } @@ -1992,7 +1992,7 @@ public async Task> GetOrchestrationStateAsync(string i /// List of public async Task> GetOrchestrationStateAsync(DateTime createdTimeFrom, DateTime? createdTimeTo, IEnumerable runtimeStatus, CancellationToken cancellationToken = default(CancellationToken)) { - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); return await this.trackingStore.GetStateAsync(createdTimeFrom, createdTimeTo, runtimeStatus, cancellationToken).ToListAsync(); } @@ -2008,7 +2008,7 @@ public async Task> GetOrchestrationStateAsync(string i /// List of public async Task GetOrchestrationStateAsync(DateTime createdTimeFrom, DateTime? createdTimeTo, IEnumerable runtimeStatus, int top, string continuationToken, CancellationToken cancellationToken = default(CancellationToken)) { - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); Page page = await this.trackingStore .GetStateAsync(createdTimeFrom, createdTimeTo, runtimeStatus, cancellationToken) .AsPages(continuationToken, top) @@ -2029,7 +2029,7 @@ public async Task> GetOrchestrationStateAsync(string i /// List of public async Task GetOrchestrationStateAsync(OrchestrationInstanceStatusQueryCondition condition, int top, string continuationToken, CancellationToken cancellationToken = default(CancellationToken)) { - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); Page page = await this.trackingStore .GetStateAsync(condition, cancellationToken) .AsPages(continuationToken, top) @@ -2063,7 +2063,7 @@ public async Task ForceTerminateTaskOrchestrationAsync(string instanceId, string /// The reason for rewinding. public async Task RewindTaskOrchestrationAsync(string instanceId, string reason) { - await this.EnsureTaskHubAsync(); + await this.EnsureClientTaskHubInitializedAsync(); List queueIds = await this.trackingStore.RewindHistoryAsync(instanceId).ToListAsync(); foreach (string id in queueIds) From cd6f68068a79a81af3dfd3ca09ef83986888df8b Mon Sep 17 00:00:00 2001 From: wangbill Date: Wed, 23 Sep 2026 14:37:43 -0700 Subject: [PATCH 08/10] Use neutral task hub initialization helper name Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 87bb5e78-73ce-4a5c-9374-09c192f3a938 --- .../AzureStorageOrchestrationService.cs | 28 +++++++++---------- 1 file changed, 14 insertions(+), 14 deletions(-) diff --git a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs index 4b1235862..72e3d7c0e 100644 --- a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs +++ b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs @@ -310,7 +310,7 @@ EntityBackendQueries IEntityOrchestrationService.EntityBackendQueries => new EntityTrackingStoreQueries( this.messageManager, this.trackingStore, - this.EnsureClientTaskHubInitializedAsync, + this.EnsureTaskHubInitializedAsync, ((IEntityOrchestrationService)this).EntityBackendProperties, this.SendTaskOrchestrationMessageAsync); @@ -366,7 +366,7 @@ public async Task CreateIfNotExistsAsync() this.clientTaskHubInitializer.Reset(Task.FromResult(this.settings.PartitionCount)); } - async Task EnsureClientTaskHubInitializedAsync() + async Task EnsureTaskHubInitializedAsync() { await this.GetClientPartitionCountAsync(); } @@ -793,7 +793,7 @@ async Task LockNextTaskOrchestrationWorkItemAsync(boo { Guid traceActivityId = StartNewLogicalTraceScope(useExisting: true); - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); using (var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, this.shutdownSource.Token)) { @@ -1634,7 +1634,7 @@ public async Task LockNextTaskActivityWorkItem( TimeSpan receiveTimeout, CancellationToken cancellationToken) { - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); using (var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, this.shutdownSource.Token)) { @@ -1850,7 +1850,7 @@ public async Task CreateTaskOrchestrationAsync(TaskMessage creationMessage, Orch Utils.ConvertDateTimeInHistoryEventsToUTC(creationMessage.Event); - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); InstanceStatus existingInstance = await this.trackingStore.FetchInstanceStatusAsync( creationMessage.OrchestrationInstance.InstanceId); @@ -1918,7 +1918,7 @@ public Task SendTaskOrchestrationMessageBatchAsync(params TaskMessage[] messages /// The message to send. public async Task SendTaskOrchestrationMessageAsync(TaskMessage message) { - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); ControlQueue controlQueue = await this.GetControlQueueAsync(message.OrchestrationInstance.InstanceId); await this.SendTaskOrchestrationMessageInternalAsync(EmptySourceInstance, controlQueue, message); } @@ -1939,7 +1939,7 @@ internal Task SendTaskOrchestrationMessageInternalAsync( /// List of objects that represent the list of orchestrations. public async Task> GetOrchestrationStateAsync(string instanceId, bool allExecutions) { - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); return new OrchestrationState[] { await this.trackingStore.GetStateAsync(instanceId, allExecutions, fetchInput: true).FirstOrDefaultAsync(), @@ -1954,7 +1954,7 @@ await this.trackingStore.GetStateAsync(instanceId, allExecutions, fetchInput: tr /// The object that represents the orchestration. public async Task GetOrchestrationStateAsync(string instanceId, string executionId) { - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); return await this.trackingStore.GetStateAsync(instanceId, executionId, fetchInput: true); } @@ -1968,7 +1968,7 @@ public async Task GetOrchestrationStateAsync(string instance /// List of objects that represent the list of orchestrations. public async Task> GetOrchestrationStateAsync(string instanceId, bool allExecutions, bool fetchInput = true) { - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); return await this.trackingStore.GetStateAsync(instanceId, allExecutions, fetchInput).ToListAsync(); } @@ -1978,7 +1978,7 @@ public async Task> GetOrchestrationStateAsync(string i /// List of public async Task> GetOrchestrationStateAsync(CancellationToken cancellationToken = default(CancellationToken)) { - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); return await this.trackingStore.GetStateAsync(cancellationToken).ToListAsync(); } @@ -1992,7 +1992,7 @@ public async Task> GetOrchestrationStateAsync(string i /// List of public async Task> GetOrchestrationStateAsync(DateTime createdTimeFrom, DateTime? createdTimeTo, IEnumerable runtimeStatus, CancellationToken cancellationToken = default(CancellationToken)) { - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); return await this.trackingStore.GetStateAsync(createdTimeFrom, createdTimeTo, runtimeStatus, cancellationToken).ToListAsync(); } @@ -2008,7 +2008,7 @@ public async Task> GetOrchestrationStateAsync(string i /// List of public async Task GetOrchestrationStateAsync(DateTime createdTimeFrom, DateTime? createdTimeTo, IEnumerable runtimeStatus, int top, string continuationToken, CancellationToken cancellationToken = default(CancellationToken)) { - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); Page page = await this.trackingStore .GetStateAsync(createdTimeFrom, createdTimeTo, runtimeStatus, cancellationToken) .AsPages(continuationToken, top) @@ -2029,7 +2029,7 @@ public async Task> GetOrchestrationStateAsync(string i /// List of public async Task GetOrchestrationStateAsync(OrchestrationInstanceStatusQueryCondition condition, int top, string continuationToken, CancellationToken cancellationToken = default(CancellationToken)) { - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); Page page = await this.trackingStore .GetStateAsync(condition, cancellationToken) .AsPages(continuationToken, top) @@ -2063,7 +2063,7 @@ public async Task ForceTerminateTaskOrchestrationAsync(string instanceId, string /// The reason for rewinding. public async Task RewindTaskOrchestrationAsync(string instanceId, string reason) { - await this.EnsureClientTaskHubInitializedAsync(); + await this.EnsureTaskHubInitializedAsync(); List queueIds = await this.trackingStore.RewindHistoryAsync(instanceId).ToListAsync(); foreach (string id in queueIds) From 5e9d1abaed348e5ade5d3ea69a1a4555334f0caa Mon Sep 17 00:00:00 2001 From: wangbill Date: Wed, 23 Sep 2026 16:14:01 -0700 Subject: [PATCH 09/10] Preserve published task hub topology during worker initialization Reject conflicting marked worker configurations before writes, retain matching metadata, and keep explicit legacy bootstrap and destructive recreation boundaries documented. Restore accurate client initialization comments. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 53306d22-2f23-4522-b14f-fc52789a4292 --- docs/providers/azure-storage.md | 4 + .../AzureStorageOrchestrationService.cs | 52 ++++++-- .../AzureStorageScaleTests.cs | 7 +- .../ClientPartitionTests.cs | 125 +++++++++++++++++- 4 files changed, 172 insertions(+), 16 deletions(-) diff --git a/docs/providers/azure-storage.md b/docs/providers/azure-storage.md index bbfa53937..c78dbf9b4 100644 --- a/docs/providers/azure-storage.md +++ b/docs/providers/azure-storage.md @@ -244,6 +244,10 @@ Both worker and client binaries must be upgraded to benefit from this behavior. The hub creation APIs remain administrative operations that apply their configured worker settings; they should not be used to discover another worker's configuration. Hub deletion removes the metadata along with the existing work-item queue, including when deletion is performed by an older provider. Unlike best-effort app-lease cleanup, work-item queue deletion failures are propagated. +When worker metadata is already published, explicit initialization and worker startup validate it before changing resources. A different configured partition count or invalid existing metadata is rejected; a matching marker, including its creation timestamp and unrelated queue metadata, is preserved. Destructive `CreateAsync()` still deletes and recreates a hub and can therefore establish a different configuration after deletion. + +This guard does not infer a legacy hub's complete topology from lease rows. Explicit initialization of an older unmarked hub still requires the operator to supply its correct existing partition configuration. All concurrently initializing hosts must agree on that configuration. Queue metadata read/validation/publication is not a compare-and-swap protocol; conflicting first-time publishers with different counts are not a supported deployment. The guard protects a count already present when read, not arbitrary conflicting first-publication races. + Client initialization retains references to every discovered control queue. Deleting the hub through that service therefore removes the full discovered topology, including queues the client has never sent to, even when the caller's configured partition count is smaller. ### Lease Management diff --git a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs index 72e3d7c0e..5950337ce 100644 --- a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs +++ b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs @@ -356,11 +356,30 @@ public async Task CreateIfNotExistsAsync() // Publish the worker's complete topology before leases are created individually. // Queue deletion removes this metadata, including when performed by older versions. Queue queue = GetWorkItemQueue(this.azureStorageClient); - await queue.CreateIfNotExistsAsync(); + if (!await queue.ExistsAsync()) + { + await queue.CreateIfNotExistsAsync(); + } + IDictionary metadata = await queue.GetMetadataAsync(); - metadata[WorkerTaskHubInfoMetadataKey] = - Utils.SerializeToJson(GetTaskHubInfo(this.settings.TaskHubName, this.settings.PartitionCount)); - await queue.SetMetadataAsync(metadata); + if (metadata.TryGetValue(WorkerTaskHubInfoMetadataKey, out string serializedHubInfo)) + { + TaskHubInfo hubInfo = this.DeserializeAndValidateTaskHubInfo(serializedHubInfo); + if (hubInfo.PartitionCount != this.settings.PartitionCount) + { + throw new InvalidOperationException( + $"Task hub '{this.settings.TaskHubName}' has published partition count {hubInfo.PartitionCount}, " + + $"but this worker is configured for {this.settings.PartitionCount}. " + + "Use the existing partition configuration, or initialize a new task hub to use a different partition count."); + } + } + else + { + metadata[WorkerTaskHubInfoMetadataKey] = + Utils.SerializeToJson(GetTaskHubInfo(this.settings.TaskHubName, this.settings.PartitionCount)); + await queue.SetMetadataAsync(metadata); + } + await this.EnsureTaskHubCreatedAsync(); // Worker admission must not be delayed by rediscovery after explicit initialization. this.clientTaskHubInitializer.Reset(Task.FromResult(this.settings.PartitionCount)); @@ -400,13 +419,7 @@ async Task InitializeClientTaskHubAsync() "using the target's partition configuration before retrying the client operation."); } - TaskHubInfo hubInfo = Utils.DeserializeFromJson(serializedHubInfo); - if (hubInfo == null || - !string.Equals(hubInfo.TaskHubName, this.settings.TaskHubName, StringComparison.OrdinalIgnoreCase) || - hubInfo.PartitionCount < 1 || hubInfo.PartitionCount > 16) - { - throw new InvalidOperationException($"Task hub '{this.settings.TaskHubName}' has invalid worker partition metadata."); - } + TaskHubInfo hubInfo = this.DeserializeAndValidateTaskHubInfo(serializedHubInfo); // Clients may finish creating resources after an interrupted worker initialization, // but must not create leases or overwrite the worker's partition configuration. @@ -426,6 +439,18 @@ async Task InitializeClientTaskHubAsync() return hubInfo.PartitionCount; } + TaskHubInfo DeserializeAndValidateTaskHubInfo(string serializedHubInfo) + { + TaskHubInfo hubInfo = Utils.DeserializeFromJson(serializedHubInfo); + if (hubInfo == null || + !string.Equals(hubInfo.TaskHubName, this.settings.TaskHubName, StringComparison.OrdinalIgnoreCase) || + hubInfo.PartitionCount < 1 || hubInfo.PartitionCount > 16) + { + throw new InvalidOperationException($"Task hub '{this.settings.TaskHubName}' has invalid worker partition metadata."); + } + return hubInfo; + } + async Task EnsureTaskHubCreatedAsync() { try @@ -1850,6 +1875,7 @@ public async Task CreateTaskOrchestrationAsync(TaskMessage creationMessage, Orch Utils.ConvertDateTimeInHistoryEventsToUTC(creationMessage.Event); + // Client operations require published worker metadata before using the task hub. await this.EnsureTaskHubInitializedAsync(); InstanceStatus existingInstance = await this.trackingStore.FetchInstanceStatusAsync( @@ -1918,6 +1944,7 @@ public Task SendTaskOrchestrationMessageBatchAsync(params TaskMessage[] messages /// The message to send. public async Task SendTaskOrchestrationMessageAsync(TaskMessage message) { + // Client operations require published worker metadata before using the task hub. await this.EnsureTaskHubInitializedAsync(); ControlQueue controlQueue = await this.GetControlQueueAsync(message.OrchestrationInstance.InstanceId); await this.SendTaskOrchestrationMessageInternalAsync(EmptySourceInstance, controlQueue, message); @@ -1939,6 +1966,7 @@ internal Task SendTaskOrchestrationMessageInternalAsync( /// List of objects that represent the list of orchestrations. public async Task> GetOrchestrationStateAsync(string instanceId, bool allExecutions) { + // Client operations require published worker metadata before using the task hub. await this.EnsureTaskHubInitializedAsync(); return new OrchestrationState[] { @@ -1954,6 +1982,7 @@ await this.trackingStore.GetStateAsync(instanceId, allExecutions, fetchInput: tr /// The object that represents the orchestration. public async Task GetOrchestrationStateAsync(string instanceId, string executionId) { + // Client operations require published worker metadata before using the task hub. await this.EnsureTaskHubInitializedAsync(); return await this.trackingStore.GetStateAsync(instanceId, executionId, fetchInput: true); } @@ -1968,6 +1997,7 @@ public async Task GetOrchestrationStateAsync(string instance /// List of objects that represent the list of orchestrations. public async Task> GetOrchestrationStateAsync(string instanceId, bool allExecutions, bool fetchInput = true) { + // Client operations require published worker metadata before using the task hub. await this.EnsureTaskHubInitializedAsync(); return await this.trackingStore.GetStateAsync(instanceId, allExecutions, fetchInput).ToListAsync(); } diff --git a/test/DurableTask.AzureStorage.Tests/AzureStorageScaleTests.cs b/test/DurableTask.AzureStorage.Tests/AzureStorageScaleTests.cs index 7d7ac8a78..42baf83a6 100644 --- a/test/DurableTask.AzureStorage.Tests/AzureStorageScaleTests.cs +++ b/test/DurableTask.AzureStorage.Tests/AzureStorageScaleTests.cs @@ -701,14 +701,14 @@ public async Task MonitorIdleTaskHubDisconnected() } [TestMethod] - public async Task UpdateTaskHubJsonWithNewPartitionCount() + public async Task UpdatePartitionCountAfterDeleteAndRecreate() { string connectionString = TestHelpers.GetTestStorageAccountConnectionString(); var settings = new AzureStorageOrchestrationServiceSettings { PartitionCount = 4, StorageAccountClientProvider = new StorageAccountClientProvider(connectionString), - TaskHubName = nameof(UpdateTaskHubJsonWithNewPartitionCount), + TaskHubName = nameof(UpdatePartitionCountAfterDeleteAndRecreate), UseAppLease = false, }; @@ -747,7 +747,8 @@ public async Task UpdateTaskHubJsonWithNewPartitionCount() Assert.IsNotNull(recommendation.Reason); } - // Change the default partition count, and start and stop the worker to try and update taskhub.json. + // Delete the empty hub before changing its count, then recreate it through worker startup. + await service.DeleteAsync(); settings.PartitionCount = 8; service = new AzureStorageOrchestrationService(settings); var worker = new TaskHubWorker(service); diff --git a/test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs b/test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs index 661cf1e12..e58651c6e 100644 --- a/test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs +++ b/test/DurableTask.AzureStorage.Tests/ClientPartitionTests.cs @@ -129,7 +129,7 @@ public async Task ClientRequiresExplicitCreationAndRecreationAfterDelete(int par [DataRow("Table")] [DataRow("Safe")] [DataRow("Legacy")] - public async Task ExplicitCreationRetainsConfiguredWorkerPartitions(string manager) + public async Task ExplicitRecreationCanChangeWorkerPartitions(string manager) { string hub = NewHub(); using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, manager)); @@ -139,7 +139,8 @@ public async Task ExplicitCreationRetainsConfiguredWorkerPartitions(string manag await target.CreateIfNotExistsAsync(); var client = new TaskHubClient(service); await client.GetOrchestrationStateAsync("missing"); - await service.CreateIfNotExistsAsync(); + await Assert.ThrowsExceptionAsync(() => service.CreateIfNotExistsAsync()); + await service.CreateAsync(); Assert.AreEqual(16, (await this.QueueNamesAsync(hub)).Length); Assert.AreEqual(16, await this.PartitionCountAsync(hub, manager)); await client.CreateOrchestrationInstanceAsync("Probe", string.Empty, "partition-probe", "hello"); @@ -152,6 +153,126 @@ public async Task ExplicitCreationRetainsConfiguredWorkerPartitions(string manag } } + [DataTestMethod] + [DataRow(4, 16, "Table")] + [DataRow(16, 4, "Table")] + [DataRow(4, 16, "Safe")] + [DataRow(16, 4, "Safe")] + [DataRow(4, 16, "Legacy")] + [DataRow(16, 4, "Legacy")] + public async Task MarkedHubRejectsWorkerPartitionMismatchWithoutPublishing( + int targetPartitions, int workerPartitions, string manager) + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, targetPartitions, manager)); + var writes = new RecordingWritePolicy(); + using var worker = new AzureStorageOrchestrationService(this.Settings(hub, workerPartitions, manager, writes)); + await target.CreateIfNotExistsAsync(); + string metadata = await this.ReadWorkerMetadataAsync(hub); + string[] inventory = await this.ResourceInventoryAsync(hub); + try + { + await Assert.ThrowsExceptionAsync(() => worker.CreateIfNotExistsAsync()); + await Assert.ThrowsExceptionAsync(() => worker.StartAsync()); + Assert.AreEqual(metadata, await this.ReadWorkerMetadataAsync(hub)); + Assert.AreEqual(0, writes.Count); + CollectionAssert.AreEqual(inventory, await this.ResourceInventoryAsync(hub)); + Assert.AreEqual(targetPartitions, await this.PartitionCountAsync(hub, manager)); + } + finally + { + Console.WriteLine($"Target={targetPartitions}, configured worker={workerPartitions}, published after call={JObject.Parse(await this.ReadWorkerMetadataAsync(hub)).Value("PartitionCount")}, partition-store count after call={await this.PartitionCountAsync(hub, manager)}."); + await this.CleanupAsync(hub, manager); + } + } + + [DataTestMethod] + [DataRow("Table")] + [DataRow("Safe")] + [DataRow("Legacy")] + public async Task MatchingWorkerPreservesPublishedMetadata(string manager) + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, manager)); + using var worker = new AzureStorageOrchestrationService(this.Settings(hub, 4, manager)); + try + { + await target.CreateIfNotExistsAsync(); + var queue = new QueueClient(this.connection, hub.ToLowerInvariant() + "-workitems"); + var metadata = (await queue.GetPropertiesAsync()).Value.Metadata; + metadata["application"] = "preserved"; + await queue.SetMetadataAsync(metadata); + string published = await this.ReadWorkerMetadataAsync(hub); + await worker.CreateIfNotExistsAsync(); + await worker.StartAsync(); + await worker.StopAsync(isForced: true); + Assert.AreEqual(published, await this.ReadWorkerMetadataAsync(hub)); + Assert.AreEqual("preserved", (await queue.GetPropertiesAsync()).Value.Metadata["application"]); + Assert.AreEqual(4, await this.PartitionCountAsync(hub, manager)); + } + finally + { + await this.CleanupAsync(hub, manager); + } + } + + [DataTestMethod] + [DataRow("null", false)] + [DataRow("{broken", true)] + [DataRow("{\"TaskHubName\":\"wrong\",\"PartitionCount\":4}", false)] + public async Task WorkerRejectsInvalidPublishedMetadataBeforeWrites(string metadata, bool invalidJson) + { + string hub = NewHub(); + using var target = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table")); + var writes = new RecordingWritePolicy(); + using var worker = new AzureStorageOrchestrationService(this.Settings(hub, 4, "Table", writes)); + try + { + await target.CreateIfNotExistsAsync(); + await this.WriteWorkerMetadataAsync(hub, metadata); + if (invalidJson) + { + await Assert.ThrowsExceptionAsync(() => worker.CreateIfNotExistsAsync()); + } + else + { + await Assert.ThrowsExceptionAsync(() => worker.CreateIfNotExistsAsync()); + } + Assert.AreEqual(metadata, await this.ReadWorkerMetadataAsync(hub)); + Assert.AreEqual(0, writes.Count); + await this.WriteWorkerMetadataAsync(hub, $"{{\"TaskHubName\":\"{hub}\",\"PartitionCount\":4}}"); + await worker.CreateIfNotExistsAsync(); + } + finally + { + await this.CleanupAsync(hub, "Table"); + } + } + + [DataTestMethod] + [DataRow("Table")] + [DataRow("Safe")] + [DataRow("Legacy")] + public async Task ExplicitLegacyBootstrapRequiresOperatorMatchedConfiguration(string manager) + { + string hub = NewHub(); + using var original = new AzureStorageOrchestrationService(this.Settings(hub, 4, manager)); + using var upgraded = new AzureStorageOrchestrationService(this.Settings(hub, 4, manager)); + try + { + await original.CreateIfNotExistsAsync(); + await this.WriteWorkerMetadataAsync(hub, null); + await upgraded.CreateIfNotExistsAsync(); + Assert.AreEqual(4, JObject.Parse(await this.ReadWorkerMetadataAsync(hub)).Value("PartitionCount")); + Assert.AreEqual(4, (await this.QueueNamesAsync(hub)).Length); + Assert.AreEqual(4, await this.PartitionCountAsync(hub, manager)); + } + finally + { + await this.CleanupAsync(hub, manager); + } + } + [DataTestMethod] [DataRow("Table")] [DataRow("Safe")] From b0548b21670303727450fe44748b2211eb70721c Mon Sep 17 00:00:00 2001 From: wangbill Date: Wed, 23 Sep 2026 17:19:19 -0700 Subject: [PATCH 10/10] Require initialized metadata for history and purge operations Gate remaining client storage entry points and include shared initialization waits in purge deadlines without cancelling other callers or resuming timed-out purges. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 53306d22-2f23-4522-b14f-fc52789a4292 --- docs/providers/azure-storage.md | 2 + .../AzureStorageOrchestrationService.cs | 51 ++- .../ClientStorageInitializationTests.cs | 322 ++++++++++++++++++ 3 files changed, 366 insertions(+), 9 deletions(-) create mode 100644 test/DurableTask.AzureStorage.Tests/ClientStorageInitializationTests.cs diff --git a/docs/providers/azure-storage.md b/docs/providers/azure-storage.md index c78dbf9b4..755f4bdaa 100644 --- a/docs/providers/azure-storage.md +++ b/docs/providers/azure-storage.md @@ -240,6 +240,8 @@ Deploy the upgraded target worker first and let it publish metadata, or explicit If the work-item queue or its `durabletask_taskhub` metadata entry is absent, the client throws an actionable exception without creating queues, tables, blobs, or leases or enqueuing the message. This includes calls racing with worker startup before metadata publication. Failed initialization is not cached: retrying the same client after publication succeeds. There is no fallback to the caller's count or a partially populated partition table. Malformed metadata and storage access failures remain explicit errors. +The same precondition applies to history reads, purge operations, and large-message downloads. A timed purge includes initialization in its deadline; if the deadline expires first, it returns zero deleted instances with `IsComplete = false` and does not later start purging. Shared initialization is not cancelled for other callers and may continue provisioning resources once valid metadata is available. This is a no-late-purge guarantee, not cancellation of all shared initialization work. + Both worker and client binaries must be upgraded to benefit from this behavior. Older client binaries ignore the metadata and are not fixed by upgrading the worker alone. Successfully initialized clients cache topology; recreate them after deleting and recreating a hub externally. Explicit worker initialization preserves its already-completed initialization path for local clients and activity admission. The hub creation APIs remain administrative operations that apply their configured worker settings; they should not be used to discover another worker's configuration. Hub deletion removes the metadata along with the existing work-item queue, including when deletion is performed by an older provider. Unlike best-effort app-lease cleanup, work-item queue deletion failures are propagated. diff --git a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs index 5950337ce..352999b10 100644 --- a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs +++ b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs @@ -2131,6 +2131,7 @@ public Task ForceChangeAppLeaseAsync() /// String with formatted JSON array representing the execution history. public async Task GetOrchestrationHistoryAsync(string instanceId, string executionId) { + await this.EnsureTaskHubInitializedAsync(); OrchestrationHistory history = await this.trackingStore.GetHistoryEventsAsync( instanceId, executionId, @@ -2143,9 +2144,10 @@ public async Task GetOrchestrationHistoryAsync(string instanceId, string /// /// Instance ID of the orchestration. /// Class containing number of storage requests sent, along with instances and rows deleted/purged - public Task PurgeInstanceHistoryAsync(string instanceId) + public async Task PurgeInstanceHistoryAsync(string instanceId) { - return this.trackingStore.PurgeInstanceHistoryAsync(instanceId); + await this.EnsureTaskHubInitializedAsync(); + return await this.trackingStore.PurgeInstanceHistoryAsync(instanceId); } /// @@ -2155,9 +2157,10 @@ public Task PurgeInstanceHistoryAsync(string instanceId) /// CreatedTime of orchestrations. Purges history less than this value. /// RuntimeStatus of orchestrations. You can specify several statuses. /// Class containing number of storage requests sent, along with instances and rows deleted/purged - public Task PurgeInstanceHistoryAsync(DateTime createdTimeFrom, DateTime? createdTimeTo, IEnumerable runtimeStatus) + public async Task PurgeInstanceHistoryAsync(DateTime createdTimeFrom, DateTime? createdTimeTo, IEnumerable runtimeStatus) { - return this.trackingStore.PurgeInstanceHistoryAsync(createdTimeFrom, createdTimeTo, runtimeStatus); + await this.EnsureTaskHubInitializedAsync(); + return await this.trackingStore.PurgeInstanceHistoryAsync(createdTimeFrom, createdTimeTo, runtimeStatus); } /// @@ -2176,6 +2179,34 @@ async Task IOrchestrationServicePurgeClient.PurgeInstanceStateAsync // Convert the timeout into a CancellationToken so that the tracking store // only needs to observe a single cancellation mechanism. using var timeoutCts = new CancellationTokenSource(purgeInstanceFilter.Timeout.Value); + if (purgeInstanceFilter.Timeout.Value == TimeSpan.Zero || timeoutCts.IsCancellationRequested) + { + return new PurgeResult(0, false); + } + Task initialization = this.EnsureTaskHubInitializedAsync(); + if (!initialization.IsCompleted) + { + var cancelled = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + using (timeoutCts.Token.Register(() => cancelled.TrySetResult(true))) + { + if (await Task.WhenAny(initialization, cancelled.Task) != initialization) + { + // Initialization is shared. Keep it available to other callers, but never + // resume this purge after its deadline. Its failures are logged by the initializer. + _ = initialization.ContinueWith( + task => { _ = task.Exception; }, + CancellationToken.None, + TaskContinuationOptions.OnlyOnFaulted | TaskContinuationOptions.ExecuteSynchronously, + TaskScheduler.Default); + return new PurgeResult(0, false); + } + } + } + await initialization; + if (timeoutCts.IsCancellationRequested) + { + return new PurgeResult(0, false); + } storagePurgeHistoryResult = await this.trackingStore.PurgeInstanceHistoryAsync( purgeInstanceFilter.CreatedTimeFrom, purgeInstanceFilter.CreatedTimeTo, @@ -2186,7 +2217,7 @@ async Task IOrchestrationServicePurgeClient.PurgeInstanceStateAsync { // No timeout: use the original code path (no CancellationToken) to preserve // backward-compatible behavior where IsComplete is null. - storagePurgeHistoryResult = await this.trackingStore.PurgeInstanceHistoryAsync( + storagePurgeHistoryResult = await this.PurgeInstanceHistoryAsync( purgeInstanceFilter.CreatedTimeFrom, purgeInstanceFilter.CreatedTimeTo, purgeInstanceFilter.RuntimeStatus); @@ -2261,9 +2292,10 @@ async Task IOrchestrationServicePurgeClient.PurgeInstanceStateAsync /// /// Threshold date time in UTC /// What to compare the threshold date time against - public Task PurgeOrchestrationHistoryAsync(DateTime thresholdDateTimeUtc, OrchestrationStateTimeRangeFilterType timeRangeFilterType) + public async Task PurgeOrchestrationHistoryAsync(DateTime thresholdDateTimeUtc, OrchestrationStateTimeRangeFilterType timeRangeFilterType) { - return this.trackingStore.PurgeHistoryAsync(thresholdDateTimeUtc, timeRangeFilterType); + await this.EnsureTaskHubInitializedAsync(); + await this.trackingStore.PurgeHistoryAsync(thresholdDateTimeUtc, timeRangeFilterType); } /// @@ -2271,9 +2303,10 @@ public Task PurgeOrchestrationHistoryAsync(DateTime thresholdDateTimeUtc, Orches /// such as the input or output status fields, and for which the blob URI was stored instead. /// /// The URI of the blob. - public Task DownloadBlobAsync(string blobUri) + public async Task DownloadBlobAsync(string blobUri) { - return this.messageManager.DownloadAndDecompressAsBytesAsync(new Uri(blobUri)); + await this.EnsureTaskHubInitializedAsync(); + return await this.messageManager.DownloadAndDecompressAsBytesAsync(new Uri(blobUri)); } #endregion diff --git a/test/DurableTask.AzureStorage.Tests/ClientStorageInitializationTests.cs b/test/DurableTask.AzureStorage.Tests/ClientStorageInitializationTests.cs new file mode 100644 index 000000000..abc8c501d --- /dev/null +++ b/test/DurableTask.AzureStorage.Tests/ClientStorageInitializationTests.cs @@ -0,0 +1,322 @@ +// Copyright Microsoft Corporation +// Licensed under the Apache License, Version 2.0. + +namespace DurableTask.AzureStorage.Tests +{ + using System; + using System.Collections.Concurrent; + using System.IO; + using System.IO.Compression; + using System.Text; + using System.Threading; + using System.Threading.Tasks; + using Azure.Core; + using Azure.Core.Pipeline; + using Azure.Data.Tables; + using Azure.Storage.Blobs; + using Azure.Storage.Queues; + using DurableTask.AzureStorage.Storage; + using DurableTask.Core; + using Microsoft.VisualStudio.TestTools.UnitTesting; + using Newtonsoft.Json; + + [TestClass] + public class ClientStorageInitializationTests + { + readonly string connection = TestHelpers.GetTestStorageAccountConnectionString(); + + [DataTestMethod] + [DataRow("history", false)] + [DataRow("history", true)] + [DataRow("instance", false)] + [DataRow("instance", true)] + [DataRow("range", false)] + [DataRow("range", true)] + [DataRow("interface-instance", false)] + [DataRow("interface-instance", true)] + [DataRow("interface-range", false)] + [DataRow("interface-range", true)] + [DataRow("timeout", false)] + [DataRow("timeout", true)] + [DataRow("legacy", false)] + [DataRow("legacy", true)] + [DataRow("download", false)] + [DataRow("download", true)] + public async Task ClientStorageOperationRequiresMetadataBeforeAccess(string operation, bool existing) + { + string hub = "ClientStorage" + Guid.NewGuid().ToString("N"); + using var target = this.CreateService(hub); + var traffic = new StorageAccessPolicy(hub); + using var client = this.CreateService(hub, traffic); + var instance = new OrchestrationInstance { InstanceId = "probe", ExecutionId = "missing" }; + string blobUri = new BlobClient(this.connection, hub.ToLowerInvariant() + "-largemessages", "probe/input").Uri.AbsoluteUri; + string history = null; + try + { + if (existing) + { + await target.CreateIfNotExistsAsync(); + instance = await new TaskHubClient(target).CreateOrchestrationInstanceAsync("Probe", "", "probe", "hello"); + history = await target.GetOrchestrationHistoryAsync(instance.InstanceId, instance.ExecutionId); + await this.RemoveMetadataAsync(hub); + } + + InvalidOperationException error = await Assert.ThrowsExceptionAsync( + () => InvokeAsync(client, operation, instance, blobUri)); + StringAssert.Contains(error.Message, "CreateIfNotExistsAsync"); + Assert.AreEqual(0, traffic.NonMetadataRequests.Count, "Rejected operation accessed hub storage beyond metadata."); + if (existing) + { + Assert.AreEqual(history, await target.GetOrchestrationHistoryAsync(instance.InstanceId, instance.ExecutionId)); + Assert.AreEqual(OrchestrationStatus.Pending, (await new TaskHubClient(target).GetOrchestrationStateAsync("probe")).OrchestrationStatus); + } + + await target.CreateIfNotExistsAsync(); + if (!existing) + { + instance = await new TaskHubClient(target).CreateOrchestrationInstanceAsync("Probe", "", "probe", "hello"); + } + if (operation == "download") + { + await this.UploadPayloadAsync(hub); + } + if (operation == "legacy") + { + await Assert.ThrowsExceptionAsync(() => InvokeAsync(client, operation, instance, blobUri)); + } + else + { + await InvokeAsync(client, operation, instance, blobUri); + } + } + finally + { + await target.DeleteAsync(); + } + } + + [TestMethod] + public async Task PurgeDeadlineIncludesInitializationWithoutCancellingOtherCallers() + { + string hub = "PurgeDeadline" + Guid.NewGuid().ToString("N"); + using var target = this.CreateService(hub); + var traffic = new StorageAccessPolicy(hub) { PauseMetadata = true }; + using var client = this.CreateService(hub, traffic); + try + { + await target.CreateIfNotExistsAsync(); + await new TaskHubClient(target).CreateOrchestrationInstanceAsync("Probe", "", "probe", "hello"); + Task purge = ((IOrchestrationServicePurgeClient)client).PurgeInstanceStateAsync( + new PurgeInstanceFilter(DateTime.UtcNow.AddDays(-1), null, null) { Timeout = TimeSpan.FromMilliseconds(100) }); + Assert.AreSame(traffic.MetadataReached, await Task.WhenAny(traffic.MetadataReached, Task.Delay(TimeSpan.FromSeconds(10)))); + Task otherCaller = new TaskHubClient(client).GetOrchestrationStateAsync("probe"); + Assert.AreSame(purge, await Task.WhenAny(purge, Task.Delay(TimeSpan.FromSeconds(5)))); + PurgeResult result = await purge; + Assert.AreEqual(0, result.DeletedInstanceCount); + Assert.AreEqual(false, result.IsComplete); + Assert.AreEqual(0, traffic.NonMetadataRequests.Count); + Assert.IsFalse(otherCaller.IsCompleted, "The shared initialization should still be waiting, not cancelled by the purge."); + + traffic.Release(); + Assert.IsNotNull(await otherCaller); + Assert.IsNotNull(await new TaskHubClient(target).GetOrchestrationStateAsync("probe"), + "A timed-out purge must not resume deletion after initialization completes."); + PurgeResult retry = await ((IOrchestrationServicePurgeClient)client).PurgeInstanceStateAsync( + new PurgeInstanceFilter(DateTime.UtcNow.AddDays(-1), null, null) { Timeout = TimeSpan.FromSeconds(10) }); + Assert.AreEqual(1, retry.DeletedInstanceCount); + Assert.AreEqual(true, retry.IsComplete); + } + finally + { + traffic.Release(); + await target.DeleteAsync(); + } + } + + [TestMethod] + public async Task ExpiredPurgeNeverStartsTrackingDeletion() + { + string hub = "ExpiredPurge" + Guid.NewGuid().ToString("N"); + using var target = this.CreateService(hub); + var traffic = new StorageAccessPolicy(hub) { PauseMetadata = true }; + using var client = this.CreateService(hub, traffic); + try + { + await target.CreateIfNotExistsAsync(); + await new TaskHubClient(target).CreateOrchestrationInstanceAsync("Probe", "", "probe", "hello"); + PurgeResult result = await ((IOrchestrationServicePurgeClient)client).PurgeInstanceStateAsync( + new PurgeInstanceFilter(DateTime.UtcNow.AddDays(-1), null, null) { Timeout = TimeSpan.Zero }); + Assert.AreEqual(0, result.DeletedInstanceCount); + Assert.AreEqual(false, result.IsComplete); + Assert.AreEqual(0, traffic.NonMetadataRequests.Count); + Assert.AreEqual(0, traffic.MetadataReads, "An already-expired purge must not start initialization."); + traffic.Release(); + Assert.IsNotNull(await new TaskHubClient(client).GetOrchestrationStateAsync("probe")); + } + finally + { + traffic.Release(); + await target.DeleteAsync(); + } + } + + [DataTestMethod] + [DataRow(false)] + [DataRow(true)] + public async Task TimedPurgePreservesInitializationErrorsBeforeDeadline(bool authorizationFailure) + { + string hub = "PurgeError" + Guid.NewGuid().ToString("N"); + using var target = this.CreateService(hub); + var traffic = new StorageAccessPolicy(hub) { FailMetadata = authorizationFailure }; + using var client = this.CreateService(hub, traffic); + try + { + await target.CreateIfNotExistsAsync(); + if (!authorizationFailure) + { + var queue = new QueueClient(this.connection, hub.ToLowerInvariant() + "-workitems"); + var metadata = (await queue.GetPropertiesAsync()).Value.Metadata; + metadata["durabletask_taskhub"] = "{broken"; + await queue.SetMetadataAsync(metadata); + } + Func purge = () => ((IOrchestrationServicePurgeClient)client).PurgeInstanceStateAsync( + new PurgeInstanceFilter(DateTime.UtcNow.AddDays(-1), null, null) { Timeout = TimeSpan.FromSeconds(10) }); + if (authorizationFailure) + { + var error = await Assert.ThrowsExceptionAsync(purge); + Assert.AreEqual(403, error.HttpStatusCode); + } + else + { + await Assert.ThrowsExceptionAsync(purge); + } + Assert.AreEqual(0, traffic.NonMetadataRequests.Count); + } + finally + { + await target.DeleteAsync(); + } + } + + static async Task InvokeAsync(AzureStorageOrchestrationService service, string operation, OrchestrationInstance instance, string blobUri) + { + switch (operation) + { + case "history": + await service.GetOrchestrationHistoryAsync(instance.InstanceId, instance.ExecutionId); + break; + case "instance": + await service.PurgeInstanceHistoryAsync(instance.InstanceId); + break; + case "range": + await service.PurgeInstanceHistoryAsync(DateTime.UtcNow.AddDays(-1), null, null); + break; + case "interface-instance": + await ((IOrchestrationServicePurgeClient)service).PurgeInstanceStateAsync(instance.InstanceId); + break; + case "interface-range": + case "timeout": + var filter = new PurgeInstanceFilter(DateTime.UtcNow.AddDays(-1), null, null); + if (operation == "timeout") + { + filter.Timeout = TimeSpan.FromSeconds(10); + } + await ((IOrchestrationServicePurgeClient)service).PurgeInstanceStateAsync(filter); + break; + case "legacy": + await service.PurgeOrchestrationHistoryAsync(DateTime.UtcNow, OrchestrationStateTimeRangeFilterType.OrchestrationCompletedTimeFilter); + break; + case "download": + Assert.AreEqual("payload", await service.DownloadBlobAsync(blobUri)); + break; + default: + throw new ArgumentException(nameof(operation)); + } + } + + AzureStorageOrchestrationService CreateService(string hub, StorageAccessPolicy policy = null) + { + var blob = new BlobClientOptions(); + var queue = new QueueClientOptions(); + var table = new TableClientOptions(); + if (policy != null) + { + blob.AddPolicy(policy, HttpPipelinePosition.PerCall); + queue.AddPolicy(policy, HttpPipelinePosition.PerCall); + table.AddPolicy(policy, HttpPipelinePosition.PerCall); + } + return new AzureStorageOrchestrationService(new AzureStorageOrchestrationServiceSettings + { + TaskHubName = hub, + PartitionCount = 4, + UseTablePartitionManagement = true, + StorageAccountClientProvider = new StorageAccountClientProvider( + StorageServiceClientProvider.ForBlob(this.connection, blob), + StorageServiceClientProvider.ForQueue(this.connection, queue), + StorageServiceClientProvider.ForTable(this.connection, table)), + }); + } + + async Task RemoveMetadataAsync(string hub) + { + var queue = new QueueClient(this.connection, hub.ToLowerInvariant() + "-workitems"); + var metadata = (await queue.GetPropertiesAsync()).Value.Metadata; + metadata.Remove("durabletask_taskhub"); + await queue.SetMetadataAsync(metadata); + } + + async Task UploadPayloadAsync(string hub) + { + var container = new BlobContainerClient(this.connection, hub.ToLowerInvariant() + "-largemessages"); + await container.CreateIfNotExistsAsync(); + using var content = new MemoryStream(); + using (var gzip = new GZipStream(content, CompressionLevel.Optimal, leaveOpen: true)) + { + byte[] bytes = Encoding.UTF8.GetBytes("payload"); + await gzip.WriteAsync(bytes, 0, bytes.Length); + } + content.Position = 0; + await container.GetBlobClient("probe/input").UploadAsync(content); + } + + class StorageAccessPolicy : HttpPipelinePolicy + { + readonly string metadataPath; + readonly TaskCompletionSource reached = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + readonly TaskCompletionSource released = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + int metadataReads; + + public StorageAccessPolicy(string hub) => this.metadataPath = "/" + hub.ToLowerInvariant() + "-workitems"; + public ConcurrentQueue NonMetadataRequests { get; } = new ConcurrentQueue(); + public bool PauseMetadata { get; set; } + public bool FailMetadata { get; set; } + public int MetadataReads => this.metadataReads; + public Task MetadataReached => this.reached.Task; + public void Release() => this.released.TrySetResult(true); + public override void Process(HttpMessage message, ReadOnlyMemory pipeline) => throw new NotSupportedException(); + public override async ValueTask ProcessAsync(HttpMessage message, ReadOnlyMemory pipeline) + { + bool metadataRead = (message.Request.Method == RequestMethod.Get || message.Request.Method == RequestMethod.Head) + && message.Request.Uri.ToUri().AbsolutePath.EndsWith(this.metadataPath, StringComparison.Ordinal); + if (metadataRead) + { + Interlocked.Increment(ref this.metadataReads); + } + if (!metadataRead) + { + this.NonMetadataRequests.Enqueue(message.Request.Method + " " + message.Request.Uri.ToUri().AbsolutePath); + } + else if (this.FailMetadata) + { + throw new Azure.RequestFailedException(403, "Injected metadata authorization failure."); + } + else if (this.PauseMetadata) + { + this.reached.TrySetResult(true); + await this.released.Task; + } + await ProcessNextAsync(message, pipeline); + } + } + } +}