From 478ebb42123e4035019f68f7577e82ea3914db0f Mon Sep 17 00:00:00 2001 From: Edouard Chevalier Date: Mon, 5 Oct 2026 14:33:06 +0200 Subject: [PATCH] Fix emulator session message snapshot race --- src/DurableTask.Emulator/AssemblyInfo.cs | 18 ++++ .../LocalOrchestrationService.cs | 2 +- .../PeekLockSessionQueue.cs | 9 +- .../PeekLockSessionQueueTests.cs | 82 +++++++++++++++++++ 4 files changed, 109 insertions(+), 2 deletions(-) create mode 100644 src/DurableTask.Emulator/AssemblyInfo.cs create mode 100644 test/DurableTask.Emulator.Tests/PeekLockSessionQueueTests.cs diff --git a/src/DurableTask.Emulator/AssemblyInfo.cs b/src/DurableTask.Emulator/AssemblyInfo.cs new file mode 100644 index 000000000..1fbc28ff0 --- /dev/null +++ b/src/DurableTask.Emulator/AssemblyInfo.cs @@ -0,0 +1,18 @@ +// ---------------------------------------------------------------------------------- +// 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. +// ---------------------------------------------------------------------------------- + +using System.Runtime.CompilerServices; + +#if !SIGN_ASSEMBLY +[assembly: InternalsVisibleTo("DurableTask.Emulator.Tests")] +#endif diff --git a/src/DurableTask.Emulator/LocalOrchestrationService.cs b/src/DurableTask.Emulator/LocalOrchestrationService.cs index 32ebea257..8821524c2 100644 --- a/src/DurableTask.Emulator/LocalOrchestrationService.cs +++ b/src/DurableTask.Emulator/LocalOrchestrationService.cs @@ -396,7 +396,7 @@ public async Task LockNextTaskOrchestrationWorkItemAs var wi = new TaskOrchestrationWorkItem { - NewMessages = taskSession.Messages.ToList(), + NewMessages = taskSession.Messages, InstanceId = taskSession.Id, LockedUntilUtc = DateTime.UtcNow.AddMinutes(5), OrchestrationRuntimeState = diff --git a/src/DurableTask.Emulator/PeekLockSessionQueue.cs b/src/DurableTask.Emulator/PeekLockSessionQueue.cs index f85184992..4a4538de5 100644 --- a/src/DurableTask.Emulator/PeekLockSessionQueue.cs +++ b/src/DurableTask.Emulator/PeekLockSessionQueue.cs @@ -177,7 +177,14 @@ public async Task AcceptSessionAsync(TimeSpan receiveTimeout, Cance ts.LockTable.Add(tm); } - return ts; + // Producers can append to the locked session as soon as this lock is released. + // Return only the messages recorded in the lock table for this delivery. + return new TaskSession + { + Id = ts.Id, + SessionState = ts.SessionState, + Messages = ts.Messages.ToList(), + }; } } } diff --git a/test/DurableTask.Emulator.Tests/PeekLockSessionQueueTests.cs b/test/DurableTask.Emulator.Tests/PeekLockSessionQueueTests.cs new file mode 100644 index 000000000..0d62d6e85 --- /dev/null +++ b/test/DurableTask.Emulator.Tests/PeekLockSessionQueueTests.cs @@ -0,0 +1,82 @@ +// ---------------------------------------------------------------------------------- +// 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.Emulator.Tests +{ + using System; + using System.Threading; + using System.Threading.Tasks; + using DurableTask.Core; + using DurableTask.Core.History; + using Microsoft.VisualStudio.TestTools.UnitTesting; + + [TestClass] + public class PeekLockSessionQueueTests + { + [TestMethod] + public async Task AcceptSessionKeepsLateMessagesForNextBatch() + { + var queue = new PeekLockSessionQueue(); + TaskMessage first = CreateMessage(1); + TaskMessage late = CreateMessage(2); + queue.SendMessage(first); + + TaskSession accepted = await queue.AcceptSessionAsync(TimeSpan.FromSeconds(1), CancellationToken.None); + + // Reproduce a producer appending after acceptance, before the consumer reads the batch. + // This must not change the delivered batch or add messages absent from its lock table. + queue.SendMessage(late); + CollectionAssert.AreEqual(new[] { first }, accepted.Messages); + + byte[] state = { 1, 2, 3 }; + queue.CompleteSession(accepted.Id, state, Array.Empty(), null); + + TaskSession next = await queue.AcceptSessionAsync(TimeSpan.FromSeconds(1), CancellationToken.None); + CollectionAssert.AreEqual(new[] { late }, next.Messages); + CollectionAssert.AreEqual(state, next.SessionState); + + queue.CompleteSession(next.Id, state, Array.Empty(), null); + Assert.IsNull(await queue.AcceptSessionAsync(TimeSpan.FromMilliseconds(1), CancellationToken.None)); + } + + [TestMethod] + public async Task AbandonSessionRedeliversAcceptedAndLateMessages() + { + var queue = new PeekLockSessionQueue(); + TaskMessage first = CreateMessage(1); + TaskMessage late = CreateMessage(2); + queue.SendMessage(first); + + TaskSession accepted = await queue.AcceptSessionAsync(TimeSpan.FromSeconds(1), CancellationToken.None); + queue.SendMessage(late); + queue.AbandonSession(accepted.Id); + + TaskSession next = await queue.AcceptSessionAsync(TimeSpan.FromSeconds(1), CancellationToken.None); + CollectionAssert.AreEqual(new[] { first, late }, next.Messages); + CollectionAssert.AreEqual(new[] { first }, accepted.Messages); + + byte[] state = { 1 }; + queue.CompleteSession(next.Id, state, Array.Empty(), null); + Assert.IsNull(await queue.AcceptSessionAsync(TimeSpan.FromMilliseconds(1), CancellationToken.None)); + } + + static TaskMessage CreateMessage(int eventId) + { + return new TaskMessage + { + OrchestrationInstance = new OrchestrationInstance { InstanceId = "instance", ExecutionId = "execution" }, + Event = new EventRaisedEvent(eventId, "payload") { Name = "event" }, + }; + } + } +}