< Summary

Information
Class: Elsa.Workflows.Management.Stores.MemoryWorkflowDefinitionStore
Assembly: Elsa.Workflows.Management
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Management/Stores/MemoryWorkflowDefinitionStore.cs
Line coverage
71%
Covered lines: 64
Uncovered lines: 26
Coverable lines: 90
Total lines: 244
Line coverage: 71.1%
Branch coverage
87%
Covered branches: 28
Total branches: 32
Branch coverage: 87.5%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Management/Stores/MemoryWorkflowDefinitionStore.cs

#LineLine coverage
 1using Elsa.Common.Models;
 2using Elsa.Common.Multitenancy;
 3using Elsa.Common.Services;
 4using Elsa.Extensions;
 5using Elsa.Workflows.Management.Entities;
 6using Elsa.Workflows.Management.Filters;
 7using Elsa.Workflows.Management.Models;
 8
 9namespace Elsa.Workflows.Management.Stores;
 10
 11/// <summary>
 12/// A memory implementation of <see cref="IWorkflowDefinitionStore"/>.
 13/// </summary>
 3614public class MemoryWorkflowDefinitionStore(MemoryStore<WorkflowDefinition> store, ITenantAccessor? tenantAccessor = null
 15{
 16    /// <inheritdoc />
 17    public Task<WorkflowDefinition?> FindAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = de
 18    {
 1619        var result = store.Query(query => Filter(query, filter)).FirstOrDefault();
 820        return Task.FromResult(result);
 21    }
 22
 23    /// <inheritdoc />
 24    public Task<WorkflowDefinition?> FindAsync<TOrderBy>(WorkflowDefinitionFilter filter, WorkflowDefinitionOrder<TOrder
 25    {
 026        var result = store.Query(query => Filter(query, filter).OrderBy(order)).FirstOrDefault();
 027        return Task.FromResult(result);
 28    }
 29
 30    /// <inheritdoc />
 31    public Task<Page<WorkflowDefinition>> FindManyAsync(WorkflowDefinitionFilter filter, PageArgs pageArgs, Cancellation
 32    {
 033        var count = store.Query(query => Filter(query, filter)).LongCount();
 034        var result = store.Query(query => Filter(query, filter).Paginate(pageArgs)).ToList();
 035        return Task.FromResult(Page.Of(result, count));
 36    }
 37
 38    /// <inheritdoc />
 39    public Task<Page<WorkflowDefinition>> FindManyAsync<TOrderBy>(WorkflowDefinitionFilter filter, WorkflowDefinitionOrd
 40    {
 041        var count = store.Query(query => Filter(query, filter).OrderBy(order)).LongCount();
 042        var result = store.Query(query => Filter(query, filter).Paginate(pageArgs)).ToList();
 043        return Task.FromResult(Page.Of(result, count));
 44    }
 45
 46    /// <inheritdoc />
 47    public Task<IEnumerable<WorkflowDefinition>> FindManyAsync(WorkflowDefinitionFilter filter, CancellationToken cancel
 48    {
 3849        var result = store.Query(query => Filter(query, filter)).ToList().AsEnumerable();
 1950        return Task.FromResult(result);
 51    }
 52
 53    /// <inheritdoc />
 54    public Task<IEnumerable<WorkflowDefinition>> FindManyAsync<TOrderBy>(WorkflowDefinitionFilter filter, WorkflowDefini
 55    {
 056        var result = store.Query(query => Filter(query, filter).OrderBy(order)).ToList().AsEnumerable();
 057        return Task.FromResult(result);
 58    }
 59
 60    /// <inheritdoc />
 61    public Task<Page<WorkflowDefinitionSummary>> FindSummariesAsync(WorkflowDefinitionFilter filter, PageArgs pageArgs, 
 62    {
 063        var count = store.Query(query => Filter(query, filter)).LongCount();
 064        var result = store.Query(query => Filter(query, filter).Paginate(pageArgs)).Select(WorkflowDefinitionSummary.Fro
 065        return Task.FromResult(Page.Of(result, count));
 66    }
 67
 68    /// <inheritdoc />
 69    public Task<Page<WorkflowDefinitionSummary>> FindSummariesAsync<TOrderBy>(WorkflowDefinitionFilter filter, WorkflowD
 70    {
 071        var count = store.Query(query => Filter(query, filter).OrderBy(order)).LongCount();
 072        var result = store.Query(query => Filter(query, filter).Paginate(pageArgs)).Select(WorkflowDefinitionSummary.Fro
 073        return Task.FromResult(Page.Of(result, count));
 74    }
 75
 76    /// <inheritdoc />
 77    public Task<IEnumerable<WorkflowDefinitionSummary>> FindSummariesAsync(WorkflowDefinitionFilter filter, Cancellation
 78    {
 079        var result = store.Query(query => Filter(query, filter)).Select(WorkflowDefinitionSummary.FromDefinition).ToList
 080        return Task.FromResult(result);
 81    }
 82
 83    /// <inheritdoc />
 84    public Task<IEnumerable<WorkflowDefinitionSummary>> FindSummariesAsync<TOrderBy>(WorkflowDefinitionFilter filter, Wo
 85    {
 086        var result = store.Query(query => Filter(query, filter).OrderBy(order)).Select(WorkflowDefinitionSummary.FromDef
 087        return Task.FromResult(result);
 88    }
 89
 90    /// <inheritdoc />
 91    public Task<WorkflowDefinition?> FindLastVersionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellatio
 92    {
 093        var result = store.Query(query => Filter(query, filter)).MaxBy(x => x.Version);
 094        return Task.FromResult(result);
 95    }
 96
 97    /// <inheritdoc />
 98    public Task SaveAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default)
 99    {
 41100        lock (store.Sync)
 101        {
 41102            EnsureVersionKeyAvailable(definition);
 40103            store.Save(definition, GetId);
 40104        }
 105
 40106        return Task.CompletedTask;
 107    }
 108
 109    /// <inheritdoc />
 110    public Task SaveManyAsync(IEnumerable<WorkflowDefinition> definitions, CancellationToken cancellationToken = default
 111    {
 3112        lock (store.Sync)
 113        {
 3114            var definitionList = definitions.ToList();
 3115            EnsureBatchVersionKeysUnique(definitionList);
 116
 3117            foreach (var definition in definitionList)
 1118                EnsureVersionKeyAvailable(definition);
 119
 0120            store.SaveMany(definitionList, GetId);
 0121        }
 122
 0123        return Task.CompletedTask;
 124    }
 125
 126    /// <inheritdoc />
 127    public Task<WorkflowDefinitionUpdateResult> TryUpdateLatestAsync(
 128        WorkflowDefinitionFilter filter,
 129        Func<WorkflowDefinition, bool> matchesExpected,
 130        Func<WorkflowDefinition, WorkflowDefinition> update,
 131        CancellationToken cancellationToken = default)
 132    {
 7133        lock (store.Sync)
 134        {
 14135            var current = store.Query(query => Filter(query, filter)).FirstOrDefault();
 136
 7137            if (current == null)
 1138                return Task.FromResult(WorkflowDefinitionUpdateResult.NotFound());
 139
 6140            if (!current.IsLatest || !matchesExpected(current))
 3141                return Task.FromResult(WorkflowDefinitionUpdateResult.Conflict());
 142
 3143            var next = update(current);
 144
 3145            if (next.Id != current.Id)
 146            {
 1147                current.IsLatest = false;
 1148                store.Save(current, GetId);
 149            }
 150
 3151            store.Save(next, GetId);
 3152            return Task.FromResult(WorkflowDefinitionUpdateResult.Updated(next));
 153        }
 7154    }
 155
 156    /// <inheritdoc />
 157    public Task<long> DeleteAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
 158    {
 3159        lock (store.Sync)
 160        {
 11161            var workflowDefinitionIds = store.Query(query => Filter(query, filter)).Select(x => x.DefinitionId).Distinct
 3162            store.DeleteWhere(x =>
 11163                workflowDefinitionIds.Contains(x.DefinitionId)
 11164                && (filter.TenantAgnostic || TenantVisibility.IsVisible(x.TenantId, CurrentTenantId)));
 3165            return Task.FromResult(workflowDefinitionIds.LongCount());
 166        }
 3167    }
 168
 169    /// <inheritdoc />
 170    public Task<bool> AnyAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default)
 171    {
 4172        var exists = store.Query(query => Filter(query, filter)).Any();
 2173        return Task.FromResult(exists);
 174    }
 175
 176    /// <inheritdoc />
 177    public Task<long> CountDistinctAsync(CancellationToken cancellationToken = default)
 178    {
 4179        var count = store.Query(query => query.WhereVisibleToTenant(CurrentTenantId))
 4180            .Select(x => x.DefinitionId)
 2181            .Distinct()
 2182            .LongCount();
 2183        return Task.FromResult(count);
 184    }
 185
 186    /// <inheritdoc />
 187    public Task<bool> GetIsNameUnique(string name, string? definitionId = default, CancellationToken cancellationToken =
 188    {
 2189        var exists = store.Any(x =>
 4190            x.Name == name
 4191            && x.DefinitionId != definitionId
 4192            && TenantVisibility.IsVisible(x.TenantId, CurrentTenantId));
 2193        return Task.FromResult(!exists);
 194    }
 195
 196    /// <remarks>
 197    /// Ambient tenant is applied here rather than in <see cref="WorkflowDefinitionFilter.Apply"/>.
 198    /// EF owns that via <c>SetTenantIdFilter</c> / <c>IgnoreQueryFilters</c>; Memory must compensate.
 199    /// </remarks>
 200    private IQueryable<WorkflowDefinition> Filter(IQueryable<WorkflowDefinition> queryable, WorkflowDefinitionFilter fil
 39201        filter.Apply(queryable.WhereVisibleToTenant(CurrentTenantId, filter.TenantAgnostic));
 202
 49203    private string CurrentTenantId => tenantAccessor?.TenantId ?? Tenant.DefaultTenantId;
 204
 205    /// <remarks>
 206    /// EF enforces <c>(DefinitionId, Version)</c> globally via
 207    /// <c>IX_WorkflowDefinition_DefinitionId_Version</c>. Memory keeps tenant in the key so
 208    /// same-tenant duplicates fail closed while cross-tenant rows remain distinct until #7539
 209    /// adds <c>TenantId</c> to that index.
 210    /// </remarks>
 211    private void EnsureVersionKeyAvailable(WorkflowDefinition definition)
 212    {
 42213        var versionKey = GetVersionKey(definition);
 73214        var existing = store.Find(x => x.Id != definition.Id && GetVersionKey(x) == versionKey);
 215
 42216        if (existing is not null)
 2217            throw CreateVersionKeyConflict(definition);
 40218    }
 219
 220    private static void EnsureBatchVersionKeysUnique(IEnumerable<WorkflowDefinition> definitions)
 221    {
 3222        var seen = new Dictionary<DefinitionVersionKey, string>();
 223
 14224        foreach (var definition in definitions)
 225        {
 5226            var versionKey = GetVersionKey(definition);
 227
 5228            if (seen.TryGetValue(versionKey, out var existingId) && existingId != definition.Id)
 2229                throw CreateVersionKeyConflict(definition);
 230
 3231            seen[versionKey] = definition.Id;
 232        }
 1233    }
 234
 235    private static InvalidOperationException CreateVersionKeyConflict(WorkflowDefinition definition) =>
 4236        new($"A workflow definition already exists for definition '{definition.DefinitionId}' version {definition.Versio
 237
 238    private static DefinitionVersionKey GetVersionKey(WorkflowDefinition definition) =>
 77239        new(definition.DefinitionId, definition.Version, definition.TenantId);
 240
 44241    private string GetId(WorkflowDefinition workflowDefinition) => workflowDefinition.Id;
 242
 0243    private readonly record struct DefinitionVersionKey(string DefinitionId, int Version, string? TenantId);
 244}

Methods/Properties

.ctor(Elsa.Common.Services.MemoryStore`1<Elsa.Workflows.Management.Entities.WorkflowDefinition>,Elsa.Common.Multitenancy.ITenantAccessor)
FindAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,System.Threading.CancellationToken)
FindAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,Elsa.Workflows.Management.Filters.WorkflowDefinitionOrder`1<TOrderBy>,System.Threading.CancellationToken)
FindManyAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,Elsa.Common.Models.PageArgs,System.Threading.CancellationToken)
FindManyAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,Elsa.Workflows.Management.Filters.WorkflowDefinitionOrder`1<TOrderBy>,Elsa.Common.Models.PageArgs,System.Threading.CancellationToken)
FindManyAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,System.Threading.CancellationToken)
FindManyAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,Elsa.Workflows.Management.Filters.WorkflowDefinitionOrder`1<TOrderBy>,System.Threading.CancellationToken)
FindSummariesAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,Elsa.Common.Models.PageArgs,System.Threading.CancellationToken)
FindSummariesAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,Elsa.Workflows.Management.Filters.WorkflowDefinitionOrder`1<TOrderBy>,Elsa.Common.Models.PageArgs,System.Threading.CancellationToken)
FindSummariesAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,System.Threading.CancellationToken)
FindSummariesAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,Elsa.Workflows.Management.Filters.WorkflowDefinitionOrder`1<TOrderBy>,System.Threading.CancellationToken)
FindLastVersionAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,System.Threading.CancellationToken)
SaveAsync(Elsa.Workflows.Management.Entities.WorkflowDefinition,System.Threading.CancellationToken)
SaveManyAsync(System.Collections.Generic.IEnumerable`1<Elsa.Workflows.Management.Entities.WorkflowDefinition>,System.Threading.CancellationToken)
TryUpdateLatestAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,System.Func`2<Elsa.Workflows.Management.Entities.WorkflowDefinition,System.Boolean>,System.Func`2<Elsa.Workflows.Management.Entities.WorkflowDefinition,Elsa.Workflows.Management.Entities.WorkflowDefinition>,System.Threading.CancellationToken)
DeleteAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,System.Threading.CancellationToken)
AnyAsync(Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter,System.Threading.CancellationToken)
CountDistinctAsync(System.Threading.CancellationToken)
GetIsNameUnique(System.String,System.String,System.Threading.CancellationToken)
Filter(System.Linq.IQueryable`1<Elsa.Workflows.Management.Entities.WorkflowDefinition>,Elsa.Workflows.Management.Filters.WorkflowDefinitionFilter)
get_CurrentTenantId()
EnsureVersionKeyAvailable(Elsa.Workflows.Management.Entities.WorkflowDefinition)
EnsureBatchVersionKeysUnique(System.Collections.Generic.IEnumerable`1<Elsa.Workflows.Management.Entities.WorkflowDefinition>)
CreateVersionKeyConflict(Elsa.Workflows.Management.Entities.WorkflowDefinition)
GetVersionKey(Elsa.Workflows.Management.Entities.WorkflowDefinition)
GetId(Elsa.Workflows.Management.Entities.WorkflowDefinition)
get_DefinitionId()