| | | 1 | | using System.Text.Json; |
| | | 2 | | using System.Text.Json.Serialization; |
| | | 3 | | using Elsa.Persistence.EFCore; |
| | | 4 | | using Elsa.UserTasks.Contracts; |
| | | 5 | | using Elsa.UserTasks.Models; |
| | | 6 | | using Elsa.UserTasks.Persistence.EFCore.Contracts; |
| | | 7 | | using Microsoft.EntityFrameworkCore; |
| | | 8 | | |
| | | 9 | | namespace Elsa.UserTasks.Persistence.EFCore.Repositories; |
| | | 10 | | |
| | | 11 | | /// <summary> |
| | | 12 | | /// EF Core implementation of the public User Tasks repository and the low-level projection adapter. |
| | | 13 | | /// Participant identities are flattened into provider/type/id columns; no identity or foreign-key |
| | | 14 | | /// dependency is introduced by this package. |
| | | 15 | | /// </summary> |
| | 7 | 16 | | public sealed class EFCoreUserTaskRepository(Store<UserTasksElsaDbContext, UserTaskRecord> store) |
| | | 17 | | : IUserTaskRepository, IUserTaskPersistenceAdapter |
| | | 18 | | { |
| | 2 | 19 | | private static readonly JsonSerializerOptions JsonOptions = new(JsonSerializerDefaults.Web) |
| | 2 | 20 | | { |
| | 2 | 21 | | Converters = { new JsonStringEnumConverter() } |
| | 2 | 22 | | }; |
| | | 23 | | |
| | | 24 | | public async Task<UserTask?> GetAsync(string tenantId, string taskId, CancellationToken cancellationToken = default) |
| | | 25 | | { |
| | 84 | 26 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 84 | 27 | | var record = await dbContext.UserTasks.AsNoTracking().FirstOrDefaultAsync(x => x.TenantId == tenantId && x.Id == |
| | 84 | 28 | | return record is null ? null : await LoadAggregateAsync(dbContext, record, cancellationToken); |
| | 84 | 29 | | } |
| | | 30 | | |
| | | 31 | | public async Task<UserTaskQueryResult> QueryAsync(UserTaskQuery query, CancellationToken cancellationToken = default |
| | | 32 | | { |
| | 111 | 33 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 111 | 34 | | var records = BuildQuery(dbContext, query); |
| | 111 | 35 | | int? totalCount = query.IncludeTotalCount ? await records.CountAsync(cancellationToken) : null; |
| | 111 | 36 | | var limit = Math.Clamp(query.Limit <= 0 ? 50 : query.Limit, 1, 200); |
| | 111 | 37 | | var rows = await ApplyCursor(ApplyOrdering(records, query), query).Take(limit + 1).ToListAsync(cancellationToken |
| | 111 | 38 | | var hasMore = rows.Count > limit; |
| | 111 | 39 | | var page = rows.Take(limit).ToList(); |
| | | 40 | | |
| | | 41 | | // The query predicate is intentionally SQL-side, but the policy layer still needs the |
| | | 42 | | // normalized candidate/snapshot/exclusion/history relationships to calculate allowed actions. |
| | | 43 | | // Hydrate those relationships in bounded batches after paging rather than issuing one query per |
| | | 44 | | // task or allowing a summary projection with empty candidate collections. |
| | 111 | 45 | | var items = await LoadAggregatesAsync(dbContext, page, cancellationToken); |
| | 111 | 46 | | return new UserTaskQueryResult(items, hasMore ? CreateCursor(page[^1], query.Sort) : null, totalCount); |
| | 111 | 47 | | } |
| | | 48 | | |
| | | 49 | | public async Task<UserTask?> FindByMaterializationKeyAsync(string tenantId, string key, CancellationToken cancellati |
| | | 50 | | { |
| | 2 | 51 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 2 | 52 | | var record = await dbContext.UserTasks.AsNoTracking().FirstOrDefaultAsync(x => x.TenantId == tenantId && x.Mater |
| | 2 | 53 | | return record is null ? null : await LoadAggregateAsync(dbContext, record, cancellationToken); |
| | 2 | 54 | | } |
| | | 55 | | |
| | | 56 | | public async Task<UserTask?> FindByBookmarkIdAsync(string tenantId, string bookmarkId, CancellationToken cancellatio |
| | | 57 | | { |
| | 2 | 58 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 2 | 59 | | var record = await dbContext.UserTasks.AsNoTracking().FirstOrDefaultAsync(x => x.TenantId == tenantId && x.Bookm |
| | 2 | 60 | | return record is null ? null : await LoadAggregateAsync(dbContext, record, cancellationToken); |
| | 2 | 61 | | } |
| | | 62 | | |
| | | 63 | | public async Task<(UserTask Task, UserTaskInvitation Invitation)?> FindByInvitationTokenHashAsync(string tokenHash, |
| | | 64 | | { |
| | 8 | 65 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | | 66 | | // The unique index on TokenHash makes this a single seek. Deliberately not tenant-filtered: an |
| | | 67 | | // anonymous holder presents only a secret and must not be trusted to name its own tenant. |
| | 8 | 68 | | var row = await dbContext.UserTaskInvitations.AsNoTracking().FirstOrDefaultAsync(x => x.TokenHash == tokenHash, |
| | 8 | 69 | | if (row is null) |
| | 2 | 70 | | return null; |
| | | 71 | | |
| | 6 | 72 | | var record = await dbContext.UserTasks.AsNoTracking().FirstOrDefaultAsync(x => x.TenantId == row.TenantId && x.I |
| | 6 | 73 | | if (record is null) |
| | 0 | 74 | | return null; |
| | | 75 | | |
| | 6 | 76 | | var task = await LoadAggregateAsync(dbContext, record, cancellationToken); |
| | 6 | 77 | | return (task, ToInvitation(row)); |
| | 8 | 78 | | } |
| | | 79 | | |
| | | 80 | | public async Task SaveAsync(UserTask task, int expectedRevision, CancellationToken cancellationToken = default) |
| | | 81 | | { |
| | 19 | 82 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 19 | 83 | | await using var transaction = await dbContext.Database.BeginTransactionAsync(cancellationToken); |
| | 19 | 84 | | var existing = await dbContext.UserTasks.FirstOrDefaultAsync(x => x.TenantId == task.TenantId && x.Id == task.Id |
| | 19 | 85 | | EnsureExpectedRevision(existing, task.Id, expectedRevision); |
| | | 86 | | |
| | 16 | 87 | | Copy(task, existing!); |
| | 16 | 88 | | existing!.Revision = expectedRevision + 1; |
| | 16 | 89 | | existing.UpdatedAt = DateTimeOffset.UtcNow; |
| | 16 | 90 | | await ReplaceChildrenAsync(dbContext, task, cancellationToken); |
| | | 91 | | |
| | | 92 | | try |
| | | 93 | | { |
| | 16 | 94 | | await dbContext.SaveChangesAsync(cancellationToken); |
| | 16 | 95 | | await transaction.CommitAsync(cancellationToken); |
| | 16 | 96 | | task.Revision = existing.Revision; |
| | 16 | 97 | | task.UpdatedAt = existing.UpdatedAt; |
| | 16 | 98 | | } |
| | 0 | 99 | | catch (DbUpdateConcurrencyException exception) |
| | | 100 | | { |
| | 0 | 101 | | await transaction.RollbackAsync(cancellationToken); |
| | | 102 | | // Translated to the repository contract's exception so callers can distinguish a lost |
| | | 103 | | // optimistic-concurrency race from a fault without depending on EF Core. |
| | 0 | 104 | | throw new UserTaskRevisionConflictException(task.Id, expectedRevision, exception); |
| | | 105 | | } |
| | 16 | 106 | | } |
| | | 107 | | |
| | | 108 | | public async Task AddProjectionAsync(UserTask task, CancellationToken cancellationToken = default) |
| | | 109 | | { |
| | 95 | 110 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 95 | 111 | | if (await ExistsByMaterializationKeyAsync(dbContext, task.TenantId, task.MaterializationKey, cancellationToken)) |
| | | 112 | | return; |
| | | 113 | | |
| | 94 | 114 | | dbContext.UserTasks.Add(ToRecord(task)); |
| | 188 | 115 | | dbContext.UserTaskCandidates.AddRange(task.CandidateUsers.Select(x => ToCandidate(task, x, UserTaskPersistenceCa |
| | 94 | 116 | | dbContext.UserTaskCandidates.AddRange(task.CandidateGroups.Select(x => ToCandidate(task, x, UserTaskPersistenceC |
| | 94 | 117 | | dbContext.UserTaskSnapshotMembers.AddRange(task.SnapshotMembers.Select(x => ToSnapshot(task, x))); |
| | 94 | 118 | | dbContext.UserTaskSnapshotMembers.AddRange(task.SnapshotGroups.Select(x => ToSnapshot(task, x))); |
| | 96 | 119 | | dbContext.UserTaskExclusions.AddRange(task.ExcludedUsers.Select(x => ToExclusion(task, x))); |
| | 94 | 120 | | dbContext.UserTaskEvents.AddRange(task.Events.Select(ToEventRecord)); |
| | 94 | 121 | | dbContext.UserTaskOperations.AddRange(task.Operations.Select(ToOperationRecord)); |
| | 94 | 122 | | dbContext.UserTaskInvitations.AddRange(task.Invitations.Select(ToInvitationRecord)); |
| | | 123 | | try |
| | | 124 | | { |
| | 94 | 125 | | await dbContext.SaveChangesAsync(cancellationToken); |
| | 94 | 126 | | } |
| | 0 | 127 | | catch (DbUpdateException) |
| | | 128 | | { |
| | 0 | 129 | | if (!await ExistsByMaterializationKeyAsync(dbContext, task.TenantId, task.MaterializationKey, cancellationTo |
| | 0 | 130 | | throw; |
| | | 131 | | // A concurrent projection won the unique materialization-key race. Projection is idempotent. |
| | | 132 | | } |
| | 95 | 133 | | } |
| | | 134 | | |
| | | 135 | | public async Task AppendEventAsync(string tenantId, string taskId, UserTaskEvent @event, CancellationToken cancellat |
| | | 136 | | { |
| | 3 | 137 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 3 | 138 | | if (!await dbContext.UserTasks.AsNoTracking().AnyAsync(x => x.TenantId == tenantId && x.Id == taskId, cancellati |
| | | 139 | | return; |
| | | 140 | | |
| | | 141 | | // A plain insert: the aggregate row, and therefore its revision, is deliberately left untouched. |
| | 2 | 142 | | dbContext.UserTaskEvents.Add(ToEventRecord(@event)); |
| | 2 | 143 | | await dbContext.SaveChangesAsync(cancellationToken); |
| | 3 | 144 | | } |
| | | 145 | | |
| | | 146 | | public async Task<bool> TryMutateAsync(string tenantId, string taskId, int expectedRevision, Func<UserTask, bool> mu |
| | | 147 | | { |
| | 11 | 148 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 11 | 149 | | await using var transaction = await dbContext.Database.BeginTransactionAsync(cancellationToken); |
| | 11 | 150 | | var record = await dbContext.UserTasks.FirstOrDefaultAsync(x => x.TenantId == tenantId && x.Id == taskId, cancel |
| | 11 | 151 | | if (record is null || record.Revision != expectedRevision) |
| | 1 | 152 | | return false; |
| | | 153 | | |
| | 10 | 154 | | var task = await LoadAggregateAsync(dbContext, record, cancellationToken); |
| | 10 | 155 | | if (!mutation(task)) |
| | 1 | 156 | | return false; |
| | | 157 | | |
| | 9 | 158 | | Copy(task, record); |
| | 9 | 159 | | record.Revision = expectedRevision + 1; |
| | 9 | 160 | | record.UpdatedAt = DateTimeOffset.UtcNow; |
| | 9 | 161 | | await ReplaceChildrenAsync(dbContext, task, cancellationToken); |
| | | 162 | | |
| | | 163 | | try |
| | | 164 | | { |
| | 9 | 165 | | await dbContext.SaveChangesAsync(cancellationToken); |
| | 9 | 166 | | await transaction.CommitAsync(cancellationToken); |
| | 9 | 167 | | return true; |
| | | 168 | | } |
| | | 169 | | catch (DbUpdateConcurrencyException) |
| | | 170 | | { |
| | 0 | 171 | | await transaction.RollbackAsync(cancellationToken); |
| | 0 | 172 | | return false; |
| | | 173 | | } |
| | 11 | 174 | | } |
| | | 175 | | |
| | | 176 | | async Task<UserTaskRecord?> IUserTaskPersistenceAdapter.GetAsync(string tenantId, string taskId, CancellationToken c |
| | | 177 | | { |
| | 0 | 178 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 0 | 179 | | return await dbContext.UserTasks.AsNoTracking().FirstOrDefaultAsync(x => x.TenantId == tenantId && x.Id == taskI |
| | 0 | 180 | | } |
| | | 181 | | |
| | | 182 | | public async Task<IReadOnlyCollection<UserTaskRecord>> QueryAsync(UserTaskPersistenceQuery query, CancellationToken |
| | | 183 | | { |
| | 0 | 184 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 0 | 185 | | var records = dbContext.UserTasks.AsNoTracking().Where(x => x.TenantId == query.TenantId); |
| | 0 | 186 | | if (query.Statuses.Count > 0) |
| | 0 | 187 | | records = records.Where(x => query.Statuses.Contains(x.Status)); |
| | 0 | 188 | | if (!string.IsNullOrWhiteSpace(query.AssigneeProvider)) |
| | 0 | 189 | | records = records.Where(x => x.AssigneeProvider == query.AssigneeProvider); |
| | 0 | 190 | | if (!string.IsNullOrWhiteSpace(query.AssigneeType)) |
| | 0 | 191 | | records = records.Where(x => x.AssigneeType == query.AssigneeType); |
| | 0 | 192 | | if (!string.IsNullOrWhiteSpace(query.AssigneeId)) |
| | 0 | 193 | | records = records.Where(x => x.AssigneeId == query.AssigneeId); |
| | 0 | 194 | | if (!string.IsNullOrWhiteSpace(query.Search)) |
| | | 195 | | { |
| | 0 | 196 | | var search = query.Search.Trim(); |
| | 0 | 197 | | records = records.Where(x => x.Title.Contains(search) || (x.Summary != null && x.Summary.Contains(search)) | |
| | | 198 | | } |
| | | 199 | | |
| | 0 | 200 | | var limit = query.Limit is > 0 ? Math.Min(query.Limit.Value, 200) : 100; |
| | 0 | 201 | | return await records.OrderBy(x => x.DueAt == null).ThenBy(x => x.DueAt).ThenByDescending(x => x.Priority).ThenBy |
| | 0 | 202 | | } |
| | | 203 | | |
| | | 204 | | public async Task<bool> TryAddProjectionAsync(UserTaskRecord task, CancellationToken cancellationToken = default) |
| | | 205 | | { |
| | 0 | 206 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 0 | 207 | | if (await ExistsByMaterializationKeyAsync(dbContext, task.TenantId, task.MaterializationKey, cancellationToken)) |
| | 0 | 208 | | return false; |
| | 0 | 209 | | dbContext.UserTasks.Add(task); |
| | | 210 | | try |
| | | 211 | | { |
| | 0 | 212 | | await dbContext.SaveChangesAsync(cancellationToken); |
| | 0 | 213 | | return true; |
| | | 214 | | } |
| | 0 | 215 | | catch (DbUpdateException) |
| | | 216 | | { |
| | 0 | 217 | | if (!await ExistsByMaterializationKeyAsync(dbContext, task.TenantId, task.MaterializationKey, cancellationTo |
| | 0 | 218 | | throw; |
| | 0 | 219 | | return false; |
| | | 220 | | } |
| | 0 | 221 | | } |
| | | 222 | | |
| | | 223 | | public async Task<bool> TrySaveAsync(UserTaskRecord task, int expectedRevision, CancellationToken cancellationToken |
| | | 224 | | { |
| | 0 | 225 | | await using var dbContext = await store.CreateDbContextAsync(cancellationToken); |
| | 0 | 226 | | var existing = await dbContext.UserTasks.FirstOrDefaultAsync(x => x.TenantId == task.TenantId && x.Id == task.Id |
| | 0 | 227 | | if (existing is null || existing.Revision != expectedRevision) |
| | 0 | 228 | | return false; |
| | 0 | 229 | | Copy(task, existing); |
| | 0 | 230 | | existing.Revision = expectedRevision + 1; |
| | 0 | 231 | | existing.UpdatedAt = DateTimeOffset.UtcNow; |
| | | 232 | | try |
| | | 233 | | { |
| | 0 | 234 | | await dbContext.SaveChangesAsync(cancellationToken); |
| | 0 | 235 | | return true; |
| | | 236 | | } |
| | 0 | 237 | | catch (DbUpdateConcurrencyException) |
| | | 238 | | { |
| | 0 | 239 | | return false; |
| | | 240 | | } |
| | 0 | 241 | | } |
| | | 242 | | |
| | | 243 | | private static IQueryable<UserTaskRecord> BuildQuery(UserTasksElsaDbContext dbContext, UserTaskQuery query) |
| | | 244 | | { |
| | 111 | 245 | | var records = dbContext.UserTasks.AsNoTracking(); |
| | 111 | 246 | | records = records.Where(x => x.TenantId == query.TenantId); |
| | 111 | 247 | | if (query.Statuses.Count > 0) |
| | | 248 | | { |
| | 0 | 249 | | var statuses = query.Statuses.ToArray(); |
| | 0 | 250 | | records = records.Where(x => statuses.Contains(x.Status)); |
| | | 251 | | } |
| | 111 | 252 | | if (query.OnlyOverdue) |
| | 0 | 253 | | records = records.Where(x => x.IsOverdue); |
| | 111 | 254 | | if (query.OnlyWithoutDueDate) |
| | 0 | 255 | | records = records.Where(x => x.DueAt == null); |
| | 111 | 256 | | if (!string.IsNullOrWhiteSpace(query.TaskType)) |
| | 0 | 257 | | records = records.Where(x => x.TaskType == query.TaskType); |
| | 111 | 258 | | if (query.PriorityFrom is not null) |
| | 0 | 259 | | records = records.Where(x => x.Priority >= query.PriorityFrom); |
| | 111 | 260 | | if (query.PriorityTo is not null) |
| | 0 | 261 | | records = records.Where(x => x.Priority <= query.PriorityTo); |
| | 111 | 262 | | if (query.DueFrom is not null) |
| | 0 | 263 | | records = records.Where(x => x.DueAt >= query.DueFrom); |
| | 111 | 264 | | if (query.DueTo is not null) |
| | 0 | 265 | | records = records.Where(x => x.DueAt <= query.DueTo); |
| | 111 | 266 | | if (!string.IsNullOrWhiteSpace(query.WorkflowDefinitionId)) |
| | 0 | 267 | | records = records.Where(x => x.WorkflowDefinitionId == query.WorkflowDefinitionId); |
| | 111 | 268 | | if (!string.IsNullOrWhiteSpace(query.WorkflowInstanceId)) |
| | 0 | 269 | | records = records.Where(x => x.WorkflowInstanceId == query.WorkflowInstanceId); |
| | 111 | 270 | | if (!string.IsNullOrWhiteSpace(query.Reference)) |
| | 0 | 271 | | records = records.Where(x => x.Reference == query.Reference); |
| | 111 | 272 | | if (!string.IsNullOrWhiteSpace(query.Search)) |
| | | 273 | | { |
| | 0 | 274 | | var search = query.Search.Trim(); |
| | 0 | 275 | | records = records.Where(x => x.Title.Contains(search) || (x.Summary != null && x.Summary.Contains(search)) | |
| | | 276 | | } |
| | 111 | 277 | | if (query.Scope is { } requestedScope && |
| | 111 | 278 | | (!string.Equals(requestedScope.TenantId, query.TenantId, StringComparison.Ordinal) || |
| | 111 | 279 | | !string.Equals(requestedScope.Subject.TenantId, query.TenantId, StringComparison.Ordinal) || |
| | 111 | 280 | | requestedScope.Groups.Any(group => !string.Equals(group.TenantId, query.TenantId, StringComparison.Ordinal) |
| | 2 | 281 | | return records.Where(_ => false); |
| | | 282 | | |
| | 109 | 283 | | if (query.Scope is not { } scope) |
| | 0 | 284 | | return records; |
| | | 285 | | |
| | 109 | 286 | | var subject = scope.Subject; |
| | 109 | 287 | | var subjectType = subject.Type.ToString(); |
| | 109 | 288 | | var groupKeys = scope.Groups.Select(GetParticipantKey).Distinct(StringComparer.Ordinal).ToArray(); |
| | 109 | 289 | | var tenantId = query.TenantId; |
| | | 290 | | |
| | 109 | 291 | | if (scope.ExcludeBlocking) |
| | 109 | 292 | | records = records.Where(task => task.HealthSeverity != UserTaskHealthSeverity.Blocking); |
| | | 293 | | |
| | | 294 | | // Manager-only scopes were already rejected by the policy for non-managers. Reaching them here means |
| | | 295 | | // the caller manages the tenant, so the tenant predicate above is the whole authorization. |
| | 109 | 296 | | if (scope.RequiresManager) |
| | | 297 | | { |
| | 0 | 298 | | if (!scope.IsManager) |
| | 0 | 299 | | return records.Where(_ => false); |
| | 0 | 300 | | if (scope.Kind == UserTaskQueryScopeKind.NeedsAttention) |
| | | 301 | | { |
| | 0 | 302 | | records = records.Where(task => |
| | 0 | 303 | | task.HealthSeverity == UserTaskHealthSeverity.Blocking |
| | 0 | 304 | | || task.IsOverdue |
| | 0 | 305 | | || (task.AssigneeId == null && task.Status != UserTaskStatus.Completed && task.Status != UserTaskSta |
| | 0 | 306 | | || task.Status == UserTaskStatus.Completing |
| | 0 | 307 | | || task.Status == UserTaskStatus.TimingOut |
| | 0 | 308 | | || task.Status == UserTaskStatus.Cancelling); |
| | | 309 | | } |
| | 0 | 310 | | return records; |
| | | 311 | | } |
| | | 312 | | |
| | | 313 | | // Correlated predicates are applied in the query path so tenant, eligibility, and exclusion are |
| | | 314 | | // evaluated before totals, cursors, and page limits — an unauthorized row can never reach the page. |
| | 109 | 315 | | return scope.Kind switch |
| | 109 | 316 | | { |
| | 1 | 317 | | UserTaskQueryScopeKind.Assigned => records.Where(task => |
| | 1 | 318 | | task.AssigneeProvider == subject.Provider && task.AssigneeType == subjectType && task.AssigneeId == subj |
| | 109 | 319 | | |
| | 108 | 320 | | UserTaskQueryScopeKind.Available => records.Where(task => |
| | 108 | 321 | | task.AssigneeId == null |
| | 108 | 322 | | && task.Status != UserTaskStatus.Completed && task.Status != UserTaskStatus.TimedOut && task.Status != U |
| | 108 | 323 | | && !dbContext.UserTaskExclusions.Any(exclusion => |
| | 108 | 324 | | exclusion.TenantId == tenantId && exclusion.TaskId == task.Id && |
| | 108 | 325 | | exclusion.ParticipantType == UserTaskParticipantType.User && exclusion.Provider == subject.Provider |
| | 108 | 326 | | && (dbContext.UserTaskCandidates.Any(candidate => |
| | 108 | 327 | | candidate.TenantId == tenantId && candidate.TaskId == task.Id && |
| | 108 | 328 | | candidate.ParticipantType == UserTaskParticipantType.User && candidate.Provider == subject.Provi |
| | 108 | 329 | | || (groupKeys.Length > 0 && dbContext.UserTaskCandidates.Any(candidate => |
| | 108 | 330 | | candidate.TenantId == tenantId && candidate.TaskId == task.Id && |
| | 108 | 331 | | candidate.ParticipantType == UserTaskParticipantType.Group && groupKeys.Contains(candidate.Parti |
| | 108 | 332 | | || dbContext.UserTaskSnapshotMembers.Any(member => |
| | 108 | 333 | | member.TenantId == tenantId && member.TaskId == task.Id && |
| | 108 | 334 | | ((member.ParticipantType == UserTaskParticipantType.User && member.Provider == subject.Provider |
| | 108 | 335 | | || (member.ParticipantType == UserTaskParticipantType.Group && groupKeys.Contains(member.Partic |
| | 109 | 336 | | |
| | 0 | 337 | | UserTaskQueryScopeKind.History => records.Where(task => |
| | 0 | 338 | | (task.Status == UserTaskStatus.Completed || task.Status == UserTaskStatus.TimedOut || task.Status == Use |
| | 0 | 339 | | && dbContext.UserTaskEvents.Any(@event => |
| | 0 | 340 | | @event.TenantId == tenantId && @event.TaskId == task.Id && |
| | 0 | 341 | | @event.ActorProvider == subject.Provider && @event.ActorType == subjectType && @event.ActorId == sub |
| | 109 | 342 | | |
| | 0 | 343 | | _ => records.Where(_ => false) |
| | 109 | 344 | | }; |
| | | 345 | | } |
| | | 346 | | |
| | | 347 | | private static IQueryable<UserTaskRecord> ApplyOrdering(IQueryable<UserTaskRecord> records, UserTaskQuery query) |
| | | 348 | | { |
| | 111 | 349 | | return query.Sort.ToLowerInvariant() switch |
| | 111 | 350 | | { |
| | 39 | 351 | | "due" when query.Descending => records.OrderBy(x => x.DueAt == null).ThenByDescending(x => x.DueAt).ThenBy(x |
| | 13 | 352 | | "due" => records.OrderBy(x => x.DueAt == null).ThenBy(x => x.DueAt).ThenBy(x => x.Id), |
| | 36 | 353 | | "priority" when query.Descending => records.OrderByDescending(x => x.Priority).ThenBy(x => x.Id), |
| | 12 | 354 | | "priority" => records.OrderBy(x => x.Priority).ThenBy(x => x.Id), |
| | 36 | 355 | | "title" when query.Descending => records.OrderByDescending(x => x.Title).ThenBy(x => x.Id), |
| | 12 | 356 | | "title" => records.OrderBy(x => x.Title).ThenBy(x => x.Id), |
| | 49 | 357 | | _ when query.Descending => records.OrderByDescending(x => x.CreatedAt).ThenBy(x => x.Id), |
| | 25 | 358 | | _ => records.OrderBy(x => x.CreatedAt).ThenBy(x => x.Id) |
| | 111 | 359 | | }; |
| | | 360 | | } |
| | | 361 | | |
| | | 362 | | private static IQueryable<UserTaskRecord> ApplyCursor(IQueryable<UserTaskRecord> records, UserTaskQuery query) |
| | | 363 | | { |
| | 111 | 364 | | if (string.IsNullOrWhiteSpace(query.Cursor) || !TryReadCursor(query.Cursor, out var cursorValue, out var cursorI |
| | 45 | 365 | | return records; |
| | | 366 | | |
| | | 367 | | return query.Sort.ToLowerInvariant() switch |
| | | 368 | | { |
| | 32 | 369 | | "priority" when int.TryParse(cursorValue, out var priority) => query.Descending |
| | 16 | 370 | | ? records.Where(x => x.Priority < priority || (x.Priority == priority && string.Compare(x.Id, cursorId) |
| | 16 | 371 | | : records.Where(x => x.Priority > priority || (x.Priority == priority && string.Compare(x.Id, cursorId) |
| | 16 | 372 | | "title" => query.Descending |
| | 16 | 373 | | ? records.Where(x => string.Compare(x.Title, cursorValue) < 0 || (x.Title == cursorValue && string.Compa |
| | 16 | 374 | | : records.Where(x => string.Compare(x.Title, cursorValue) > 0 || (x.Title == cursorValue && string.Compa |
| | 18 | 375 | | "due" when cursorValue == "~null" => records.Where(x => x.DueAt == null && string.Compare(x.Id, cursorId) > |
| | 28 | 376 | | "due" when DateTimeOffset.TryParse(cursorValue, out var dueAt) => query.Descending |
| | 14 | 377 | | ? records.Where(x => x.DueAt == null || x.DueAt < dueAt || (x.DueAt == dueAt && string.Compare(x.Id, cur |
| | 14 | 378 | | : records.Where(x => x.DueAt == null || x.DueAt > dueAt || (x.DueAt == dueAt && string.Compare(x.Id, cur |
| | 36 | 379 | | _ when DateTimeOffset.TryParse(cursorValue, out var createdAt) => query.Descending |
| | 18 | 380 | | ? records.Where(x => x.CreatedAt < createdAt || (x.CreatedAt == createdAt && string.Compare(x.Id, cursor |
| | 18 | 381 | | : records.Where(x => x.CreatedAt > createdAt || (x.CreatedAt == createdAt && string.Compare(x.Id, cursor |
| | 0 | 382 | | _ => records |
| | | 383 | | }; |
| | | 384 | | } |
| | | 385 | | |
| | | 386 | | private static string CreateCursor(UserTaskRecord record, string sort) |
| | | 387 | | { |
| | 67 | 388 | | var value = sort.ToLowerInvariant() switch |
| | 67 | 389 | | { |
| | 16 | 390 | | "priority" => record.Priority.ToString(System.Globalization.CultureInfo.InvariantCulture), |
| | 16 | 391 | | "title" => record.Title, |
| | 16 | 392 | | "due" => record.DueAt?.ToString("O", System.Globalization.CultureInfo.InvariantCulture) ?? "~null", |
| | 19 | 393 | | _ => record.CreatedAt.ToString("O", System.Globalization.CultureInfo.InvariantCulture) |
| | 67 | 394 | | }; |
| | 67 | 395 | | return Convert.ToBase64String(JsonSerializer.SerializeToUtf8Bytes(new[] { value, record.Id }, JsonOptions)) |
| | 67 | 396 | | .TrimEnd('=').Replace('+', '-').Replace('/', '_'); |
| | | 397 | | } |
| | | 398 | | |
| | | 399 | | private static bool TryReadCursor(string cursor, out string value, out string id) |
| | | 400 | | { |
| | 67 | 401 | | value = id = ""; |
| | | 402 | | try |
| | | 403 | | { |
| | 67 | 404 | | var padded = cursor.Replace('-', '+').Replace('_', '/') + new string('=', (4 - cursor.Length % 4) % 4); |
| | 67 | 405 | | var parts = JsonSerializer.Deserialize<string[]>(Convert.FromBase64String(padded), JsonOptions); |
| | 66 | 406 | | if (parts is not [var cursorValue, var cursorId] || string.IsNullOrWhiteSpace(cursorId)) |
| | 0 | 407 | | return false; |
| | 66 | 408 | | value = cursorValue; |
| | 66 | 409 | | id = cursorId; |
| | 66 | 410 | | return true; |
| | | 411 | | } |
| | 0 | 412 | | catch (FormatException) |
| | | 413 | | { |
| | 0 | 414 | | return false; |
| | | 415 | | } |
| | 1 | 416 | | catch (JsonException) |
| | | 417 | | { |
| | 1 | 418 | | return false; |
| | | 419 | | } |
| | 67 | 420 | | } |
| | | 421 | | |
| | | 422 | | private static async Task<UserTask> LoadAggregateAsync(UserTasksElsaDbContext dbContext, UserTaskRecord record, Canc |
| | | 423 | | { |
| | 99 | 424 | | return (await LoadAggregatesAsync(dbContext, [record], cancellationToken))[0]; |
| | 99 | 425 | | } |
| | | 426 | | |
| | | 427 | | private static async Task<List<UserTask>> LoadAggregatesAsync(UserTasksElsaDbContext dbContext, IReadOnlyCollection< |
| | | 428 | | { |
| | 210 | 429 | | var tasks = records.Select(ToModel).ToList(); |
| | 210 | 430 | | if (tasks.Count == 0) |
| | 3 | 431 | | return tasks; |
| | | 432 | | |
| | 207 | 433 | | var tenantId = records.First().TenantId; |
| | 533 | 434 | | var taskIds = tasks.Select(x => x.Id).ToArray(); |
| | 533 | 435 | | var taskById = tasks.ToDictionary(x => x.Id, StringComparer.Ordinal); |
| | 207 | 436 | | var candidates = await dbContext.UserTaskCandidates.AsNoTracking().Where(x => x.TenantId == tenantId && taskIds. |
| | 1392 | 437 | | foreach (var (task, group) in GroupByLoadedTask(candidates, x => x.TaskId, taskById)) |
| | | 438 | | { |
| | 652 | 439 | | task.CandidateUsers = group.Where(x => x.ParticipantType == UserTaskParticipantType.User).Select(ToParticipa |
| | 652 | 440 | | task.CandidateGroups = group.Where(x => x.ParticipantType == UserTaskParticipantType.Group).Select(ToPartici |
| | | 441 | | } |
| | | 442 | | |
| | 207 | 443 | | var snapshots = await dbContext.UserTaskSnapshotMembers.AsNoTracking().Where(x => x.TenantId == tenantId && task |
| | 414 | 444 | | foreach (var (task, group) in GroupByLoadedTask(snapshots, x => x.TaskId, taskById)) |
| | | 445 | | { |
| | 0 | 446 | | task.SnapshotMembers = group.Where(x => x.ParticipantType == UserTaskParticipantType.User).Select(ToParticip |
| | 0 | 447 | | task.SnapshotGroups = group.Where(x => x.ParticipantType == UserTaskParticipantType.Group).Select(ToParticip |
| | | 448 | | } |
| | | 449 | | |
| | 207 | 450 | | var exclusions = await dbContext.UserTaskExclusions.AsNoTracking().Where(x => x.TenantId == tenantId && taskIds. |
| | 414 | 451 | | foreach (var (task, group) in GroupByLoadedTask(exclusions, x => x.TaskId, taskById)) |
| | 0 | 452 | | task.ExcludedUsers = group.Select(ToParticipant).ToList(); |
| | | 453 | | |
| | 207 | 454 | | var events = await dbContext.UserTaskEvents.AsNoTracking().Where(x => x.TenantId == tenantId && taskIds.Contains |
| | 603 | 455 | | foreach (var (task, group) in GroupByLoadedTask(events, x => x.TaskId, taskById)) |
| | 49 | 456 | | task.Events = group.Select(ToEvent).ToList(); |
| | | 457 | | |
| | 207 | 458 | | var operations = await dbContext.UserTaskOperations.AsNoTracking().Where(x => x.TenantId == tenantId && taskIds. |
| | 438 | 459 | | foreach (var (task, group) in GroupByLoadedTask(operations, x => x.TaskId, taskById)) |
| | 7 | 460 | | task.Operations = group.Select(ToOperation).ToList(); |
| | | 461 | | |
| | 207 | 462 | | var invitations = await dbContext.UserTaskInvitations.AsNoTracking().Where(x => x.TenantId == tenantId && taskId |
| | 540 | 463 | | foreach (var (task, group) in GroupByLoadedTask(invitations, x => x.TaskId, taskById)) |
| | 42 | 464 | | task.Invitations = group.Select(ToInvitation).ToList(); |
| | | 465 | | |
| | 207 | 466 | | return tasks; |
| | 210 | 467 | | } |
| | | 468 | | |
| | | 469 | | /// <summary> |
| | | 470 | | /// Pairs each group of child rows with the task it belongs to, dropping groups whose task is not on the |
| | | 471 | | /// current page. The filter is explicit and each group still costs a single dictionary probe. |
| | | 472 | | /// </summary> |
| | | 473 | | private static IEnumerable<(UserTask Task, IGrouping<string, TRow> Rows)> GroupByLoadedTask<TRow>( |
| | | 474 | | IEnumerable<TRow> rows, Func<TRow, string> taskIdSelector, Dictionary<string, UserTask> taskById) => |
| | 1242 | 475 | | rows.GroupBy(taskIdSelector) |
| | 424 | 476 | | .Select(group => (Task: taskById.GetValueOrDefault(group.Key), Rows: group)) |
| | 424 | 477 | | .Where(pair => pair.Task is not null) |
| | 1666 | 478 | | .Select(pair => (pair.Task!, pair.Rows)); |
| | | 479 | | |
| | 326 | 480 | | private static UserTask ToModel(UserTaskRecord record) => new() |
| | 326 | 481 | | { |
| | 326 | 482 | | Id = record.Id, |
| | 326 | 483 | | TenantId = record.TenantId, |
| | 326 | 484 | | WorkflowDefinitionId = record.WorkflowDefinitionId, |
| | 326 | 485 | | WorkflowDefinitionName = record.WorkflowDefinitionName, |
| | 326 | 486 | | WorkflowDefinitionVersion = record.WorkflowDefinitionVersion, |
| | 326 | 487 | | WorkflowInstanceId = record.WorkflowInstanceId, |
| | 326 | 488 | | WorkflowInstanceReference = record.WorkflowInstanceReference, |
| | 326 | 489 | | ActivityInstanceId = record.ActivityInstanceId, |
| | 326 | 490 | | BookmarkId = record.BookmarkId, |
| | 326 | 491 | | MaterializationKey = record.MaterializationKey, |
| | 326 | 492 | | Title = record.Title, |
| | 326 | 493 | | Summary = record.Summary, |
| | 326 | 494 | | Reference = record.Reference, |
| | 326 | 495 | | Tags = Deserialize<HashSet<string>>(record.TagsJson) ?? new(StringComparer.OrdinalIgnoreCase), |
| | 326 | 496 | | TaskType = record.TaskType, |
| | 326 | 497 | | Requester = ToParticipant(record.RequesterProvider, record.RequesterType, record.RequesterId, record.RequesterDi |
| | 326 | 498 | | Assignee = ToParticipant(record.AssigneeProvider, record.AssigneeType, record.AssigneeId, record.AssigneeDisplay |
| | 326 | 499 | | MembershipResolutionMode = record.MembershipResolutionMode ?? UserTaskMembershipResolutionMode.Live, |
| | 326 | 500 | | AllowManagerExclusionOverride = record.AllowManagerExclusionOverride, |
| | 326 | 501 | | Priority = record.Priority, |
| | 326 | 502 | | DueAt = record.DueAt, |
| | 326 | 503 | | IsOverdue = record.IsOverdue, |
| | 326 | 504 | | Instructions = Deserialize<string>(record.InstructionsJson), |
| | 326 | 505 | | TaskData = DeserializeJson(record.TaskDataJson), |
| | 326 | 506 | | RequestedForm = Deserialize<UserTaskFormReference>(record.FormReferenceJson), |
| | 326 | 507 | | PinnedForm = Deserialize<ResolvedUserTaskForm>(record.PinnedFormJson), |
| | 326 | 508 | | Actions = Deserialize<List<UserTaskAction>>(record.ActionsJson) ?? [new UserTaskAction("Complete", "Complete")], |
| | 326 | 509 | | InvitationDefinitions = Deserialize<List<UserTaskInvitationDefinition>>(record.InvitationDefinitionsJson) ?? [], |
| | 326 | 510 | | EnableTimeoutOutcome = record.TimeoutEnabled, |
| | 326 | 511 | | EnableCancellationOutcome = record.CancellationEnabled, |
| | 326 | 512 | | Status = record.Status, |
| | 326 | 513 | | HealthSeverity = record.HealthSeverity, |
| | 326 | 514 | | HealthCode = record.HealthCode, |
| | 326 | 515 | | HealthMessage = record.HealthMessage, |
| | 326 | 516 | | CompletionActionKey = record.CompletionActionKey, |
| | 326 | 517 | | CompletionData = DeserializeJson(record.CompletionDataJson), |
| | 326 | 518 | | CompletedBy = Deserialize<ParticipantReference>(record.CompletionActorJson), |
| | 326 | 519 | | CreatedAt = record.CreatedAt, |
| | 326 | 520 | | UpdatedAt = record.UpdatedAt, |
| | 326 | 521 | | AssignedAt = record.AssignedAt, |
| | 326 | 522 | | CompletedAt = record.CompletedAt, |
| | 326 | 523 | | Revision = record.Revision |
| | 326 | 524 | | }; |
| | | 525 | | |
| | | 526 | | private static UserTaskRecord ToRecord(UserTask task) |
| | | 527 | | { |
| | 94 | 528 | | var record = new UserTaskRecord { Id = string.IsNullOrWhiteSpace(task.Id) ? Guid.NewGuid().ToString("N") : task. |
| | 94 | 529 | | Copy(task, record); |
| | 94 | 530 | | return record; |
| | | 531 | | } |
| | | 532 | | |
| | | 533 | | private static void Copy(UserTask source, UserTaskRecord target) |
| | | 534 | | { |
| | 119 | 535 | | target.Id = source.Id; |
| | 119 | 536 | | target.TenantId = source.TenantId; |
| | 119 | 537 | | target.WorkflowDefinitionId = source.WorkflowDefinitionId; |
| | 119 | 538 | | target.WorkflowDefinitionName = source.WorkflowDefinitionName; |
| | 119 | 539 | | target.WorkflowDefinitionVersion = source.WorkflowDefinitionVersion; |
| | 119 | 540 | | target.WorkflowInstanceId = source.WorkflowInstanceId; |
| | 119 | 541 | | target.WorkflowInstanceReference = source.WorkflowInstanceReference; |
| | 119 | 542 | | target.ActivityInstanceId = source.ActivityInstanceId; |
| | 119 | 543 | | target.BookmarkId = source.BookmarkId; |
| | 119 | 544 | | target.MaterializationKey = source.MaterializationKey; |
| | 119 | 545 | | target.Title = source.Title; |
| | 119 | 546 | | target.Summary = source.Summary; |
| | 119 | 547 | | target.Reference = source.Reference; |
| | 119 | 548 | | target.TaskType = source.TaskType; |
| | 119 | 549 | | target.TagsJson = Serialize(source.Tags); |
| | 119 | 550 | | target.RequesterProvider = source.Requester?.Provider; |
| | 119 | 551 | | target.RequesterType = source.Requester?.Type.ToString(); |
| | 119 | 552 | | target.RequesterId = source.Requester?.Id; |
| | 119 | 553 | | target.RequesterDisplayName = source.Requester?.DisplayName; |
| | 119 | 554 | | target.Priority = source.Priority; |
| | 119 | 555 | | target.DueAt = source.DueAt; |
| | 119 | 556 | | target.IsOverdue = source.IsOverdue; |
| | 119 | 557 | | target.Status = source.Status; |
| | 119 | 558 | | target.TimeoutEnabled = source.EnableTimeoutOutcome; |
| | 119 | 559 | | target.CancellationEnabled = source.EnableCancellationOutcome; |
| | 119 | 560 | | target.AllowManagerExclusionOverride = source.AllowManagerExclusionOverride; |
| | 119 | 561 | | target.MembershipResolutionMode = source.MembershipResolutionMode; |
| | 119 | 562 | | target.AssigneeProvider = source.Assignee?.Provider; |
| | 119 | 563 | | target.AssigneeType = source.Assignee?.Type.ToString(); |
| | 119 | 564 | | target.AssigneeId = source.Assignee?.Id; |
| | 119 | 565 | | target.AssigneeDisplayName = source.Assignee?.DisplayName; |
| | 119 | 566 | | target.InstructionsJson = Serialize(source.Instructions); |
| | 119 | 567 | | target.TaskDataJson = Serialize(source.TaskData); |
| | 119 | 568 | | target.FormReferenceJson = Serialize(source.RequestedForm); |
| | 119 | 569 | | target.PinnedFormJson = Serialize(source.PinnedForm); |
| | 119 | 570 | | target.ActionsJson = Serialize(source.Actions); |
| | 119 | 571 | | target.InvitationDefinitionsJson = Serialize(source.InvitationDefinitions); |
| | 119 | 572 | | target.HealthIssuesJson = Serialize(new { source.HealthSeverity, source.HealthCode, source.HealthMessage }); |
| | 119 | 573 | | target.HealthSeverity = source.HealthSeverity; |
| | 119 | 574 | | target.HealthCode = source.HealthCode; |
| | 119 | 575 | | target.HealthMessage = source.HealthMessage; |
| | 119 | 576 | | target.CompletionActionKey = source.CompletionActionKey; |
| | 119 | 577 | | target.CompletionDataJson = Serialize(source.CompletionData); |
| | 119 | 578 | | target.CompletionActorJson = Serialize(source.CompletedBy); |
| | 119 | 579 | | target.CreatedAt = source.CreatedAt; |
| | 119 | 580 | | target.UpdatedAt = source.UpdatedAt; |
| | 119 | 581 | | target.AssignedAt = source.AssignedAt; |
| | 119 | 582 | | target.CompletedAt = source.CompletedAt; |
| | 119 | 583 | | target.Revision = source.Revision; |
| | 119 | 584 | | } |
| | | 585 | | |
| | | 586 | | private static void Copy(UserTaskRecord source, UserTaskRecord target) |
| | | 587 | | { |
| | 0 | 588 | | target.TenantId = source.TenantId; |
| | 0 | 589 | | target.WorkflowDefinitionId = source.WorkflowDefinitionId; |
| | 0 | 590 | | target.WorkflowDefinitionName = source.WorkflowDefinitionName; |
| | 0 | 591 | | target.WorkflowDefinitionVersion = source.WorkflowDefinitionVersion; |
| | 0 | 592 | | target.WorkflowInstanceId = source.WorkflowInstanceId; |
| | 0 | 593 | | target.WorkflowInstanceReference = source.WorkflowInstanceReference; |
| | 0 | 594 | | target.ActivityInstanceId = source.ActivityInstanceId; |
| | 0 | 595 | | target.BookmarkId = source.BookmarkId; |
| | 0 | 596 | | target.MaterializationKey = source.MaterializationKey; |
| | 0 | 597 | | target.Title = source.Title; |
| | 0 | 598 | | target.Summary = source.Summary; |
| | 0 | 599 | | target.Reference = source.Reference; |
| | 0 | 600 | | target.TaskType = source.TaskType; |
| | 0 | 601 | | target.TagsJson = source.TagsJson; |
| | 0 | 602 | | target.RequesterProvider = source.RequesterProvider; |
| | 0 | 603 | | target.RequesterType = source.RequesterType; |
| | 0 | 604 | | target.RequesterId = source.RequesterId; |
| | 0 | 605 | | target.RequesterDisplayName = source.RequesterDisplayName; |
| | 0 | 606 | | target.Priority = source.Priority; |
| | 0 | 607 | | target.DueAt = source.DueAt; |
| | 0 | 608 | | target.IsOverdue = source.IsOverdue; |
| | 0 | 609 | | target.Status = source.Status; |
| | 0 | 610 | | target.TimeoutEnabled = source.TimeoutEnabled; |
| | 0 | 611 | | target.CancellationEnabled = source.CancellationEnabled; |
| | 0 | 612 | | target.AllowManagerExclusionOverride = source.AllowManagerExclusionOverride; |
| | 0 | 613 | | target.MembershipResolutionMode = source.MembershipResolutionMode; |
| | 0 | 614 | | target.AssigneeProvider = source.AssigneeProvider; |
| | 0 | 615 | | target.AssigneeType = source.AssigneeType; |
| | 0 | 616 | | target.AssigneeId = source.AssigneeId; |
| | 0 | 617 | | target.AssigneeDisplayName = source.AssigneeDisplayName; |
| | 0 | 618 | | target.InstructionsJson = source.InstructionsJson; |
| | 0 | 619 | | target.TaskDataJson = source.TaskDataJson; |
| | 0 | 620 | | target.FormReferenceJson = source.FormReferenceJson; |
| | 0 | 621 | | target.PinnedFormJson = source.PinnedFormJson; |
| | 0 | 622 | | target.ActionsJson = source.ActionsJson; |
| | 0 | 623 | | target.InvitationDefinitionsJson = source.InvitationDefinitionsJson; |
| | 0 | 624 | | target.HealthIssuesJson = source.HealthIssuesJson; |
| | 0 | 625 | | target.HealthSeverity = source.HealthSeverity; |
| | 0 | 626 | | target.HealthCode = source.HealthCode; |
| | 0 | 627 | | target.HealthMessage = source.HealthMessage; |
| | 0 | 628 | | target.CompletionActionKey = source.CompletionActionKey; |
| | 0 | 629 | | target.CompletionDataJson = source.CompletionDataJson; |
| | 0 | 630 | | target.CompletionActorJson = source.CompletionActorJson; |
| | 0 | 631 | | target.CreatedAt = source.CreatedAt; |
| | 0 | 632 | | target.UpdatedAt = source.UpdatedAt; |
| | 0 | 633 | | target.AssignedAt = source.AssignedAt; |
| | 0 | 634 | | target.CompletedAt = source.CompletedAt; |
| | 0 | 635 | | target.CreatedFromBookmarkRevision = source.CreatedFromBookmarkRevision; |
| | 0 | 636 | | } |
| | | 637 | | |
| | | 638 | | private static async Task ReplaceChildrenAsync(UserTasksElsaDbContext dbContext, UserTask task, CancellationToken ca |
| | | 639 | | { |
| | 25 | 640 | | dbContext.UserTaskCandidates.RemoveRange(await dbContext.UserTaskCandidates.Where(x => x.TenantId == task.Tenant |
| | 25 | 641 | | dbContext.UserTaskSnapshotMembers.RemoveRange(await dbContext.UserTaskSnapshotMembers.Where(x => x.TenantId == t |
| | 25 | 642 | | dbContext.UserTaskExclusions.RemoveRange(await dbContext.UserTaskExclusions.Where(x => x.TenantId == task.Tenant |
| | 25 | 643 | | var existingEvents = await dbContext.UserTaskEvents.Where(x => x.TenantId == task.TenantId && x.TaskId == task.I |
| | 41 | 644 | | var existingEventIds = existingEvents.Select(x => x.Id).ToHashSet(StringComparer.Ordinal); |
| | 41 | 645 | | var existingEventRevisions = existingEvents.Select(x => x.Revision).ToHashSet(); |
| | 56 | 646 | | dbContext.UserTaskEvents.AddRange(task.Events.Where(x => !existingEventIds.Contains(x.Id) && !existingEventRevis |
| | | 647 | | |
| | 25 | 648 | | var existingOperations = await dbContext.UserTaskOperations.Where(x => x.TenantId == task.TenantId && x.TaskId = |
| | 27 | 649 | | var existingOperationIds = existingOperations.Select(x => x.Id).ToHashSet(StringComparer.Ordinal); |
| | 58 | 650 | | foreach (var operation in task.Operations) |
| | | 651 | | { |
| | 7 | 652 | | var entity = existingOperations.FirstOrDefault(x => x.Id == operation.Id || x.OperationId == operation.Opera |
| | 4 | 653 | | if (entity is null) |
| | | 654 | | { |
| | 2 | 655 | | if (!existingOperationIds.Contains(operation.Id)) |
| | 2 | 656 | | dbContext.UserTaskOperations.Add(ToOperationRecord(operation)); |
| | 2 | 657 | | continue; |
| | | 658 | | } |
| | | 659 | | |
| | 2 | 660 | | entity.Kind = operation.Kind.ToString(); |
| | 2 | 661 | | entity.ExpectedRevision = operation.ExpectedRevision; |
| | 2 | 662 | | entity.RequestHash = operation.RequestHash; |
| | 2 | 663 | | entity.Status = operation.Status switch |
| | 2 | 664 | | { |
| | 2 | 665 | | UserTaskOperationStatus.Completed => UserTaskPersistenceOperationStatus.Completed, |
| | 0 | 666 | | UserTaskOperationStatus.Failed => UserTaskPersistenceOperationStatus.Failed, |
| | 0 | 667 | | _ => UserTaskPersistenceOperationStatus.Enqueued |
| | 2 | 668 | | }; |
| | 2 | 669 | | entity.UpdatedAt = operation.UpdatedAt; |
| | 2 | 670 | | entity.ActionKey = operation.ActionKey; |
| | 2 | 671 | | entity.ProtectedPayloadJson = Serialize(operation.Data); |
| | 2 | 672 | | entity.ErrorCode = operation.ErrorCode; |
| | | 673 | | } |
| | | 674 | | |
| | 25 | 675 | | var existingInvitations = await dbContext.UserTaskInvitations.Where(x => x.TenantId == task.TenantId && x.TaskId |
| | 33 | 676 | | var existingInvitationIds = existingInvitations.Select(x => x.Id).ToHashSet(StringComparer.Ordinal); |
| | 74 | 677 | | foreach (var invitation in task.Invitations) |
| | | 678 | | { |
| | 20 | 679 | | var entity = existingInvitations.FirstOrDefault(x => x.Id == invitation.Id); |
| | 12 | 680 | | if (entity is null) |
| | | 681 | | { |
| | 4 | 682 | | if (!existingInvitationIds.Contains(invitation.Id)) |
| | 4 | 683 | | dbContext.UserTaskInvitations.Add(ToInvitationRecord(invitation)); |
| | 4 | 684 | | continue; |
| | | 685 | | } |
| | | 686 | | |
| | 8 | 687 | | entity.Status = invitation.Status; |
| | 8 | 688 | | entity.VerifiedAt = invitation.VerifiedAt; |
| | 8 | 689 | | entity.ConsumedAt = invitation.ConsumedAt; |
| | 8 | 690 | | entity.RevokedAt = invitation.RevokedAt; |
| | | 691 | | } |
| | 50 | 692 | | dbContext.UserTaskCandidates.AddRange(task.CandidateUsers.Select(x => ToCandidate(task, x, UserTaskPersistenceCa |
| | 25 | 693 | | dbContext.UserTaskCandidates.AddRange(task.CandidateGroups.Select(x => ToCandidate(task, x, UserTaskPersistenceC |
| | 25 | 694 | | dbContext.UserTaskSnapshotMembers.AddRange(task.SnapshotMembers.Select(x => ToSnapshot(task, x))); |
| | 25 | 695 | | dbContext.UserTaskSnapshotMembers.AddRange(task.SnapshotGroups.Select(x => ToSnapshot(task, x))); |
| | 25 | 696 | | dbContext.UserTaskExclusions.AddRange(task.ExcludedUsers.Select(x => ToExclusion(task, x))); |
| | 25 | 697 | | } |
| | | 698 | | |
| | 119 | 699 | | private static UserTaskCandidateRecord ToCandidate(UserTask task, ParticipantReference participant, UserTaskPersiste |
| | 119 | 700 | | { |
| | 119 | 701 | | TenantId = task.TenantId, TaskId = task.Id, Provider = participant.Provider, ParticipantKey = GetParticipantKey( |
| | 119 | 702 | | ParticipantId = participant.Id, DisplayName = participant.DisplayName, Source = source |
| | 119 | 703 | | }; |
| | | 704 | | |
| | 0 | 705 | | private static UserTaskSnapshotMemberRecord ToSnapshot(UserTask task, ParticipantReference participant) => new() |
| | 0 | 706 | | { |
| | 0 | 707 | | TenantId = task.TenantId, TaskId = task.Id, Provider = participant.Provider, ParticipantKey = GetParticipantKey( |
| | 0 | 708 | | ParticipantId = participant.Id, CreatedAt = task.CreatedAt |
| | 0 | 709 | | }; |
| | | 710 | | |
| | 2 | 711 | | private static UserTaskExclusionRecord ToExclusion(UserTask task, ParticipantReference participant) => new() |
| | 2 | 712 | | { |
| | 2 | 713 | | TenantId = task.TenantId, TaskId = task.Id, Provider = participant.Provider, ParticipantKey = GetParticipantKey( |
| | 2 | 714 | | ParticipantId = participant.Id, CreatedAt = task.CreatedAt |
| | 2 | 715 | | }; |
| | | 716 | | |
| | 326 | 717 | | private static ParticipantReference ToParticipant(UserTaskCandidateRecord row) => new(row.TenantId, row.Provider, ro |
| | 0 | 718 | | private static ParticipantReference ToParticipant(UserTaskSnapshotMemberRecord row) => new(row.TenantId, row.Provide |
| | 0 | 719 | | private static ParticipantReference ToParticipant(UserTaskExclusionRecord row) => new(row.TenantId, row.Provider, ro |
| | | 720 | | |
| | | 721 | | private static ParticipantReference? ToParticipant(string? provider, string? type, string? id, string? displayName, |
| | | 722 | | { |
| | 656 | 723 | | if (string.IsNullOrWhiteSpace(provider) || string.IsNullOrWhiteSpace(type) || string.IsNullOrWhiteSpace(id) || ! |
| | 623 | 724 | | return null; |
| | 33 | 725 | | return new ParticipantReference(tenantId, provider, participantType, id, displayName); |
| | | 726 | | } |
| | | 727 | | |
| | 121 | 728 | | private static string GetParticipantKey(ParticipantReference participant) => $"{participant.Provider}|{participant.T |
| | | 729 | | |
| | 91 | 730 | | private static UserTaskEvent ToEvent(UserTaskEventRecord row) => new(row.Id, row.TenantId, row.TaskId, row.Revision, |
| | | 731 | | |
| | 18 | 732 | | private static UserTaskEventRecord ToEventRecord(UserTaskEvent value) => new() |
| | 18 | 733 | | { |
| | 18 | 734 | | Id = value.Id, TenantId = value.TenantId, TaskId = value.TaskId, Revision = value.Revision, EventType = value.Ev |
| | 18 | 735 | | OccurredAt = value.OccurredAt, ActorProvider = value.Actor?.Provider, ActorType = value.Actor?.Type.ToString(), |
| | 18 | 736 | | ActorJson = Serialize(value.Actor), OperationId = value.OperationId, Reason = value.Reason, |
| | 18 | 737 | | MetadataJson = Serialize(value.Metadata ?? new Dictionary<string, object?>()) |
| | 18 | 738 | | }; |
| | | 739 | | |
| | 10 | 740 | | private static UserTaskOperation ToOperation(UserTaskOperationRecord row) => new( |
| | 10 | 741 | | row.Id, row.TenantId, row.TaskId, row.OperationId, |
| | 10 | 742 | | Enum.TryParse<UserTaskOperationKind>(row.Kind, true, out var kind) ? kind : UserTaskOperationKind.Claim, |
| | 10 | 743 | | row.ExpectedRevision, row.RequestHash, |
| | 10 | 744 | | row.Status switch |
| | 10 | 745 | | { |
| | 6 | 746 | | UserTaskPersistenceOperationStatus.Completed => UserTaskOperationStatus.Completed, |
| | 0 | 747 | | UserTaskPersistenceOperationStatus.Failed => UserTaskOperationStatus.Failed, |
| | 4 | 748 | | _ => UserTaskOperationStatus.Accepted |
| | 10 | 749 | | }, |
| | 10 | 750 | | row.CreatedAt, row.UpdatedAt, row.ActionKey, DeserializeJson(row.ProtectedPayloadJson), row.ErrorCode); |
| | | 751 | | |
| | 3 | 752 | | private static UserTaskOperationRecord ToOperationRecord(UserTaskOperation value) => new() |
| | 3 | 753 | | { |
| | 3 | 754 | | Id = value.Id, TenantId = value.TenantId, TaskId = value.TaskId, OperationId = value.OperationId, Kind = value.K |
| | 3 | 755 | | ExpectedRevision = value.ExpectedRevision, RequestHash = value.RequestHash, |
| | 3 | 756 | | Status = value.Status switch |
| | 3 | 757 | | { |
| | 1 | 758 | | UserTaskOperationStatus.Completed => UserTaskPersistenceOperationStatus.Completed, |
| | 0 | 759 | | UserTaskOperationStatus.Failed => UserTaskPersistenceOperationStatus.Failed, |
| | 2 | 760 | | _ => UserTaskPersistenceOperationStatus.Enqueued |
| | 3 | 761 | | }, |
| | 3 | 762 | | CreatedAt = value.CreatedAt, UpdatedAt = value.UpdatedAt, ActionKey = value.ActionKey, ProtectedPayloadJson = Se |
| | 3 | 763 | | }; |
| | | 764 | | |
| | 48 | 765 | | private static UserTaskInvitation ToInvitation(UserTaskInvitationRecord row) => new( |
| | 48 | 766 | | row.Id, row.TenantId, row.TaskId, Deserialize<string>(row.RecipientJson), row.TokenHash, row.Status, |
| | 48 | 767 | | row.IssuedAt, row.ExpiresAt, row.VerifierProvider, row.VerifiedAt, row.ConsumedAt, row.RevokedAt, row.SiblingGro |
| | 48 | 768 | | { |
| | 48 | 769 | | AllowedActions = Deserialize<List<string>>(row.AllowedActionsJson) ?? [] |
| | 48 | 770 | | }; |
| | | 771 | | |
| | 6 | 772 | | private static UserTaskInvitationRecord ToInvitationRecord(UserTaskInvitation value) => new() |
| | 6 | 773 | | { |
| | 6 | 774 | | Id = value.Id, TenantId = value.TenantId, TaskId = value.TaskId, RecipientJson = Serialize(value.Recipient), Tok |
| | 6 | 775 | | VerifierProvider = value.VerifierName ?? "default", Status = value.Status, IssuedAt = value.IssuedAt, ExpiresAt |
| | 6 | 776 | | VerifiedAt = value.VerifiedAt, ConsumedAt = value.ConsumedAt, RevokedAt = value.RevokedAt, SiblingGroupId = valu |
| | 6 | 777 | | AllowedActionsJson = Serialize(value.AllowedActions) |
| | 6 | 778 | | }; |
| | | 779 | | |
| | 95 | 780 | | private static async Task<bool> ExistsByMaterializationKeyAsync(UserTasksElsaDbContext dbContext, string tenantId, s |
| | | 781 | | |
| | | 782 | | private static void EnsureExpectedRevision(UserTaskRecord? existing, string taskId, int expectedRevision) |
| | | 783 | | { |
| | 19 | 784 | | if (existing is null) |
| | 0 | 785 | | throw new KeyNotFoundException($"User task '{taskId}' was not found."); |
| | 19 | 786 | | if (existing.Revision != expectedRevision) |
| | 3 | 787 | | throw new UserTaskRevisionConflictException(taskId, expectedRevision); |
| | 16 | 788 | | } |
| | | 789 | | |
| | 1000 | 790 | | private static string Serialize<T>(T value) => JsonSerializer.Serialize(value, JsonOptions); |
| | 243 | 791 | | private static string Serialize(JsonElement? value) => !value.HasValue || value.Value.ValueKind == JsonValueKind.Und |
| | 3222 | 792 | | private static T? Deserialize<T>(string? value) => string.IsNullOrWhiteSpace(value) || string.Equals(value, "null", |
| | 662 | 793 | | private static JsonElement? DeserializeJson(string? value) => Deserialize<JsonElement>(value); |
| | | 794 | | } |