using Microsoft.EntityFrameworkCore; using Nexus.Api.Data; using Nexus.Api.Models; using System.Text.Json; namespace Nexus.Api.Repositories; public sealed class TaskRepository(NexusDbContext db) : ITaskRepository { public Task> GetAllAsync(CancellationToken ct = default) => db.Tasks.AsNoTracking().OrderByDescending(x => x.UpdatedAt).ToListAsync(ct); public ValueTask GetByIdAsync(Guid id, CancellationToken ct = default) => db.Tasks.FindAsync([id], ct); public Task> GetPendingApprovalAsync(CancellationToken ct = default) { var threshold = DateTimeOffset.UtcNow.AddHours(-1); return db.Tasks.AsNoTracking() .Where(x => x.State == TaskStateHelper.ToStateString(TaskState.InProgress) && x.UpdatedAt <= threshold) .OrderByDescending(x => x.UpdatedAt) .ToListAsync(ct); } public async Task AddAsync(WorkTask task, CancellationToken ct = default) { db.Tasks.Add(task); db.OutboxEvents.Add(CreateTaskEvent("task.created", task)); await db.SaveChangesAsync(ct); return task; } public async Task TryResetStaleInProgressToBacklogAsync( Guid id, DateTimeOffset staleBefore, DateTimeOffset updatedAt, CancellationToken ct = default) { if (!db.Database.IsRelational()) { var task = await db.Tasks .FirstOrDefaultAsync(task => task.Id == id && task.State == TaskStateHelper.ToStateString(TaskState.InProgress) && task.UpdatedAt < staleBefore, ct); if (task is null) { return false; } task.State = TaskStateHelper.ToStateString(TaskState.Backlog); task.UpdatedAt = updatedAt; db.OutboxEvents.Add(CreateTaskEvent("task.updated", task)); await db.SaveChangesAsync(ct); return true; } await using var transaction = await db.Database.BeginTransactionAsync(ct); var affectedRows = await db.Tasks .Where(task => task.Id == id && task.State == TaskStateHelper.ToStateString(TaskState.InProgress) && task.UpdatedAt < staleBefore) .ExecuteUpdateAsync(setters => setters .SetProperty(task => task.State, TaskStateHelper.ToStateString(TaskState.Backlog)) .SetProperty(task => task.UpdatedAt, updatedAt), ct); if (affectedRows > 0) { var updatedTask = await db.Tasks .AsNoTracking() .SingleAsync(task => task.Id == id, ct); db.OutboxEvents.Add(CreateTaskEvent("task.updated", updatedTask)); await db.SaveChangesAsync(ct); } await transaction.CommitAsync(ct); return affectedRows > 0; } public async Task UpdateAsync(WorkTask task, CancellationToken ct = default) { task.UpdatedAt = DateTimeOffset.UtcNow; db.Tasks.Update(task); db.OutboxEvents.Add(CreateTaskEvent("task.updated", task)); await db.SaveChangesAsync(ct); } public async Task DeleteAsync(WorkTask task, CancellationToken ct = default) { db.OutboxEvents.Add(CreateTaskEvent("task.deleted", task)); db.Tasks.Remove(task); await db.SaveChangesAsync(ct); } public Task CountAsync(CancellationToken ct = default) => db.Tasks.CountAsync(ct); public Task CountByStateAsync(string state, CancellationToken ct = default) => db.Tasks.CountAsync(x => x.State == state, ct); public Task GetLastBlockedAsync(CancellationToken ct = default) => db.Tasks.AsNoTracking() .Where(x => x.State == TaskStateHelper.ToStateString(TaskState.Blocked)) .OrderByDescending(x => x.UpdatedAt) .FirstOrDefaultAsync(ct); public async Task GetBoardPageAsync( int doneLimit, DateTimeOffset? doneBeforeUpdatedAt, Guid? doneBeforeId, CancellationToken ct = default) { var doneState = TaskStateHelper.ToStateString(TaskState.Done); var activeQuery = db.Tasks .AsNoTracking() .Where(task => task.State != doneState) .OrderBy(task => task.State == "Backlog" ? 0 : task.State == "In progress" ? 1 : task.State == "Review" ? 2 : task.State == "Blocked" ? 3 : 4) .ThenByDescending(task => task.Priority == "High" ? 3 : task.Priority == "Medium" || task.Priority == "Normal" ? 2 : task.Priority == "Low" ? 1 : 2) .ThenBy(task => task.CreatedAt) .ThenBy(task => task.Id); var doneQuery = db.Tasks .AsNoTracking() .Where(task => task.State == doneState); if (doneBeforeUpdatedAt.HasValue && doneBeforeId.HasValue) { var cursorUpdatedAt = doneBeforeUpdatedAt.Value; var cursorId = doneBeforeId.Value; doneQuery = doneQuery.Where(task => task.UpdatedAt < cursorUpdatedAt || (task.UpdatedAt == cursorUpdatedAt && task.Id.CompareTo(cursorId) < 0)); } doneQuery = doneQuery .OrderByDescending(task => task.UpdatedAt) .ThenByDescending(task => task.Id); // Three SQL statements total. Child counts and latest activity are // correlated scalar subqueries inside each projected statement rather // than per-card round trips. var isDoneContinuation = doneBeforeUpdatedAt.HasValue && doneBeforeId.HasValue; var activeTasks = isDoneContinuation ? [] : await ProjectBoardCards(activeQuery).ToListAsync(ct); var doneTasks = await ProjectBoardCards(doneQuery) .Take(doneLimit + 1) .ToListAsync(ct); var revision = await db.Tasks .AsNoTracking() .OrderByDescending(task => task.UpdatedAt) .ThenByDescending(task => task.Id) .Select(task => new TaskBoardRevisionPoint(task.UpdatedAt, task.Id)) .FirstOrDefaultAsync(ct); var hasMoreDone = doneTasks.Count > doneLimit; if (hasMoreDone) doneTasks.RemoveAt(doneLimit); return new TaskBoardQueryPage(activeTasks, doneTasks, hasMoreDone, revision); } public Task GetBoardCardAsync( Guid id, CancellationToken ct = default) => ProjectBoardCards( db.Tasks .AsNoTracking() .Where(task => task.Id == id)) .SingleOrDefaultAsync(ct); private IQueryable ProjectBoardCards(IQueryable query) => query.Select(task => new TaskBoardCardDto( task.Id, task.Title, task.Detail, task.Source, task.State, task.Priority, task.AssignedTo, task.ParentTaskId, task.ProjectId, task.DueDate, task.CreatedAt, task.UpdatedAt, task.IsAgentTask, task.ExpectedFrom, db.Activity .Where(activity => activity.TaskId == task.Id) .OrderByDescending(activity => activity.CreatedAt) .ThenByDescending(activity => activity.Id) .Select(activity => activity.Message) .FirstOrDefault(), db.Activity .Where(activity => activity.TaskId == task.Id) .OrderByDescending(activity => activity.CreatedAt) .ThenByDescending(activity => activity.Id) .Select(activity => (DateTimeOffset?)activity.CreatedAt) .FirstOrDefault(), db.Tasks.Count(child => child.ParentTaskId == task.Id), db.Tasks.Count(child => child.ParentTaskId == task.Id && child.State != "Done"), task.ParentTaskId.HasValue || task.IsAgentTask || db.Tasks.Any(child => child.ParentTaskId == task.Id))); private static OutboxEvent CreateTaskEvent(string type, WorkTask task) => new() { Type = type, AggregateType = "task", AggregateId = task.Id.ToString(), AggregateRevision = 0, PayloadJson = JsonSerializer.Serialize(new { entityType = "task", id = task.Id, state = task.State, projectId = task.ProjectId, assignedAgentId = task.AssignedTo }) }; }