From db3a4283b30392440aa711a55b52b978fdf899fd Mon Sep 17 00:00:00 2001 From: Bernd Verst Date: Mon, 27 Jul 2026 13:36:25 -0700 Subject: [PATCH 1/3] Remove Service Bus inline message Task.Run Return completed stream tasks for inline message bodies while retaining asynchronous blob-store loading. Add cross-target coverage for round trips, synchronous completion, and stream ownership. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> Copilot-Session: 1347830a-17c5-4248-8e41-fee30d3066b4 --- .../ServiceBusUtilsTests.cs | 133 ++++++++++++++++++ .../Common/ServiceBusUtils.cs | 4 +- 2 files changed, 135 insertions(+), 2 deletions(-) create mode 100644 Test/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs diff --git a/Test/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs b/Test/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs new file mode 100644 index 000000000..500110338 --- /dev/null +++ b/Test/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs @@ -0,0 +1,133 @@ +// ---------------------------------------------------------------------------------- +// Copyright Microsoft Corporation +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// http://www.apache.org/licenses/LICENSE-2.0 +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. +// ---------------------------------------------------------------------------------- + +namespace DurableTask.ServiceBus.Tests +{ + using System; + using System.IO; + using System.Threading.Tasks; + using DurableTask.Core; + using DurableTask.Core.Common; + using DurableTask.Core.Settings; + using DurableTask.Core.Tracking; + using DurableTask.ServiceBus.Common.Abstraction; + using Microsoft.VisualStudio.TestTools.UnitTesting; + + [TestClass] + public class ServiceBusUtilsTests + { + const string Payload = "inline message payload"; + + [TestMethod] + public Task UncompressedInlineMessageDeserializationCompletesSynchronously() + { + return AssertInlineMessageRoundTrips(CompressionStyle.Never, expectSynchronousCompletion: true); + } + + [TestMethod] + public Task CompressedInlineMessageDeserializationRoundTrips() + { + return AssertInlineMessageRoundTrips(CompressionStyle.Always, expectSynchronousCompletion: false); + } + + [TestMethod] + public async Task ExternalMessageDeserializationAwaitsBlobStore() + { + const string blobKey = "message-blob"; + var message = new Message(); + message.UserProperties[ServiceBusConstants.MessageBlobKey] = blobKey; + message.UserProperties[FrameworkConstants.CompressionTypePropertyName] = + FrameworkConstants.CompressionTypeNonePropertyValue; + + var stream = new MemoryStream(); + Utils.WriteObjectToStream(stream, Payload); + stream.Position = stream.Length; + + var blobStore = new DeferredBlobStore(); + Task deserializationTask = + ServiceBusUtils.GetObjectFromBrokeredMessageAsync(message, blobStore); + + Assert.AreEqual(blobKey, blobStore.LoadedBlobKey); + Assert.IsFalse(deserializationTask.IsCompleted); + + blobStore.CompleteLoad(stream); + + Assert.AreEqual(Payload, await deserializationTask); + Assert.ThrowsException(() => stream.ReadByte()); + } + + static async Task AssertInlineMessageRoundTrips( + CompressionStyle compressionStyle, + bool expectSynchronousCompletion) + { + Message message = await ServiceBusUtils.GetBrokeredMessageFromObjectAsync( + Payload, + new CompressionSettings { Style = compressionStyle }); + + Task deserializationTask = + ServiceBusUtils.GetObjectFromBrokeredMessageAsync(message, null); + + if (expectSynchronousCompletion) + { + Assert.IsTrue( + deserializationTask.IsCompleted, + "Inline message deserialization should not require a thread-pool continuation."); + } + + Assert.AreEqual(Payload, await deserializationTask); + } + + sealed class DeferredBlobStore : IOrchestrationServiceBlobStore + { + readonly TaskCompletionSource loadCompletionSource = new TaskCompletionSource(); + + public string LoadedBlobKey { get; private set; } + + public string BuildMessageBlobKey(OrchestrationInstance orchestrationInstance, DateTime messageFireTime) + { + throw new NotSupportedException(); + } + + public string BuildSessionBlobKey(string sessionId) + { + throw new NotSupportedException(); + } + + public Task SaveStreamAsync(string blobKey, Stream stream) + { + throw new NotSupportedException(); + } + + public Task LoadStreamAsync(string blobKey) + { + this.LoadedBlobKey = blobKey; + return this.loadCompletionSource.Task; + } + + public Task DeleteStoreAsync() + { + throw new NotSupportedException(); + } + + public Task PurgeExpiredBlobsAsync(DateTime thresholdDateTimeUtc) + { + throw new NotSupportedException(); + } + + public void CompleteLoad(Stream stream) + { + this.loadCompletionSource.SetResult(stream); + } + } + } +} diff --git a/src/DurableTask.ServiceBus/Common/ServiceBusUtils.cs b/src/DurableTask.ServiceBus/Common/ServiceBusUtils.cs index cc9472a1a..861b16852 100644 --- a/src/DurableTask.ServiceBus/Common/ServiceBusUtils.cs +++ b/src/DurableTask.ServiceBus/Common/ServiceBusUtils.cs @@ -301,9 +301,9 @@ static Task LoadMessageStreamAsync(Message message, IOrchestrationServic // load the stream from the message directly if the blob key property is not set, // i.e., it is not stored externally #if NETSTANDARD2_0 - return Task.Run(() => new System.IO.MemoryStream(message.Body) as Stream); + return Task.FromResult(new MemoryStream(message.Body)); #else - return Task.Run(() => message.GetBody()); + return Task.FromResult(message.GetBody()); #endif } From ea16dafb0e103729367c57057ce605b072db4160 Mon Sep 17 00:00:00 2001 From: Bernd Verst Date: Tue, 29 Sep 2026 17:36:12 -0700 Subject: [PATCH 2/3] Address Service Bus deserialization test review feedback Move ServiceBusUtilsTests into the lowercase test project directory and ensure the external stream is disposed even if the test fails. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../ServiceBusUtilsTests.cs | 24 ++++++++++--------- 1 file changed, 13 insertions(+), 11 deletions(-) rename {Test => test}/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs (85%) diff --git a/Test/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs b/test/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs similarity index 85% rename from Test/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs rename to test/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs index 500110338..5332364db 100644 --- a/Test/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs +++ b/test/DurableTask.ServiceBus.Tests/ServiceBusUtilsTests.cs @@ -49,21 +49,23 @@ public async Task ExternalMessageDeserializationAwaitsBlobStore() message.UserProperties[FrameworkConstants.CompressionTypePropertyName] = FrameworkConstants.CompressionTypeNonePropertyValue; - var stream = new MemoryStream(); - Utils.WriteObjectToStream(stream, Payload); - stream.Position = stream.Length; + using (var stream = new MemoryStream()) + { + Utils.WriteObjectToStream(stream, Payload); + stream.Position = stream.Length; - var blobStore = new DeferredBlobStore(); - Task deserializationTask = - ServiceBusUtils.GetObjectFromBrokeredMessageAsync(message, blobStore); + var blobStore = new DeferredBlobStore(); + Task deserializationTask = + ServiceBusUtils.GetObjectFromBrokeredMessageAsync(message, blobStore); - Assert.AreEqual(blobKey, blobStore.LoadedBlobKey); - Assert.IsFalse(deserializationTask.IsCompleted); + Assert.AreEqual(blobKey, blobStore.LoadedBlobKey); + Assert.IsFalse(deserializationTask.IsCompleted); - blobStore.CompleteLoad(stream); + blobStore.CompleteLoad(stream); - Assert.AreEqual(Payload, await deserializationTask); - Assert.ThrowsException(() => stream.ReadByte()); + Assert.AreEqual(Payload, await deserializationTask); + Assert.ThrowsException(() => stream.ReadByte()); + } } static async Task AssertInlineMessageRoundTrips( From 170eaed9b2a966eb5f43f3bf77f1c3a0997213fd Mon Sep 17 00:00:00 2001 From: Bernd Verst Date: Thu, 1 Oct 2026 16:06:25 -0700 Subject: [PATCH 3/3] Bump Service Bus patch version to 4.0.4 Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- src/DurableTask.ServiceBus/DurableTask.ServiceBus.csproj | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/DurableTask.ServiceBus/DurableTask.ServiceBus.csproj b/src/DurableTask.ServiceBus/DurableTask.ServiceBus.csproj index 0152c0860..3e2a23474 100644 --- a/src/DurableTask.ServiceBus/DurableTask.ServiceBus.csproj +++ b/src/DurableTask.ServiceBus/DurableTask.ServiceBus.csproj @@ -11,7 +11,7 @@ 4 0 - 3 + 4 $(MajorVersion).$(MinorVersion).$(PatchVersion) $(VersionPrefix).0