| | | 1 | | using Elsa.Alterations.Core.Contracts; |
| | | 2 | | using Elsa.Alterations.Core.Entities; |
| | | 3 | | using Elsa.Alterations.Core.Filters; |
| | | 4 | | using Elsa.Common.Entities; |
| | | 5 | | using Elsa.Common.Multitenancy; |
| | | 6 | | using Elsa.Common.Services; |
| | | 7 | | |
| | | 8 | | namespace Elsa.Alterations.Core.Stores; |
| | | 9 | | |
| | | 10 | | /// <summary> |
| | | 11 | | /// A memory-based store for alteration jobs. |
| | | 12 | | /// </summary> |
| | | 13 | | /// <remarks> |
| | | 14 | | /// Ambient tenant is applied here rather than in callers. |
| | | 15 | | /// EF owns that via <c>SetTenantIdFilter</c> / <c>ApplyTenantId</c>; Memory must compensate. |
| | | 16 | | /// Alterations contracts have no TenantAgnostic flag, so isolation always applies (EF query filter). |
| | | 17 | | /// </remarks> |
| | | 18 | | public class MemoryAlterationJobStore : IAlterationJobStore |
| | | 19 | | { |
| | | 20 | | private readonly MemoryStore<AlterationJob> _store; |
| | | 21 | | private readonly ITenantAccessor? _tenantAccessor; |
| | | 22 | | |
| | | 23 | | /// <summary> |
| | | 24 | | /// Initializes a new instance of the <see cref="MemoryAlterationJobStore"/> class. |
| | | 25 | | /// </summary> |
| | 36 | 26 | | public MemoryAlterationJobStore(MemoryStore<AlterationJob> store, ITenantAccessor? tenantAccessor = null) |
| | | 27 | | { |
| | 36 | 28 | | _store = store; |
| | 36 | 29 | | _tenantAccessor = tenantAccessor; |
| | 36 | 30 | | } |
| | | 31 | | |
| | | 32 | | /// <inheritdoc /> |
| | | 33 | | public Task SaveAsync(AlterationJob job, CancellationToken cancellationToken = default) |
| | | 34 | | { |
| | 49 | 35 | | lock (_store.Sync) |
| | | 36 | | { |
| | 49 | 37 | | ApplyCurrentTenant(job); |
| | 49 | 38 | | EnsureIdAvailable(job); |
| | 90 | 39 | | _store.Save(job, x => x.Id); |
| | 45 | 40 | | } |
| | | 41 | | |
| | 45 | 42 | | return Task.CompletedTask; |
| | | 43 | | } |
| | | 44 | | |
| | | 45 | | /// <inheritdoc /> |
| | | 46 | | public Task SaveManyAsync(IEnumerable<AlterationJob> jobs, CancellationToken cancellationToken = default) |
| | | 47 | | { |
| | 11 | 48 | | var list = jobs.ToList(); |
| | | 49 | | |
| | 11 | 50 | | lock (_store.Sync) |
| | | 51 | | { |
| | 56 | 52 | | foreach (var job in list) |
| | 17 | 53 | | ApplyCurrentTenant(job); |
| | | 54 | | |
| | 11 | 55 | | var tenantIdsById = new Dictionary<string, string?>(StringComparer.Ordinal); |
| | 50 | 56 | | foreach (var job in list) |
| | 17 | 57 | | EnsureIdAvailable(job, tenantIdsById); |
| | | 58 | | |
| | 14 | 59 | | _store.SaveMany(list, x => x.Id); |
| | 5 | 60 | | } |
| | | 61 | | |
| | 5 | 62 | | return Task.CompletedTask; |
| | | 63 | | } |
| | | 64 | | |
| | | 65 | | /// <inheritdoc /> |
| | | 66 | | public Task<AlterationJob?> FindAsync(AlterationJobFilter filter, CancellationToken cancellationToken = default) |
| | | 67 | | { |
| | 52 | 68 | | var entity = _store.Query(query => Filter(query, filter)).FirstOrDefault(); |
| | 26 | 69 | | return Task.FromResult(entity); |
| | | 70 | | } |
| | | 71 | | |
| | | 72 | | /// <inheritdoc /> |
| | | 73 | | public Task<IEnumerable<AlterationJob>> FindManyAsync(AlterationJobFilter filter, CancellationToken cancellationToke |
| | | 74 | | { |
| | 26 | 75 | | var entities = _store.Query(query => Filter(query, filter)).ToList().AsEnumerable(); |
| | 13 | 76 | | return Task.FromResult(entities); |
| | | 77 | | } |
| | | 78 | | |
| | | 79 | | /// <inheritdoc /> |
| | | 80 | | public Task<IEnumerable<string>> FindManyIdsAsync(AlterationJobFilter filter, CancellationToken cancellationToken = |
| | | 81 | | { |
| | 14 | 82 | | var ids = _store.Query(query => Filter(query, filter)).Select(x => x.Id).ToList().AsEnumerable(); |
| | 3 | 83 | | return Task.FromResult(ids); |
| | | 84 | | } |
| | | 85 | | |
| | | 86 | | /// <inheritdoc /> |
| | | 87 | | public Task<long> CountAsync(AlterationJobFilter filter, CancellationToken cancellationToken = default) |
| | | 88 | | { |
| | 8 | 89 | | var count = _store.Query(query => Filter(query, filter)).LongCount(); |
| | 4 | 90 | | return Task.FromResult(count); |
| | | 91 | | } |
| | | 92 | | |
| | | 93 | | /// <remarks> |
| | | 94 | | /// Ambient tenant is applied here rather than in <see cref="AlterationJobFilter.Apply"/>. |
| | | 95 | | /// EF owns that via <c>SetTenantIdFilter</c>; Memory must compensate. |
| | | 96 | | /// </remarks> |
| | | 97 | | private IQueryable<AlterationJob> Filter(IQueryable<AlterationJob> query, AlterationJobFilter filter) => |
| | 46 | 98 | | filter.Apply(query.WhereVisibleToTenant(CurrentTenantId)); |
| | | 99 | | |
| | 62 | 100 | | private string CurrentTenantId => _tenantAccessor?.TenantId ?? Tenant.DefaultTenantId; |
| | | 101 | | |
| | | 102 | | private void EnsureIdAvailable(AlterationJob job, IDictionary<string, string?>? tenantIdsById = null) |
| | | 103 | | { |
| | 66 | 104 | | if (tenantIdsById?.TryGetValue(job.Id, out var stagedTenantId) == true) |
| | | 105 | | { |
| | 4 | 106 | | if (!CanReplace(stagedTenantId, job.TenantId)) |
| | 2 | 107 | | throw AlterationStoreConflict.HiddenJobId(job.Id); |
| | | 108 | | |
| | | 109 | | // A repeated ID in one batch is an update of the row established by the |
| | | 110 | | // earlier item, even when that row was absent before the batch started. |
| | 2 | 111 | | job.TenantId = stagedTenantId; |
| | 2 | 112 | | return; |
| | | 113 | | } |
| | | 114 | | |
| | 128 | 115 | | var existing = _store.Find(x => x.Id == job.Id); |
| | | 116 | | |
| | 62 | 117 | | if (existing is null) |
| | | 118 | | { |
| | 50 | 119 | | tenantIdsById?.Add(job.Id, job.TenantId); |
| | 50 | 120 | | return; |
| | | 121 | | } |
| | | 122 | | |
| | 12 | 123 | | if (!CanReplace(existing.TenantId, job.TenantId)) |
| | 8 | 124 | | throw AlterationStoreConflict.HiddenJobId(job.Id); |
| | | 125 | | |
| | | 126 | | // An accepted update may change the payload, but it must not rehome the row. |
| | 4 | 127 | | job.TenantId = existing.TenantId; |
| | 4 | 128 | | tenantIdsById?.Add(job.Id, existing.TenantId); |
| | 4 | 129 | | } |
| | | 130 | | |
| | | 131 | | /// <summary> |
| | | 132 | | /// <c>*</c> is visible to every tenant, but only an agnostic writer may replace it. |
| | | 133 | | /// Named tenants may upsert their own visible rows. |
| | | 134 | | /// </summary> |
| | | 135 | | private bool CanReplace(string? existingTenantId, string? incomingTenantId) => |
| | 16 | 136 | | TenantVisibility.CanReplaceOwnedRow(existingTenantId, incomingTenantId, CurrentTenantId); |
| | | 137 | | |
| | | 138 | | private void ApplyCurrentTenant(Entity entity) |
| | | 139 | | { |
| | 66 | 140 | | if (entity.TenantId == Tenant.AgnosticTenantId || _tenantAccessor is null) |
| | 15 | 141 | | return; |
| | | 142 | | |
| | 51 | 143 | | entity.TenantId ??= _tenantAccessor.TenantId; |
| | 51 | 144 | | } |
| | | 145 | | } |