< Summary

Information
Class: Elsa.Persistence.EFCore.Modules.Management.EFCoreWorkflowInstanceStore
Assembly: Elsa.Persistence.EFCore
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Persistence.EFCore/Modules/Management/WorkflowInstanceStore.cs
Line coverage
48%
Covered lines: 62
Uncovered lines: 65
Coverable lines: 127
Total lines: 286
Line coverage: 48.8%
Branch coverage
31%
Covered branches: 7
Total branches: 22
Branch coverage: 31.8%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
FindAsync()100%11100%
FindManyAsync()100%210%
FindManyAsync()100%210%
FindManyAsync()100%11100%
FindManyAsync()100%11100%
CountAsync()100%210%
SummarizeManyAsync()100%11100%
SummarizeManyAsync()100%22100%
SummarizeManyAsync()100%11100%
SummarizeManyAsync()100%11100%
FindManyIdsAsync()100%11100%
FindManyIdsAsync()100%210%
FindManyIdsAsync()0%2040%
DeleteAsync()100%11100%
UpdateUpdatedTimestampAsync()0%7280%
TryMarkInterruptedAsync()0%620%
SaveAsync()100%11100%
AddAsync()100%210%
UpdateAsync()100%210%
SaveManyAsync()100%210%
OnSaveAsync()50%22100%
OnLoadAsync()66.66%7673.33%
Filter(...)100%11100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Persistence.EFCore/Modules/Management/WorkflowInstanceStore.cs

#LineLine coverage
 1using System.Diagnostics.CodeAnalysis;
 2using Elsa.Common;
 3using Elsa.Common.Codecs;
 4using Elsa.Common.Entities;
 5using Elsa.Common.Models;
 6using Elsa.Extensions;
 7using Elsa.Workflows;
 8using Elsa.Workflows.Management;
 9using Elsa.Workflows.Management.Entities;
 10using Elsa.Workflows.Management.Filters;
 11using Elsa.Workflows.Management.Models;
 12using Elsa.Workflows.Management.Options;
 13using JetBrains.Annotations;
 14using Microsoft.EntityFrameworkCore;
 15using Microsoft.Extensions.Logging;
 16using Microsoft.Extensions.Options;
 17using Open.Linq.AsyncExtensions;
 18
 19namespace Elsa.Persistence.EFCore.Modules.Management;
 20
 21/// <summary>
 22/// An EF Core implementation of <see cref="IWorkflowInstanceStore"/>.
 23/// </summary>
 24[UsedImplicitly]
 25public class EFCoreWorkflowInstanceStore : IWorkflowInstanceStore
 26{
 27    private readonly EntityStore<ManagementElsaDbContext, WorkflowInstance> _store;
 28    private readonly IWorkflowStateSerializer _workflowStateSerializer;
 29    private readonly ICompressionCodecResolver _compressionCodecResolver;
 30    private readonly IOptions<ManagementOptions> _options;
 31    private readonly ILogger<EFCoreWorkflowInstanceStore> _logger;
 32
 33    /// <summary>
 34    /// Constructor.
 35    /// </summary>
 41036    public EFCoreWorkflowInstanceStore(
 41037        EntityStore<ManagementElsaDbContext, WorkflowInstance> store,
 41038        IWorkflowStateSerializer workflowStateSerializer,
 41039        ICompressionCodecResolver compressionCodecResolver,
 41040        IOptions<ManagementOptions> options,
 41041        ILogger<EFCoreWorkflowInstanceStore> logger)
 42    {
 41043        _store = store;
 41044        _workflowStateSerializer = workflowStateSerializer;
 41045        _compressionCodecResolver = compressionCodecResolver;
 41046        _options = options;
 41047        _logger = logger;
 41048    }
 49
 50    /// <inheritdoc />
 51    public async ValueTask<WorkflowInstance?> FindAsync(WorkflowInstanceFilter filter, CancellationToken cancellationTok
 52    {
 23453        return await _store.QueryAsync(query => Filter(query, filter), OnLoadAsync, cancellationToken).FirstOrDefault();
 11754    }
 55
 56    /// <inheritdoc />
 57    public async ValueTask<Page<WorkflowInstance>> FindManyAsync(WorkflowInstanceFilter filter, PageArgs pageArgs, Cance
 58    {
 059        var orderBy = new WorkflowInstanceOrder<DateTimeOffset>(x => x.CreatedAt, OrderDirection.Ascending);
 060        return await FindManyAsync(filter, pageArgs, orderBy, cancellationToken);
 061    }
 62
 63    /// <inheritdoc />
 64    public async ValueTask<Page<WorkflowInstance>> FindManyAsync<TOrderBy>(WorkflowInstanceFilter filter, PageArgs pageA
 65    {
 066        var count = await _store.QueryAsync(query => Filter(query, filter), x => x.Id, cancellationToken).LongCount();
 067        var entities = await _store.QueryAsync(query => Filter(query, filter).OrderBy(order).Paginate(pageArgs), OnLoadA
 068        return Page.Of(entities, count);
 069    }
 70
 71    /// <inheritdoc />
 72    public async ValueTask<IEnumerable<WorkflowInstance>> FindManyAsync(WorkflowInstanceFilter filter, CancellationToken
 73    {
 1074        var orderBy = new WorkflowInstanceOrder<DateTimeOffset>(x => x.CreatedAt, OrderDirection.Ascending);
 1075        return await FindManyAsync(filter, orderBy, cancellationToken);
 1076    }
 77
 78    /// <inheritdoc />
 79    public async ValueTask<IEnumerable<WorkflowInstance>> FindManyAsync<TOrderBy>(WorkflowInstanceFilter filter, Workflo
 80    {
 2081        return await _store.QueryAsync(query => Filter(query, filter).OrderBy(order), OnLoadAsync, cancellationToken).To
 1082    }
 83
 84    /// <inheritdoc />
 85    [RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.SerializeAsync(WorkflowState, Can
 86    public async ValueTask<long> CountAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = default
 87    {
 088        return await _store.CountAsync(filter.Apply, cancellationToken);
 089    }
 90
 91    /// <inheritdoc />
 92    public async ValueTask<Page<WorkflowInstanceSummary>> SummarizeManyAsync(WorkflowInstanceFilter filter, PageArgs pag
 93    {
 10994        var orderBy = new WorkflowInstanceOrder<DateTimeOffset>(x => x.CreatedAt, OrderDirection.Ascending);
 10995        return await SummarizeManyAsync(filter, pageArgs, orderBy, cancellationToken);
 10996    }
 97
 98    /// <inheritdoc />
 99    public async ValueTask<Page<WorkflowInstanceSummary>> SummarizeManyAsync<TOrderBy>(WorkflowInstanceFilter filter, Pa
 100    {
 109101        await using var dbContext = await _store.CreateDbContextAsync(cancellationToken);
 109102        var set = dbContext.WorkflowInstances;
 109103        var queryable = Filter(set.AsQueryable(), filter).OrderBy(order);
 109104        var count = await queryable.LongCountAsync(cancellationToken);
 109105        queryable = queryable.Paginate(pageArgs);
 109106        var entities = await queryable.Select(WorkflowInstanceSummary.FromInstanceExpression()).ToListAsync(cancellation
 107
 109108        return Page.Of(entities, count);
 109109    }
 110
 111    /// <inheritdoc />
 112    public async ValueTask<IEnumerable<WorkflowInstanceSummary>> SummarizeManyAsync(WorkflowInstanceFilter filter, Cance
 113    {
 9114        var orderBy = new WorkflowInstanceOrder<DateTimeOffset>(x => x.CreatedAt, OrderDirection.Ascending);
 9115        return await SummarizeManyAsync(filter, orderBy, cancellationToken);
 9116    }
 117
 118    /// <inheritdoc />
 119    public async ValueTask<IEnumerable<WorkflowInstanceSummary>> SummarizeManyAsync<TOrderBy>(WorkflowInstanceFilter fil
 120    {
 27121        return await _store.QueryAsync(query => Filter(query, filter).OrderBy(order), WorkflowInstanceSummary.FromInstan
 9122    }
 123
 124    /// <inheritdoc />
 125    public async ValueTask<IEnumerable<string>> FindManyIdsAsync(WorkflowInstanceFilter filter, CancellationToken cancel
 126    {
 6127        var entities = await _store.QueryAsync(query => Filter(query, filter).OrderBy(x => x.CreatedAt), WorkflowInstanc
 4128        return entities.Select(x => x.Id).ToList();
 2129    }
 130
 131    /// <inheritdoc />
 132    public async ValueTask<Page<string>> FindManyIdsAsync(WorkflowInstanceFilter filter, PageArgs pageArgs, Cancellation
 133    {
 0134        var orderBy = new WorkflowInstanceOrder<DateTimeOffset>(x => x.CreatedAt, OrderDirection.Ascending);
 0135        return await FindManyIdsAsync(filter, pageArgs, orderBy, cancellationToken);
 0136    }
 137
 138    /// <inheritdoc />
 139    public async ValueTask<Page<string>> FindManyIdsAsync<TOrderBy>(WorkflowInstanceFilter filter, PageArgs pageArgs, Wo
 140    {
 0141        await using var dbContext = await _store.CreateDbContextAsync(cancellationToken);
 0142        var set = dbContext.WorkflowInstances;
 0143        var queryable = Filter(set.AsQueryable(), filter).OrderBy(order);
 0144        var count = await queryable.LongCountAsync(cancellationToken);
 0145        queryable = queryable.Paginate(pageArgs);
 0146        var entities = await queryable.Select(WorkflowInstanceId.FromInstanceExpression()).ToListAsync(cancellationToken
 0147        var ids = entities.Select(x => x.Id).ToList();
 148
 0149        return Page.Of(ids, count);
 0150    }
 151
 152    /// <inheritdoc />
 153    public async ValueTask<long> DeleteAsync(WorkflowInstanceFilter filter, CancellationToken cancellationToken = defaul
 154    {
 32155        return await _store.DeleteWhereAsync(query => Filter(query, filter), cancellationToken);
 16156    }
 157
 158    public async Task UpdateUpdatedTimestampAsync(string workflowInstanceId, DateTimeOffset value, CancellationToken can
 159    {
 0160        var entity = new WorkflowInstance
 0161        {
 0162            Id = workflowInstanceId,
 0163            UpdatedAt = value
 0164        };
 165
 0166        await using var dbContext = await _store.CreateDbContextAsync(cancellationToken);
 0167        dbContext.Attach(entity);
 0168        dbContext.Entry(entity).Property(x => x.UpdatedAt).IsModified = true;
 169
 170        try
 171        {
 0172            await dbContext.SaveChangesAsync(cancellationToken);
 0173        }
 0174        catch (DbUpdateConcurrencyException e)
 175        {
 0176            foreach (var entry in e.Entries)
 177            {
 0178                var proposedValues = entry.CurrentValues;
 0179                var databaseValues = await entry.GetDatabaseValuesAsync(cancellationToken);
 180
 0181                if(databaseValues == null)
 182                    continue;
 183
 0184                var updatedAtProperty = entry.Metadata.GetProperty(nameof(WorkflowInstance.UpdatedAt));
 0185                var proposedValue = (DateTimeOffset)proposedValues[updatedAtProperty]!;
 0186                var databaseValue = (DateTimeOffset)databaseValues[updatedAtProperty]!;
 187
 0188                if (proposedValue > databaseValue)
 0189                    proposedValues[updatedAtProperty] = proposedValue;
 190
 0191                entry.OriginalValues.SetValues(databaseValues);
 0192            }
 193        }
 0194    }
 195
 196    /// <inheritdoc />
 197    /// <remarks>
 198    /// <paramref name="allowFinishedCancelled"/> is unused: drain no longer promotes Finished/Cancelled (#8419).
 199    /// The parameter remains so the 3.8.4 signature stays binary-compatible.
 200    /// </remarks>
 201    public async ValueTask<bool> TryMarkInterruptedAsync(string workflowInstanceId, CancellationToken cancellationToken 
 202    {
 0203        await using var dbContext = await _store.CreateDbContextAsync(cancellationToken);
 204        // SetTenantIdFilter's global query filter applies to ExecuteUpdateAsync unless
 205        // IgnoreQueryFilters is used. This store does not ignore it, so a tenant-B
 206        // caller cannot mark a tenant-A row.
 0207        var updated = await dbContext.WorkflowInstances
 0208            .Where(x => x.Id == workflowInstanceId && x.Status != WorkflowStatus.Finished)
 0209            .ExecuteUpdateAsync(
 0210                setters => setters
 0211                    .SetProperty(x => x.Status, WorkflowStatus.Running)
 0212                    .SetProperty(x => x.SubStatus, WorkflowSubStatus.Interrupted)
 0213                    .SetProperty(x => x.IsExecuting, false),
 0214                cancellationToken);
 215
 0216        return updated > 0;
 0217    }
 218
 219    /// <inheritdoc />
 220    [RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.SerializeAsync(WorkflowState, Can
 221    public async ValueTask SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken = default)
 222    {
 420223        await _store.SaveAsync(instance, OnSaveAsync, cancellationToken);
 420224    }
 225
 226    /// <inheritdoc />
 227    public async ValueTask AddAsync(WorkflowInstance instance, CancellationToken cancellationToken = default)
 228    {
 0229        await _store.AddAsync(instance, OnSaveAsync, cancellationToken);
 0230    }
 231
 232    /// <inheritdoc />
 233    public async ValueTask UpdateAsync(WorkflowInstance instance, CancellationToken cancellationToken = default)
 234    {
 0235        await _store.UpdateAsync(instance, OnSaveAsync, cancellationToken);
 0236    }
 237
 238    /// <inheritdoc />
 239    public async ValueTask SaveManyAsync(IEnumerable<WorkflowInstance> instances, CancellationToken cancellationToken = 
 240    {
 0241        await _store.SaveManyAsync(instances, OnSaveAsync, cancellationToken);
 0242    }
 243
 244    [RequiresUnreferencedCode("Calls Elsa.Workflows.Contracts.IWorkflowStateSerializer.SerializeAsync(WorkflowState, Can
 245    private async ValueTask OnSaveAsync(ManagementElsaDbContext managementElsaDbContext, WorkflowInstance entity, Cancel
 246    {
 420247        var data = entity.WorkflowState;
 420248        var json = _workflowStateSerializer.Serialize(data);
 420249        var compressionAlgorithm = _options.Value.CompressionAlgorithm ?? nameof(None);
 420250        var compressionCodec = _compressionCodecResolver.Resolve(compressionAlgorithm);
 420251        var compressedJson = await compressionCodec.CompressAsync(json, cancellationToken);
 252
 420253        managementElsaDbContext.Entry(entity).Property("Data").CurrentValue = compressedJson;
 420254        managementElsaDbContext.Entry(entity).Property("DataCompressionAlgorithm").CurrentValue = compressionAlgorithm;
 420255    }
 256
 257    private async ValueTask OnLoadAsync(ManagementElsaDbContext managementElsaDbContext, WorkflowInstance? entity, Cance
 258    {
 114259        if (entity == null)
 0260            return;
 261
 114262        var data = entity.WorkflowState;
 114263        var json = (string?)managementElsaDbContext.Entry(entity).Property("Data").CurrentValue;
 114264        var compressionAlgorithm = (string?)managementElsaDbContext.Entry(entity).Property("DataCompressionAlgorithm").C
 114265        var compressionStrategy = _compressionCodecResolver.Resolve(compressionAlgorithm);
 266
 267        try
 268        {
 114269            if (!string.IsNullOrWhiteSpace(json))
 270            {
 114271                json = await compressionStrategy.DecompressAsync(json, cancellationToken);
 114272                data = _workflowStateSerializer.Deserialize(json);
 273            }
 114274        }
 0275        catch (Exception exp)
 276        {
 0277            _logger.LogWarning(exp, "Exception while deserializing workflow instance state: {InstanceId}. Reverting to d
 0278        }
 114279        entity.WorkflowState = data;
 114280    }
 281
 282    private static IQueryable<WorkflowInstance> Filter(IQueryable<WorkflowInstance> query, WorkflowInstanceFilter filter
 283    {
 274284        return filter.Apply(query);
 285    }
 286}