< Summary

Information
Class: Elsa.UserTasks.Persistence.VNext.Repositories.VNextUserTaskRepository
Assembly: Elsa.UserTasks.Persistence.VNext
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.UserTasks.Persistence.VNext/Repositories/VNextUserTaskRepository.cs
Line coverage
87%
Covered lines: 183
Uncovered lines: 27
Coverable lines: 210
Total lines: 344
Line coverage: 87.1%
Branch coverage
73%
Covered branches: 206
Total branches: 282
Branch coverage: 73%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
.cctor()100%11100%
GetAsync()100%22100%
QueryAsync()87.5%88100%
FindByMaterializationKeyAsync(...)100%11100%
FindByBookmarkIdAsync(...)100%11100%
FindByInvitationTokenHashAsync()100%66100%
SaveAsync()75%4476.92%
AddProjectionAsync()100%2271.42%
AppendEventAsync()100%2280%
TryMutateAsync()66.66%12646.15%
LoadAllAsync()100%22100%
FindByIndexAsync()100%44100%
LoadDocumentAsync()100%22100%
CreateRequest(...)58.33%1212100%
Deserialize(...)50%22100%
DocumentId(...)100%11100%
Matches(...)55.88%6868100%
IsVisible(...)38.23%773466.66%
IsEligible(...)100%66100%
NeedsAttention(...)0%110100%
ApplyOrdering(...)100%1818100%
ApplyCursor(...)95.58%746888.88%
TitleIsAfterCursor(...)100%1010100%
CreateCursor(...)100%1212100%
TryReadCursor(...)66.66%6684.61%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.UserTasks.Persistence.VNext/Repositories/VNextUserTaskRepository.cs

#LineLine coverage
 1using System.Text.Json;
 2using System.Text.Json.Serialization;
 3using Elsa.Persistence.VNext.Document;
 4using Elsa.UserTasks.Contracts;
 5using Elsa.UserTasks.Models;
 6
 7namespace Elsa.UserTasks.Persistence.VNext.Repositories;
 8
 9/// <summary>
 10/// Provider-neutral document implementation. The aggregate is stored as one document while the schema
 11/// provider advertises the same logical units and indexes as the relational providers.
 12/// </summary>
 213public sealed class VNextUserTaskRepository(IDocumentStore documentStore) : IUserTaskRepository
 14{
 15    public const string StorageUnitName = "UserTasks";
 116    private static readonly JsonSerializerOptions JsonOptions = new(JsonSerializerDefaults.Web)
 117    {
 118        Converters = { new JsonStringEnumConverter() }
 119    };
 20
 21    /// <summary>
 22    /// Title OrderBy, ties, and cursors share this comparer. Title cursors are not portable to EF
 23    /// (column collation) or across databases; recreate the list after a provider change.
 24    /// </summary>
 125    private static readonly StringComparer TitleComparer = StringComparer.Ordinal;
 26
 27    public async Task<UserTask?> GetAsync(string tenantId, string taskId, CancellationToken cancellationToken = default)
 28    {
 2229        var document = await documentStore.LoadAsync(StorageUnitName, DocumentId(tenantId, taskId), cancellationToken);
 2230        return document is null ? null : Deserialize(document);
 2231    }
 32
 33    public async Task<UserTaskQueryResult> QueryAsync(UserTaskQuery query, CancellationToken cancellationToken = default
 34    {
 14035        var tasks = new List<UserTask>();
 280036        foreach (var status in Enum.GetValues<UserTaskStatus>())
 37        {
 126038            var documents = await documentStore.QueryAsync(
 126039                new DocumentQuery(StorageUnitName, new Dictionary<string, string?>
 126040                {
 126041                    ["TenantId"] = query.TenantId,
 126042                    ["Status"] = status.ToString()
 126043                }), cancellationToken);
 126044            tasks.AddRange(documents.Select(Deserialize));
 45        }
 46
 171647        var filtered = tasks.Where(x => Matches(x, query)).Where(x => query.Scope is null || IsVisible(x, query.Scope)).
 14048        int? totalCount = query.IncludeTotalCount ? filtered.Count : null;
 14049        filtered = ApplyOrdering(filtered, query).ToList();
 14050        filtered = ApplyCursor(filtered, query).ToList();
 14051        var limit = Math.Clamp(query.Limit <= 0 ? 50 : query.Limit, 1, 200);
 14052        var hasMore = filtered.Count > limit;
 14053        var page = filtered.Take(limit).ToList();
 14054        return new UserTaskQueryResult(page, hasMore ? CreateCursor(page[^1], query.Sort) : null, totalCount);
 14055    }
 56
 33257    public Task<UserTask?> FindByMaterializationKeyAsync(string tenantId, string key, CancellationToken cancellationToke
 58
 359    public Task<UserTask?> FindByBookmarkIdAsync(string tenantId, string bookmarkId, CancellationToken cancellationToken
 60
 61    public async Task<(UserTask Task, UserTaskInvitation Invitation)?> FindByInvitationTokenHashAsync(string tokenHash, 
 62    {
 63        // The document store has no cross-tenant secondary index for invitation hashes, so this scans the
 64        // storage unit. Verification is rate limited and rare; correctness matters more than the scan here.
 28165        foreach (var task in await LoadAllAsync(cancellationToken))
 66        {
 14167            var invitation = task.Invitations.FirstOrDefault(x => string.Equals(x.TokenHash, tokenHash, StringComparison
 13968            if (invitation != null)
 169                return (task, invitation);
 70        }
 171        return null;
 272    }
 73
 74    public async Task SaveAsync(UserTask task, int expectedRevision, CancellationToken cancellationToken = default)
 75    {
 676        var existing = await LoadDocumentAsync(task.TenantId, task.Id, cancellationToken);
 677        if (existing is null)
 078            throw new KeyNotFoundException($"User task '{task.Id}' was not found.");
 679        var loaded = existing.Value;
 680        if (loaded.Task.Revision != expectedRevision)
 281            throw new UserTaskRevisionConflictException(task.Id, expectedRevision);
 82
 483        task.Revision = expectedRevision + 1;
 484        task.UpdatedAt = DateTimeOffset.UtcNow;
 85        try
 86        {
 487            await documentStore.SaveAsync(CreateRequest(task, loaded.Document.Version), cancellationToken);
 488        }
 089        catch (DocumentStoreConcurrencyException exception)
 90        {
 91            // A writer won the race between the read above and this save. Same contract as the other
 92            // providers so the caller sees one exception type regardless of the installed store.
 093            throw new UserTaskRevisionConflictException(task.Id, expectedRevision, exception);
 94        }
 495    }
 96
 97    public async Task AddProjectionAsync(UserTask task, CancellationToken cancellationToken = default)
 98    {
 10799        if (await FindByMaterializationKeyAsync(task.TenantId, task.MaterializationKey, cancellationToken) is not null)
 1100            return;
 101        try
 102        {
 106103            await documentStore.SaveAsync(CreateRequest(task, expectedVersion: 0), cancellationToken);
 106104        }
 0105        catch (DocumentStoreConcurrencyException)
 106        {
 107            // Projection is idempotent when the same aggregate ID was committed concurrently.
 0108        }
 107109    }
 110
 111    public async Task AppendEventAsync(string tenantId, string taskId, UserTaskEvent @event, CancellationToken cancellat
 112    {
 3113        var existing = await LoadDocumentAsync(tenantId, taskId, cancellationToken);
 3114        if (existing is null)
 1115            return;
 116
 117        // The document store has no separate audit stream, so the entry is written back with the aggregate.
 118        // The revision is deliberately left as-is: an audited read must not consume the concurrency token.
 2119        var loaded = existing.Value;
 2120        loaded.Task.Events.Add(@event);
 121        try
 122        {
 2123            await documentStore.SaveAsync(CreateRequest(loaded.Task, loaded.Document.Version), cancellationToken);
 2124        }
 0125        catch (DocumentStoreConcurrencyException)
 126        {
 127            // A concurrent writer won. The audit entry is advisory, so losing this race is not an error.
 0128        }
 3129    }
 130
 131    public async Task<bool> TryMutateAsync(string tenantId, string taskId, int expectedRevision, Func<UserTask, bool> mu
 132    {
 2133        var existing = await LoadDocumentAsync(tenantId, taskId, cancellationToken);
 2134        if (existing is null)
 0135            return false;
 2136        var loaded = existing.Value;
 2137        if (loaded.Task.Revision != expectedRevision || !mutation(loaded.Task))
 2138            return false;
 0139        loaded.Task.Revision = expectedRevision + 1;
 0140        loaded.Task.UpdatedAt = DateTimeOffset.UtcNow;
 141        try
 142        {
 0143            await documentStore.SaveAsync(CreateRequest(loaded.Task, loaded.Document.Version), cancellationToken);
 0144            return true;
 145        }
 0146        catch (DocumentStoreConcurrencyException)
 147        {
 0148            return false;
 149        }
 2150    }
 151
 152    /// <summary>
 153    /// Loads every stored task across tenants. Only the invitation-hash lookup uses this: an anonymous
 154    /// holder presents a secret and no tenant, so the scan cannot be narrowed by an index.
 155    /// </summary>
 156    private async Task<IReadOnlyCollection<UserTask>> LoadAllAsync(CancellationToken cancellationToken)
 157    {
 2158        var tasks = new List<UserTask>();
 40159        foreach (var status in Enum.GetValues<UserTaskStatus>())
 160        {
 18161            var documents = await documentStore.QueryAsync(new DocumentQuery(StorageUnitName, new Dictionary<string, str
 18162            {
 18163                ["Status"] = status.ToString()
 18164            }), cancellationToken);
 18165            tasks.AddRange(documents.Select(Deserialize));
 166        }
 2167        return tasks;
 2168    }
 169
 170    private async Task<UserTask?> FindByIndexAsync(string tenantId, Func<UserTask, bool> predicate, CancellationToken ca
 171    {
 2175172        foreach (var status in Enum.GetValues<UserTaskStatus>())
 173        {
 978174            var documents = await documentStore.QueryAsync(new DocumentQuery(StorageUnitName, new Dictionary<string, str
 978175            {
 978176                ["TenantId"] = tenantId,
 978177                ["Status"] = status.ToString()
 978178            }), cancellationToken);
 978179            var task = documents.Select(Deserialize).FirstOrDefault(predicate);
 978180            if (task is not null)
 3181                return task;
 182        }
 108183        return null;
 111184    }
 185
 186    private async Task<(StoredDocument Document, UserTask Task)?> LoadDocumentAsync(string tenantId, string taskId, Canc
 187    {
 11188        var document = await documentStore.LoadAsync(StorageUnitName, DocumentId(tenantId, taskId), cancellationToken);
 11189        return document is null ? null : (document, Deserialize(document));
 11190    }
 191
 112192    private static SaveDocumentRequest CreateRequest(UserTask task, long expectedVersion) => new(
 112193        StorageUnitName,
 112194        DocumentId(task.TenantId, task.Id),
 112195        JsonSerializer.Serialize(task, JsonOptions),
 112196        new Dictionary<string, string?>
 112197        {
 112198            ["TenantId"] = task.TenantId,
 112199            ["Status"] = task.Status.ToString(),
 112200            ["MaterializationKey"] = task.MaterializationKey,
 112201            ["BookmarkId"] = task.BookmarkId,
 112202            ["TaskType"] = task.TaskType,
 112203            // Every index the schema provider declares must be supplied on save; the store rejects the
 112204            // write outright when one is absent, so an omission here disables the provider entirely
 112205            // rather than merely losing an index.
 112206            ["WorkflowDefinitionId"] = task.WorkflowDefinitionId,
 112207            ["WorkflowInstanceId"] = task.WorkflowInstanceId,
 112208            ["ActivityInstanceId"] = task.ActivityInstanceId,
 112209            ["CreatedAt"] = task.CreatedAt.ToString("O", System.Globalization.CultureInfo.InvariantCulture),
 112210            ["CompletedAt"] = task.CompletedAt?.ToString("O", System.Globalization.CultureInfo.InvariantCulture),
 112211            ["AssigneeProvider"] = task.Assignee?.Provider,
 112212            ["AssigneeType"] = task.Assignee?.Type.ToString(),
 112213            ["AssigneeId"] = task.Assignee?.Id,
 112214            ["HealthSeverity"] = task.HealthSeverity?.ToString(),
 112215            ["Priority"] = task.Priority.ToString(System.Globalization.CultureInfo.InvariantCulture),
 112216            ["DueAt"] = task.DueAt?.ToString("O", System.Globalization.CultureInfo.InvariantCulture)
 112217        }, expectedVersion);
 218
 1208219    private static UserTask Deserialize(StoredDocument document) => JsonSerializer.Deserialize<UserTask>(document.Conten
 1208220        ?? throw new DocumentStoreValidationException($"Stored User Task document '{document.Id}' could not be deseriali
 221
 145222    private static string DocumentId(string tenantId, string taskId) => $"{tenantId}:{taskId}";
 223
 224    private static bool Matches(UserTask task, UserTaskQuery query)
 225    {
 791226        var search = query.Search?.Trim();
 791227        return task.TenantId == query.TenantId && (query.Statuses.Count == 0 || query.Statuses.Contains(task.Status)) &&
 791228               (!query.OnlyOverdue || task.IsOverdue) && (!query.OnlyWithoutDueDate || task.DueAt is null) &&
 791229               (string.IsNullOrWhiteSpace(query.TaskType) || task.TaskType == query.TaskType) &&
 791230               (!query.PriorityFrom.HasValue || task.Priority >= query.PriorityFrom) && (!query.PriorityTo.HasValue || t
 791231               (!query.DueFrom.HasValue || task.DueAt >= query.DueFrom) && (!query.DueTo.HasValue || task.DueAt <= query
 791232               (string.IsNullOrWhiteSpace(query.WorkflowDefinitionId) || task.WorkflowDefinitionId == query.WorkflowDefi
 791233               (string.IsNullOrWhiteSpace(query.WorkflowInstanceId) || task.WorkflowInstanceId == query.WorkflowInstance
 791234               (string.IsNullOrWhiteSpace(query.Reference) || task.Reference == query.Reference) &&
 802235               (string.IsNullOrWhiteSpace(search) || task.Title.Contains(search, StringComparison.OrdinalIgnoreCase) || 
 236    }
 237
 238    private static bool IsVisible(UserTask task, UserTaskQueryScope scope)
 239    {
 785240        if (!string.Equals(scope.TenantId, task.TenantId, StringComparison.Ordinal) ||
 785241            !string.Equals(scope.Subject.TenantId, task.TenantId, StringComparison.Ordinal) ||
 786242            scope.Groups.Any(group => !string.Equals(group.TenantId, task.TenantId, StringComparison.Ordinal)))
 1243            return false;
 784244        if (scope.ExcludeBlocking && task.HealthSeverity == UserTaskHealthSeverity.Blocking)
 0245            return false;
 246        // Manager-only scopes were already rejected by the policy for non-managers.
 784247        if (scope.RequiresManager)
 0248            return scope.IsManager && (scope.Kind != UserTaskQueryScopeKind.NeedsAttention || NeedsAttention(task));
 249
 784250        var subject = scope.Subject;
 784251        var groups = scope.Groups;
 784252        return scope.Kind switch
 784253        {
 0254            UserTaskQueryScopeKind.Assigned => task.Assignee?.Matches(subject) == true,
 784255            UserTaskQueryScopeKind.Available => task.IsOpen && task.Assignee is null && IsEligible(task, subject, groups
 0256            UserTaskQueryScopeKind.History => task.IsTerminal &&
 0257                                              (task.CompletedBy?.Matches(subject) == true || task.Events.Any(x => x.Acto
 0258            _ => false
 784259        };
 260    }
 261
 262    private static bool IsEligible(UserTask task, ParticipantReference subject, IReadOnlyCollection<ParticipantReference
 263    {
 788264        if (task.ExcludedUsers.Any(x => x.Matches(subject)))
 4265            return false;
 266        // SnapshotGroups are the original group refs. Eligibility is the expanded SnapshotMembers,
 267        // matching DefaultUserTaskAccessPolicy.IsCandidate and InMemory IsEligible.
 780268        if (task.MembershipResolutionMode == UserTaskMembershipResolutionMode.Snapshot)
 5269            return task.SnapshotMembers.Any(x => x.Matches(subject));
 1554270        return task.CandidateUsers.Any(x => x.Matches(subject)) || task.CandidateGroups.Any(x => groups.Any(x.Matches));
 271    }
 272
 273    private static bool NeedsAttention(UserTask task) =>
 0274        task.HealthSeverity == UserTaskHealthSeverity.Blocking
 0275        || task.IsOverdue
 0276        || (task.IsOpen && task.Assignee is null)
 0277        || task.Status is UserTaskStatus.Completing or UserTaskStatus.TimingOut or UserTaskStatus.Cancelling;
 278
 279    // Id is always ThenBy ascending so a page of ties is the same in both directions and across providers.
 140280    private static IEnumerable<UserTask> ApplyOrdering(IEnumerable<UserTask> tasks, UserTaskQuery query) => query.Sort.T
 140281    {
 312282        "priority" => query.Descending ? tasks.OrderByDescending(x => x.Priority).ThenBy(x => x.Id) : tasks.OrderBy(x =>
 312283        "title" => query.Descending ? tasks.OrderByDescending(x => x.Title, TitleComparer).ThenBy(x => x.Id) : tasks.Ord
 494284        "due" => query.Descending ? tasks.OrderBy(x => x.DueAt == null).ThenByDescending(x => x.DueAt).ThenBy(x => x.Id)
 338285        "updated" => query.Descending ? tasks.OrderByDescending(x => x.UpdatedAt).ThenBy(x => x.Id) : tasks.OrderBy(x =>
 376286        _ => query.Descending ? tasks.OrderByDescending(x => x.CreatedAt).ThenBy(x => x.Id) : tasks.OrderBy(x => x.Creat
 140287    };
 288
 289    private static IEnumerable<UserTask> ApplyCursor(IEnumerable<UserTask> tasks, UserTaskQuery query)
 290    {
 140291        if (!TryReadCursor(query.Cursor, out var value, out var id))
 59292            return tasks;
 293        return query.Sort.ToLowerInvariant() switch
 294        {
 128295            "priority" when int.TryParse(value, out var priority) => tasks.Where(x => query.Descending ? x.Priority < pr
 112296            "title" => tasks.Where(x => TitleIsAfterCursor(x.Title, value, x.Id, id, query.Descending)),
 30297            "due" when value == "~null" => tasks.Where(x => x.DueAt == null && string.Compare(x.Id, id) > 0),
 112298            "due" when DateTimeOffset.TryParse(value, out var due) => tasks.Where(x => x.DueAt == null || query.Descendi
 128299            "updated" when DateTimeOffset.TryParse(value, out var updated) => tasks.Where(x => query.Descending ? x.Upda
 132300            _ when DateTimeOffset.TryParse(value, out var created) => tasks.Where(x => query.Descending ? x.CreatedAt < 
 0301            _ => tasks
 302        };
 303    }
 304
 305    private static bool TitleIsAfterCursor(string title, string cursorTitle, string id, string cursorId, bool descending
 306    {
 96307        var comparison = TitleComparer.Compare(title, cursorTitle);
 96308        return descending
 96309            ? comparison < 0 || comparison == 0 && string.Compare(id, cursorId) > 0
 96310            : comparison > 0 || comparison == 0 && string.Compare(id, cursorId) > 0;
 311    }
 312
 313    private static string CreateCursor(UserTask task, string sort)
 314    {
 82315        var value = sort.ToLowerInvariant() switch
 82316        {
 16317            "priority" => task.Priority.ToString(System.Globalization.CultureInfo.InvariantCulture),
 16318            "title" => task.Title,
 16319            "due" => task.DueAt?.ToString("O", System.Globalization.CultureInfo.InvariantCulture) ?? "~null",
 16320            "updated" => task.UpdatedAt.ToString("O", System.Globalization.CultureInfo.InvariantCulture),
 18321            _ => task.CreatedAt.ToString("O", System.Globalization.CultureInfo.InvariantCulture)
 82322        };
 82323        return Convert.ToBase64String(JsonSerializer.SerializeToUtf8Bytes(new[] { value, task.Id }, JsonOptions)).TrimEn
 324    }
 325
 326    private static bool TryReadCursor(string? cursor, out string value, out string id)
 327    {
 140328        value = id = "";
 140329        if (string.IsNullOrWhiteSpace(cursor))
 58330            return false;
 331        try
 332        {
 82333            var padded = cursor.Replace('-', '+').Replace('_', '/') + new string('=', (4 - cursor.Length % 4) % 4);
 82334            var values = JsonSerializer.Deserialize<string[]>(Convert.FromBase64String(padded), JsonOptions);
 81335            if (values is not [var parsedValue, var parsedId])
 0336                return false;
 81337            value = parsedValue;
 81338            id = parsedId;
 81339            return true;
 340        }
 0341        catch (FormatException) { return false; }
 2342        catch (JsonException) { return false; }
 82343    }
 344}