From f11b7d5a21e374471d660bc3b6a968bc03e8915c Mon Sep 17 00:00:00 2001 From: n84ck <43924278+n84ck@users.noreply.github.com> Date: Sat, 12 Jul 2025 18:40:52 +0200 Subject: [PATCH 01/10] Update BackgroundStimulusDispatcher.cs Include tenant headers during command dispatch. --- .../Services/BackgroundStimulusDispatcher.cs | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundStimulusDispatcher.cs b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundStimulusDispatcher.cs index 58fad4747..614f63a8d 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/BackgroundStimulusDispatcher.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/BackgroundStimulusDispatcher.cs @@ -1,5 +1,7 @@ +using Elsa.Common.Multitenancy; using Elsa.Mediator; using Elsa.Mediator.Contracts; +using Elsa.Tenants.Mediator; using Elsa.Workflows.Runtime.Commands; using Elsa.Workflows.Runtime.Requests; using Elsa.Workflows.Runtime.Responses; @@ -9,13 +11,18 @@ namespace Elsa.Workflows.Runtime; /// /// A simple implementation that queues the specified request for delivering stimuli on a non-durable background worker. /// -public class BackgroundStimulusDispatcher(ICommandSender commandSender) : IStimulusDispatcher +public class BackgroundStimulusDispatcher(ICommandSender commandSender, ITenantAccessor tenantAccessor) : IStimulusDispatcher { /// public async Task SendAsync(DispatchStimulusRequest request, CancellationToken cancellationToken = default) { var command = new DispatchStimulusCommand(request); - await commandSender.SendAsync(command, CommandStrategy.Background, cancellationToken); + await commandSender.SendAsync(command, CommandStrategy.Background, CreateHeaders(), cancellationToken); return DispatchStimulusResponse.Empty; } -} \ No newline at end of file + + private IDictionary CreateHeaders() + { + return TenantHeaders.CreateHeaders(tenantAccessor.Tenant?.Id); + } +} From b85ed3b32966c778d1128975c10112b76f175b8e Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 22 Jul 2025 21:37:23 +0200 Subject: [PATCH 02/10] Remove unused `GetCapturedActivityExecutionRecord` method from `ActivityExecutionContextRecordExtensions`. --- .../Extensions/ActivityExecutionContextRecordExtensions.cs | 5 ----- 1 file changed, 5 deletions(-) diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs index c2706cc1e..9f0810fc1 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs @@ -16,11 +16,6 @@ public static class ActivityExecutionContextRecordExtensions context.TransientProperties[ActivityExecutionRecordKey] = record; } - public static ActivityExecutionRecord? GetCapturedActivityExecutionRecord(this ActivityExecutionContext context) - { - return context.TransientProperties.TryGetValue(ActivityExecutionRecordKey, out var record) ? (ActivityExecutionRecord?)record : null; - } - public static async Task GetOrMapCapturedActivityExecutionRecordAsync(this ActivityExecutionContext context) { if(context.TransientProperties.TryGetValue(ActivityExecutionRecordKey, out var record)) From 18ee694d2360afb64a5359bd4e969f15ce751797 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 24 Jul 2025 19:47:08 +0200 Subject: [PATCH 03/10] Refactor variable reading and workflow state extraction logic. Revised method parameters, simplified object initializations in `DefaultWorkflowInstanceVariableReader`, and enhanced hierarchy reconstruction in `WorkflowStateExtractor` with logging for missing parent contexts. --- .../DefaultWorkflowInstanceVariableReader.cs | 6 +++--- .../Services/WorkflowStateExtractor.cs | 13 ++++++++++--- 2 files changed, 13 insertions(+), 6 deletions(-) diff --git a/src/modules/Elsa.Workflows.Core/Services/DefaultWorkflowInstanceVariableReader.cs b/src/modules/Elsa.Workflows.Core/Services/DefaultWorkflowInstanceVariableReader.cs index c3557bf10..3eca2c3a5 100644 --- a/src/modules/Elsa.Workflows.Core/Services/DefaultWorkflowInstanceVariableReader.cs +++ b/src/modules/Elsa.Workflows.Core/Services/DefaultWorkflowInstanceVariableReader.cs @@ -2,11 +2,11 @@ namespace Elsa.Workflows; public class DefaultWorkflowInstanceVariableReader(IVariablePersistenceManager variablePersistenceManager) : IWorkflowInstanceVariableReader { - public async Task> GetVariables(WorkflowExecutionContext workflowExecutionContext, IEnumerable? excludeTags = default, CancellationToken cancellationToken = default) + public async Task> GetVariables(WorkflowExecutionContext workflowExecutionContext, IEnumerable? excludeTags = null, CancellationToken cancellationToken = default) { var workflow = workflowExecutionContext.Workflow; var workflowVariables = workflow.Variables; - var rootWorkflowActivityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x.ParentActivityExecutionContext == null); + var rootWorkflowActivityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.FirstOrDefault(x => x.Activity == workflow); if (rootWorkflowActivityExecutionContext == null) return []; @@ -17,7 +17,7 @@ public class DefaultWorkflowInstanceVariableReader(IVariablePersistenceManager v foreach (var workflowVariable in workflowVariables) { var value = workflowVariable.Get(rootWorkflowActivityExecutionContext.ExpressionExecutionContext); - resolvedVariables.Add(new ResolvedVariable(workflowVariable, value)); + resolvedVariables.Add(new(workflowVariable, value)); } return resolvedVariables; diff --git a/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs b/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs index 40297f297..f37d2a6b4 100644 --- a/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs +++ b/src/modules/Elsa.Workflows.Core/Services/WorkflowStateExtractor.cs @@ -2,11 +2,12 @@ using Elsa.Extensions; using Elsa.Workflows.Models; using Elsa.Workflows.Services; using Elsa.Workflows.State; +using Microsoft.Extensions.Logging; namespace Elsa.Workflows; /// -public class WorkflowStateExtractor : IWorkflowStateExtractor +public class WorkflowStateExtractor(ILogger logger) : IWorkflowStateExtractor { /// public WorkflowState Extract(WorkflowExecutionContext workflowExecutionContext) @@ -100,7 +101,7 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor workflowExecutionContext.Properties[property.Key] = property.Value; } - private static async Task ApplyActivityExecutionContextsAsync(WorkflowState state, WorkflowExecutionContext workflowExecutionContext) + private async Task ApplyActivityExecutionContextsAsync(WorkflowState state, WorkflowExecutionContext workflowExecutionContext) { var activityExecutionContexts = (await Task.WhenAll( state.ActivityExecutionContexts.Select(async item => await CreateActivityExecutionContextAsync(item)))) @@ -113,7 +114,13 @@ public class WorkflowStateExtractor : IWorkflowStateExtractor // Reconstruct hierarchy. foreach (var contextState in state.ActivityExecutionContexts.Where(x => !string.IsNullOrWhiteSpace(x.ParentContextId))) { - var parentContext = lookup[contextState.ParentContextId!]; + var parentContextId = contextState.ParentContextId; + if (parentContextId == null || !lookup.TryGetValue(parentContextId, out var parentContext)) + { + logger.LogWarning("Parent context with ID '{ParentContextId}' not found for context with ID '{ContextId}'.", parentContextId, contextState.Id); + continue; // Skip if parent context is not found. + } + var contextId = contextState.Id; if (lookup.TryGetValue(contextId, out var context)) From 61179bcb53e70dbd9787e09473e46586fffab0dd Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 24 Jul 2025 19:47:26 +0200 Subject: [PATCH 04/10] Remove redundant activity metadata properties from `ActivityExecutionRecordSnapshot` and streamline mapping logic Deleted unused metadata properties to simplify `ActivityExecutionRecordSnapshot`. Updated `GetOrMapCapturedActivityExecutionRecordAsync` to maintain serialized snapshots when mapping, ensuring consistency in activity execution records. --- .../ActivityExecutionContextRecordExtensions.cs | 10 ++++++---- .../Models/ActivityExecutionRecordSnapshot.cs | 13 ------------- .../Services/DefaultActivityExecutionMapper.cs | 13 ------------- 3 files changed, 6 insertions(+), 30 deletions(-) diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs index 9f0810fc1..950e8bfdf 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs @@ -18,10 +18,12 @@ public static class ActivityExecutionContextRecordExtensions public static async Task GetOrMapCapturedActivityExecutionRecordAsync(this ActivityExecutionContext context) { - if(context.TransientProperties.TryGetValue(ActivityExecutionRecordKey, out var record)) - return (ActivityExecutionRecord)record; - var mapper = context.GetRequiredService(); - return await mapper.MapAsync(context); + var record = await mapper.MapAsync(context); + + if (context.TransientProperties.TryGetValue(ActivityExecutionRecordKey, out var capturedRecord)) + record.SerializedSnapshot = ((ActivityExecutionRecord)capturedRecord).SerializedSnapshot; + + return record; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs b/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs index 377850153..3f0031f2a 100644 --- a/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs +++ b/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs @@ -2,19 +2,6 @@ namespace Elsa.Workflows.Runtime; public class ActivityExecutionRecordSnapshot { - public string Id { get; set; } = null!; - public string? TenantId { get; set; } - public string WorkflowInstanceId { get; set; } = null!; - public string ActivityId { get; set; } = null!; - public string ActivityNodeId { get; set; } = null!; - public string ActivityType { get; set; } = null!; - public int ActivityTypeVersion { get; set; } - public string? ActivityName { get; set; } - public DateTimeOffset StartedAt { get; set; } - public bool HasBookmarks { get; set; } - public ActivityStatus Status { get; set; } - public int AggregateFaultCount { get; set; } - public DateTimeOffset? CompletedAt { get; set; } public string? SerializedActivityState { get; set; } public string? SerializedOutputs { get; set; } public string? SerializedProperties { get; set; } diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs index 223cd55f3..3779f253d 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs @@ -58,19 +58,6 @@ public class DefaultActivityExecutionMapper( var serializedMetadata = record.Metadata != null ? payloadSerializer.Serialize(record.Metadata) : null; record.SerializedSnapshot = new() { - Id = record.Id, - TenantId = record.TenantId, - WorkflowInstanceId = record.WorkflowInstanceId, - ActivityId = record.ActivityId, - ActivityNodeId = record.ActivityNodeId, - ActivityType = record.ActivityType, - ActivityTypeVersion = record.ActivityTypeVersion, - ActivityName = record.ActivityName, - StartedAt = record.StartedAt, - HasBookmarks = record.HasBookmarks, - Status = record.Status, - AggregateFaultCount = record.AggregateFaultCount, - CompletedAt = record.CompletedAt, SerializedActivityState = compressedSerializedActivityState, SerializedActivityStateCompressionAlgorithm = compressionAlgorithm, SerializedOutputs = record.Outputs?.Any() == true ? safeSerializer.Serialize(record.Outputs) : null, From 9b1a76d047eef7f5244e61527399e12f396e2d04 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 24 Jul 2025 21:38:18 +0200 Subject: [PATCH 05/10] Extend `ActivityExecutionRecordSnapshot` and update `DefaultActivityExecutionMapper` to include additional activity execution details. --- .../Models/ActivityExecutionRecordSnapshot.cs | 13 +++++++++++++ .../Services/DefaultActivityExecutionMapper.cs | 13 +++++++++++++ 2 files changed, 26 insertions(+) diff --git a/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs b/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs index 3f0031f2a..377850153 100644 --- a/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs +++ b/src/modules/Elsa.Workflows.Runtime/Models/ActivityExecutionRecordSnapshot.cs @@ -2,6 +2,19 @@ namespace Elsa.Workflows.Runtime; public class ActivityExecutionRecordSnapshot { + public string Id { get; set; } = null!; + public string? TenantId { get; set; } + public string WorkflowInstanceId { get; set; } = null!; + public string ActivityId { get; set; } = null!; + public string ActivityNodeId { get; set; } = null!; + public string ActivityType { get; set; } = null!; + public int ActivityTypeVersion { get; set; } + public string? ActivityName { get; set; } + public DateTimeOffset StartedAt { get; set; } + public bool HasBookmarks { get; set; } + public ActivityStatus Status { get; set; } + public int AggregateFaultCount { get; set; } + public DateTimeOffset? CompletedAt { get; set; } public string? SerializedActivityState { get; set; } public string? SerializedOutputs { get; set; } public string? SerializedProperties { get; set; } diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs index 3779f253d..223cd55f3 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultActivityExecutionMapper.cs @@ -58,6 +58,19 @@ public class DefaultActivityExecutionMapper( var serializedMetadata = record.Metadata != null ? payloadSerializer.Serialize(record.Metadata) : null; record.SerializedSnapshot = new() { + Id = record.Id, + TenantId = record.TenantId, + WorkflowInstanceId = record.WorkflowInstanceId, + ActivityId = record.ActivityId, + ActivityNodeId = record.ActivityNodeId, + ActivityType = record.ActivityType, + ActivityTypeVersion = record.ActivityTypeVersion, + ActivityName = record.ActivityName, + StartedAt = record.StartedAt, + HasBookmarks = record.HasBookmarks, + Status = record.Status, + AggregateFaultCount = record.AggregateFaultCount, + CompletedAt = record.CompletedAt, SerializedActivityState = compressedSerializedActivityState, SerializedActivityStateCompressionAlgorithm = compressionAlgorithm, SerializedOutputs = record.Outputs?.Any() == true ? safeSerializer.Serialize(record.Outputs) : null, From 3fd8a252cdf7a0f009c0140cc551c32acdacfcb2 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 25 Jul 2025 07:46:35 +0200 Subject: [PATCH 06/10] Add cancellation handling for executing scheduled tasks (#6819) * Add cancellation handling for executing scheduled tasks Enhanced `ScheduledRecurringTask`, `ScheduledCronTask`, and `ScheduledSpecificInstantTask` to properly handle cancellation scenarios when tasks are executing. Introduced `_executing` and `_cancellationRequested` flags to ensure clean cancellation processes. * Refactor `ScheduledRecurringTask` to improve readability and fix formatting issues. * Add `SemaphoreSlim` for concurrent task execution control in scheduled tasks Introduce `SemaphoreSlim` to manage and safeguard concurrent executions in `ScheduledRecurringTask`, `ScheduledCronTask`, and `ScheduledSpecificInstantTask`. Enhances thread safety and prevents overlapping executions. Added exception handling and proper semaphore release to ensure robustness. * Add ILogger to scheduled tasks and improve error logging Integrated `ILogger` into `ScheduledRecurringTask` to enhance logging capabilities and replaced the generic comment-based error handling with proper logging for better traceability. Cleaned up and formatted `ScheduledCronTask` for improved code readability. * Dispose of semaphore fields in scheduled task classes to ensure proper resource cleanup. * Update src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Set `_executing` to `false` in task `finally` blocks to ensure state reset after execution. * Update src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledSpecificInstantTask.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Remove redundant `_executing` assignment after `SendAsync` in scheduled task classes. * Update src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledCronTask.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledSpecificInstantTask.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update semaphore logic to prevent blocked tasks when cancellation is requested Refactored `_executionSemaphore.WaitAsync` usage in `ScheduledRecurringTask`, `ScheduledCronTask`, and `ScheduledSpecificInstantTask` to use non-blocking semaphore acquisition with cancellation token support. This ensures graceful handling of pending tasks during cancellation scenarios. * Refactor delay condition checks to use `TimeSpan.Zero` for improved readability and precision. --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- .../ScheduledTasks/ScheduledCronTask.cs | 54 ++++++++++++++-- .../ScheduledTasks/ScheduledRecurringTask.cs | 61 +++++++++++++++---- .../ScheduledSpecificInstantTask.cs | 44 ++++++++++--- 3 files changed, 134 insertions(+), 25 deletions(-) diff --git a/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledCronTask.cs b/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledCronTask.cs index 72b87444c..2b61acc78 100644 --- a/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledCronTask.cs +++ b/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledCronTask.cs @@ -20,6 +20,9 @@ public class ScheduledCronTask : IScheduledTask, IDisposable private readonly ICronParser _cronParser; private readonly IServiceScopeFactory _scopeFactory; private readonly CancellationTokenSource _cancellationTokenSource; + private readonly SemaphoreSlim _executionSemaphore = new(1, 1); + private bool _executing; + private bool _cancellationRequested; /// /// Initializes a new instance of . @@ -32,7 +35,7 @@ public class ScheduledCronTask : IScheduledTask, IDisposable _scopeFactory = scopeFactory; _systemClock = systemClock; _logger = logger; - _cancellationTokenSource = new CancellationTokenSource(); + _cancellationTokenSource = new(); Schedule(); } @@ -41,6 +44,13 @@ public class ScheduledCronTask : IScheduledTask, IDisposable public void Cancel() { _timer?.Dispose(); + + if (_executing) + { + _cancellationRequested = true; + return; + } + _cancellationTokenSource.Cancel(); } @@ -54,7 +64,7 @@ public class ScheduledCronTask : IScheduledTask, IDisposable var nextOccurence = _cronParser.GetNextOccurrence(_cronExpression); var delay = nextOccurence - now; - if (!adjusted && delay.Milliseconds <= 0) + if (!adjusted && delay <= TimeSpan.Zero) { adjusted = true; continue; @@ -67,7 +77,7 @@ public class ScheduledCronTask : IScheduledTask, IDisposable private void TrySetupTimer(TimeSpan delay) { - if (delay.Milliseconds <= 0) + if (delay <= TimeSpan.Zero) return; try @@ -82,7 +92,10 @@ public class ScheduledCronTask : IScheduledTask, IDisposable private void SetupTimer(TimeSpan delay) { - _timer = new Timer(delay.TotalMilliseconds) { Enabled = true }; + _timer = new(delay.TotalMilliseconds) + { + Enabled = true + }; _timer.Elapsed += async (_, _) => { @@ -93,8 +106,36 @@ public class ScheduledCronTask : IScheduledTask, IDisposable var commandSender = scope.ServiceProvider.GetRequiredService(); var cancellationToken = _cancellationTokenSource.Token; - if (!cancellationToken.IsCancellationRequested) await commandSender.SendAsync(new RunScheduledTask(_task), cancellationToken); - if (!cancellationToken.IsCancellationRequested) Schedule(); + + if (!cancellationToken.IsCancellationRequested) + { + try + { + var acquired = await _executionSemaphore.WaitAsync(0, cancellationToken); + if (!acquired) return; + + _executing = true; + await commandSender.SendAsync(new RunScheduledTask(_task), cancellationToken); + + if (_cancellationRequested) + { + _cancellationRequested = false; + _cancellationTokenSource.Cancel(); + } + } + catch (Exception e) + { + _logger.LogError(e, "Error executing scheduled task"); + } + finally + { + _executing = false; + _executionSemaphore.Release(); + } + } + + if (!cancellationToken.IsCancellationRequested) + Schedule(); }; } @@ -102,5 +143,6 @@ public class ScheduledCronTask : IScheduledTask, IDisposable { _timer?.Dispose(); _cancellationTokenSource.Dispose(); + _executionSemaphore.Dispose(); } } \ No newline at end of file diff --git a/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs b/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs index 21d7fd380..89a75575d 100644 --- a/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs +++ b/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledRecurringTask.cs @@ -2,6 +2,7 @@ using Elsa.Common; using Elsa.Mediator.Contracts; using Elsa.Scheduling.Commands; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; using Timer = System.Timers.Timer; namespace Elsa.Scheduling.ScheduledTasks; @@ -14,24 +15,24 @@ public class ScheduledRecurringTask : IScheduledTask, IDisposable private readonly ITask _task; private readonly ISystemClock _systemClock; private readonly IServiceScopeFactory _scopeFactory; + private readonly ILogger _logger; private readonly TimeSpan _interval; private readonly CancellationTokenSource _cancellationTokenSource; + private readonly SemaphoreSlim _executionSemaphore = new(1, 1); private DateTimeOffset _startAt; private Timer? _timer; + private bool _executing; + private bool _cancellationRequested; /// /// Initializes a new instance of . /// - /// The task to execute. - /// The instant at which to start executing the task. - /// The interval at which to execute the task. - /// The system clock. - /// Scope factory to create the scope and get dependancies. - public ScheduledRecurringTask(ITask task, DateTimeOffset startAt, TimeSpan interval, ISystemClock systemClock, IServiceScopeFactory scopeFactory) + public ScheduledRecurringTask(ITask task, DateTimeOffset startAt, TimeSpan interval, ISystemClock systemClock, IServiceScopeFactory scopeFactory, ILogger logger) { _task = task; _systemClock = systemClock; _scopeFactory = scopeFactory; + _logger = logger; _startAt = startAt; _interval = interval; _cancellationTokenSource = new(); @@ -43,6 +44,13 @@ public class ScheduledRecurringTask : IScheduledTask, IDisposable public void Cancel() { _timer?.Dispose(); + + if (_executing) + { + _cancellationRequested = true; + return; + } + _cancellationTokenSource.Cancel(); } @@ -56,7 +64,7 @@ public class ScheduledRecurringTask : IScheduledTask, IDisposable var now = _systemClock.UtcNow; var delay = startAt - now; - if (!adjusted && delay.Milliseconds <= 0) + if (!adjusted && delay <= TimeSpan.Zero) { adjusted = true; continue; @@ -71,7 +79,10 @@ public class ScheduledRecurringTask : IScheduledTask, IDisposable { if (delay < TimeSpan.Zero) delay = TimeSpan.FromSeconds(1); - _timer = new(delay.TotalMilliseconds) { Enabled = true }; + _timer = new(delay.TotalMilliseconds) + { + Enabled = true + }; _timer.Elapsed += async (_, _) => { @@ -81,10 +92,35 @@ public class ScheduledRecurringTask : IScheduledTask, IDisposable using var scope = _scopeFactory.CreateScope(); var commandSender = scope.ServiceProvider.GetRequiredService(); - var cancellationToken = _cancellationTokenSource.Token; - if (!cancellationToken.IsCancellationRequested) await commandSender.SendAsync(new RunScheduledTask(_task), cancellationToken); - if (!cancellationToken.IsCancellationRequested) Schedule(); + if (!cancellationToken.IsCancellationRequested) + { + try + { + var acquired = await _executionSemaphore.WaitAsync(0, cancellationToken); + if (!acquired) return; + _executing = true; + await commandSender.SendAsync(new RunScheduledTask(_task), cancellationToken); + + if (_cancellationRequested) + { + _cancellationRequested = false; + _cancellationTokenSource.Cancel(); + } + } + catch (Exception e) + { + _logger.LogError(e, "Error executing scheduled task"); + } + finally + { + _executing = false; + _executionSemaphore.Release(); + } + } + + if (!cancellationToken.IsCancellationRequested) + Schedule(); }; } @@ -92,5 +128,6 @@ public class ScheduledRecurringTask : IScheduledTask, IDisposable { _cancellationTokenSource.Dispose(); _timer?.Dispose(); + _executionSemaphore.Dispose(); } -} +} \ No newline at end of file diff --git a/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledSpecificInstantTask.cs b/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledSpecificInstantTask.cs index 3d2e9c61e..6979dd410 100644 --- a/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledSpecificInstantTask.cs +++ b/src/modules/Elsa.Scheduling/ScheduledTasks/ScheduledSpecificInstantTask.cs @@ -18,7 +18,10 @@ public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable private readonly ILogger _logger; private readonly DateTimeOffset _startAt; private readonly CancellationTokenSource _cancellationTokenSource; + private readonly SemaphoreSlim _executionSemaphore = new(1, 1); private Timer? _timer; + private bool _executing; + private bool _cancellationRequested; /// /// Initializes a new instance of . @@ -30,23 +33,37 @@ public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable _scopeFactory = scopeFactory; _logger = logger; _startAt = startAt; - _cancellationTokenSource = new CancellationTokenSource(); + _cancellationTokenSource = new(); Schedule(); } /// - public void Cancel() => _timer?.Dispose(); + public void Cancel() + { + _timer?.Dispose(); + + if (_executing) + { + _cancellationRequested = true; + return; + } + + _cancellationTokenSource.Cancel(); + } private void Schedule() { var now = _systemClock.UtcNow; var delay = _startAt - now; - if (delay.Milliseconds <= 0) + if (delay <= TimeSpan.Zero) delay = TimeSpan.FromMilliseconds(1); - _timer = new Timer(delay.TotalMilliseconds) { Enabled = true }; + _timer = new(delay.TotalMilliseconds) + { + Enabled = true + }; _timer.Elapsed += async (_, _) => { @@ -55,19 +72,31 @@ public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable using var scope = _scopeFactory.CreateScope(); var commandSender = scope.ServiceProvider.GetRequiredService(); - var cancellationToken = _cancellationTokenSource.Token; if (!cancellationToken.IsCancellationRequested) { try { + var acquired = await _executionSemaphore.WaitAsync(0, cancellationToken); + if (!acquired) return; + _executing = true; await commandSender.SendAsync(new RunScheduledTask(_task), cancellationToken); + + if (_cancellationRequested) + { + _cancellationRequested = false; + _cancellationTokenSource.Cancel(); + } } catch (Exception e) { - _logger.LogError(e, "Error scheduled task"); + _logger.LogError(e, "Error executing scheduled task"); + } + finally + { + _executing = false; + _executionSemaphore.Release(); } - } }; } @@ -76,5 +105,6 @@ public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable { _cancellationTokenSource.Dispose(); _timer?.Dispose(); + _executionSemaphore.Dispose(); } } \ No newline at end of file From deb42d7f19011e478b4ba834169481cd6f9b8c1f Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 25 Jul 2025 14:32:12 +0200 Subject: [PATCH 07/10] Update `ActivityExecutionContextRecordExtensions` to preserve and update serialized snapshots (#6823) Refactor the extension method to merge existing serialized snapshots with updated activity execution properties, ensuring the latest context state is retained without overwriting prior data. --- .../ActivityExecutionContextRecordExtensions.cs | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs index 950e8bfdf..b9a49106c 100644 --- a/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs +++ b/src/modules/Elsa.Workflows.Runtime/Extensions/ActivityExecutionContextRecordExtensions.cs @@ -21,8 +21,20 @@ public static class ActivityExecutionContextRecordExtensions var mapper = context.GetRequiredService(); var record = await mapper.MapAsync(context); - if (context.TransientProperties.TryGetValue(ActivityExecutionRecordKey, out var capturedRecord)) - record.SerializedSnapshot = ((ActivityExecutionRecord)capturedRecord).SerializedSnapshot; + if (context.TransientProperties.TryGetValue(ActivityExecutionRecordKey, out var capturedRecord)) + { + var serializedSnapshot = ((ActivityExecutionRecord)capturedRecord).SerializedSnapshot!; + + // Take the existing serialized snapshot. + record.SerializedSnapshot = serializedSnapshot; + + // Update the serialized snapshot with the current record's properties. + // This will reflect the latest state of the activity execution context without losing the existing serialized snapshot representing e.g., variable values at the time of the record capture. + serializedSnapshot.HasBookmarks = record.HasBookmarks; + serializedSnapshot.Status = record.Status; + serializedSnapshot.AggregateFaultCount = record.AggregateFaultCount; + serializedSnapshot.CompletedAt = record.CompletedAt; + } return record; } From 7467b6347dddacc78a27baf877c819edb382759d Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 29 Jul 2025 20:22:39 +0200 Subject: [PATCH 08/10] Add support for flow authorization activities and bookmark trigger URL generation (#6828) * Add support for flow authorization activities and bookmark trigger URL generation - Introduced `AuthorizeFlow` activity for configurable policy-based flow authorization. - Added extensions for generating bookmark trigger URLs. - Created `BookmarkTokenPayload` and updated APIs to handle bookmark resumption with SAS tokens. - Refactored and consolidated related code for improved modularity and clarity. * Update src/modules/Elsa.Http/Extensions/BookmarkExpressionExecutionContextExtensions.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update src/modules/Elsa.Http/Extensions/BookmarkExpressionExecutionContextExtensions.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * Update src/apps/Elsa.Server.Web/Activities/AuthorizeFlow.cs Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> --- .../Activities/AuthorizeFlow.cs | 49 ++++++++++ src/modules/Elsa.Http/Elsa.Http.csproj | 1 + ...arkExpressionExecutionContextExtensions.cs | 75 ++++++++++++++++ ...ntExpressionExecutionContextExtensions.cs} | 13 +-- .../Elsa.Http/Options/HttpActivityOptions.cs | 2 +- .../Elsa.Workflows.Api.csproj | 2 +- .../Endpoints/Bookmarks/Resume/Endpoint.cs | 89 +++++++++++++++++++ .../Events/TriggerPublic/Endpoint.cs | 1 - .../Features/WorkflowsApiFeature.cs | 10 +-- .../Activities/Flowchart/Models/Outcomes.cs | 10 +-- .../Models/BookmarkTokenPayload.cs | 6 ++ .../Models/EventTokenPayload.cs | 2 +- 12 files changed, 236 insertions(+), 24 deletions(-) create mode 100644 src/apps/Elsa.Server.Web/Activities/AuthorizeFlow.cs create mode 100644 src/modules/Elsa.Http/Extensions/BookmarkExpressionExecutionContextExtensions.cs rename src/modules/Elsa.Http/Extensions/{ExpressionExecutionContextExtensions.cs => EventExpressionExecutionContextExtensions.cs} (88%) create mode 100644 src/modules/Elsa.Workflows.Api/Endpoints/Bookmarks/Resume/Endpoint.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Models/BookmarkTokenPayload.cs rename src/modules/{Elsa.Http => Elsa.Workflows.Runtime}/Models/EventTokenPayload.cs (90%) diff --git a/src/apps/Elsa.Server.Web/Activities/AuthorizeFlow.cs b/src/apps/Elsa.Server.Web/Activities/AuthorizeFlow.cs new file mode 100644 index 000000000..1761e811f --- /dev/null +++ b/src/apps/Elsa.Server.Web/Activities/AuthorizeFlow.cs @@ -0,0 +1,49 @@ +using Elsa.Extensions; +using Elsa.Workflows; +using Elsa.Workflows.Activities.Flowchart.Attributes; +using Elsa.Workflows.Attributes; + +namespace Elsa.Server.Web.Activities; + +[Activity("Elsa", "Authorization", "Authorizes a flow based on the configured policies.")] +[FlowNode("Authorized", "Unauthorized", "Error")] +public class AuthorizeFlow : Activity +{ + protected override ValueTask ExecuteAsync(ActivityExecutionContext context) + { + var httpContext = context.GetRequiredService().HttpContext; + + if (httpContext == null) + throw new InvalidOperationException("HttpContext is not available. Ensure that the activity is executed within an HTTP request context."); + + var bookmark = context.CreateBookmark(new AuthorizeStimulus(), OnResumeAsync); + var redirectUrl = context.ExpressionExecutionContext.GenerateBookmarkTriggerUrl(bookmark.Id); + + Result.Set(context, redirectUrl); + return ValueTask.CompletedTask; + } + + private async ValueTask OnResumeAsync(ActivityExecutionContext context) + { + if (!context.TryGetWorkflowInput("Answer", out var response)) + { + await context.CompleteActivityWithOutcomesAsync("Unauthorized"); + return; + } + + switch (response) + { + case "Authorized": + await context.CompleteActivityWithOutcomesAsync("Authorized"); + return; + case "Error": + await context.CompleteActivityWithOutcomesAsync("Error"); + return; + default: + await context.CompleteActivityWithOutcomesAsync("Unauthorized"); + break; + } + } +} + +public record AuthorizeStimulus; \ No newline at end of file diff --git a/src/modules/Elsa.Http/Elsa.Http.csproj b/src/modules/Elsa.Http/Elsa.Http.csproj index 7901bd00d..06a602fb2 100644 --- a/src/modules/Elsa.Http/Elsa.Http.csproj +++ b/src/modules/Elsa.Http/Elsa.Http.csproj @@ -19,6 +19,7 @@ + diff --git a/src/modules/Elsa.Http/Extensions/BookmarkExpressionExecutionContextExtensions.cs b/src/modules/Elsa.Http/Extensions/BookmarkExpressionExecutionContextExtensions.cs new file mode 100644 index 000000000..1cdcdbd92 --- /dev/null +++ b/src/modules/Elsa.Http/Extensions/BookmarkExpressionExecutionContextExtensions.cs @@ -0,0 +1,75 @@ +using Elsa.Expressions.Models; +using Elsa.Http; +using Elsa.Http.Options; +using Elsa.SasTokens.Contracts; +using Elsa.Workflows.Api; +using Elsa.Workflows.Runtime; +using Microsoft.Extensions.Options; + +// ReSharper disable once CheckNamespace +namespace Elsa.Extensions; + +/// +/// Provides extension methods for working with and generating bookmark trigger URLs. +/// +public static class BookmarkExpressionExecutionContextExtensions +{ + /// + /// Generates a URL that can be used to resume a bookmarked workflow. + /// + /// The expression execution context. + /// The ID of the bookmark to resume. + /// The lifetime of the bookmark trigger token. + /// A URL that can be used to resume a bookmarked workflow. + public static string GenerateBookmarkTriggerUrl(this ExpressionExecutionContext context, string bookmarkId, TimeSpan lifetime) + { + var token = context.GenerateBookmarkTriggerTokenInternal(bookmarkId, lifetime); + return context.GenerateBookmarkTriggerUrlInternal(token); + } + + /// + /// Generates a URL that can be used to resume a bookmarked workflow. + /// + /// The expression execution context. + /// The ID of the bookmark to resume. + /// The expiration date of the bookmark trigger token. + /// A URL that can be used to resume a bookmarked workflow. + public static string GenerateBookmarkTriggerUrl(this ExpressionExecutionContext context, string bookmarkId, DateTimeOffset expiresAt) + { + var token = context.GenerateBookmarkTriggerTokenInternal(bookmarkId, expiresAt: expiresAt); + return context.GenerateBookmarkTriggerUrlInternal(token); + } + + /// + /// Generates a URL that can be used to resume a bookmarked workflow. + /// + /// The expression execution context. + /// The ID of the bookmark to resume. + /// A URL that can be used to trigger an event. + public static string GenerateBookmarkTriggerUrl(this ExpressionExecutionContext context, string bookmarkId) + { + var token = context.GenerateBookmarkTriggerTokenInternal(bookmarkId); + return context.GenerateBookmarkTriggerUrlInternal(token); + } + + private static string GenerateBookmarkTriggerUrlInternal(this ExpressionExecutionContext context, string token) + { + var options = context.GetRequiredService>().Value; + var url = $"{options.RoutePrefix}/bookmarks/resume?t={token}"; + var absoluteUrlProvider = context.GetRequiredService(); + return absoluteUrlProvider.ToAbsoluteUrl(url).ToString(); + } + + private static string GenerateBookmarkTriggerTokenInternal(this ExpressionExecutionContext context, string bookmarkId, TimeSpan? lifetime = null, DateTimeOffset? expiresAt = null) + { + var workflowInstanceId = context.GetWorkflowExecutionContext().Id; + var payload = new BookmarkTokenPayload(bookmarkId, workflowInstanceId); + var tokenService = context.GetRequiredService(); + + return lifetime != null + ? tokenService.CreateToken(payload, lifetime.Value) + : expiresAt != null + ? tokenService.CreateToken(payload, expiresAt.Value) + : tokenService.CreateToken(payload); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Http/Extensions/ExpressionExecutionContextExtensions.cs b/src/modules/Elsa.Http/Extensions/EventExpressionExecutionContextExtensions.cs similarity index 88% rename from src/modules/Elsa.Http/Extensions/ExpressionExecutionContextExtensions.cs rename to src/modules/Elsa.Http/Extensions/EventExpressionExecutionContextExtensions.cs index a4eac875f..9f882d7bb 100644 --- a/src/modules/Elsa.Http/Extensions/ExpressionExecutionContextExtensions.cs +++ b/src/modules/Elsa.Http/Extensions/EventExpressionExecutionContextExtensions.cs @@ -1,16 +1,17 @@ using Elsa.Expressions.Models; using Elsa.Http; -using Elsa.Http.Options; using Elsa.SasTokens.Contracts; +using Elsa.Workflows.Api; +using Elsa.Workflows.Runtime; using Microsoft.Extensions.Options; // ReSharper disable once CheckNamespace namespace Elsa.Extensions; /// -/// +/// Provides extension methods for working with . /// -public static class ExpressionExecutionContextExtensions +public static class EventExpressionExecutionContextExtensions { /// /// Generates a URL that can be used to trigger an event. @@ -52,13 +53,13 @@ public static class ExpressionExecutionContextExtensions private static string GenerateEventTriggerUrlInternal(this ExpressionExecutionContext context, string token) { - var options = context.GetRequiredService>().Value; - var url = $"{options.ApiRoutePrefix}/events/trigger?t={token}"; + var options = context.GetRequiredService>().Value; + var url = $"{options.RoutePrefix}/events/trigger?t={token}"; var absoluteUrlProvider = context.GetRequiredService(); return absoluteUrlProvider.ToAbsoluteUrl(url).ToString(); } - private static string GenerateEventTriggerTokenInternal(this ExpressionExecutionContext context, string eventName, TimeSpan? lifetime = default, DateTimeOffset? expiresAt = default) + private static string GenerateEventTriggerTokenInternal(this ExpressionExecutionContext context, string eventName, TimeSpan? lifetime = null, DateTimeOffset? expiresAt = null) { var workflowInstanceId = context.GetWorkflowExecutionContext().Id; var payload = new EventTokenPayload(eventName, workflowInstanceId); diff --git a/src/modules/Elsa.Http/Options/HttpActivityOptions.cs b/src/modules/Elsa.Http/Options/HttpActivityOptions.cs index 2080c0068..1d0c75f49 100644 --- a/src/modules/Elsa.Http/Options/HttpActivityOptions.cs +++ b/src/modules/Elsa.Http/Options/HttpActivityOptions.cs @@ -15,7 +15,7 @@ public class HttpActivityOptions /// /// The base URL of the server. This should be set to the same value at which the Elsa Server is publicly available. It will be used when generating absolute URLs need to be generated by activities such as SendEmail. /// - public Uri BaseUrl { get; set; } = default!; + public Uri BaseUrl { get; set; } = null!; /// /// The prefix used for API routes. diff --git a/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj b/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj index 94828c0a3..9a86f64aa 100644 --- a/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj +++ b/src/modules/Elsa.Workflows.Api/Elsa.Workflows.Api.csproj @@ -9,8 +9,8 @@ - + diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/Bookmarks/Resume/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/Bookmarks/Resume/Endpoint.cs new file mode 100644 index 000000000..95ddfabd8 --- /dev/null +++ b/src/modules/Elsa.Workflows.Api/Endpoints/Bookmarks/Resume/Endpoint.cs @@ -0,0 +1,89 @@ +using Elsa.Abstractions; +using Elsa.SasTokens.Contracts; +using Elsa.Workflows.Runtime; +using FastEndpoints; +using JetBrains.Annotations; +using Microsoft.AspNetCore.Http; + +namespace Elsa.Workflows.Api.Endpoints.Bookmarks.Resume; + +/// +/// Resumes a bookmarked workflow instance with the bookmark ID specified in the provided SAS token. +/// +[PublicAPI] +internal class Resume(ITokenService tokenService, IBookmarkQueue bookmarkQueue, IPayloadSerializer serializer) : ElsaEndpoint +{ + /// + public override void Configure() + { + Routes("/bookmarks/resume"); + Verbs(Http.GET, Http.POST); + AllowAnonymous(); + } + + /// + public override async Task HandleAsync(Request request, CancellationToken cancellationToken) + { + var token = Query("t")!; + + if (!tokenService.TryDecryptToken(token, out var payload)) + AddError("Invalid token."); + + var input = HttpContext.Request.Method == HttpMethods.Post ? request.Input : GetInputFromQueryString(); + + if (ValidationFailed) + { + await SendErrorsAsync(cancellation: cancellationToken); + return; + } + + await ResumeBookmarkedWorkflowAsync(payload, input, cancellationToken); + + if (!HttpContext.Response.HasStarted) + await SendOkAsync(cancellationToken); + } + + private IDictionary? GetInputFromQueryString() + { + var inputJson = Query("in", false); + if (string.IsNullOrWhiteSpace(inputJson)) + return null; + + try + { + return serializer.Deserialize>(inputJson); + } + catch + { + AddError("Invalid input format. Expected a valid JSON string."); + return null; + } + } + + private async Task ResumeBookmarkedWorkflowAsync(BookmarkTokenPayload tokenPayload, IDictionary? input, CancellationToken cancellationToken) + { + var bookmarkId = tokenPayload.BookmarkId; + var workflowInstanceId = tokenPayload.WorkflowInstanceId; + var item = new NewBookmarkQueueItem + { + BookmarkId = bookmarkId, + WorkflowInstanceId = workflowInstanceId, + Options = new() + { + Input = input + } + }; + await bookmarkQueue.EnqueueAsync(item, cancellationToken); + } +} + +/// +/// The request model for the Resume endpoint. +/// +internal class Request +{ + /// + /// The input to provide to the workflow when resuming. + /// + public IDictionary? Input { get; set; } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Api/Endpoints/Events/TriggerPublic/Endpoint.cs b/src/modules/Elsa.Workflows.Api/Endpoints/Events/TriggerPublic/Endpoint.cs index 6685e7c97..0f3e50a40 100644 --- a/src/modules/Elsa.Workflows.Api/Endpoints/Events/TriggerPublic/Endpoint.cs +++ b/src/modules/Elsa.Workflows.Api/Endpoints/Events/TriggerPublic/Endpoint.cs @@ -1,5 +1,4 @@ using Elsa.Abstractions; -using Elsa.Http; using Elsa.SasTokens.Contracts; using Elsa.Workflows.Runtime; using JetBrains.Annotations; diff --git a/src/modules/Elsa.Workflows.Api/Features/WorkflowsApiFeature.cs b/src/modules/Elsa.Workflows.Api/Features/WorkflowsApiFeature.cs index 6fb8a3003..721a36827 100644 --- a/src/modules/Elsa.Workflows.Api/Features/WorkflowsApiFeature.cs +++ b/src/modules/Elsa.Workflows.Api/Features/WorkflowsApiFeature.cs @@ -2,7 +2,6 @@ using Elsa.Extensions; using Elsa.Features.Abstractions; using Elsa.Features.Attributes; using Elsa.Features.Services; -using Elsa.Http.Features; using Elsa.SasTokens.Features; using Elsa.Workflows.Api.Constants; using Elsa.Workflows.Api.Requirements; @@ -21,15 +20,9 @@ namespace Elsa.Workflows.Api.Features; [DependsOn(typeof(WorkflowInstancesFeature))] [DependsOn(typeof(WorkflowManagementFeature))] [DependsOn(typeof(WorkflowRuntimeFeature))] -[DependsOn(typeof(HttpFeature))] [DependsOn(typeof(SasTokensFeature))] -public class WorkflowsApiFeature : FeatureBase +public class WorkflowsApiFeature(IModule module) : FeatureBase(module) { - /// - public WorkflowsApiFeature(IModule module) : base(module) - { - } - /// public override void Configure() { @@ -43,7 +36,6 @@ public class WorkflowsApiFeature : FeatureBase Module.AddFastEndpointsFromModule(); Services.AddScoped(); - Services.AddScoped(); Services.Configure(options => { diff --git a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Outcomes.cs b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Outcomes.cs index 9cfa3cfde..c0154ab4b 100644 --- a/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Outcomes.cs +++ b/src/modules/Elsa.Workflows.Core/Activities/Flowchart/Models/Outcomes.cs @@ -4,8 +4,8 @@ namespace Elsa.Workflows.Activities.Flowchart.Models; /// Represents a list of outcomes that can be send when completing an activity. This information is used by . /// /// A list of outcome names. -public record Outcomes(params string[] Names) -{ - public static readonly Outcomes Default = new([null!, "Done"]); - public static readonly Outcomes Empty = new(); -} +public record Outcomes(params string[] Names) +{ + public static readonly Outcomes Default = new(null!, "Done"); + public static readonly Outcomes Empty = new(); +} diff --git a/src/modules/Elsa.Workflows.Runtime/Models/BookmarkTokenPayload.cs b/src/modules/Elsa.Workflows.Runtime/Models/BookmarkTokenPayload.cs new file mode 100644 index 000000000..038959449 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Models/BookmarkTokenPayload.cs @@ -0,0 +1,6 @@ +namespace Elsa.Workflows.Runtime; + +/// +/// Represents the payload for a bookmark token, including the bookmark identifier and the associated workflow instance identifier. +/// +public record BookmarkTokenPayload(string BookmarkId, string WorkflowInstanceId); \ No newline at end of file diff --git a/src/modules/Elsa.Http/Models/EventTokenPayload.cs b/src/modules/Elsa.Workflows.Runtime/Models/EventTokenPayload.cs similarity index 90% rename from src/modules/Elsa.Http/Models/EventTokenPayload.cs rename to src/modules/Elsa.Workflows.Runtime/Models/EventTokenPayload.cs index 22d5fd214..e6981adde 100644 --- a/src/modules/Elsa.Http/Models/EventTokenPayload.cs +++ b/src/modules/Elsa.Workflows.Runtime/Models/EventTokenPayload.cs @@ -1,4 +1,4 @@ -namespace Elsa.Http; +namespace Elsa.Workflows.Runtime; /// /// Represents the payload of an event, serialized as a secured token. From cbcd9c3dc9fb1502ddf6632d9ea49956abb00514 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 30 Jul 2025 09:13:18 +0200 Subject: [PATCH 09/10] Refactor `BookmarkExpressionExecutionContextExtensions` to `BookmarkExecutionContextExtensions` and add new helper methods for generating bookmark trigger URLs. --- ...xtensions.cs => BookmarkExecutionContextExtensions.cs} | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) rename src/modules/Elsa.Http/Extensions/{BookmarkExpressionExecutionContextExtensions.cs => BookmarkExecutionContextExtensions.cs} (83%) diff --git a/src/modules/Elsa.Http/Extensions/BookmarkExpressionExecutionContextExtensions.cs b/src/modules/Elsa.Http/Extensions/BookmarkExecutionContextExtensions.cs similarity index 83% rename from src/modules/Elsa.Http/Extensions/BookmarkExpressionExecutionContextExtensions.cs rename to src/modules/Elsa.Http/Extensions/BookmarkExecutionContextExtensions.cs index 1cdcdbd92..c509b1f84 100644 --- a/src/modules/Elsa.Http/Extensions/BookmarkExpressionExecutionContextExtensions.cs +++ b/src/modules/Elsa.Http/Extensions/BookmarkExecutionContextExtensions.cs @@ -1,7 +1,7 @@ using Elsa.Expressions.Models; using Elsa.Http; -using Elsa.Http.Options; using Elsa.SasTokens.Contracts; +using Elsa.Workflows; using Elsa.Workflows.Api; using Elsa.Workflows.Runtime; using Microsoft.Extensions.Options; @@ -12,8 +12,12 @@ namespace Elsa.Extensions; /// /// Provides extension methods for working with and generating bookmark trigger URLs. /// -public static class BookmarkExpressionExecutionContextExtensions +public static class BookmarkExecutionContextExtensions { + public static string GenerateBookmarkTriggerUrl(this ActivityExecutionContext context, string bookmarkId, TimeSpan lifetime) => context.ExpressionExecutionContext.GenerateBookmarkTriggerUrl(bookmarkId, lifetime); + public static string GenerateBookmarkTriggerUrl(this ActivityExecutionContext context, string bookmarkId, DateTimeOffset expiresAt) => context.ExpressionExecutionContext.GenerateBookmarkTriggerUrl(bookmarkId, expiresAt); + public static string GenerateBookmarkTriggerUrl(this ActivityExecutionContext context, string bookmarkId) => context.ExpressionExecutionContext.GenerateBookmarkTriggerUrl(bookmarkId); + /// /// Generates a URL that can be used to resume a bookmarked workflow. /// From 207356cf5b6d4fded5ae13c4c2d79df44b7634d6 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 30 Jul 2025 13:11:02 +0200 Subject: [PATCH 10/10] Add `Bookmarks` property to workflow state mapping in `LocalWorkflowClient` --- .../Elsa.Workflows.Runtime/Services/LocalWorkflowClient.cs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/src/modules/Elsa.Workflows.Runtime/Services/LocalWorkflowClient.cs b/src/modules/Elsa.Workflows.Runtime/Services/LocalWorkflowClient.cs index 1146d7c88..f7446ee99 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/LocalWorkflowClient.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/LocalWorkflowClient.cs @@ -140,7 +140,8 @@ public class LocalWorkflowClient( WorkflowInstanceId = WorkflowInstanceId, Status = workflowState.Status, SubStatus = workflowState.SubStatus, - Incidents = workflowState.Incidents + Incidents = workflowState.Incidents, + Bookmarks = workflowState.Bookmarks }; }