| | | 1 | | using System.Data.Common; |
| | | 2 | | using System.Text.Json; |
| | | 3 | | using Elsa.Alterations.Core.Contracts; |
| | | 4 | | using Elsa.Alterations.Core.Entities; |
| | | 5 | | using Elsa.Alterations.Core.Filters; |
| | | 6 | | using Elsa.Alterations.Core.Models; |
| | | 7 | | using Elsa.Alterations.Core.Stores; |
| | | 8 | | using Elsa.Tenants.Options; |
| | | 9 | | using Microsoft.EntityFrameworkCore; |
| | | 10 | | using Microsoft.Extensions.DependencyInjection; |
| | | 11 | | using Microsoft.Extensions.Options; |
| | | 12 | | using Open.Linq.AsyncExtensions; |
| | | 13 | | |
| | | 14 | | namespace Elsa.Persistence.EFCore.Modules.Alterations; |
| | | 15 | | |
| | | 16 | | /// <summary> |
| | | 17 | | /// An EF Core implementation of <see cref="IAlterationJobStore"/>. |
| | | 18 | | /// </summary> |
| | | 19 | | public class EFCoreAlterationJobStore : IAlterationJobStore |
| | | 20 | | { |
| | | 21 | | private readonly EntityStore<AlterationsElsaDbContext, AlterationJob> _store; |
| | | 22 | | private readonly bool _tenantEnabled; |
| | | 23 | | |
| | | 24 | | /// <summary> |
| | | 25 | | /// Constructor. |
| | | 26 | | /// </summary> |
| | | 27 | | public EFCoreAlterationJobStore(EntityStore<AlterationsElsaDbContext, AlterationJob> store) |
| | 0 | 28 | | : this(store, Options.Create(new TenantsOptions())) |
| | | 29 | | { |
| | 0 | 30 | | } |
| | | 31 | | |
| | | 32 | | /// <summary> |
| | | 33 | | /// Constructor used by dependency injection. Direct construction through the legacy |
| | | 34 | | /// overload keeps tenancy-aware upsert disabled for compatibility. |
| | | 35 | | /// </summary> |
| | | 36 | | [ActivatorUtilitiesConstructor] |
| | 445 | 37 | | public EFCoreAlterationJobStore(EntityStore<AlterationsElsaDbContext, AlterationJob> store, IOptions<TenantsOptions> |
| | | 38 | | { |
| | 445 | 39 | | _store = store; |
| | 445 | 40 | | _tenantEnabled = tenantsOptions.Value.IsEnabled; |
| | 445 | 41 | | } |
| | | 42 | | |
| | | 43 | | /// <inheritdoc /> |
| | | 44 | | public async Task SaveAsync(AlterationJob record, CancellationToken cancellationToken = default) |
| | | 45 | | { |
| | 31 | 46 | | if (!_tenantEnabled) |
| | | 47 | | { |
| | 2 | 48 | | await _store.SaveAsync(record, OnSaveAsync, cancellationToken); |
| | 2 | 49 | | return; |
| | | 50 | | } |
| | | 51 | | |
| | 29 | 52 | | await _store.ExecuteWithDbExceptionHandlingAsync( |
| | 29 | 53 | | async () => |
| | 29 | 54 | | { |
| | 29 | 55 | | await using var dbContext = await _store.CreateDbContextAsync(cancellationToken); |
| | 29 | 56 | | await UpsertAsync(dbContext, record, cancellationToken, handleDbExceptions: false); |
| | 25 | 57 | | return true; |
| | 29 | 58 | | }, |
| | 29 | 59 | | cancellationToken, |
| | 29 | 60 | | IsDatabaseException); |
| | 52 | 61 | | } |
| | | 62 | | |
| | | 63 | | /// <inheritdoc /> |
| | | 64 | | public async Task SaveManyAsync(IEnumerable<AlterationJob> jobs, CancellationToken cancellationToken = default) |
| | | 65 | | { |
| | 14 | 66 | | if (!_tenantEnabled) |
| | | 67 | | { |
| | 1 | 68 | | await _store.SaveManyAsync(jobs, OnSaveAsync, cancellationToken); |
| | 1 | 69 | | return; |
| | | 70 | | } |
| | | 71 | | |
| | 17 | 72 | | var list = jobs.OrderBy(job => job.Id, StringComparer.Ordinal).ToList(); |
| | 13 | 73 | | if (list.Count == 0) |
| | 0 | 74 | | return; |
| | | 75 | | |
| | 13 | 76 | | await _store.ExecuteWithDbExceptionHandlingAsync( |
| | 13 | 77 | | async () => |
| | 13 | 78 | | { |
| | 13 | 79 | | await _store.ExecuteWriteWithRetryAsync(async (dbContext, ct) => |
| | 13 | 80 | | { |
| | 13 | 81 | | await using var transaction = await dbContext.Database.BeginTransactionAsync(ct); |
| | 13 | 82 | | |
| | 49 | 83 | | foreach (var job in list) |
| | 15 | 84 | | await UpsertAsync(dbContext, job, ct, handleDbExceptions: false); |
| | 13 | 85 | | |
| | 6 | 86 | | await transaction.CommitAsync(ct); |
| | 13 | 87 | | }, cancellationToken); |
| | 6 | 88 | | return true; |
| | 6 | 89 | | }, |
| | 13 | 90 | | cancellationToken, |
| | 13 | 91 | | IsDatabaseException); |
| | 13 | 92 | | } |
| | | 93 | | |
| | | 94 | | /// <inheritdoc /> |
| | | 95 | | public async Task<AlterationJob?> FindAsync(AlterationJobFilter filter, CancellationToken cancellationToken = defaul |
| | | 96 | | { |
| | 22 | 97 | | return await _store.FindAsync(filter.Apply, OnLoadAsync, cancellationToken); |
| | 22 | 98 | | } |
| | | 99 | | |
| | | 100 | | /// <inheritdoc /> |
| | | 101 | | public async Task<IEnumerable<AlterationJob>> FindManyAsync(AlterationJobFilter filter, CancellationToken cancellati |
| | | 102 | | { |
| | 9 | 103 | | return await _store.QueryAsync(filter.Apply, OnLoadAsync, cancellationToken).ToList(); |
| | 9 | 104 | | } |
| | | 105 | | |
| | | 106 | | /// <inheritdoc /> |
| | | 107 | | public async Task<IEnumerable<string>> FindManyIdsAsync(AlterationJobFilter filter, CancellationToken cancellationTo |
| | | 108 | | { |
| | 2 | 109 | | return await _store.QueryAsync(filter.Apply, x => x.Id, cancellationToken); |
| | 2 | 110 | | } |
| | | 111 | | |
| | | 112 | | /// <inheritdoc /> |
| | | 113 | | public async Task<long> CountAsync(AlterationJobFilter filter, CancellationToken cancellationToken = default) |
| | | 114 | | { |
| | 9 | 115 | | return await _store.CountAsync(queryable => Filter(queryable, filter), cancellationToken); |
| | 3 | 116 | | } |
| | | 117 | | |
| | | 118 | | private async Task UpsertAsync( |
| | | 119 | | AlterationsElsaDbContext dbContext, |
| | | 120 | | AlterationJob record, |
| | | 121 | | CancellationToken cancellationToken, |
| | | 122 | | bool handleDbExceptions = true) |
| | | 123 | | { |
| | 44 | 124 | | var ambientTenantId = AlterationTenantOwnedUpsert.AmbientTenantId(dbContext); |
| | 44 | 125 | | AlterationTenantOwnedUpsert.StampTenantId(record, ambientTenantId); |
| | 44 | 126 | | OnSave(dbContext, record); |
| | | 127 | | |
| | 44 | 128 | | var planId = record.PlanId; |
| | 44 | 129 | | var workflowInstanceId = record.WorkflowInstanceId; |
| | 44 | 130 | | var status = record.Status; |
| | 44 | 131 | | var createdAt = record.CreatedAt; |
| | 44 | 132 | | var startedAt = record.StartedAt; |
| | 44 | 133 | | var completedAt = record.CompletedAt; |
| | 44 | 134 | | var serializedLog = dbContext.Entry(record).Property<string>("SerializedLog").CurrentValue; |
| | | 135 | | |
| | 44 | 136 | | var query = dbContext.Set<AlterationJob>() |
| | 44 | 137 | | .IgnoreQueryFilters() |
| | 44 | 138 | | .Where(AlterationTenantOwnedUpsert.OwnedId<AlterationJob>(record.Id, record.TenantId, ambientTenantId)); |
| | | 139 | | |
| | | 140 | | // Inline lambda so net8/net9 bind SetPropertyCalls and net10 binds UpdateSettersBuilder. |
| | 53 | 141 | | Task<int> UpdateOwnedAsync() => query.ExecuteUpdateAsync( |
| | 53 | 142 | | setters => setters |
| | 53 | 143 | | .SetProperty(job => job.PlanId, planId) |
| | 53 | 144 | | .SetProperty(job => job.WorkflowInstanceId, workflowInstanceId) |
| | 53 | 145 | | .SetProperty(job => job.Status, status) |
| | 53 | 146 | | .SetProperty(job => job.CreatedAt, createdAt) |
| | 53 | 147 | | .SetProperty(job => job.StartedAt, startedAt) |
| | 53 | 148 | | .SetProperty(job => job.CompletedAt, completedAt) |
| | 53 | 149 | | .SetProperty(job => EF.Property<string>(job, "SerializedLog"), serializedLog), |
| | 53 | 150 | | cancellationToken); |
| | | 151 | | |
| | | 152 | | Task<TResult> ExecuteWriteAsync<TResult>(Func<Task<TResult>> operation) => |
| | 93 | 153 | | handleDbExceptions |
| | 93 | 154 | | ? _store.ExecuteWithDbExceptionHandlingAsync(operation, cancellationToken) |
| | 93 | 155 | | : operation(); |
| | | 156 | | |
| | 44 | 157 | | var updated = await ExecuteWriteAsync(UpdateOwnedAsync); |
| | | 158 | | |
| | 43 | 159 | | if (updated == 0) |
| | | 160 | | { |
| | 40 | 161 | | var inserted = await ExecuteWriteAsync( |
| | 80 | 162 | | () => AlterationTenantOwnedUpsert.InsertIfAbsentAsync(dbContext, record, cancellationToken)); |
| | 39 | 163 | | if (!inserted) |
| | | 164 | | { |
| | 9 | 165 | | var retried = await ExecuteWriteAsync(UpdateOwnedAsync); |
| | | 166 | | |
| | 9 | 167 | | if (retried == 0) |
| | 9 | 168 | | throw AlterationStoreConflict.HiddenJobId(record.Id); |
| | | 169 | | } |
| | | 170 | | } |
| | 33 | 171 | | } |
| | | 172 | | |
| | | 173 | | private static void OnSave(AlterationsElsaDbContext elsaDbContext, AlterationJob entity) |
| | | 174 | | { |
| | 48 | 175 | | elsaDbContext.Entry(entity).Property("SerializedLog").CurrentValue = JsonSerializer.Serialize(entity.Log); |
| | 48 | 176 | | } |
| | | 177 | | |
| | | 178 | | private static ValueTask OnSaveAsync(AlterationsElsaDbContext dbContext, AlterationJob entity, CancellationToken can |
| | | 179 | | { |
| | 4 | 180 | | OnSave(dbContext, entity); |
| | 4 | 181 | | return default; |
| | | 182 | | } |
| | | 183 | | |
| | | 184 | | private static ValueTask OnLoadAsync(AlterationsElsaDbContext elsaDbContext, AlterationJob? entity, CancellationToke |
| | | 185 | | { |
| | 38 | 186 | | if (entity is null) |
| | 0 | 187 | | return default; |
| | | 188 | | |
| | 38 | 189 | | var logJson = elsaDbContext.Entry(entity).Property<string>("SerializedLog").CurrentValue; |
| | 38 | 190 | | entity.Log = JsonSerializer.Deserialize<AlterationLogEntry[]>(logJson)!; |
| | | 191 | | |
| | 38 | 192 | | return default; |
| | | 193 | | } |
| | | 194 | | |
| | 6 | 195 | | private static IQueryable<AlterationJob> Filter(IQueryable<AlterationJob> queryable, AlterationJobFilter filter) => |
| | | 196 | | |
| | | 197 | | private static bool IsDatabaseException(Exception exception) => |
| | 11 | 198 | | exception is DbException or DbUpdateException || exception.InnerException is DbException; |
| | | 199 | | } |