using ClaudeDo.Data; using ClaudeDo.Data.Git; using ClaudeDo.Data.Models; using ClaudeDo.Data.Repositories; using ClaudeDo.Worker.Config; using ClaudeDo.Worker.External; using ClaudeDo.Worker.Hub; using ClaudeDo.Worker.Queue; using ClaudeDo.Worker.Runner; using ClaudeDo.Worker.Tests.Infrastructure; using ClaudeDo.Worker.Usage; using Microsoft.Extensions.Logging.Abstractions; using TaskStatus = ClaudeDo.Data.Models.TaskStatus; namespace ClaudeDo.Worker.Tests.External; public sealed class QueueStateMcpToolsTests : IDisposable { private readonly DbFixture _db = new(); private readonly ClaudeDoDbContext _ctx; private readonly TaskRepository _taskRepo; private readonly ListRepository _listRepo; private readonly WorkerConfig _cfg; private readonly string _tempDir; public QueueStateMcpToolsTests() { _ctx = _db.CreateContext(); _taskRepo = new TaskRepository(_ctx); _listRepo = new ListRepository(_ctx); _tempDir = Path.Combine(Path.GetTempPath(), $"claudedo_test_{Guid.NewGuid():N}"); Directory.CreateDirectory(_tempDir); _cfg = new WorkerConfig { SandboxRoot = Path.Combine(_tempDir, "sandbox"), LogRoot = Path.Combine(_tempDir, "logs"), QueueBackstopIntervalMs = 50, }; } public void Dispose() { _ctx.Dispose(); _db.Dispose(); try { Directory.Delete(_tempDir, true); } catch { } } private (QueueService queue, QueueStateMcpTools sut) CreateSut( Func, Func, CancellationToken, Task>? handler = null, UsageState? usageState = null) { var dbFactory = _db.CreateFactory(); var fake = new FakeClaudeProcess(handler); var broadcaster = new HubBroadcaster(new CapturingHubContext()); var wtManager = new WorktreeManager(new GitService(), dbFactory, _cfg, NullLogger.Instance); var built = TaskStateServiceBuilder.Build(dbFactory); var runner = new TaskRunner(fake, dbFactory, broadcaster, wtManager, new ClaudeArgsBuilder(), _cfg, NullLogger.Instance, built.State, new TaskRunTokenRegistry(), new AttachmentStore(), new FakeSessionSkillSeeder(), new FakeTranscriptUsageReader()); var picker = new QueuePicker(dbFactory); var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger.Instance, built.RunCancels); var queue = new QueueService(dbFactory, runner, _cfg, NullLogger.Instance, new QueueWaker(), picker, overrideSlot, built.State, built.RunCancels, new FakeUsageGate(), usageState ?? new UsageState(), broadcaster); return (queue, new QueueStateMcpTools(queue, dbFactory)); } private async Task SeedListAsync() { var listId = Guid.NewGuid().ToString(); await _listRepo.AddAsync(new ListEntity { Id = listId, Name = "Test", CreatedAt = DateTime.UtcNow }); return listId; } private async Task SeedTaskAsync( string listId, TaskStatus status, int sortOrder = 0, DateTime? createdAt = null, bool isManual = false, string? blockedByTaskId = null, DateTime? scheduledFor = null) { var task = new TaskEntity { Id = Guid.NewGuid().ToString(), ListId = listId, Title = "Test task", Status = status, SortOrder = sortOrder, CreatedAt = createdAt ?? DateTime.UtcNow, IsManual = isManual, BlockedByTaskId = blockedByTaskId, ScheduledFor = scheduledFor, }; // Bypass TaskRepository.AddAsync, which overwrites SortOrder with max(listId)+1 -- // these tests need to control SortOrder directly to exercise queue pick order. _ctx.Tasks.Add(task); await _ctx.SaveChangesAsync(); return task; } private async Task SetMaxParallelAsync(int maxParallel, int softPct = 0, int hardPct = 0) { using var ctx = _db.CreateContext(); var repo = new AppSettingsRepository(ctx); var settings = await repo.GetAsync(); settings.MaxParallelExecutions = maxParallel; settings.UsageThrottleSoftPct = softPct; settings.UsageThrottleHardPct = hardPct; await repo.UpdateAsync(settings); } [Fact] public async Task GetQueueState_NoThrottle_EffectiveEqualsConfigured() { await SetMaxParallelAsync(maxParallel: 3); var (_, sut) = CreateSut(); var result = await sut.GetQueueState(CancellationToken.None); Assert.Equal(3, result.ConfiguredSlots); Assert.Equal(3, result.EffectiveSlots); } [Fact] public async Task GetQueueState_UsageThrottleActive_EffectiveBelowConfigured() { await SetMaxParallelAsync(maxParallel: 3, softPct: 50, hardPct: 65); var usageState = new UsageState(); usageState.ReportSuccess(new UsageSnapshot( new UsageBucket(60, null), new UsageBucket(0, null), Array.Empty(), DateTime.UtcNow)); var (_, sut) = CreateSut(usageState: usageState); var result = await sut.GetQueueState(CancellationToken.None); Assert.Equal(3, result.ConfiguredSlots); Assert.Equal(2, result.EffectiveSlots); } [Fact] public async Task GetQueueState_ActiveSlots_ReflectsOverrideSlot() { var listId = await SeedListAsync(); var tcs = new TaskCompletionSource(); var (queue, sut) = CreateSut((_, _, _, _, _) => tcs.Task); var task = await SeedTaskAsync(listId, TaskStatus.Queued); await queue.RunNow(task.Id); var result = await sut.GetQueueState(CancellationToken.None); var slot = Assert.Single(result.ActiveSlots); Assert.Equal("override", slot.Slot); Assert.Equal(task.Id, slot.TaskId); tcs.SetResult(new RunResult { ExitCode = 0, ResultMarkdown = "ok" }); } [Fact] public async Task GetQueueState_WaitingTaskIds_OrderedBySortOrderThenCreatedAt_AndFiltersIneligible() { var listId = await SeedListAsync(); var (_, sut) = CreateSut(); var second = await SeedTaskAsync(listId, TaskStatus.Queued, sortOrder: 1, createdAt: DateTime.UtcNow); var first = await SeedTaskAsync(listId, TaskStatus.Queued, sortOrder: 0, createdAt: DateTime.UtcNow.AddMinutes(1)); await SeedTaskAsync(listId, TaskStatus.Queued, sortOrder: 2, isManual: true); await SeedTaskAsync(listId, TaskStatus.Queued, sortOrder: 3, blockedByTaskId: first.Id); await SeedTaskAsync(listId, TaskStatus.Queued, sortOrder: 4, scheduledFor: DateTime.UtcNow.AddHours(1)); await SeedTaskAsync(listId, TaskStatus.Running, sortOrder: -1); var result = await sut.GetQueueState(CancellationToken.None); Assert.Equal(new[] { first.Id, second.Id }, result.WaitingTaskIds); } }