| | | 1 | | using System.Text.Json; |
| | | 2 | | using System.Text.Json.Serialization; |
| | | 3 | | using Elsa.UserTasks.Contracts; |
| | | 4 | | using Elsa.UserTasks.Models; |
| | | 5 | | |
| | | 6 | | namespace Elsa.UserTasks.Repositories; |
| | | 7 | | |
| | | 8 | | /// <summary> |
| | | 9 | | /// A deterministic repository for the Core module and development hosts. It uses a single lock to |
| | | 10 | | /// provide the same compare-and-swap and projection idempotency guarantees expected from durable stores. |
| | | 11 | | /// </summary> |
| | | 12 | | public sealed class InMemoryUserTaskRepository : IUserTaskRepository |
| | | 13 | | { |
| | 2 | 14 | | private static readonly JsonSerializerOptions JsonOptions = new(JsonSerializerDefaults.Web) |
| | 2 | 15 | | { |
| | 2 | 16 | | Converters = { new JsonStringEnumConverter() } |
| | 2 | 17 | | }; |
| | | 18 | | |
| | | 19 | | /// <summary> |
| | | 20 | | /// Title OrderBy, ties, and cursors share this comparer. Title cursors are not portable to EF |
| | | 21 | | /// (column collation) or across databases; recreate the list after a provider change. |
| | | 22 | | /// </summary> |
| | 2 | 23 | | private static readonly StringComparer TitleComparer = StringComparer.Ordinal; |
| | | 24 | | |
| | 93 | 25 | | private readonly object _sync = new(); |
| | 93 | 26 | | private readonly Dictionary<string, UserTask> _tasks = new(StringComparer.Ordinal); |
| | | 27 | | |
| | | 28 | | public Task<UserTask?> GetAsync(string tenantId, string taskId, CancellationToken cancellationToken = default) |
| | | 29 | | { |
| | 246 | 30 | | lock (_sync) |
| | 246 | 31 | | return Task.FromResult(_tasks.TryGetValue(Key(tenantId, taskId), out var task) ? Clone(task) : null); |
| | 246 | 32 | | } |
| | | 33 | | |
| | | 34 | | public Task<UserTaskQueryResult> QueryAsync(UserTaskQuery query, CancellationToken cancellationToken = default) |
| | | 35 | | { |
| | 174 | 36 | | lock (_sync) |
| | | 37 | | { |
| | 174 | 38 | | var scope = query.Scope; |
| | 174 | 39 | | var items = _tasks.Values |
| | 8176 | 40 | | .Where(x => string.Equals(x.TenantId, query.TenantId, StringComparison.Ordinal)) |
| | 174 | 41 | | // Visibility is evaluated before counting, cursoring, and paging so an unauthorized row can |
| | 174 | 42 | | // never influence a total or push an authorized row off the page. |
| | 963 | 43 | | .Where(x => scope == null || IsVisible(x, scope)) |
| | 943 | 44 | | .Where(x => query.Statuses.Count == 0 || query.Statuses.Contains(x.Status)) |
| | 943 | 45 | | .Where(x => !query.OnlyOverdue || x.IsOverdue) |
| | 943 | 46 | | .Where(x => !query.OnlyWithoutDueDate || x.DueAt == null) |
| | 943 | 47 | | .Where(x => query.PriorityFrom == null || x.Priority >= query.PriorityFrom) |
| | 943 | 48 | | .Where(x => query.PriorityTo == null || x.Priority <= query.PriorityTo) |
| | 943 | 49 | | .Where(x => query.DueFrom == null || x.DueAt >= query.DueFrom) |
| | 943 | 50 | | .Where(x => query.DueTo == null || x.DueAt <= query.DueTo) |
| | 943 | 51 | | .Where(x => query.WorkflowDefinitionId == null || x.WorkflowDefinitionId == query.WorkflowDefinitionId) |
| | 943 | 52 | | .Where(x => query.WorkflowInstanceId == null || x.WorkflowInstanceId == query.WorkflowInstanceId) |
| | 943 | 53 | | .Where(x => query.Reference == null || string.Equals(x.Reference, query.Reference, StringComparison.Ordi |
| | 943 | 54 | | .Where(x => query.TaskType == null || string.Equals(x.TaskType, query.TaskType, StringComparison.Ordinal |
| | 1117 | 55 | | .Where(x => MatchesSearch(x, query.Search)); |
| | | 56 | | |
| | 174 | 57 | | var filteredCount = query.IncludeTotalCount ? items.Count() : 0; |
| | 174 | 58 | | var materialized = ApplyOrdering(items, query).ToList(); |
| | 174 | 59 | | materialized = ApplyCursor(materialized, query).ToList(); |
| | | 60 | | |
| | 174 | 61 | | int? total = query.IncludeTotalCount ? filteredCount : null; |
| | 174 | 62 | | var limit = Math.Clamp(query.Limit, 1, 200); |
| | 174 | 63 | | var page = materialized.Take(limit).Select(Clone).ToArray(); |
| | 174 | 64 | | var next = materialized.Count > limit ? CreateCursor(page[^1], query.Sort) : null; |
| | 174 | 65 | | return Task.FromResult(new UserTaskQueryResult(page, next, total)); |
| | | 66 | | } |
| | 174 | 67 | | } |
| | | 68 | | |
| | | 69 | | public Task<UserTask?> FindByMaterializationKeyAsync(string tenantId, string key, CancellationToken cancellationToke |
| | | 70 | | { |
| | 88 | 71 | | lock (_sync) |
| | 307 | 72 | | return Task.FromResult(_tasks.Values.FirstOrDefault(x => x.TenantId == tenantId && x.MaterializationKey == k |
| | 88 | 73 | | } |
| | | 74 | | |
| | | 75 | | public Task<UserTask?> FindByBookmarkIdAsync(string tenantId, string bookmarkId, CancellationToken cancellationToken |
| | | 76 | | { |
| | 2 | 77 | | lock (_sync) |
| | 176 | 78 | | return Task.FromResult(_tasks.Values.FirstOrDefault(x => x.TenantId == tenantId && x.BookmarkId == bookmarkI |
| | 2 | 79 | | } |
| | | 80 | | |
| | | 81 | | public Task<(UserTask Task, UserTaskInvitation Invitation)?> FindByInvitationTokenHashAsync(string tokenHash, Cancel |
| | | 82 | | { |
| | 35 | 83 | | lock (_sync) |
| | | 84 | | { |
| | 469 | 85 | | foreach (var task in _tasks.Values) |
| | | 86 | | { |
| | | 87 | | // Ordinal comparison over a fixed-length hex hash; the secret itself is never stored. |
| | 263 | 88 | | var invitation = task.Invitations.FirstOrDefault(x => string.Equals(x.TokenHash, tokenHash, StringCompar |
| | 215 | 89 | | if (invitation != null) |
| | 31 | 90 | | return Task.FromResult<(UserTask, UserTaskInvitation)?>((Clone(task), invitation)); |
| | | 91 | | } |
| | 4 | 92 | | return Task.FromResult<(UserTask, UserTaskInvitation)?>(null); |
| | | 93 | | } |
| | 35 | 94 | | } |
| | | 95 | | |
| | | 96 | | public Task SaveAsync(UserTask task, int expectedRevision, CancellationToken cancellationToken = default) |
| | | 97 | | { |
| | 55 | 98 | | lock (_sync) |
| | | 99 | | { |
| | 55 | 100 | | var key = Key(task.TenantId, task.Id); |
| | 55 | 101 | | if (!_tasks.TryGetValue(key, out var current) || current.Revision != expectedRevision) |
| | 2 | 102 | | throw new UserTaskRevisionConflictException(task.Id, expectedRevision); |
| | | 103 | | |
| | 53 | 104 | | var copy = Clone(task); |
| | 53 | 105 | | copy.Revision = expectedRevision + 1; |
| | 53 | 106 | | copy.UpdatedAt = DateTimeOffset.UtcNow; |
| | 53 | 107 | | _tasks[key] = copy; |
| | 53 | 108 | | return Task.CompletedTask; |
| | | 109 | | } |
| | 53 | 110 | | } |
| | | 111 | | |
| | | 112 | | public Task AddProjectionAsync(UserTask task, CancellationToken cancellationToken = default) |
| | | 113 | | { |
| | 200 | 114 | | lock (_sync) |
| | | 115 | | { |
| | 6803 | 116 | | if (_tasks.Values.Any(x => x.TenantId == task.TenantId && |
| | 6803 | 117 | | ((!string.IsNullOrEmpty(task.MaterializationKey) && x.MaterializationKey == task.MaterializationKey) || |
| | 6803 | 118 | | (!string.IsNullOrEmpty(task.BookmarkId) && x.BookmarkId == task.BookmarkId)))) |
| | 1 | 119 | | return Task.CompletedTask; |
| | | 120 | | |
| | 199 | 121 | | _tasks[Key(task.TenantId, task.Id)] = Clone(task); |
| | 199 | 122 | | return Task.CompletedTask; |
| | | 123 | | } |
| | 200 | 124 | | } |
| | | 125 | | |
| | | 126 | | public Task AppendEventAsync(string tenantId, string taskId, UserTaskEvent @event, CancellationToken cancellationTok |
| | | 127 | | { |
| | 5 | 128 | | lock (_sync) |
| | | 129 | | { |
| | 5 | 130 | | if (_tasks.TryGetValue(Key(tenantId, taskId), out var current)) |
| | 4 | 131 | | current.Events.Add(@event); |
| | 5 | 132 | | return Task.CompletedTask; |
| | | 133 | | } |
| | 5 | 134 | | } |
| | | 135 | | |
| | | 136 | | public Task<bool> TryMutateAsync(string tenantId, string taskId, int expectedRevision, Func<UserTask, bool> mutation |
| | | 137 | | { |
| | 39 | 138 | | lock (_sync) |
| | | 139 | | { |
| | 39 | 140 | | var key = Key(tenantId, taskId); |
| | 39 | 141 | | if (!_tasks.TryGetValue(key, out var current) || current.Revision != expectedRevision) |
| | 1 | 142 | | return Task.FromResult(false); |
| | | 143 | | |
| | 38 | 144 | | var copy = Clone(current); |
| | 38 | 145 | | if (!mutation(copy)) |
| | 1 | 146 | | return Task.FromResult(false); |
| | | 147 | | |
| | 37 | 148 | | copy.Revision = expectedRevision + 1; |
| | 37 | 149 | | copy.UpdatedAt = DateTimeOffset.UtcNow; |
| | 37 | 150 | | _tasks[key] = copy; |
| | 37 | 151 | | return Task.FromResult(true); |
| | | 152 | | } |
| | 39 | 153 | | } |
| | | 154 | | |
| | | 155 | | private static bool IsVisible(UserTask task, UserTaskQueryScope scope) |
| | | 156 | | { |
| | 816 | 157 | | if (!string.Equals(task.TenantId, scope.TenantId, StringComparison.Ordinal)) |
| | 2 | 158 | | return false; |
| | 814 | 159 | | if (scope.ExcludeBlocking && task.HealthSeverity == UserTaskHealthSeverity.Blocking) |
| | 0 | 160 | | return false; |
| | | 161 | | // Manager-only scopes were already rejected by the policy for non-managers, so reaching them here |
| | | 162 | | // means the caller manages the tenant. |
| | 814 | 163 | | if (scope.RequiresManager) |
| | 0 | 164 | | return scope.IsManager && (scope.Kind != UserTaskQueryScopeKind.NeedsAttention || NeedsAttention(task)); |
| | | 165 | | |
| | 814 | 166 | | return scope.Kind switch |
| | 814 | 167 | | { |
| | 1 | 168 | | UserTaskQueryScopeKind.Assigned => task.Assignee?.Matches(scope.Subject) == true, |
| | 813 | 169 | | UserTaskQueryScopeKind.Available => task.IsOpen && task.Assignee == null && IsEligible(task, scope), |
| | 0 | 170 | | UserTaskQueryScopeKind.History => task.IsTerminal && |
| | 0 | 171 | | (task.CompletedBy?.Matches(scope.Subject) == true || task.Events.Any(x => |
| | 0 | 172 | | _ => false |
| | 814 | 173 | | }; |
| | | 174 | | } |
| | | 175 | | |
| | | 176 | | private static bool IsEligible(UserTask task, UserTaskQueryScope scope) |
| | | 177 | | { |
| | 820 | 178 | | if (task.ExcludedUsers.Any(x => x.Matches(scope.Subject))) |
| | 7 | 179 | | return false; |
| | 805 | 180 | | if (task.MembershipResolutionMode == UserTaskMembershipResolutionMode.Snapshot) |
| | 10 | 181 | | return task.SnapshotMembers.Any(x => x.Matches(scope.Subject)); |
| | 1598 | 182 | | return task.CandidateUsers.Any(x => x.Matches(scope.Subject)) |
| | 799 | 183 | | || task.CandidateGroups.Any(candidate => scope.Groups.Any(candidate.Matches)); |
| | | 184 | | } |
| | | 185 | | |
| | | 186 | | private static bool NeedsAttention(UserTask task) => |
| | 0 | 187 | | task.HealthSeverity == UserTaskHealthSeverity.Blocking |
| | 0 | 188 | | || task.IsOverdue |
| | 0 | 189 | | || (task.IsOpen && task.Assignee == null) |
| | 0 | 190 | | || task.Status is UserTaskStatus.Completing or UserTaskStatus.TimingOut or UserTaskStatus.Cancelling; |
| | | 191 | | |
| | | 192 | | private static bool MatchesSearch(UserTask task, string? search) |
| | | 193 | | { |
| | 943 | 194 | | if (string.IsNullOrWhiteSpace(search)) |
| | 925 | 195 | | return true; |
| | 18 | 196 | | var value = search.Trim(); |
| | 18 | 197 | | return (task.Title?.Contains(value, StringComparison.OrdinalIgnoreCase) == true) |
| | 18 | 198 | | || (task.Summary?.Contains(value, StringComparison.OrdinalIgnoreCase) == true) |
| | 18 | 199 | | || (task.Reference?.Contains(value, StringComparison.OrdinalIgnoreCase) == true) |
| | 18 | 200 | | || (task.TaskType?.Contains(value, StringComparison.OrdinalIgnoreCase) == true) |
| | 40 | 201 | | || task.Tags.Any(x => x.Contains(value, StringComparison.OrdinalIgnoreCase)); |
| | | 202 | | } |
| | | 203 | | |
| | | 204 | | // Same contract as EF/VNext: REST sorts only, Id always ThenBy ascending, JSON base64url cursors. |
| | 174 | 205 | | private static IEnumerable<UserTask> ApplyOrdering(IEnumerable<UserTask> tasks, UserTaskQuery query) => query.Sort.T |
| | 174 | 206 | | { |
| | 345 | 207 | | "priority" => query.Descending ? tasks.OrderByDescending(x => x.Priority).ThenBy(x => x.Id) : tasks.OrderBy(x => |
| | 360 | 208 | | "title" => query.Descending ? tasks.OrderByDescending(x => x.Title, TitleComparer).ThenBy(x => x.Id) : tasks.Ord |
| | 535 | 209 | | "due" => query.Descending ? tasks.OrderBy(x => x.DueAt == null).ThenByDescending(x => x.DueAt).ThenBy(x => x.Id) |
| | 366 | 210 | | "updated" => query.Descending ? tasks.OrderByDescending(x => x.UpdatedAt).ThenBy(x => x.Id) : tasks.OrderBy(x => |
| | 432 | 211 | | _ => query.Descending ? tasks.OrderByDescending(x => x.CreatedAt).ThenBy(x => x.Id) : tasks.OrderBy(x => x.Creat |
| | 174 | 212 | | }; |
| | | 213 | | |
| | | 214 | | private static IEnumerable<UserTask> ApplyCursor(IEnumerable<UserTask> tasks, UserTaskQuery query) |
| | | 215 | | { |
| | 174 | 216 | | if (string.IsNullOrWhiteSpace(query.Cursor) || !TryReadCursor(query.Cursor, out var value, out var id)) |
| | 80 | 217 | | return tasks; |
| | | 218 | | return query.Sort.ToLowerInvariant() switch |
| | | 219 | | { |
| | 138 | 220 | | "priority" when int.TryParse(value, out var priority) => tasks.Where(x => query.Descending ? x.Priority < pr |
| | 126 | 221 | | "title" => tasks.Where(x => TitleIsAfterCursor(x.Title, value, x.Id, id, query.Descending)), |
| | 32 | 222 | | "due" when value == "~null" => tasks.Where(x => x.DueAt == null && string.Compare(x.Id, id) > 0), |
| | 122 | 223 | | "due" when DateTimeOffset.TryParse(value, out var due) => tasks.Where(x => x.DueAt == null || query.Descendi |
| | 138 | 224 | | "updated" when DateTimeOffset.TryParse(value, out var updated) => tasks.Where(x => query.Descending ? x.Upda |
| | 147 | 225 | | _ when DateTimeOffset.TryParse(value, out var created) => tasks.Where(x => query.Descending ? x.CreatedAt < |
| | 0 | 226 | | _ => tasks |
| | | 227 | | }; |
| | | 228 | | } |
| | | 229 | | |
| | | 230 | | private static bool TitleIsAfterCursor(string title, string cursorTitle, string id, string cursorId, bool descending |
| | | 231 | | { |
| | 106 | 232 | | var comparison = TitleComparer.Compare(title, cursorTitle); |
| | 106 | 233 | | return descending |
| | 106 | 234 | | ? comparison < 0 || comparison == 0 && string.Compare(id, cursorId) > 0 |
| | 106 | 235 | | : comparison > 0 || comparison == 0 && string.Compare(id, cursorId) > 0; |
| | | 236 | | } |
| | | 237 | | |
| | | 238 | | private static string CreateCursor(UserTask task, string sort) |
| | | 239 | | { |
| | 95 | 240 | | var value = sort.ToLowerInvariant() switch |
| | 95 | 241 | | { |
| | 18 | 242 | | "priority" => task.Priority.ToString(System.Globalization.CultureInfo.InvariantCulture), |
| | 20 | 243 | | "title" => task.Title, |
| | 18 | 244 | | "due" => task.DueAt?.ToString("O", System.Globalization.CultureInfo.InvariantCulture) ?? "~null", |
| | 18 | 245 | | "updated" => task.UpdatedAt.ToString("O", System.Globalization.CultureInfo.InvariantCulture), |
| | 21 | 246 | | _ => task.CreatedAt.ToString("O", System.Globalization.CultureInfo.InvariantCulture) |
| | 95 | 247 | | }; |
| | 95 | 248 | | return Convert.ToBase64String(JsonSerializer.SerializeToUtf8Bytes(new[] { value, task.Id }, JsonOptions)).TrimEn |
| | | 249 | | } |
| | | 250 | | |
| | | 251 | | private static bool TryReadCursor(string? cursor, out string value, out string id) |
| | | 252 | | { |
| | 95 | 253 | | value = id = ""; |
| | 95 | 254 | | if (string.IsNullOrWhiteSpace(cursor)) |
| | 0 | 255 | | return false; |
| | | 256 | | try |
| | | 257 | | { |
| | 95 | 258 | | var padded = cursor.Replace('-', '+').Replace('_', '/') + new string('=', (4 - cursor.Length % 4) % 4); |
| | 95 | 259 | | var values = JsonSerializer.Deserialize<string[]>(Convert.FromBase64String(padded), JsonOptions); |
| | 94 | 260 | | if (values is not [var parsedValue, var parsedId] || string.IsNullOrWhiteSpace(parsedId)) |
| | 0 | 261 | | return false; |
| | 94 | 262 | | value = parsedValue; |
| | 94 | 263 | | id = parsedId; |
| | 94 | 264 | | return true; |
| | | 265 | | } |
| | 0 | 266 | | catch (FormatException) |
| | | 267 | | { |
| | 0 | 268 | | return false; |
| | | 269 | | } |
| | 1 | 270 | | catch (JsonException) |
| | | 271 | | { |
| | 1 | 272 | | return false; |
| | | 273 | | } |
| | 95 | 274 | | } |
| | | 275 | | |
| | 544 | 276 | | private static string Key(string tenantId, string taskId) => tenantId + "\0" + taskId; |
| | | 277 | | |
| | 943 | 278 | | private static UserTask Clone(UserTask task) => new() |
| | 943 | 279 | | { |
| | 943 | 280 | | Id = task.Id, |
| | 943 | 281 | | TenantId = task.TenantId, |
| | 943 | 282 | | WorkflowDefinitionId = task.WorkflowDefinitionId, |
| | 943 | 283 | | WorkflowDefinitionName = task.WorkflowDefinitionName, |
| | 943 | 284 | | WorkflowDefinitionVersion = task.WorkflowDefinitionVersion, |
| | 943 | 285 | | WorkflowInstanceId = task.WorkflowInstanceId, |
| | 943 | 286 | | WorkflowInstanceReference = task.WorkflowInstanceReference, |
| | 943 | 287 | | ActivityInstanceId = task.ActivityInstanceId, |
| | 943 | 288 | | BookmarkId = task.BookmarkId, |
| | 943 | 289 | | MaterializationKey = task.MaterializationKey, |
| | 943 | 290 | | Title = task.Title, |
| | 943 | 291 | | Summary = task.Summary, |
| | 943 | 292 | | Reference = task.Reference, |
| | 943 | 293 | | Tags = new HashSet<string>(task.Tags, StringComparer.OrdinalIgnoreCase), |
| | 943 | 294 | | TaskType = task.TaskType, |
| | 943 | 295 | | Requester = task.Requester, |
| | 943 | 296 | | Assignee = task.Assignee, |
| | 943 | 297 | | CandidateUsers = [..task.CandidateUsers], |
| | 943 | 298 | | CandidateGroups = [..task.CandidateGroups], |
| | 943 | 299 | | SnapshotMembers = [..task.SnapshotMembers], |
| | 943 | 300 | | SnapshotGroups = [..task.SnapshotGroups], |
| | 943 | 301 | | ExcludedUsers = [..task.ExcludedUsers], |
| | 943 | 302 | | MembershipResolutionMode = task.MembershipResolutionMode, |
| | 943 | 303 | | AllowManagerExclusionOverride = task.AllowManagerExclusionOverride, |
| | 943 | 304 | | Priority = task.Priority, |
| | 943 | 305 | | DueAt = task.DueAt, |
| | 943 | 306 | | IsOverdue = task.IsOverdue, |
| | 943 | 307 | | Instructions = task.Instructions, |
| | 943 | 308 | | TaskData = Clone(task.TaskData), |
| | 943 | 309 | | RequestedForm = task.RequestedForm, |
| | 943 | 310 | | PinnedForm = task.PinnedForm, |
| | 943 | 311 | | Actions = [..task.Actions], |
| | 943 | 312 | | InvitationDefinitions = [..task.InvitationDefinitions], |
| | 943 | 313 | | EnableTimeoutOutcome = task.EnableTimeoutOutcome, |
| | 943 | 314 | | EnableCancellationOutcome = task.EnableCancellationOutcome, |
| | 943 | 315 | | Status = task.Status, |
| | 943 | 316 | | HealthSeverity = task.HealthSeverity, |
| | 943 | 317 | | HealthCode = task.HealthCode, |
| | 943 | 318 | | HealthMessage = task.HealthMessage, |
| | 943 | 319 | | CompletionActionKey = task.CompletionActionKey, |
| | 943 | 320 | | CompletionData = Clone(task.CompletionData), |
| | 943 | 321 | | CompletedBy = task.CompletedBy, |
| | 943 | 322 | | CreatedAt = task.CreatedAt, |
| | 943 | 323 | | UpdatedAt = task.UpdatedAt, |
| | 943 | 324 | | AssignedAt = task.AssignedAt, |
| | 943 | 325 | | CompletedAt = task.CompletedAt, |
| | 943 | 326 | | Revision = task.Revision, |
| | 943 | 327 | | Events = [..task.Events], |
| | 943 | 328 | | Operations = [..task.Operations], |
| | 943 | 329 | | Invitations = [..task.Invitations] |
| | 943 | 330 | | }; |
| | | 331 | | |
| | 1886 | 332 | | private static JsonElement? Clone(JsonElement? value) => value is { } element ? element.Clone() : null; |
| | | 333 | | } |