diff --git a/doc/wiki/bpmn-workflows.md b/doc/wiki/bpmn-workflows.md index edfaa4a24..325c54ec9 100644 --- a/doc/wiki/bpmn-workflows.md +++ b/doc/wiki/bpmn-workflows.md @@ -147,11 +147,19 @@ whole-definition import and resets them from the document, same as it always has - **Missing, or the wildcard `*`** — `428 Precondition Required`. Neither says which revision the caller is replacing (`*` matches whatever is stored), so the endpoint refuses rather than overwrite blindly, before doing any import work or persisting anything. -- **Anything other than exactly the definition's current `ETag`** — `412 Precondition Failed`, checked before any - import work and before anything is persisted. The comparison is exact: a weak (`W/`) tag or a list of tags never - matches. The definition was written since the caller last read it; `GET` the document again and reapply the edit. +- **Anything other than exactly the definition's current `ETag`** — `412 Precondition Failed`. The comparison is + exact: a weak (`W/`) tag or a list of tags never matches. Checked once on arrival so a stale client is refused + before import work runs, and again in the same compare-and-swap that reads the definition's non-BPMN metadata and + saves — a write that lands in that window is also `412`, not a silent overwrite. The definition was written since + the caller last read it; `GET` the document again and reapply the edit. - **Exactly the current `ETag`** — the request proceeds exactly as before, and the response carries the `ETag` of - the draft as this `PUT` stored it. + the draft as this `PUT` stored it. The draft is prepared first and `WorkflowDefinitionDraftSaving` is + dispatched before persist so a rejecting handler fails the request before anything is written; handlers + may mutate that draft, and the compare-and-swap writes that same instance. If name, description or the + ETag-covered graph moved in that window the swap is `412` rather than silently overwriting the other + write. A lost race is still `412`. The swap is implemented on the in-memory and EF Core definition stores + (`IWorkflowDefinitionStore.TryUpdateLatestAsync`). Mongo, Dapper and Event Sourcing providers need the same + method before this endpoint is concurrency-safe on those stores. ### The document endpoints and the JSON payload format diff --git a/src/modules/Elsa.Bpmn.Interchange/BpmnErrorCodes.cs b/src/modules/Elsa.Bpmn.Interchange/BpmnErrorCodes.cs index ae25cc2e1..88b5efb18 100644 --- a/src/modules/Elsa.Bpmn.Interchange/BpmnErrorCodes.cs +++ b/src/modules/Elsa.Bpmn.Interchange/BpmnErrorCodes.cs @@ -74,7 +74,8 @@ public static class BpmnErrorCodes /// /// The document PUT refuses an If-Match header that does not match the workflow definition's - /// current ETag: the definition was written since the caller last read it. + /// current ETag — including when the header matched on arrival but another writer saved before this PUT's + /// compare-and-swap completed: the definition was written since the caller last read it. /// public const string DocumentPreconditionFailed = "bpmn.document.precondition-failed"; } diff --git a/src/modules/Elsa.Bpmn.Interchange/Endpoints/Bpmn/BpmnImportErrorResponses.cs b/src/modules/Elsa.Bpmn.Interchange/Endpoints/Bpmn/BpmnImportErrorResponses.cs index 094e259a1..bb0abf914 100644 --- a/src/modules/Elsa.Bpmn.Interchange/Endpoints/Bpmn/BpmnImportErrorResponses.cs +++ b/src/modules/Elsa.Bpmn.Interchange/Endpoints/Bpmn/BpmnImportErrorResponses.cs @@ -39,6 +39,11 @@ internal static class BpmnImportErrorResponses await BpmnErrorResponse.SendAsync(httpResponse, NotFoundResponseFor(exception), cancellationToken); return null; } + catch (BpmnDocumentPreconditionFailedException exception) + { + await BpmnErrorResponse.SendAsync(httpResponse, PreconditionFailedResponseFor(exception), cancellationToken); + return null; + } catch (BpmnInterchangeException exception) { addError(exception.Message); @@ -84,6 +89,13 @@ internal static class BpmnImportErrorResponses internal static BpmnErrorResponse NotFoundResponseFor(BpmnDefinitionNotFoundException exception) => BpmnErrorResponse.Create(exception.Message, BpmnErrorCodes.DocumentNotFound, StatusCodes.Status404NotFound); + /// + /// The response for a lost compare-and-swap — + /// the same code and status the PUT already sends when If-Match is stale on arrival. + /// + internal static BpmnErrorResponse PreconditionFailedResponseFor(BpmnDocumentPreconditionFailedException exception) => + BpmnErrorResponse.Create(exception.Message, BpmnErrorCodes.DocumentPreconditionFailed, StatusCodes.Status412PreconditionFailed); + /// The response for . internal static BpmnErrorResponse BindingInvalidResponseFor(BpmnBindingException exception) => BpmnErrorResponse.Create(exception.Message, BpmnErrorCodes.ImportBindingInvalid, StatusCodes.Status422UnprocessableEntity); diff --git a/src/modules/Elsa.Bpmn.Interchange/Endpoints/Bpmn/Document/Put/Endpoint.cs b/src/modules/Elsa.Bpmn.Interchange/Endpoints/Bpmn/Document/Put/Endpoint.cs index 204e1df20..15eabfc1c 100644 --- a/src/modules/Elsa.Bpmn.Interchange/Endpoints/Bpmn/Document/Put/Endpoint.cs +++ b/src/modules/Elsa.Bpmn.Interchange/Endpoints/Bpmn/Document/Put/Endpoint.cs @@ -54,7 +54,9 @@ internal sealed class Put(IWorkflowDefinitionStore store, BpmnInterchangeDocumen // ETag its GET returned as If-Match. A missing header cannot express "I know what I'm overwriting" at all, and // neither can "*", which matches whatever is stored — so both are refused as 428 rather than honoured. Anything // else must be exactly the current strong ETag (a weak W/ tag or a list never is), or it proves the client's copy - // is no longer current. All of this is checked before any import work runs or anything is persisted. + // is no longer current. The header is checked here so a stale client is refused before import work runs, and + // again inside ImportDocumentAsync's compare-and-swap so a write that lands in that window is 412, not a silent + // overwrite. Metadata carried forward is read in that same swap, not from this lookup. var ifMatch = HttpContext.Request.Headers.IfMatch.ToString().Trim(); if (ifMatch is "" or "*") @@ -114,7 +116,7 @@ internal sealed class Put(IWorkflowDefinitionStore store, BpmnInterchangeDocumen : null; var result = await BpmnImportErrorResponses.RunAsync( - () => documentService.ImportDocumentAsync(document, definitionId, processId, cancellationToken), + () => documentService.ImportDocumentAsync(document, definitionId, processId, cancellationToken, ifMatch), HttpContext.Response, message => AddError(message), Send.ErrorsAsync, diff --git a/src/modules/Elsa.Bpmn.Interchange/Exceptions/BpmnDocumentPreconditionFailedException.cs b/src/modules/Elsa.Bpmn.Interchange/Exceptions/BpmnDocumentPreconditionFailedException.cs new file mode 100644 index 000000000..0cd37a361 --- /dev/null +++ b/src/modules/Elsa.Bpmn.Interchange/Exceptions/BpmnDocumentPreconditionFailedException.cs @@ -0,0 +1,12 @@ +namespace Elsa.Bpmn.Interchange.Exceptions; + +/// +/// Thrown when loses the +/// compare-and-swap against the caller's If-Match: another writer saved the definition after +/// the PUT's early ETag check (or between that check and this save). +/// +/// +/// Mapped to the same 412 / bpmn.document.precondition-failed response the PUT already +/// sends when If-Match is stale on arrival, so a lost race and a stale header are one refusal. +/// +public class BpmnDocumentPreconditionFailedException(string message) : Exception(message); diff --git a/src/modules/Elsa.Bpmn.Interchange/Services/BpmnInterchangeDocumentService.cs b/src/modules/Elsa.Bpmn.Interchange/Services/BpmnInterchangeDocumentService.cs index 73c05a9ee..b2486a877 100644 --- a/src/modules/Elsa.Bpmn.Interchange/Services/BpmnInterchangeDocumentService.cs +++ b/src/modules/Elsa.Bpmn.Interchange/Services/BpmnInterchangeDocumentService.cs @@ -5,13 +5,19 @@ using Bpmn.Semantics; using Elsa.Bpmn.Activities; using Elsa.Bpmn.Hosting; using Elsa.Bpmn.Interchange.Binding; +using Elsa.Bpmn.Interchange.Endpoints.Bpmn.Document; using Elsa.Bpmn.Interchange.Exceptions; +using Elsa.Common; using Elsa.Common.Models; using Elsa.Extensions; +using Elsa.Mediator.Contracts; +using Elsa.Workflows; using Elsa.Workflows.Management; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Mappers; +using Elsa.Workflows.Management.Materializers; using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Management.Notifications; using Elsa.Workflows.Models; namespace Elsa.Bpmn.Interchange.Services; @@ -91,7 +97,11 @@ public sealed class BpmnInterchangeDocumentService( BpmnWorkBinder binder, IWorkflowDefinitionImporter importer, IWorkflowDefinitionStore store, - VariableDefinitionMapper variableDefinitionMapper) + VariableDefinitionMapper variableDefinitionMapper, + IActivitySerializer activitySerializer, + IIdentityGenerator identityGenerator, + ISystemClock systemClock, + IMediator mediator) { /// The workflow definition custom property the original BPMN XML is carried under, for . public const string SourceXmlCustomPropertyKey = "Bpmn:SourceXml"; @@ -190,7 +200,7 @@ public sealed class BpmnInterchangeDocumentService( /// The document needs a host capability this deployment does not declare. /// A work binding cannot be turned into an Elsa activity. public Task ImportAsync(string xml, string? definitionId, string? name, string? processId, CancellationToken cancellationToken) => - ImportCoreAsync(xml, definitionId, name, processId, preserveMetadataFrom: null, cancellationToken); + ImportCoreAsync(xml, definitionId, name, processId, preserveMetadataFrom: null, expectedETag: null, compareAndSwap: false, cancellationToken); /// /// The shared import logic behind both the public and : @@ -208,10 +218,17 @@ public sealed class BpmnInterchangeDocumentService( /// inputs, outputs, outcomes, options, tool version, read-only flag and custom properties are carried onto the /// imported definition as-is, and only the bound activity graph and the , /// , and - /// custom properties this method owns change. This is what - /// passes so the document PUT - /// edits the BPMN document without silently resetting the rest of the definition; left null for a - /// whole-definition import, where the model built from the document alone is the intended contract. + /// custom properties this method owns change. Left null + /// for a whole-definition import and for a document edit, where metadata is copied onto the prepared draft + /// and the compare-and-swap refuses if that snapshot has moved. + /// + /// + /// When is true, the If-Match value the save must still equal, or + /// null to accept any ETag while still refusing if the prepared draft's snapshot has moved. + /// + /// + /// True for a document edit: persist through so + /// the precondition and the metadata-preserving save are one step. /// /// The cancellation token. /// The document cannot be read, or declares more than one process and does not pick one. @@ -223,6 +240,8 @@ public sealed class BpmnInterchangeDocumentService( string? name, string? processId, WorkflowDefinition? preserveMetadataFrom, + string? expectedETag, + bool compareAndSwap, CancellationToken cancellationToken) { var result = reader.Read(xml, new BpmnImportOptions { ProcessId = processId }); @@ -265,6 +284,11 @@ public sealed class BpmnInterchangeDocumentService( model.CustomProperties = new Dictionary(preserveMetadataFrom.CustomProperties); } + if (compareAndSwap) + { + return await PersistDocumentEditAsync(xml, definitionId!, process, rootDefinition, result.Analysis, expectedETag, cancellationToken); + } + var importResult = await importer.ImportAsync(new SaveWorkflowDefinitionRequest { Model = model, Publish = false }, cancellationToken); // The definition's final Version is only known once the importer/publisher has assigned and persisted it — @@ -285,6 +309,125 @@ public sealed class BpmnInterchangeDocumentService( return new BpmnDocumentImportResult(importResult, result.Analysis); } + /// + /// The document-edit persist: prepare the draft once, dispatch + /// so a rejecting handler fails the request before + /// anything is written, then compare-and-swap that same draft. The swap accepts the write only + /// when If-Match and the loaded snapshot (id, version, graph, name, description, IsLatest) are + /// still the row the draft was built from — so a metadata-only save in the window is 412, not a + /// silent overwrite. A lost race is . + /// + private async Task PersistDocumentEditAsync( + string xml, + string definitionId, + BpmnProcess process, + BpmnProcessDefinition rootDefinition, + BpmnImportAnalysis analysis, + string? expectedETag, + CancellationToken cancellationToken) + { + var filter = WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Latest).ToFilter(); + var current = await store.FindAsync(filter, cancellationToken); + + if (current is null) + { + throw new BpmnDefinitionNotFoundException( + $"Workflow definition '{definitionId}' does not exist, so its BPMN document cannot be edited."); + } + + if (expectedETag is not null && !string.Equals(BpmnDocumentETag.From(current), expectedETag, StringComparison.Ordinal)) + { + throw new BpmnDocumentPreconditionFailedException( + "The workflow definition has been written since the ETag in If-Match was issued. GET the document again, reapply the edit, and PUT it with the new ETag."); + } + + var expectedId = current.Id; + var expectedVersion = current.Version; + var expectedName = current.Name; + var expectedDescription = current.Description; + var expectedStringData = current.StringData; + + var draft = ApplyDocumentEdit(current, process, xml, rootDefinition); + await mediator.SendAsync(new WorkflowDefinitionDraftSaving(draft), cancellationToken); + + var result = await store.TryUpdateLatestAsync( + filter, + loaded => loaded.IsLatest + && loaded.Id == expectedId + && loaded.Version == expectedVersion + && loaded.StringData == expectedStringData + && loaded.Name == expectedName + && loaded.Description == expectedDescription + && (expectedETag is null || string.Equals(BpmnDocumentETag.From(loaded), expectedETag, StringComparison.Ordinal)), + _ => draft, + cancellationToken); + + if (result.Outcome == WorkflowDefinitionUpdateOutcome.NotFound) + { + throw new BpmnDefinitionNotFoundException( + $"Workflow definition '{definitionId}' does not exist, so its BPMN document cannot be edited."); + } + + if (result.Outcome == WorkflowDefinitionUpdateOutcome.Conflict) + { + throw new BpmnDocumentPreconditionFailedException( + "The workflow definition has been written since the ETag in If-Match was issued. GET the document again, reapply the edit, and PUT it with the new ETag."); + } + + await mediator.SendAsync(new WorkflowDefinitionDraftSaved(result.Definition!), cancellationToken); + return new BpmnDocumentImportResult(new ImportWorkflowResult(true, result.Definition!, []), analysis); + } + + /// + /// Builds the draft the compare-and-swap will save from : + /// metadata comes from that row; only the bound graph and the BPMN source properties this service owns change. + /// + private WorkflowDefinition ApplyDocumentEdit( + WorkflowDefinition current, + BpmnProcess process, + string xml, + BpmnProcessDefinition rootDefinition) + { + var draft = current.IsPublished ? NewDraftFrom(current) : current.ShallowClone(); + var stringData = activitySerializer.Serialize(process); + + draft.StringData = stringData; + draft.MaterializerName = JsonWorkflowMaterializer.MaterializerName; + draft.Name = current.Name; + draft.Description = current.Description; + draft.Variables = current.Variables; + draft.Inputs = current.Inputs; + draft.Outputs = current.Outputs; + draft.Outcomes = current.Outcomes; + draft.Options = current.Options; + draft.ToolVersion = current.ToolVersion; + draft.IsReadonly = current.IsReadonly; + draft.CustomProperties = new Dictionary(current.CustomProperties) + { + [SourceXmlCustomPropertyKey] = xml, + [SourceVersionCustomPropertyKey] = draft.Version, + [SourceProcessIdCustomPropertyKey] = rootDefinition.ProcessId, + [SourceGraphHashCustomPropertyKey] = BpmnContentHash.OfGraph(stringData) + }; + + return draft; + } + + /// + /// The unpublished draft would + /// return for a published latest row, without a store read: new id, next version, not published. + /// + private WorkflowDefinition NewDraftFrom(WorkflowDefinition published) + { + var draft = published.ShallowClone(); + draft.Id = identityGenerator.GenerateId(); + draft.Version = published.Version + 1; + draft.CreatedAt = systemClock.UtcNow; + draft.IsLatest = true; + draft.IsPublished = false; + return draft; + } + /// /// Writes the document a workflow definition was imported from back out as BPMN 2.0 XML, through the same /// reader-then-writer path used, so retained extension elements, foreign attributes and @@ -349,9 +492,9 @@ public sealed class BpmnInterchangeDocumentService( /// Unlike a whole-definition import, this edits the BPMN document of an existing definition: the caller /// is changing a binding, not replacing the definition. So 's current metadata — /// name, description, variables, inputs, outputs, outcomes, options, tool version, read-only flag and custom - /// properties other than the ones this service owns — is carried onto the result unchanged; see the shared - /// import logic's preserveMetadataFrom parameter, which this passes the existing definition to. Only the activity graph - /// and the // + /// properties other than the ones this service owns — is copied onto the prepared draft and the compare-and-swap + /// refuses if that snapshot has moved, rather than being taken from the lookup that restores nested scopes. Only + /// the activity graph and the // /// / custom properties move. /// /// Nested scopes come from the stored source, not from , which cannot carry them (see @@ -375,6 +518,12 @@ public sealed class BpmnInterchangeDocumentService( /// See for where a caller re-importing an existing definition /// finds the value that was used the first time. /// + /// + /// The If-Match value the document PUT already checked. When set, the compare-and-swap that persists + /// the prepared draft also refuses with if the stored + /// content that ETag covers — or the name/description snapshot the draft was built from — has moved. When + /// omitted, any ETag is accepted but a moved metadata snapshot is still refused. + /// /// The cancellation token. /// /// The workflow definition to edit no longer exists — e.g. it was deleted between the PUT endpoint's own @@ -382,13 +531,21 @@ public sealed class BpmnInterchangeDocumentService( /// whole-definition import path, which would silently create a definition under /// with reset metadata instead of reporting that this PUT's target disappeared. /// + /// + /// no longer matches the stored definition — another writer saved first. + /// /// /// The document declares more than one process and does not pick one, or it declares a /// subprocess element that has a stored body but no bindingRef to write that body back under. /// /// The document needs a host capability this deployment does not declare. /// A work binding cannot be turned into an Elsa activity. - public async Task ImportDocumentAsync(BpmnDefinitions document, string definitionId, string? processId, CancellationToken cancellationToken) + public async Task ImportDocumentAsync( + BpmnDefinitions document, + string definitionId, + string? processId, + CancellationToken cancellationToken, + string? expectedETag = null) { var filter = WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Latest).ToFilter(); var existingDefinition = await store.FindAsync(filter, cancellationToken); @@ -414,7 +571,7 @@ public sealed class BpmnInterchangeDocumentService( EnsureElementIdsUnique(document.Processes, bindingsToCarryAcross); var xml = writer.Write(document, bindingsToCarryAcross); - return await ImportCoreAsync(xml, definitionId, name: null, processId, preserveMetadataFrom: existingDefinition, cancellationToken); + return await ImportCoreAsync(xml, definitionId, name: null, processId, preserveMetadataFrom: null, expectedETag, compareAndSwap: true, cancellationToken); } /// diff --git a/src/modules/Elsa.Common/Services/MemoryStore.cs b/src/modules/Elsa.Common/Services/MemoryStore.cs index cbbb7f84f..a47590809 100644 --- a/src/modules/Elsa.Common/Services/MemoryStore.cs +++ b/src/modules/Elsa.Common/Services/MemoryStore.cs @@ -10,6 +10,13 @@ namespace Elsa.Common.Services; public class MemoryStore { private IDictionary Entities { get; set; } = new ConcurrentDictionary(); + + /// + /// One lock for every scoped wrapper that shares this store. Compare-and-swap + /// load/match/write and the other Save/Delete paths take it so two scopes cannot + /// race on the same in-memory set. + /// + public object Sync { get; } = new(); /// /// Gets a queryable of all entities. diff --git a/src/modules/Elsa.Persistence.EFCore/Modules/Management/WorkflowDefinitionStore.cs b/src/modules/Elsa.Persistence.EFCore/Modules/Management/WorkflowDefinitionStore.cs index e805caeaa..e7ae405ed 100644 --- a/src/modules/Elsa.Persistence.EFCore/Modules/Management/WorkflowDefinitionStore.cs +++ b/src/modules/Elsa.Persistence.EFCore/Modules/Management/WorkflowDefinitionStore.cs @@ -1,4 +1,5 @@ using System.Diagnostics.CodeAnalysis; +using System.Linq.Expressions; using System.Text.Json.Serialization; using Elsa.Common.Entities; using Elsa.Common.Models; @@ -128,6 +129,97 @@ public class EFCoreWorkflowDefinitionStore(EntityStore + public async Task TryUpdateLatestAsync( + WorkflowDefinitionFilter filter, + Func matchesExpected, + Func update, + CancellationToken cancellationToken = default) + { + await using var dbContext = await store.CreateDbContextAsync(cancellationToken); + await using var transaction = await dbContext.Database.BeginTransactionAsync(cancellationToken); + + var queryable = Filter(dbContext.WorkflowDefinitions.AsNoTracking(), filter); + + if (filter.TenantAgnostic) + queryable = queryable.IgnoreQueryFilters(); + + var current = await queryable.OrderBy(x => x.CreatedAt).FirstOrDefaultAsync(cancellationToken); + + if (current == null) + return WorkflowDefinitionUpdateResult.NotFound(); + + await OnLoadAsync(dbContext, current, cancellationToken); + + if (!matchesExpected(current)) + return WorkflowDefinitionUpdateResult.Conflict(); + + var expectedId = current.Id; + var expectedVersion = current.Version; + var expectedStringData = current.StringData; + var expectedName = current.Name; + var expectedDescription = current.Description; + var expectedData = (string?)dbContext.Entry(current).Property("Data").CurrentValue; + + var next = update(current); + var nextData = SerializeState(next); + var nextUsableAsActivity = next.Options.UsableAsActivity; + + if (next.Id != expectedId) + { + var unmarked = await dbContext.WorkflowDefinitions + .Where(MatchesLoadedSnapshot(expectedId, expectedVersion, expectedStringData, expectedName, expectedDescription, expectedData)) + .ExecuteUpdateAsync(setters => setters.SetProperty(x => x.IsLatest, false), cancellationToken); + + if (unmarked == 0) + return WorkflowDefinitionUpdateResult.Conflict(); + + dbContext.WorkflowDefinitions.Add(next); + dbContext.Entry(next).Property("Data").CurrentValue = nextData; + dbContext.Entry(next).Property("UsableAsActivity").CurrentValue = nextUsableAsActivity; + await dbContext.SaveChangesAsync(cancellationToken); + await transaction.CommitAsync(cancellationToken); + return WorkflowDefinitionUpdateResult.Updated(next); + } + + var updated = await dbContext.WorkflowDefinitions + .Where(MatchesLoadedSnapshot(expectedId, expectedVersion, expectedStringData, expectedName, expectedDescription, expectedData)) + .ExecuteUpdateAsync( + setters => setters + .SetProperty(x => x.StringData, next.StringData) + .SetProperty(x => EF.Property(x, "Data"), nextData) + .SetProperty(x => EF.Property(x, "UsableAsActivity"), nextUsableAsActivity), + cancellationToken); + + if (updated == 0) + return WorkflowDefinitionUpdateResult.Conflict(); + + await transaction.CommitAsync(cancellationToken); + return WorkflowDefinitionUpdateResult.Updated(next); + } + + /// + /// The loaded snapshot will only overwrite: identity, the ETag-covered + /// graph, the serialized Data blob (source XML, variables, options) and the name/description columns + /// a metadata-only save changes. The row must still be IsLatest so a published→draft loser + /// is Conflict rather than a unique-key failure on (DefinitionId, Version). Zero rows means + /// another writer got there first. + /// + private static Expression> MatchesLoadedSnapshot( + string expectedId, + int expectedVersion, + string? expectedStringData, + string? expectedName, + string? expectedDescription, + string? expectedData) => + x => x.Id == expectedId + && x.Version == expectedVersion + && x.IsLatest + && x.StringData == expectedStringData + && x.Name == expectedName + && x.Description == expectedDescription + && EF.Property(x, "Data") == expectedData; + /// public async Task DeleteAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { @@ -163,14 +255,19 @@ public class EFCoreWorkflowDefinitionStore(EntityStore Task SaveAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default); + /// + /// Compare-and-swap the latest definition matching . + /// + /// + /// + /// Loads the matching row, and only if is true applies + /// to that just-loaded row and saves the result. The load, the match, + /// the update and the save are one critical section (memory) or one conditional write (EF: + /// ExecuteUpdate against the loaded snapshot). A lost race returns + /// — it does not wait. + /// + /// + /// sees the definition as stored at the moment of the swap, so + /// metadata copied from it (name, variables, options, custom properties) is current — not a + /// snapshot taken by the caller before this call. Treat that argument as read-only and return + /// a new or cloned definition; mutating it in place can tear a shared in-memory instance. + /// + /// + /// When returns a definition with a different Id (a new draft + /// of a published version), the previously latest row is unmarked in the same step. + /// + /// + /// Persistence providers outside this repository (Mongo, Dapper, Event Sourcing) must implement + /// this the same way before the BPMN document PUT is concurrency-safe on those stores. + /// Until they do, that endpoint's atomic precondition is not available there. + /// + /// + /// Typically the latest version of one definition id. + /// + /// True when the loaded row is still the snapshot the caller is allowed to overwrite + /// (for the document PUT, the current ETag still equals If-Match). + /// + /// Builds the definition to save from the just-loaded row. + /// The cancellation token. + Task TryUpdateLatestAsync( + WorkflowDefinitionFilter filter, + Func matchesExpected, + Func update, + CancellationToken cancellationToken = default); + /// /// Adds the specified set of objects to te persistence store. /// diff --git a/src/modules/Elsa.Workflows.Management/Models/WorkflowDefinitionUpdateResult.cs b/src/modules/Elsa.Workflows.Management/Models/WorkflowDefinitionUpdateResult.cs new file mode 100644 index 000000000..3b98d4437 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Models/WorkflowDefinitionUpdateResult.cs @@ -0,0 +1,38 @@ +using Elsa.Workflows.Management.Entities; + +namespace Elsa.Workflows.Management.Models; + +/// +/// The outcome of . +/// +public enum WorkflowDefinitionUpdateOutcome +{ + /// The update ran and the definition was saved. + Updated, + + /// No definition matched the filter. + NotFound, + + /// + /// The loaded definition no longer matched what the caller expected — another writer saved first. + /// + Conflict +} + +/// +/// The result of a compare-and-swap update of the latest workflow definition. +/// +/// Whether the swap saved, found nothing, or lost the race. +/// The saved definition when is . +public sealed record WorkflowDefinitionUpdateResult(WorkflowDefinitionUpdateOutcome Outcome, WorkflowDefinition? Definition = null) +{ + /// The update ran and is what was saved. + public static WorkflowDefinitionUpdateResult Updated(WorkflowDefinition definition) => + new(WorkflowDefinitionUpdateOutcome.Updated, definition); + + /// No definition matched the filter. + public static WorkflowDefinitionUpdateResult NotFound() => new(WorkflowDefinitionUpdateOutcome.NotFound); + + /// Another writer saved first; the caller's expected snapshot is no longer current. + public static WorkflowDefinitionUpdateResult Conflict() => new(WorkflowDefinitionUpdateOutcome.Conflict); +} diff --git a/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs b/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs index c7facf34d..a8f6dfc25 100644 --- a/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs +++ b/src/modules/Elsa.Workflows.Management/Stores/CachingWorkflowDefinitionStore.cs @@ -99,6 +99,21 @@ public class CachingWorkflowDefinitionStore(IWorkflowDefinitionStore decoratedSt await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); } + /// + public async Task TryUpdateLatestAsync( + WorkflowDefinitionFilter filter, + Func matchesExpected, + Func update, + CancellationToken cancellationToken = default) + { + var result = await decoratedStore.TryUpdateLatestAsync(filter, matchesExpected, update, cancellationToken); + + if (result.Outcome == WorkflowDefinitionUpdateOutcome.Updated) + await cacheManager.TriggerTokenAsync(CacheInvalidationTokenKey, cancellationToken); + + return result; + } + /// public async Task SaveManyAsync(IEnumerable definitions, CancellationToken cancellationToken = default) { diff --git a/src/modules/Elsa.Workflows.Management/Stores/MemoryWorkflowDefinitionStore.cs b/src/modules/Elsa.Workflows.Management/Stores/MemoryWorkflowDefinitionStore.cs index d30153a23..ad702e3c1 100644 --- a/src/modules/Elsa.Workflows.Management/Stores/MemoryWorkflowDefinitionStore.cs +++ b/src/modules/Elsa.Workflows.Management/Stores/MemoryWorkflowDefinitionStore.cs @@ -96,23 +96,60 @@ public class MemoryWorkflowDefinitionStore(MemoryStore store /// public Task SaveAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { - store.Save(definition, GetId); + lock (store.Sync) + store.Save(definition, GetId); + return Task.CompletedTask; } /// public Task SaveManyAsync(IEnumerable definitions, CancellationToken cancellationToken = default) { - store.SaveMany(definitions, GetId); + lock (store.Sync) + store.SaveMany(definitions, GetId); + return Task.CompletedTask; } + /// + public Task TryUpdateLatestAsync( + WorkflowDefinitionFilter filter, + Func matchesExpected, + Func update, + CancellationToken cancellationToken = default) + { + lock (store.Sync) + { + var current = store.Query(query => Filter(query, filter)).FirstOrDefault(); + + if (current == null) + return Task.FromResult(WorkflowDefinitionUpdateResult.NotFound()); + + if (!current.IsLatest || !matchesExpected(current)) + return Task.FromResult(WorkflowDefinitionUpdateResult.Conflict()); + + var next = update(current); + + if (next.Id != current.Id) + { + current.IsLatest = false; + store.Save(current, GetId); + } + + store.Save(next, GetId); + return Task.FromResult(WorkflowDefinitionUpdateResult.Updated(next)); + } + } + /// public Task DeleteAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { - var workflowDefinitionIds = store.Query(query => Filter(query, filter)).Select(x => x.DefinitionId).Distinct().ToList(); - store.DeleteWhere(x => workflowDefinitionIds.Contains(x.DefinitionId)); - return Task.FromResult(workflowDefinitionIds.LongCount()); + lock (store.Sync) + { + var workflowDefinitionIds = store.Query(query => Filter(query, filter)).Select(x => x.DefinitionId).Distinct().ToList(); + store.DeleteWhere(x => workflowDefinitionIds.Contains(x.DefinitionId)); + return Task.FromResult(workflowDefinitionIds.LongCount()); + } } /// diff --git a/test/integration/Elsa.AI.IntegrationTests/AIWorkflowGroundingToolTests.cs b/test/integration/Elsa.AI.IntegrationTests/AIWorkflowGroundingToolTests.cs index 7c484a205..4fe905cd3 100644 --- a/test/integration/Elsa.AI.IntegrationTests/AIWorkflowGroundingToolTests.cs +++ b/test/integration/Elsa.AI.IntegrationTests/AIWorkflowGroundingToolTests.cs @@ -105,6 +105,32 @@ public class AIWorkflowGroundingToolTests return Task.CompletedTask; } + public Task TryUpdateLatestAsync( + WorkflowDefinitionFilter filter, + Func matchesExpected, + Func update, + CancellationToken cancellationToken = default) + { + var current = Apply(filter).FirstOrDefault(); + + if (current is null) + return Task.FromResult(WorkflowDefinitionUpdateResult.NotFound()); + + if (!matchesExpected(current)) + return Task.FromResult(WorkflowDefinitionUpdateResult.Conflict()); + + var next = update(current); + _definitions.RemoveAll(x => x.Id == current.Id || x.Id == next.Id); + if (next.Id != current.Id) + { + current.IsLatest = false; + _definitions.Add(current); + } + + _definitions.Add(next); + return Task.FromResult(WorkflowDefinitionUpdateResult.Updated(next)); + } + public Task SaveManyAsync(IEnumerable definitions, CancellationToken cancellationToken = default) { foreach (var definition in definitions) diff --git a/test/integration/Elsa.Bpmn.Interchange.IntegrationTests/Scenarios/Interchange/BpmnDocumentPutCompareAndSwapTests.cs b/test/integration/Elsa.Bpmn.Interchange.IntegrationTests/Scenarios/Interchange/BpmnDocumentPutCompareAndSwapTests.cs new file mode 100644 index 000000000..52b3b2818 --- /dev/null +++ b/test/integration/Elsa.Bpmn.Interchange.IntegrationTests/Scenarios/Interchange/BpmnDocumentPutCompareAndSwapTests.cs @@ -0,0 +1,363 @@ +using Bpmn.Interchange; +using Bpmn.Model; +using Elsa.Bpmn.Interchange.Endpoints.Bpmn.Document; +using Elsa.Bpmn.Interchange.Exceptions; +using Elsa.Bpmn.Interchange.Services; +using Elsa.Common.Models; +using Elsa.Extensions; +using Elsa.Mediator.Contracts; +using Elsa.Testing.Shared; +using Elsa.Workflows.Management; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Management.Notifications; +using Elsa.Workflows.Models; +using Microsoft.Extensions.DependencyInjection; +using Xunit.Abstractions; + +namespace Elsa.Bpmn.Interchange.IntegrationTests.Scenarios.Interchange; + +/// +/// The document PUT's precondition and save are one compare-and-swap (#8064). A store double pauses the first +/// writer after its match and before its save so a second writer can finish in that window; the first then gets +/// the same refusal a stale If-Match gets, and the second writer's change is what stays stored. +/// +public class BpmnDocumentPutCompareAndSwapTests(ITestOutputHelper testOutputHelper) +{ + [Fact(DisplayName = "A second document edit that wins the race leaves its change stored; the first writer is refused")] + public async Task ImportDocumentAsync_WhenASecondWriterSavesAfterTheFirstMatched_RefusesTheFirstAndKeepsTheSecond() + { + var services = new TestApplicationBuilder(testOutputHelper) + .ConfigureElsa(elsa => elsa.UseBpmnInterchange()) + .Build(); + await services.PopulateRegistriesAsync(); + + var innerStore = services.GetRequiredService(); + var setup = services.GetRequiredService(); + + var imported = await setup.ImportAsync(ReadAsset("camunda-order-process.bpmn"), definitionId: null, name: null, processId: null, CancellationToken.None); + var definitionId = imported.ImportResult.WorkflowDefinition.DefinitionId; + var stored = await FindLatestAsync(innerStore, definitionId); + var expectedETag = BpmnDocumentETag.From(stored); + var reader = services.GetRequiredService(); + var sourceXml = (string)stored.CustomProperties[BpmnInterchangeDocumentService.SourceXmlCustomPropertyKey]; + var firstEdit = reader.Read(sourceXml.Replace("Order Handled", "First writer"), new BpmnImportOptions()).Definitions; + var secondEdit = reader.Read(sourceXml.Replace("Order Handled", "Second writer"), new BpmnImportOptions()).Definitions; + + var gate = new CompareAndSwapPauseGate(); + var pausingStore = new PausingCompareAndSwapStore(innerStore, gate); + var firstWriter = ActivatorUtilities.CreateInstance(services, pausingStore); + + var first = firstWriter.ImportDocumentAsync(firstEdit, definitionId, processId: null, CancellationToken.None, expectedETag); + await gate.Checked.Task.WaitAsync(TimeSpan.FromSeconds(10)); + + var second = await setup.ImportDocumentAsync(secondEdit, definitionId, processId: null, CancellationToken.None, expectedETag); + Assert.True(second.ImportResult.Succeeded); + + gate.Release.TrySetResult(); + + var lost = await Assert.ThrowsAsync(() => first); + + Assert.Contains("written since the ETag in If-Match was issued", lost.Message); + + var after = await FindLatestAsync(innerStore, definitionId); + Assert.Equal(BpmnDocumentETag.From(second.ImportResult.WorkflowDefinition), BpmnDocumentETag.From(after)); + Assert.NotEqual(expectedETag, BpmnDocumentETag.From(after)); + } + + [Fact(DisplayName = "A metadata-only save in the window is 412; the rename stays and the document PUT does not overwrite it")] + public async Task ImportDocumentAsync_WhenMetadataIsSavedAfterTheFirstMatched_RefusesRatherThanRevertingTheRename() + { + var services = new TestApplicationBuilder(testOutputHelper) + .ConfigureElsa(elsa => elsa.UseBpmnInterchange()) + .Build(); + await services.PopulateRegistriesAsync(); + + var innerStore = services.GetRequiredService(); + var setup = services.GetRequiredService(); + + var imported = await setup.ImportAsync(ReadAsset("camunda-order-process.bpmn"), definitionId: null, name: "Order", processId: null, CancellationToken.None); + var definitionId = imported.ImportResult.WorkflowDefinition.DefinitionId; + var stored = await FindLatestAsync(innerStore, definitionId); + var expectedETag = BpmnDocumentETag.From(stored); + var reader = services.GetRequiredService(); + var sourceXml = (string)stored.CustomProperties[BpmnInterchangeDocumentService.SourceXmlCustomPropertyKey]; + var firstEdit = reader.Read(sourceXml.Replace("Order Handled", "First writer"), new BpmnImportOptions()).Definitions; + + var gate = new CompareAndSwapPauseGate(); + var pausingStore = new PausingCompareAndSwapStore(innerStore, gate); + var firstWriter = ActivatorUtilities.CreateInstance(services, pausingStore); + + var first = firstWriter.ImportDocumentAsync(firstEdit, definitionId, processId: null, CancellationToken.None, expectedETag); + await gate.Checked.Task.WaitAsync(TimeSpan.FromSeconds(10)); + + var latest = await FindLatestAsync(innerStore, definitionId); + latest.Name = "Renamed-in-window"; + await innerStore.SaveAsync(latest); + + gate.Release.TrySetResult(); + + var lost = await Assert.ThrowsAsync(() => first); + + Assert.Contains("written since the ETag in If-Match was issued", lost.Message); + + var after = await FindLatestAsync(innerStore, definitionId); + Assert.Equal("Renamed-in-window", after.Name); + Assert.Equal(stored.StringData, after.StringData); + } + + [Fact(DisplayName = "A rejecting DraftSaving handler fails the document PUT before persist; the stored definition is unchanged")] + public async Task ImportDocumentAsync_WhenDraftSavingHandlerRejects_DoesNotPersist() + { + var probe = new DraftNotificationProbe(); + var services = new TestApplicationBuilder(testOutputHelper) + .ConfigureElsa(elsa => elsa.UseBpmnInterchange()) + .ConfigureServices(s => + { + s.AddSingleton(probe); + s.AddNotificationHandler(); + s.AddNotificationHandler(); + }) + .Build(); + await services.PopulateRegistriesAsync(); + + var store = services.GetRequiredService(); + var setup = services.GetRequiredService(); + + var imported = await setup.ImportAsync(ReadAsset("camunda-order-process.bpmn"), definitionId: null, name: "Order", processId: null, CancellationToken.None); + var definitionId = imported.ImportResult.WorkflowDefinition.DefinitionId; + var before = await FindLatestAsync(store, definitionId); + var expectedETag = BpmnDocumentETag.From(before); + var reader = services.GetRequiredService(); + var sourceXml = (string)before.CustomProperties[BpmnInterchangeDocumentService.SourceXmlCustomPropertyKey]; + var edit = reader.Read(sourceXml.Replace("Order Handled", "Must not persist"), new BpmnImportOptions()).Definitions; + + var savingBefore = probe.SavingCount; + var savedBefore = probe.SavedCount; + probe.Reject = true; + + var rejected = await Assert.ThrowsAsync( + () => setup.ImportDocumentAsync(edit, definitionId, processId: null, CancellationToken.None, expectedETag)); + + Assert.Equal("Draft save rejected.", rejected.Message); + Assert.Equal(savingBefore + 1, probe.SavingCount); + Assert.Equal(savedBefore, probe.SavedCount); + + var after = await FindLatestAsync(store, definitionId); + Assert.Equal(expectedETag, BpmnDocumentETag.From(after)); + Assert.Equal(before.StringData, after.StringData); + Assert.Equal(before.Name, after.Name); + } + + [Fact(DisplayName = "A published→draft document PUT keeps the same id, version and created-at from DraftSaving through persist and DraftSaved")] + public async Task ImportDocumentAsync_WhenLatestIsPublished_ReusesTheAnnouncedDraftIdentity() + { + var probe = new DraftNotificationProbe(); + var services = new TestApplicationBuilder(testOutputHelper) + .ConfigureElsa(elsa => elsa.UseBpmnInterchange()) + .ConfigureServices(s => + { + s.AddSingleton(probe); + s.AddNotificationHandler(); + s.AddNotificationHandler(); + }) + .Build(); + await services.PopulateRegistriesAsync(); + + var store = services.GetRequiredService(); + var setup = services.GetRequiredService(); + var publisher = services.GetRequiredService(); + + var imported = await setup.ImportAsync(ReadAsset("camunda-order-process.bpmn"), definitionId: null, name: "Order", processId: null, CancellationToken.None); + var definitionId = imported.ImportResult.WorkflowDefinition.DefinitionId; + await Support.DefinitionPublishing.PublishLatestAsync(publisher, definitionId); + + var published = await FindLatestAsync(store, definitionId); + Assert.True(published.IsPublished); + var expectedETag = BpmnDocumentETag.From(published); + var reader = services.GetRequiredService(); + var sourceXml = (string)published.CustomProperties[BpmnInterchangeDocumentService.SourceXmlCustomPropertyKey]; + var edit = reader.Read(sourceXml.Replace("Order Handled", "After publish"), new BpmnImportOptions()).Definitions; + + probe.StampHandlerMarker = true; + + var result = await setup.ImportDocumentAsync(edit, definitionId, processId: null, CancellationToken.None, expectedETag); + + Assert.True(result.ImportResult.Succeeded); + Assert.NotNull(probe.SavingId); + Assert.Equal(probe.SavingId, probe.SavedId); + Assert.Equal(probe.SavingVersion, probe.SavedVersion); + Assert.Equal(probe.SavingCreatedAt, probe.SavedCreatedAt); + Assert.NotEqual(published.Id, probe.SavingId); + Assert.Equal(published.Version + 1, probe.SavingVersion); + + var persisted = result.ImportResult.WorkflowDefinition; + Assert.Equal(probe.SavingId, persisted.Id); + Assert.Equal(probe.SavingVersion, persisted.Version); + Assert.Equal(probe.SavingCreatedAt, persisted.CreatedAt); + + var stored = await FindLatestAsync(store, definitionId); + Assert.Equal(probe.SavingId, stored.Id); + Assert.Equal(probe.SavingVersion, stored.Version); + Assert.Equal(probe.SavingCreatedAt, stored.CreatedAt); + Assert.False(stored.IsPublished); + Assert.Equal("kept", stored.CustomProperties[DraftNotificationProbe.HandlerMarkerKey]); + Assert.Equal("Order", stored.Name); + } + + private static async Task FindLatestAsync(IWorkflowDefinitionStore store, string definitionId) + { + var found = await store.FindAsync(WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Latest).ToFilter()); + return found ?? throw new InvalidOperationException($"Latest definition '{definitionId}' was not stored."); + } + + private static string ReadAsset(string fileName) => Support.BpmnAssetReader.Read(fileName); + + private sealed class DraftNotificationProbe + { + public const string HandlerMarkerKey = "test:draft-saving-marker"; + + public bool Reject { get; set; } + public bool StampHandlerMarker { get; set; } + public int SavingCount { get; set; } + public int SavedCount { get; set; } + public string? SavingId { get; set; } + public int SavingVersion { get; set; } + public DateTimeOffset SavingCreatedAt { get; set; } + public string? SavedId { get; set; } + public int SavedVersion { get; set; } + public DateTimeOffset SavedCreatedAt { get; set; } + } + + private sealed class RejectingDraftSavingHandler(DraftNotificationProbe probe) : INotificationHandler + { + public Task HandleAsync(WorkflowDefinitionDraftSaving notification, CancellationToken cancellationToken) + { + probe.SavingCount++; + probe.SavingId = notification.WorkflowDefinition.Id; + probe.SavingVersion = notification.WorkflowDefinition.Version; + probe.SavingCreatedAt = notification.WorkflowDefinition.CreatedAt; + + if (probe.StampHandlerMarker) + notification.WorkflowDefinition.CustomProperties[DraftNotificationProbe.HandlerMarkerKey] = "kept"; + + if (probe.Reject) + throw new InvalidOperationException("Draft save rejected."); + + return Task.CompletedTask; + } + } + + private sealed class CountingDraftSavedHandler(DraftNotificationProbe probe) : INotificationHandler + { + public Task HandleAsync(WorkflowDefinitionDraftSaved notification, CancellationToken cancellationToken) + { + probe.SavedCount++; + probe.SavedId = notification.WorkflowDefinition.Id; + probe.SavedVersion = notification.WorkflowDefinition.Version; + probe.SavedCreatedAt = notification.WorkflowDefinition.CreatedAt; + return Task.CompletedTask; + } + } + + /// + /// Signals the test can run the second writer () and then lets the first writer resume + /// (). Not a lock: the first writer is parked between match and save so the race is + /// deterministic. + /// + private sealed class CompareAndSwapPauseGate + { + public TaskCompletionSource Checked { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously); + public TaskCompletionSource Release { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously); + } + + /// + /// A store double that implements the compare-and-swap contract with an explicit pause after the match: + /// the first writer stops there, the second writer finishes on the inner store, then this re-reads and + /// re-checks so a lost race is Conflict instead of an overwrite. + /// + private sealed class PausingCompareAndSwapStore(IWorkflowDefinitionStore inner, CompareAndSwapPauseGate gate) : IWorkflowDefinitionStore + { + public Task FindAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) => + inner.FindAsync(filter, cancellationToken); + + public Task FindAsync(WorkflowDefinitionFilter filter, WorkflowDefinitionOrder order, CancellationToken cancellationToken = default) => + inner.FindAsync(filter, order, cancellationToken); + + public Task> FindManyAsync(WorkflowDefinitionFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) => + inner.FindManyAsync(filter, pageArgs, cancellationToken); + + public Task> FindManyAsync(WorkflowDefinitionFilter filter, WorkflowDefinitionOrder order, PageArgs pageArgs, CancellationToken cancellationToken = default) => + inner.FindManyAsync(filter, order, pageArgs, cancellationToken); + + public Task> FindManyAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) => + inner.FindManyAsync(filter, cancellationToken); + + public Task> FindManyAsync(WorkflowDefinitionFilter filter, WorkflowDefinitionOrder order, CancellationToken cancellationToken = default) => + inner.FindManyAsync(filter, order, cancellationToken); + + public Task> FindSummariesAsync(WorkflowDefinitionFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) => + inner.FindSummariesAsync(filter, pageArgs, cancellationToken); + + public Task> FindSummariesAsync(WorkflowDefinitionFilter filter, WorkflowDefinitionOrder order, PageArgs pageArgs, CancellationToken cancellationToken = default) => + inner.FindSummariesAsync(filter, order, pageArgs, cancellationToken); + + public Task> FindSummariesAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) => + inner.FindSummariesAsync(filter, cancellationToken); + + public Task> FindSummariesAsync(WorkflowDefinitionFilter filter, WorkflowDefinitionOrder order, CancellationToken cancellationToken = default) => + inner.FindSummariesAsync(filter, order, cancellationToken); + + public Task FindLastVersionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken) => + inner.FindLastVersionAsync(filter, cancellationToken); + + public Task SaveAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) => + inner.SaveAsync(definition, cancellationToken); + + public async Task TryUpdateLatestAsync( + WorkflowDefinitionFilter filter, + Func matchesExpected, + Func update, + CancellationToken cancellationToken = default) + { + var current = await inner.FindAsync(filter, cancellationToken); + + if (current is null) + return WorkflowDefinitionUpdateResult.NotFound(); + + if (!matchesExpected(current)) + return WorkflowDefinitionUpdateResult.Conflict(); + + gate.Checked.TrySetResult(); + await gate.Release.Task.WaitAsync(cancellationToken); + + current = await inner.FindAsync(filter, cancellationToken); + + if (current is null) + return WorkflowDefinitionUpdateResult.NotFound(); + + if (!matchesExpected(current)) + return WorkflowDefinitionUpdateResult.Conflict(); + + var next = update(current); + await inner.SaveAsync(next, cancellationToken); + return WorkflowDefinitionUpdateResult.Updated(next); + } + + public Task SaveManyAsync(IEnumerable definitions, CancellationToken cancellationToken = default) => + inner.SaveManyAsync(definitions, cancellationToken); + + public Task DeleteAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) => + inner.DeleteAsync(filter, cancellationToken); + + public Task AnyAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) => + inner.AnyAsync(filter, cancellationToken); + + public Task CountDistinctAsync(CancellationToken cancellationToken = default) => + inner.CountDistinctAsync(cancellationToken); + + public Task GetIsNameUnique(string name, string? definitionId = null, CancellationToken cancellationToken = default) => + inner.GetIsNameUnique(name, definitionId, cancellationToken); + } +} diff --git a/test/integration/Elsa.Bpmn.Interchange.IntegrationTests/Scenarios/Interchange/BpmnExportAvailabilityTests.cs b/test/integration/Elsa.Bpmn.Interchange.IntegrationTests/Scenarios/Interchange/BpmnExportAvailabilityTests.cs index b2051654e..1411c2abe 100644 --- a/test/integration/Elsa.Bpmn.Interchange.IntegrationTests/Scenarios/Interchange/BpmnExportAvailabilityTests.cs +++ b/test/integration/Elsa.Bpmn.Interchange.IntegrationTests/Scenarios/Interchange/BpmnExportAvailabilityTests.cs @@ -6,7 +6,10 @@ using Elsa.Bpmn.Interchange.IntegrationTests.Support; using Elsa.Bpmn.Interchange.Services; using Elsa.Common.Models; using Elsa.Extensions; +using Elsa.Common; +using Elsa.Mediator.Contracts; using Elsa.Testing.Shared; +using Elsa.Workflows; using Elsa.Workflows.Activities; using Elsa.Workflows.Management; using Elsa.Workflows.Management.Entities; @@ -277,7 +280,11 @@ public class BpmnExportAvailabilityTests(ITestOutputHelper testOutputHelper) : B services.GetRequiredService(), importer, store, - services.GetRequiredService()); + services.GetRequiredService(), + services.GetRequiredService(), + services.GetRequiredService(), + services.GetRequiredService(), + services.GetRequiredService()); var xml = BpmnAssetReader.Read("camunda-order-process.bpmn"); @@ -374,6 +381,13 @@ public class BpmnExportAvailabilityTests(ITestOutputHelper testOutputHelper) : B public Task SaveAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) => throw new InvalidOperationException("Simulated save failure for testing atomic BPMN import."); + public Task TryUpdateLatestAsync( + WorkflowDefinitionFilter filter, + Func matchesExpected, + Func update, + CancellationToken cancellationToken = default) => + throw new InvalidOperationException("Simulated save failure for testing atomic BPMN import."); + public Task FindAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) => throw new NotSupportedException(); public Task FindAsync(WorkflowDefinitionFilter filter, WorkflowDefinitionOrder order, CancellationToken cancellationToken = default) => throw new NotSupportedException(); public Task> FindManyAsync(WorkflowDefinitionFilter filter, PageArgs pageArgs, CancellationToken cancellationToken = default) => throw new NotSupportedException(); diff --git a/test/unit/Elsa.Bpmn.Interchange.UnitTests/BpmnErrorResponseMappingTests.cs b/test/unit/Elsa.Bpmn.Interchange.UnitTests/BpmnErrorResponseMappingTests.cs index 80eb7be5d..fc03b874e 100644 --- a/test/unit/Elsa.Bpmn.Interchange.UnitTests/BpmnErrorResponseMappingTests.cs +++ b/test/unit/Elsa.Bpmn.Interchange.UnitTests/BpmnErrorResponseMappingTests.cs @@ -84,6 +84,20 @@ public class BpmnErrorResponseMappingTests Assert.Null(response.Data); } + [Fact(DisplayName = "A lost compare-and-swap is coded bpmn.document.precondition-failed, the same 412 If-Match already uses")] + public void PreconditionFailedResponseFor_CarriesTheCode() + { + var exception = new BpmnDocumentPreconditionFailedException( + "The workflow definition has been written since the ETag in If-Match was issued. GET the document again, reapply the edit, and PUT it with the new ETag."); + + var response = BpmnImportErrorResponses.PreconditionFailedResponseFor(exception); + + Assert.Equal(BpmnErrorCodes.DocumentPreconditionFailed, response.Code); + Assert.Equal(StatusCodes.Status412PreconditionFailed, response.StatusCode); + Assert.Equal(exception.Message, Assert.Single(response.Errors["generalErrors"])); + Assert.Null(response.Data); + } + [Theory(DisplayName = "Each BpmnExportUnavailableReason maps to its own code")] [InlineData(BpmnExportUnavailableReason.NotImported, BpmnErrorCodes.ExportNotImported)] [InlineData(BpmnExportUnavailableReason.SourceVersionUnknown, BpmnErrorCodes.ExportSourceVersionUnknown)] diff --git a/test/unit/Elsa.Workflows.Management.UnitTests/Stores/MemoryWorkflowDefinitionStoreCompareAndSwapTests.cs b/test/unit/Elsa.Workflows.Management.UnitTests/Stores/MemoryWorkflowDefinitionStoreCompareAndSwapTests.cs new file mode 100644 index 000000000..bf51641db --- /dev/null +++ b/test/unit/Elsa.Workflows.Management.UnitTests/Stores/MemoryWorkflowDefinitionStoreCompareAndSwapTests.cs @@ -0,0 +1,193 @@ +using Elsa.Common.Models; +using Elsa.Common.Services; +using Elsa.Extensions; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Management.Models; +using Elsa.Workflows.Management.Stores; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Management.UnitTests.Stores; + +/// +/// is the compare-and-swap the BPMN document +/// PUT uses: load, match, apply, save in one critical section. Lost match is Conflict, not an overwrite. +/// +public class MemoryWorkflowDefinitionStoreCompareAndSwapTests +{ + [Fact(DisplayName = "A matching latest row is updated and the callback sees that just-loaded row")] + public async Task TryUpdateLatestAsync_WhenTheRowMatches_SavesTheUpdateBuiltFromTheLoadedRow() + { + var store = new MemoryWorkflowDefinitionStore(new MemoryStore()); + var current = Definition("def-1", "id-1", name: "Original", stringData: "graph-v1"); + await store.SaveAsync(current); + + var result = await store.TryUpdateLatestAsync( + LatestOf("def-1"), + loaded => loaded.StringData == "graph-v1", + loaded => + { + var next = loaded.ShallowClone(); + next.StringData = "graph-v2"; + next.Name = loaded.Name + "-kept"; + return next; + }); + + Assert.Equal(WorkflowDefinitionUpdateOutcome.Updated, result.Outcome); + Assert.Equal("graph-v2", result.Definition!.StringData); + Assert.Equal("Original-kept", result.Definition.Name); + + var stored = await store.FindAsync(LatestOf("def-1")); + Assert.Equal("graph-v2", stored!.StringData); + Assert.Equal("Original-kept", stored.Name); + } + + [Fact(DisplayName = "A match that fails is Conflict and the stored row is unchanged")] + public async Task TryUpdateLatestAsync_WhenTheRowDoesNotMatch_ReturnsConflictAndWritesNothing() + { + var store = new MemoryWorkflowDefinitionStore(new MemoryStore()); + await store.SaveAsync(Definition("def-1", "id-1", name: "Original", stringData: "graph-v2")); + + var result = await store.TryUpdateLatestAsync( + LatestOf("def-1"), + loaded => loaded.StringData == "graph-v1", + loaded => + { + var next = loaded.ShallowClone(); + next.StringData = "should-not-be-saved"; + return next; + }); + + Assert.Equal(WorkflowDefinitionUpdateOutcome.Conflict, result.Outcome); + Assert.Null(result.Definition); + + var stored = await store.FindAsync(LatestOf("def-1")); + Assert.Equal("graph-v2", stored!.StringData); + Assert.Equal("Original", stored.Name); + } + + [Fact(DisplayName = "A missing definition is NotFound")] + public async Task TryUpdateLatestAsync_WhenNothingMatchesTheFilter_ReturnsNotFound() + { + var store = new MemoryWorkflowDefinitionStore(new MemoryStore()); + + var result = await store.TryUpdateLatestAsync( + LatestOf("missing"), + _ => true, + current => current); + + Assert.Equal(WorkflowDefinitionUpdateOutcome.NotFound, result.Outcome); + } + + [Fact(DisplayName = "Two scoped wrappers share the backing-store lock: the loser is Conflict and the winner's graph stays")] + public async Task TryUpdateLatestAsync_WhenTwoWrappersShareTheBackingStore_LoserIsConflictAndWinnerGraphStays() + { + var backing = new MemoryStore(); + var writerA = new MemoryWorkflowDefinitionStore(backing); + var writerB = new MemoryWorkflowDefinitionStore(backing); + await writerA.SaveAsync(Definition("def-1", "id-1", name: "Original", stringData: "graph-v1")); + + var firstHoldsLock = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var releaseFirst = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + var first = Task.Run(() => writerA.TryUpdateLatestAsync( + LatestOf("def-1"), + loaded => loaded.StringData == "graph-v1", + loaded => + { + firstHoldsLock.TrySetResult(); + releaseFirst.Task.GetAwaiter().GetResult(); + var next = loaded.ShallowClone(); + next.StringData = "winner-graph"; + return next; + })); + + await firstHoldsLock.Task; + + var second = Task.Run(() => writerB.TryUpdateLatestAsync( + LatestOf("def-1"), + loaded => loaded.StringData == "graph-v1", + loaded => + { + var next = loaded.ShallowClone(); + next.StringData = "stale-overwrite"; + return next; + })); + + // Writer B must still be waiting on the shared lock. A per-wrapper lock would have + // already written the stale graph and completed as Updated. + await Task.Delay(100); + Assert.False(second.IsCompleted); + + releaseFirst.TrySetResult(); + + var firstResult = await first; + var secondResult = await second; + + Assert.Equal(WorkflowDefinitionUpdateOutcome.Updated, firstResult.Outcome); + Assert.Equal(WorkflowDefinitionUpdateOutcome.Conflict, secondResult.Outcome); + + var stored = await writerA.FindAsync(LatestOf("def-1")); + Assert.Equal("winner-graph", stored!.StringData); + } + + [Fact(DisplayName = "A loaded row that is no longer IsLatest is Conflict and writes nothing")] + public async Task TryUpdateLatestAsync_WhenTheLoadedRowIsNoLongerLatest_ReturnsConflictAndWritesNothing() + { + var backing = new MemoryStore(); + var store = new MemoryWorkflowDefinitionStore(backing); + var published = Definition("def-1", "id-1", name: "Published", stringData: "graph-v1"); + published.IsPublished = true; + await store.SaveAsync(published); + + var winner = await store.TryUpdateLatestAsync( + LatestOf("def-1"), + _ => true, + loaded => + { + var draft = loaded.ShallowClone(); + draft.Id = "id-2"; + draft.Version = loaded.Version + 1; + draft.IsPublished = false; + draft.StringData = "winner-draft"; + return draft; + }); + + Assert.Equal(WorkflowDefinitionUpdateOutcome.Updated, winner.Outcome); + + var loser = await store.TryUpdateLatestAsync( + new WorkflowDefinitionFilter { Id = published.Id }, + _ => true, + loaded => + { + var draft = loaded.ShallowClone(); + draft.Id = "id-3"; + draft.Version = loaded.Version + 1; + draft.IsPublished = false; + draft.StringData = "should-not-be-saved"; + return draft; + }); + + Assert.Equal(WorkflowDefinitionUpdateOutcome.Conflict, loser.Outcome); + Assert.Null(loser.Definition); + + var stored = await store.FindAsync(LatestOf("def-1")); + Assert.Equal("id-2", stored!.Id); + Assert.Equal("winner-draft", stored.StringData); + } + + private static WorkflowDefinitionFilter LatestOf(string definitionId) => + WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Latest).ToFilter(); + + private static WorkflowDefinition Definition(string definitionId, string id, string name, string stringData) => + new() + { + Id = id, + DefinitionId = definitionId, + Name = name, + StringData = stringData, + Version = 1, + IsLatest = true, + MaterializerName = "Json" + }; +}