diff --git a/Test/DurableTask.AzureStorage.Tests/LoggingTests.cs b/Test/DurableTask.AzureStorage.Tests/LoggingTests.cs new file mode 100644 index 000000000..f52ba5a92 --- /dev/null +++ b/Test/DurableTask.AzureStorage.Tests/LoggingTests.cs @@ -0,0 +1,219 @@ +// ---------------------------------------------------------------------------------- +// 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.AzureStorage.Tests +{ + using System; + using System.Collections.Generic; + using System.Diagnostics.Tracing; + using System.Linq; + using System.Reflection; + using DurableTask.AzureStorage.Logging; + using DurableTask.Core.Logging; + using Microsoft.Extensions.Logging; + using Microsoft.VisualStudio.TestTools.UnitTesting; + + [TestClass] + public class LoggingTests + { + [TestMethod] + public void AbandoningMessage_HasExpectedStructuredFieldsAndPreservesMessage() + { + var logEvent = new LogEvents.AbandoningMessage( + "test-account", + "test-hub", + "TaskScheduled", + 42, + "message-id", + "instance-id", + "execution-id", + "control-queue", + 17, + "pop-receipt", + 30, + "The activity work item could not be processed."); + + var fields = (IReadOnlyDictionary)logEvent; + Assert.AreEqual("test-account", fields["Account"]); + Assert.AreEqual("test-hub", fields["TaskHub"]); + Assert.AreEqual("TaskScheduled", fields["EventType"]); + Assert.AreEqual(42, fields["TaskEventId"]); + Assert.AreEqual("message-id", fields["MessageId"]); + Assert.AreEqual("instance-id", fields["InstanceId"]); + Assert.AreEqual("execution-id", fields["ExecutionId"]); + Assert.AreEqual("control-queue", fields["PartitionId"]); + Assert.AreEqual(17L, fields["SequenceNumber"]); + Assert.AreEqual("pop-receipt", fields["PopReceipt"]); + Assert.AreEqual(30, fields["VisibilityTimeoutSeconds"]); + Assert.AreEqual("The activity work item could not be processed.", fields["Details"]); + Assert.AreEqual(LogLevel.Warning, logEvent.Level); + Assert.AreEqual(EventIds.AbandoningMessage, logEvent.EventId.Id); + Assert.AreEqual(nameof(EventIds.AbandoningMessage), logEvent.EventId.Name); + Assert.AreEqual( + "instance-id: Abandoning [TaskScheduled#42] message back to control-queue and setting a visibility delay of 30ms", + ((ILogEvent)logEvent).FormattedMessage); + } + + [TestMethod] + public void AbandoningMessage_EventSourceSchemaAppendsDetailsAndUsesVersionEight() + { + MethodInfo method = typeof(AnalyticsEventSource).GetMethod(nameof(AnalyticsEventSource.AbandoningMessage)); + EventAttribute eventAttribute = method.GetCustomAttribute(); + string[] parameterNames = method.GetParameters().Select(parameter => parameter.Name).ToArray(); + + Assert.AreEqual(8, eventAttribute.Version); + CollectionAssert.AreEqual( + new[] + { + "Account", + "TaskHub", + "EventType", + "TaskEventId", + "MessageId", + "InstanceId", + "ExecutionId", + "PartitionId", + "SequenceNumber", + "PopReceipt", + "VisibilityTimeoutSeconds", + "AppName", + "ExtensionVersion", + "Details", + }, + parameterNames); + } + + [TestMethod] + public void AbandoningMessage_WriteEventSourceWritesDetailsToFinalPayloadSlot() + { + const string details = "The dispatcher abandoned the work item."; + const string messageId = "event-source-test-message-id"; + var logEvent = new LogEvents.AbandoningMessage( + "test-account", + "test-hub", + "TaskScheduled", + 42, + messageId, + "instance-id", + "execution-id", + "control-queue", + 17, + "pop-receipt", + 30, + details); + + using (var listener = new AbandoningMessageEventListener(messageId)) + { + listener.Enable(); + ((IEventSourceEvent)logEvent).WriteEventSource(); + + Assert.AreEqual(EventIds.AbandoningMessage, listener.EventId); + Assert.AreEqual("Details", listener.PayloadNames.Last()); + Assert.AreEqual(Utils.ExtensionVersion, listener.Payload[listener.Payload.Count - 2]); + Assert.AreEqual(details, listener.Payload.Last()); + } + } + + [TestMethod] + public void AbandoningMessage_LogHelperPropagatesDetails() + { + var logger = new CapturingLogger(); + var logHelper = new LogHelper(logger); + + logHelper.AbandoningMessage( + "test-account", + "test-hub", + "TaskScheduled", + 42, + "message-id", + "instance-id", + "execution-id", + "control-queue", + 17, + "pop-receipt", + 30, + "The orchestration work item could not be processed."); + + var fields = (IReadOnlyDictionary)logger.State; + Assert.AreEqual("The orchestration work item could not be processed.", fields["Details"]); + } + + sealed class CapturingLogger : ILogger + { + public object State { get; private set; } + + public IDisposable BeginScope(TState state) => NullScope.Instance; + + public bool IsEnabled(LogLevel logLevel) => true; + + public void Log( + LogLevel logLevel, + EventId eventId, + TState state, + Exception exception, + Func formatter) + { + this.State = state; + } + } + + sealed class AbandoningMessageEventListener : EventListener + { + readonly string messageId; + + public AbandoningMessageEventListener(string messageId) + { + this.messageId = messageId; + } + + public int EventId { get; private set; } = -1; + + public IReadOnlyList PayloadNames { get; private set; } = Array.Empty(); + + public IReadOnlyList Payload { get; private set; } = Array.Empty(); + + public void Enable() + { + this.EnableEvents(AnalyticsEventSource.Log, EventLevel.Verbose); + } + + public override void Dispose() + { + this.DisableEvents(AnalyticsEventSource.Log); + base.Dispose(); + } + + protected override void OnEventWritten(EventWrittenEventArgs eventData) + { + if (eventData.EventSource == AnalyticsEventSource.Log && + eventData.EventId == EventIds.AbandoningMessage && + eventData.Payload.Count == 14 && + Equals(eventData.Payload[4], this.messageId)) + { + this.EventId = eventData.EventId; + this.PayloadNames = eventData.PayloadNames.ToArray(); + this.Payload = eventData.Payload.ToArray(); + } + } + } + + sealed class NullScope : IDisposable + { + public static readonly NullScope Instance = new NullScope(); + + public void Dispose() + { + } + } + } +} diff --git a/src/DurableTask.AzureStorage/AnalyticsEventSource.cs b/src/DurableTask.AzureStorage/AnalyticsEventSource.cs index 38b6b2ba6..db0d445ee 100644 --- a/src/DurableTask.AzureStorage/AnalyticsEventSource.cs +++ b/src/DurableTask.AzureStorage/AnalyticsEventSource.cs @@ -165,7 +165,7 @@ public void DeletingMessage( ExtensionVersion); } - [Event(EventIds.AbandoningMessage, Level = EventLevel.Warning, Version = 7)] + [Event(EventIds.AbandoningMessage, Level = EventLevel.Warning, Version = 8)] public void AbandoningMessage( string Account, string TaskHub, @@ -179,7 +179,8 @@ public void AbandoningMessage( string PopReceipt, int VisibilityTimeoutSeconds, string AppName, - string ExtensionVersion) + string ExtensionVersion, + string Details) { this.WriteEvent( EventIds.AbandoningMessage, @@ -195,7 +196,8 @@ public void AbandoningMessage( PopReceipt ?? string.Empty, VisibilityTimeoutSeconds, AppName, - ExtensionVersion); + ExtensionVersion, + Details); } [Event(EventIds.AssertFailure, Level = EventLevel.Warning, Message = "An unexpected condition was detected: {2}", Version = 2)] diff --git a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs index 74798cb45..62f3b260e 100644 --- a/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs +++ b/src/DurableTask.AzureStorage/AzureStorageOrchestrationService.cs @@ -726,7 +726,9 @@ async Task LockNextTaskOrchestrationWorkItemAsync(boo // Make sure we still own the partition. If not, abandon the session. if (session.ControlQueue.IsReleased) { - await this.AbandonAndReleaseSessionAsync(session); + await this.AbandonAndReleaseSessionAsync( + session, + "The control queue was released."); return null; } @@ -771,13 +773,18 @@ async Task LockNextTaskOrchestrationWorkItemAsync(boo if (outOfOrderMessages?.Count > 0) { // This will also remove the messages from the current batch. - await this.AbandonMessagesAsync(session, outOfOrderMessages); + await this.AbandonMessagesAsync( + session, + outOfOrderMessages, + "Message was received out of order."); } if (session.CurrentMessageBatch.Count == 0) { // All messages were removed. Release the work item. - await this.AbandonAndReleaseSessionAsync(session); + await this.AbandonAndReleaseSessionAsync( + session, + "No processable messages remained in the session."); return null; } @@ -870,7 +877,9 @@ async Task LockNextTaskOrchestrationWorkItemAsync(boo if (session != null) { // host is shutting down - release any queued messages - await this.AbandonAndReleaseSessionAsync(session); + await this.AbandonAndReleaseSessionAsync( + session, + "Message processing was canceled during shutdown or listener cancellation."); } return null; @@ -1125,11 +1134,11 @@ await this.trackingStore.UpdateInstanceStatusForCompletedOrchestrationAsync( return null; } - async Task AbandonAndReleaseSessionAsync(OrchestrationSession session) + async Task AbandonAndReleaseSessionAsync(OrchestrationSession session, string details) { try { - await this.AbandonSessionAsync(session); + await this.AbandonSessionAsync(session, details); } finally { @@ -1488,20 +1497,25 @@ public Task AbandonTaskOrchestrationWorkItemAsync(TaskOrchestrationWorkItem work return Utils.CompletedTask; } - return this.AbandonSessionAsync(session); + return this.AbandonSessionAsync( + session, + "The orchestration work item was abandoned by the dispatcher."); } - Task AbandonSessionAsync(OrchestrationSession session) + Task AbandonSessionAsync(OrchestrationSession session, string details) { session.StartNewLogicalTraceScope(); - return this.AbandonMessagesAsync(session, session.CurrentMessageBatch.ToList()); + return this.AbandonMessagesAsync(session, session.CurrentMessageBatch.ToList(), details); } - async Task AbandonMessagesAsync(OrchestrationSession session, IList messages) + async Task AbandonMessagesAsync( + OrchestrationSession session, + IList messages, + string details) { await messages.ParallelForEachAsync( this.settings.MaxStorageOperationConcurrency, - message => session.ControlQueue.AbandonMessageAsync(message, session)); + message => session.ControlQueue.AbandonMessageAsync(message, details, session)); // Remove the messages from the current batch. The remaining messages // may still be able to be processed @@ -1680,7 +1694,10 @@ public async Task AbandonTaskActivityWorkItemAsync(TaskActivityWorkItem workItem session.StartNewLogicalTraceScope(); - await this.workItemQueue.AbandonMessageAsync(session.MessageData, session); + await this.workItemQueue.AbandonMessageAsync( + session.MessageData, + "The activity work item was abandoned by the dispatcher.", + session); if (this.activeActivitySessions.TryRemove(workItem.Id, out _)) { diff --git a/src/DurableTask.AzureStorage/Logging/LogEvents.cs b/src/DurableTask.AzureStorage/Logging/LogEvents.cs index cd63f405a..e4376f84d 100644 --- a/src/DurableTask.AzureStorage/Logging/LogEvents.cs +++ b/src/DurableTask.AzureStorage/Logging/LogEvents.cs @@ -336,7 +336,8 @@ public AbandoningMessage( string partitionId, long sequenceNumber, string popReceipt, - int visibilityTimeoutSeconds) + int visibilityTimeoutSeconds, + string details) { this.Account = account; this.TaskHub = taskHub; @@ -349,6 +350,7 @@ public AbandoningMessage( this.SequenceNumber = sequenceNumber; this.PopReceipt = popReceipt; this.VisibilityTimeoutSeconds = visibilityTimeoutSeconds; + this.Details = details; } [StructuredLogField] @@ -384,6 +386,9 @@ public AbandoningMessage( [StructuredLogField] public int VisibilityTimeoutSeconds { get; } + [StructuredLogField] + public string Details { get; } + public override EventId EventId => new EventId( EventIds.AbandoningMessage, nameof(EventIds.AbandoningMessage)); @@ -410,7 +415,8 @@ void IEventSourceEvent.WriteEventSource() => AnalyticsEventSource.Log.Abandoning this.PopReceipt, this.VisibilityTimeoutSeconds, Utils.AppName, - Utils.ExtensionVersion); + Utils.ExtensionVersion, + this.Details); } internal class AssertFailure : StructuredLogEvent, IEventSourceEvent diff --git a/src/DurableTask.AzureStorage/Logging/LogHelper.cs b/src/DurableTask.AzureStorage/Logging/LogHelper.cs index e4ecaae14..bb815e826 100644 --- a/src/DurableTask.AzureStorage/Logging/LogHelper.cs +++ b/src/DurableTask.AzureStorage/Logging/LogHelper.cs @@ -134,7 +134,8 @@ internal void AbandoningMessage( string partitionId, long sequenceNumber, string popReceipt, - int visibilityTimeoutSeconds) + int visibilityTimeoutSeconds, + string details) { var logEvent = new LogEvents.AbandoningMessage( account, @@ -147,7 +148,8 @@ internal void AbandoningMessage( partitionId, sequenceNumber, popReceipt, - visibilityTimeoutSeconds); + visibilityTimeoutSeconds, + details); this.WriteStructuredLog(logEvent); } diff --git a/src/DurableTask.AzureStorage/Messaging/ControlQueue.cs b/src/DurableTask.AzureStorage/Messaging/ControlQueue.cs index 9f1c0d2ba..1c56f7fa9 100644 --- a/src/DurableTask.AzureStorage/Messaging/ControlQueue.cs +++ b/src/DurableTask.AzureStorage/Messaging/ControlQueue.cs @@ -127,7 +127,9 @@ await batch.ParallelForEachAsync(async delegate (QueueMessage queueMessage) // Abandon the message so we can try it again later. // Note: We will fetch the message again from the queue before retrying, so no need to read the receipt - _ = await this.AbandonMessageAsync(queueMessage); + _ = await this.AbandonMessageAsync( + queueMessage, + "Message deserialization failed."); return; } @@ -192,7 +194,7 @@ await batch.ParallelForEachAsync(async delegate (QueueMessage queueMessage) } // This overload is intended for cases where we aren't able to deserialize an instance of MessageData. - public Task AbandonMessageAsync(QueueMessage queueMessage) + public Task AbandonMessageAsync(QueueMessage queueMessage, string details) { this.stats.PendingOrchestratorMessages.TryRemove(queueMessage.MessageId, out _); return base.AbandonMessageAsync( @@ -200,13 +202,17 @@ await batch.ParallelForEachAsync(async delegate (QueueMessage queueMessage) taskMessage: null, instance: null, traceActivityId: null, - sequenceNumber: -1); + sequenceNumber: -1, + details: details); } - public override Task AbandonMessageAsync(MessageData message, SessionBase? session = null) + public override Task AbandonMessageAsync( + MessageData message, + string abandonmentDetails, + SessionBase? session = null) { this.stats.PendingOrchestratorMessages.TryRemove(message.OriginalQueueMessage.MessageId, out _); - return base.AbandonMessageAsync(message, session); + return base.AbandonMessageAsync(message, abandonmentDetails, session); } public override Task DeleteMessageAsync(MessageData message, SessionBase? session = null) diff --git a/src/DurableTask.AzureStorage/Messaging/TaskHubQueue.cs b/src/DurableTask.AzureStorage/Messaging/TaskHubQueue.cs index 221789566..4441f76bb 100644 --- a/src/DurableTask.AzureStorage/Messaging/TaskHubQueue.cs +++ b/src/DurableTask.AzureStorage/Messaging/TaskHubQueue.cs @@ -212,7 +212,10 @@ await this.storageQueue.AddMessageAsync( return initialVisibilityDelay; } - public virtual async Task AbandonMessageAsync(MessageData message, SessionBase? session = null) + public virtual async Task AbandonMessageAsync( + MessageData message, + string abandonmentDetails, + SessionBase? session = null) { QueueMessage queueMessage = message.OriginalQueueMessage; TaskMessage taskMessage = message.TaskMessage; @@ -224,7 +227,8 @@ public virtual async Task AbandonMessageAsync(MessageData message, SessionBase? taskMessage, instance, session?.TraceActivityId, - sequenceNumber); + sequenceNumber, + abandonmentDetails); // If we've successfully abandoned the message, update the pop receipt // (even though we'll likely no longer interact with this message) @@ -239,7 +243,8 @@ public virtual async Task AbandonMessageAsync(MessageData message, SessionBase? TaskMessage? taskMessage, OrchestrationInstance? instance, Guid? traceActivityId, - long sequenceNumber) + long sequenceNumber, + string details) { string instanceId = instance?.InstanceId ?? string.Empty; string executionId = instance?.ExecutionId ?? string.Empty; @@ -277,7 +282,8 @@ public virtual async Task AbandonMessageAsync(MessageData message, SessionBase? this.storageQueue.Name, sequenceNumber, queueMessage.PopReceipt, - numSecondsToWait); + numSecondsToWait, + details); try { diff --git a/src/DurableTask.AzureStorage/OrchestrationSessionManager.cs b/src/DurableTask.AzureStorage/OrchestrationSessionManager.cs index abf7a58b2..fafd0e4d6 100644 --- a/src/DurableTask.AzureStorage/OrchestrationSessionManager.cs +++ b/src/DurableTask.AzureStorage/OrchestrationSessionManager.cs @@ -326,7 +326,11 @@ async Task> DedupeExecutionStartedMessagesAsync( filteredMessages = filteredMessages.Except(messagesToDefer); // Defer messages on a background thread to avoid blocking the dequeue loop - _ = Task.Run(() => messagesToDefer.ParallelForEachAsync(msg => controlQueue.AbandonMessageAsync(msg, session: null))); + _ = Task.Run(() => messagesToDefer.ParallelForEachAsync( + msg => controlQueue.AbandonMessageAsync( + msg, + "Execution-start message was deferred pending instance status reconciliation.", + session: null))); } if (messagesToDiscard?.Count > 0)