From 0ac71842264cdce876e1b8dc2a7be8f18989d38c Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sun, 13 Sep 2026 01:40:43 -0700 Subject: [PATCH] fix(bpmn): make document PUT If-Match and save a compare-and-swap (#8092) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(bpmn): make document PUT If-Match and save a compare-and-swap The document PUT checked If-Match, reloaded metadata, then saved through the importer as separate steps. Two writers could both pass If-Match, and a metadata-only save in that window was silently reverted. Add IWorkflowDefinitionStore.TryUpdateLatestAsync — load, match, apply, save as one critical section (memory) or ExecuteUpdate against the loaded snapshot (EF). ImportDocumentAsync reads metadata inside that swap. A lost race throws the same 412 the stale If-Match already returns. Mongo/Dapper/ES stores need the same method before the endpoint is concurrency-safe on those providers. Co-authored-by: Sipke Schoorstra * fix(bpmn): resolve document-service DI and interleave test compile Drop the cache-manager constructor dependency (only registered when definition caching is on) and evict via DraftSaving/DraftSaved instead. Update the export-availability stub construction and the CAS interleave fixture to parse edited XML through the reader. Co-authored-by: Sipke Schoorstra * fix(bpmn): close Greptile P1s on document PUT CAS Require IsLatest in the EF ExecuteUpdate WHERE so a published-to-draft loser is Conflict instead of a unique-key failure. Lock Memory CAS on the shared MemoryStore so scoped wrappers cannot stale-overwrite. Dispatch WorkflowDefinitionDraftSaving before the CAS persist so a rejecting handler fails the request before commit. Co-authored-by: Sipke Schoorstra * fix(bpmn): reuse published draft identity across DraftSaving and CAS Allocate the published-to-draft id, version and created-at once before WorkflowDefinitionDraftSaving. The compare-and-swap still rebuilds from the just-loaded row so metadata is not frozen from the outer Find, then reuses that announced identity and keeps handler-added custom properties. Co-authored-by: Sipke Schoorstra * test(management): run Memory CAS lock holder off the test thread The shared-lock test blocked inside TryUpdateLatestAsync on the test thread, so it never reached the release signal and hung. Co-authored-by: Sipke Schoorstra * test(efcore): drop the SQLite CAS harness that cannot match ExecuteUpdate The in-memory SQLite fixture could not satisfy the store's DateTimeOffset ORDER BY plus Data snapshot WHERE, so the winner CAS returned Conflict before the IsLatest loser path ran. Memory already covers that contract. Co-authored-by: Sipke Schoorstra * fix(bpmn): persist the prepared DraftSaving draft on document PUT CAS Prepare the draft, dispatch WorkflowDefinitionDraftSaving so handlers can mutate or reject, then TryUpdateLatestAsync with If-Match plus the loaded snapshot and update: _ => draft. A metadata change in the window is 412 instead of overwriting the other write. DraftSaved still fires after CAS. Co-authored-by: Sipke Schoorstra * docs(bpmn): align document PUT remarks with prepared-draft CAS Co-authored-by: Sipke Schoorstra --------- Co-authored-by: Cursor Agent --- doc/wiki/bpmn-workflows.md | 16 +- .../Elsa.Bpmn.Interchange/BpmnErrorCodes.cs | 3 +- .../Bpmn/BpmnImportErrorResponses.cs | 12 + .../Endpoints/Bpmn/Document/Put/Endpoint.cs | 6 +- ...BpmnDocumentPreconditionFailedException.cs | 12 + .../BpmnInterchangeDocumentService.cs | 179 ++++++++- .../Elsa.Common/Services/MemoryStore.cs | 7 + .../Management/WorkflowDefinitionStore.cs | 103 ++++- .../Contracts/IWorkflowDefinitionStore.cs | 40 ++ .../Models/WorkflowDefinitionUpdateResult.cs | 38 ++ .../Stores/CachingWorkflowDefinitionStore.cs | 15 + .../Stores/MemoryWorkflowDefinitionStore.cs | 47 ++- .../AIWorkflowGroundingToolTests.cs | 26 ++ .../BpmnDocumentPutCompareAndSwapTests.cs | 363 ++++++++++++++++++ .../BpmnExportAvailabilityTests.cs | 16 +- .../BpmnErrorResponseMappingTests.cs | 14 + ...kflowDefinitionStoreCompareAndSwapTests.cs | 193 ++++++++++ 17 files changed, 1063 insertions(+), 27 deletions(-) create mode 100644 src/modules/Elsa.Bpmn.Interchange/Exceptions/BpmnDocumentPreconditionFailedException.cs create mode 100644 src/modules/Elsa.Workflows.Management/Models/WorkflowDefinitionUpdateResult.cs create mode 100644 test/integration/Elsa.Bpmn.Interchange.IntegrationTests/Scenarios/Interchange/BpmnDocumentPutCompareAndSwapTests.cs create mode 100644 test/unit/Elsa.Workflows.Management.UnitTests/Stores/MemoryWorkflowDefinitionStoreCompareAndSwapTests.cs 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" + }; +}