## Kontext: Limits sind Fenster, nicht Summen Die Runs laufen ueber das Claude-Abo. Limits greifen pro 5h-Fenster und pro 7 Tage. Nicht die Wochensumme tut weh, sondern dass ein Agent-Burst ein Fenster leerraeumt, in dem Mika selbst interaktiv arbeiten will. ## Messgrundlage (alle Transcripts unter ~/.claude/projects) Agent-Runs sind ueber die ganze Historie nur **18,4 %** des Account-Verbrauch ClaudeDo-Task: 87105f5e-c4f4-4af4-ae60-89cd2e153e3c
621 lines
21 KiB
C#
621 lines
21 KiB
C#
using ClaudeDo.Data;
|
|
using ClaudeDo.Data.Git;
|
|
using ClaudeDo.Data.Models;
|
|
using ClaudeDo.Data.Repositories;
|
|
using ClaudeDo.Worker.Config;
|
|
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.Services;
|
|
|
|
public sealed class QueueServiceTests : 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 QueueServiceTests()
|
|
{
|
|
_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, // fast for tests
|
|
};
|
|
}
|
|
|
|
public void Dispose()
|
|
{
|
|
_ctx.Dispose();
|
|
_db.Dispose();
|
|
try { Directory.Delete(_tempDir, true); } catch { }
|
|
}
|
|
|
|
private QueueWaker _waker = null!;
|
|
private FakeUsageGate _usageGate = null!;
|
|
private CapturingHubContext _hubContext = null!;
|
|
|
|
private (QueueService service, FakeClaudeProcess fakeProcess) CreateService(
|
|
Func<string, string, IReadOnlyList<string>, Func<string, Task>, CancellationToken, Task<RunResult>>? handler = null,
|
|
FakeUsageGate? usageGate = null,
|
|
UsageState? usageState = null)
|
|
{
|
|
var fake = new FakeClaudeProcess(handler);
|
|
_hubContext = new CapturingHubContext();
|
|
var broadcaster = new HubBroadcaster(_hubContext);
|
|
var dbFactory = _db.CreateFactory();
|
|
var wtManager = new WorktreeManager(new GitService(), dbFactory, _cfg, NullLogger<WorktreeManager>.Instance);
|
|
var argsBuilder = new ClaudeArgsBuilder();
|
|
var built = TaskStateServiceBuilder.Build(dbFactory);
|
|
var state = built.State;
|
|
var runner = new TaskRunner(fake, dbFactory, broadcaster, wtManager, argsBuilder, _cfg,
|
|
NullLogger<TaskRunner>.Instance, state, new TaskRunTokenRegistry(), new AttachmentStore(), new FakeSessionSkillSeeder());
|
|
_waker = new QueueWaker();
|
|
var picker = new QueuePicker(dbFactory);
|
|
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, built.RunCancels);
|
|
_usageGate = usageGate ?? new FakeUsageGate();
|
|
var service = new QueueService(dbFactory, runner, _cfg, NullLogger<QueueService>.Instance, _waker, picker, overrideSlot, state, built.RunCancels,
|
|
_usageGate, usageState ?? new UsageState(), broadcaster);
|
|
return (service, fake);
|
|
}
|
|
|
|
private async Task SetAppSettingsAsync(
|
|
int maxParallel, int softPct = 50, int hardPct = 65, int gateFive = 80, int gateSeven = 90)
|
|
{
|
|
using var ctx = _db.CreateContext();
|
|
var repo = new AppSettingsRepository(ctx);
|
|
var settings = await repo.GetAsync();
|
|
settings.MaxParallelExecutions = maxParallel;
|
|
settings.UsageThrottleSoftPct = softPct;
|
|
settings.UsageThrottleHardPct = hardPct;
|
|
settings.UsageGateFiveHourPct = gateFive;
|
|
settings.UsageGateSevenDayPct = gateSeven;
|
|
await repo.UpdateAsync(settings);
|
|
}
|
|
|
|
private async Task<string> SeedListAsync()
|
|
{
|
|
var listId = Guid.NewGuid().ToString();
|
|
await _listRepo.AddAsync(new ListEntity { Id = listId, Name = "Test", CreatedAt = DateTime.UtcNow });
|
|
return listId;
|
|
}
|
|
|
|
private async Task<TaskEntity> SeedQueuedTask(string listId, DateTime? scheduledFor = null, DateTime? createdAt = null)
|
|
{
|
|
var task = new TaskEntity
|
|
{
|
|
Id = Guid.NewGuid().ToString(),
|
|
ListId = listId,
|
|
Title = "Test task",
|
|
Description = "Do something",
|
|
Status = TaskStatus.Queued,
|
|
ScheduledFor = scheduledFor,
|
|
CreatedAt = createdAt ?? DateTime.UtcNow,
|
|
};
|
|
await _taskRepo.AddAsync(task);
|
|
return task;
|
|
}
|
|
|
|
[Fact]
|
|
public async Task RunNow_Throws_When_Override_Slot_Busy()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
var tcs = new TaskCompletionSource<RunResult>();
|
|
|
|
var (service, _) = CreateService((_, _, _, _, ct) => tcs.Task);
|
|
|
|
var task1 = await SeedQueuedTask(listId);
|
|
var task2 = await SeedQueuedTask(listId);
|
|
|
|
await service.RunNow(task1.Id);
|
|
|
|
var ex = await Assert.ThrowsAsync<InvalidOperationException>(() => service.RunNow(task2.Id));
|
|
Assert.Equal("override slot busy", ex.Message);
|
|
|
|
tcs.SetResult(new RunResult { ExitCode = 0, ResultMarkdown = "ok" });
|
|
}
|
|
|
|
[Fact]
|
|
public async Task RunNow_Throws_For_Unknown_Task()
|
|
{
|
|
var (service, _) = CreateService();
|
|
await Assert.ThrowsAsync<KeyNotFoundException>(() => service.RunNow("nonexistent"));
|
|
}
|
|
|
|
[Fact]
|
|
public async Task ReQueuedReviewTask_ResumesSession_WithFeedbackPrompt_AndClearsFeedback()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
|
|
IReadOnlyList<string>? capturedArgs = null;
|
|
string? capturedPrompt = null;
|
|
var done = new TaskCompletionSource();
|
|
|
|
var (service, _) = CreateService((prompt, _, args, _, _) =>
|
|
{
|
|
capturedPrompt = prompt;
|
|
capturedArgs = args;
|
|
done.TrySetResult();
|
|
return Task.FromResult(new RunResult { ExitCode = 0, SessionId = "sess-2", ResultMarkdown = "ok" });
|
|
});
|
|
|
|
// A task that was reviewed and rejected: Queued + ReviewFeedback, with a prior run carrying a session id.
|
|
var task = new TaskEntity
|
|
{
|
|
Id = Guid.NewGuid().ToString(),
|
|
ListId = listId,
|
|
Title = "Reviewed task",
|
|
Status = TaskStatus.Queued,
|
|
ReviewFeedback = "fix the bug",
|
|
CreatedAt = DateTime.UtcNow,
|
|
};
|
|
await _taskRepo.AddAsync(task);
|
|
await new TaskRunRepository(_ctx).AddAsync(new TaskRunEntity
|
|
{
|
|
Id = Guid.NewGuid().ToString(),
|
|
TaskId = task.Id,
|
|
RunNumber = 1,
|
|
IsRetry = false,
|
|
Prompt = "original",
|
|
SessionId = "sess-1",
|
|
StartedAt = DateTime.UtcNow.AddMinutes(-1),
|
|
});
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
await service.StartAsync(cts.Token);
|
|
_waker.Wake();
|
|
await done.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
|
|
Assert.NotNull(capturedArgs);
|
|
Assert.Contains("--resume", capturedArgs);
|
|
Assert.Contains("sess-1", capturedArgs);
|
|
Assert.Equal("fix the bug", capturedPrompt);
|
|
|
|
// Feedback is cleared after the run reaches a successful terminal state (post-run),
|
|
// so poll rather than asserting on the handler-fired instant.
|
|
var deadline = DateTime.UtcNow.AddSeconds(5);
|
|
TaskEntity? reloaded;
|
|
do
|
|
{
|
|
reloaded = await new TaskRepository(_db.CreateContext()).GetByIdAsync(task.Id);
|
|
if (reloaded?.ReviewFeedback is null) break;
|
|
await Task.Delay(25);
|
|
} while (DateTime.UtcNow < deadline);
|
|
cts.Cancel();
|
|
|
|
Assert.Equal(TaskStatus.WaitingForReview, reloaded!.Status);
|
|
Assert.Null(reloaded.ReviewFeedback);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Schedule_Filter_Skips_Future_Tasks()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
await SeedQueuedTask(listId, scheduledFor: DateTime.UtcNow.AddHours(1));
|
|
|
|
var (service, fake) = CreateService((_, _, _, _, _) =>
|
|
Task.FromResult(new RunResult { ExitCode = 0, ResultMarkdown = "ok" }));
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
|
|
// Start the service loop, wake it, give it time.
|
|
await service.StartAsync(cts.Token);
|
|
_waker.Wake();
|
|
await Task.Delay(200);
|
|
cts.Cancel();
|
|
|
|
// The fake should never have been called because the task is scheduled in the future.
|
|
Assert.Equal(0, fake.CallCount);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Queue_FIFO_Sequentiality()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
|
|
var order = new List<string>();
|
|
var gate1 = new TaskCompletionSource();
|
|
var gate2 = new TaskCompletionSource();
|
|
var callCount = 0;
|
|
|
|
var (service, _) = CreateService(async (_, _, _, _, ct) =>
|
|
{
|
|
var n = Interlocked.Increment(ref callCount);
|
|
lock (order) { order.Add(n.ToString()); }
|
|
if (n == 1) await gate1.Task;
|
|
if (n == 2) gate2.SetResult();
|
|
return new RunResult { ExitCode = 0, ResultMarkdown = "ok" };
|
|
});
|
|
|
|
await SeedQueuedTask(listId, createdAt: DateTime.UtcNow.AddSeconds(-2));
|
|
await SeedQueuedTask(listId, createdAt: DateTime.UtcNow.AddSeconds(-1));
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
await service.StartAsync(cts.Token);
|
|
_waker.Wake();
|
|
|
|
// Wait until task1 has been picked up (poll instead of fixed delay to avoid flake under load).
|
|
var deadline = DateTime.UtcNow.AddSeconds(5);
|
|
while (order.Count == 0 && DateTime.UtcNow < deadline)
|
|
await Task.Delay(20);
|
|
|
|
// Only task1 should be running (task2 waiting on the queue slot).
|
|
Assert.Single(order);
|
|
Assert.Equal("1", order[0]);
|
|
|
|
// Release first task.
|
|
gate1.SetResult();
|
|
|
|
// Wait for second task to complete.
|
|
await gate2.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
|
|
Assert.Equal(2, order.Count);
|
|
Assert.Equal("2", order[1]);
|
|
|
|
cts.Cancel();
|
|
}
|
|
|
|
[Fact]
|
|
public async Task CancelTask_Triggers_Cancellation()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
|
|
var running = new TaskCompletionSource();
|
|
var cancelled = false;
|
|
|
|
var (service, _) = CreateService(async (_, _, _, _, ct) =>
|
|
{
|
|
running.SetResult();
|
|
try
|
|
{
|
|
await Task.Delay(Timeout.Infinite, ct);
|
|
}
|
|
catch (OperationCanceledException)
|
|
{
|
|
cancelled = true;
|
|
throw;
|
|
}
|
|
return new RunResult { ExitCode = 0, ResultMarkdown = "ok" };
|
|
});
|
|
|
|
var task = await SeedQueuedTask(listId);
|
|
await service.RunNow(task.Id);
|
|
|
|
await running.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
|
|
var result = service.CancelTask(task.Id);
|
|
Assert.True(result);
|
|
|
|
await Task.Delay(200);
|
|
Assert.True(cancelled);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task RunNow_AutoRetries_On_Failure_With_SessionId()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
var task = await SeedQueuedTask(listId);
|
|
|
|
var callCount = 0;
|
|
var (service, fake) = CreateService((prompt, dir, args, onLine, ct) =>
|
|
{
|
|
callCount++;
|
|
if (callCount == 1)
|
|
{
|
|
return Task.FromResult(new RunResult
|
|
{
|
|
ExitCode = 1,
|
|
ErrorMarkdown = "something broke",
|
|
SessionId = "sess-retry-test",
|
|
});
|
|
}
|
|
return Task.FromResult(new RunResult
|
|
{
|
|
ExitCode = 0,
|
|
ResultMarkdown = "fixed it",
|
|
SessionId = "sess-retry-test",
|
|
});
|
|
});
|
|
|
|
await service.StartAsync(CancellationToken.None);
|
|
await service.RunNow(task.Id);
|
|
|
|
// Wait for both runs to complete.
|
|
await Task.Delay(2000);
|
|
|
|
await service.StopAsync(CancellationToken.None);
|
|
|
|
Assert.Equal(2, callCount);
|
|
|
|
var finalTask = await _taskRepo.GetByIdAsync(task.Id);
|
|
Assert.NotNull(finalTask);
|
|
// A standalone task that completes successfully now gates on review.
|
|
Assert.Equal(TaskStatus.WaitingForReview, finalTask.Status);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task GetActive_Returns_Running_Slots()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
var tcs = new TaskCompletionSource<RunResult>();
|
|
|
|
var (service, _) = CreateService((_, _, _, _, _) => tcs.Task);
|
|
|
|
var task = await SeedQueuedTask(listId);
|
|
await service.RunNow(task.Id);
|
|
|
|
var active = service.GetActive();
|
|
Assert.Single(active);
|
|
Assert.Equal("override", active[0].slot);
|
|
Assert.Equal(task.Id, active[0].taskId);
|
|
|
|
tcs.SetResult(new RunResult { ExitCode = 0, ResultMarkdown = "ok" });
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Blocked_UsageGate_Skips_Queue_Refill()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
await SeedQueuedTask(listId);
|
|
|
|
var gate = new FakeUsageGate { Decision = new UsageGateDecision(true, "5h-Limit 90% >= 80%") };
|
|
var (service, fake) = CreateService(
|
|
(_, _, _, _, _) => Task.FromResult(new RunResult { ExitCode = 0, ResultMarkdown = "ok" }), gate);
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
await service.StartAsync(cts.Token);
|
|
_waker.Wake();
|
|
await Task.Delay(200);
|
|
cts.Cancel();
|
|
|
|
Assert.Equal(0, fake.CallCount);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task UsageGate_Clearing_Resumes_QueueRefill_OnNextTick()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
await SeedQueuedTask(listId);
|
|
|
|
var gate = new FakeUsageGate { Decision = new UsageGateDecision(true, "5h-Limit 90% >= 80%") };
|
|
var done = new TaskCompletionSource();
|
|
var (service, fake) = CreateService((_, _, _, _, _) =>
|
|
{
|
|
done.TrySetResult();
|
|
return Task.FromResult(new RunResult { ExitCode = 0, ResultMarkdown = "ok" });
|
|
}, gate);
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
await service.StartAsync(cts.Token);
|
|
_waker.Wake();
|
|
await Task.Delay(150);
|
|
Assert.Equal(0, fake.CallCount);
|
|
|
|
// Clear the gate; the 50ms backstop timer in this test's config picks it up.
|
|
gate.Decision = new UsageGateDecision(false, null);
|
|
await done.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
cts.Cancel();
|
|
|
|
Assert.Equal(1, fake.CallCount);
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Blocked_UsageGate_Does_Not_Cancel_AlreadyRunning_Slot()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
await SeedQueuedTask(listId);
|
|
|
|
var running = new TaskCompletionSource();
|
|
var cancelled = false;
|
|
var gate = new FakeUsageGate();
|
|
var (service, _) = CreateService(async (_, _, _, _, ct) =>
|
|
{
|
|
running.SetResult();
|
|
try
|
|
{
|
|
await Task.Delay(Timeout.Infinite, ct);
|
|
}
|
|
catch (OperationCanceledException)
|
|
{
|
|
cancelled = true;
|
|
throw;
|
|
}
|
|
return new RunResult { ExitCode = 0, ResultMarkdown = "ok" };
|
|
}, gate);
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
await service.StartAsync(cts.Token);
|
|
_waker.Wake();
|
|
await running.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
|
|
// Block after the slot is already running — several backstop ticks pass.
|
|
gate.Decision = new UsageGateDecision(true, "5h-Limit 90% >= 80%");
|
|
await Task.Delay(200);
|
|
|
|
Assert.False(cancelled);
|
|
cts.Cancel();
|
|
}
|
|
|
|
[Fact]
|
|
public async Task UsageGate_TransitionLogging_FiresOncePerChange()
|
|
{
|
|
var gate = new FakeUsageGate { Decision = new UsageGateDecision(true, "5h-Limit 90% >= 80%") };
|
|
var (service, _) = CreateService(
|
|
(_, _, _, _, _) => Task.FromResult(new RunResult { ExitCode = 0, ResultMarkdown = "ok" }), gate);
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
await service.StartAsync(cts.Token);
|
|
|
|
// Several backstop ticks (50ms interval) all observe the same blocked state.
|
|
await Task.Delay(200);
|
|
cts.Cancel();
|
|
|
|
var warnCalls = _hubContext.Proxy.Calls
|
|
.Count(c => c.Method == "WorkerLog" && (WorkerLogLevel)c.Args[1]! == WorkerLogLevel.Warn);
|
|
Assert.Equal(1, warnCalls);
|
|
}
|
|
|
|
// Polls until `read()` reaches `expected` (or times out), then waits a further grace period
|
|
// to make sure the count doesn't keep climbing past it — needed because slot fills happen
|
|
// concurrently and a fixed sleep is either flaky (too short) or slow (too long).
|
|
private static async Task AssertStableCountAsync(Func<int> read, int expected)
|
|
{
|
|
var deadline = DateTime.UtcNow.AddSeconds(5);
|
|
while (read() < expected && DateTime.UtcNow < deadline)
|
|
await Task.Delay(20);
|
|
|
|
await Task.Delay(250);
|
|
Assert.Equal(expected, read());
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Throttle_StepsDownEffectiveSlots_BelowConfiguredMax()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
await SeedQueuedTask(listId);
|
|
await SeedQueuedTask(listId);
|
|
await SeedQueuedTask(listId);
|
|
|
|
await SetAppSettingsAsync(maxParallel: 3);
|
|
|
|
var usageState = new UsageState();
|
|
usageState.ReportSuccess(new UsageSnapshot(
|
|
new UsageBucket(60, null), new UsageBucket(0, null), Array.Empty<UsageLimitRow>(), DateTime.UtcNow));
|
|
|
|
var startedCount = 0;
|
|
var block = new TaskCompletionSource();
|
|
var (service, _) = CreateService(async (_, _, _, _, _) =>
|
|
{
|
|
Interlocked.Increment(ref startedCount);
|
|
await block.Task;
|
|
return new RunResult { ExitCode = 0, ResultMarkdown = "ok" };
|
|
}, usageState: usageState);
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
await service.StartAsync(cts.Token);
|
|
_waker.Wake();
|
|
|
|
// 60% is between the soft (50) and hard (65) thresholds — capped at 2 slots even
|
|
// though 3 are configured and 3 tasks are queued.
|
|
await AssertStableCountAsync(() => Volatile.Read(ref startedCount), 2);
|
|
|
|
block.SetResult();
|
|
cts.Cancel();
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Throttle_AtHardThreshold_CapsToOneSlot()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
await SeedQueuedTask(listId);
|
|
await SeedQueuedTask(listId);
|
|
|
|
await SetAppSettingsAsync(maxParallel: 3);
|
|
|
|
var usageState = new UsageState();
|
|
usageState.ReportSuccess(new UsageSnapshot(
|
|
new UsageBucket(70, null), new UsageBucket(0, null), Array.Empty<UsageLimitRow>(), DateTime.UtcNow));
|
|
|
|
var startedCount = 0;
|
|
var block = new TaskCompletionSource();
|
|
var (service, _) = CreateService(async (_, _, _, _, _) =>
|
|
{
|
|
Interlocked.Increment(ref startedCount);
|
|
await block.Task;
|
|
return new RunResult { ExitCode = 0, ResultMarkdown = "ok" };
|
|
}, usageState: usageState);
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
await service.StartAsync(cts.Token);
|
|
_waker.Wake();
|
|
|
|
await AssertStableCountAsync(() => Volatile.Read(ref startedCount), 1);
|
|
|
|
block.SetResult();
|
|
cts.Cancel();
|
|
}
|
|
|
|
[Fact]
|
|
public async Task NoUsageSnapshot_FallsBackToFullConfiguredParallelism()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
await SeedQueuedTask(listId);
|
|
await SeedQueuedTask(listId);
|
|
await SeedQueuedTask(listId);
|
|
|
|
await SetAppSettingsAsync(maxParallel: 3);
|
|
|
|
// No snapshot has landed yet (fresh UsageState) — throttle must fail open.
|
|
var startedCount = 0;
|
|
var block = new TaskCompletionSource();
|
|
var (service, _) = CreateService(async (_, _, _, _, _) =>
|
|
{
|
|
Interlocked.Increment(ref startedCount);
|
|
await block.Task;
|
|
return new RunResult { ExitCode = 0, ResultMarkdown = "ok" };
|
|
}, usageState: new UsageState());
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
await service.StartAsync(cts.Token);
|
|
_waker.Wake();
|
|
|
|
await AssertStableCountAsync(() => Volatile.Read(ref startedCount), 3);
|
|
|
|
block.SetResult();
|
|
cts.Cancel();
|
|
}
|
|
|
|
[Fact]
|
|
public async Task Throttle_Engaging_Does_Not_Cancel_AlreadyRunning_Slot()
|
|
{
|
|
var listId = await SeedListAsync();
|
|
await SeedQueuedTask(listId);
|
|
|
|
await SetAppSettingsAsync(maxParallel: 3);
|
|
var usageState = new UsageState();
|
|
|
|
var running = new TaskCompletionSource();
|
|
var cancelled = false;
|
|
var (service, _) = CreateService(async (_, _, _, _, ct) =>
|
|
{
|
|
running.SetResult();
|
|
try
|
|
{
|
|
await Task.Delay(Timeout.Infinite, ct);
|
|
}
|
|
catch (OperationCanceledException)
|
|
{
|
|
cancelled = true;
|
|
throw;
|
|
}
|
|
return new RunResult { ExitCode = 0, ResultMarkdown = "ok" };
|
|
}, usageState: usageState);
|
|
|
|
using var cts = new CancellationTokenSource();
|
|
await service.StartAsync(cts.Token);
|
|
_waker.Wake();
|
|
await running.Task.WaitAsync(TimeSpan.FromSeconds(5));
|
|
|
|
// Throttle engages hard after the slot is already running — several backstop ticks pass.
|
|
usageState.ReportSuccess(new UsageSnapshot(
|
|
new UsageBucket(70, null), new UsageBucket(0, null), Array.Empty<UsageLimitRow>(), DateTime.UtcNow));
|
|
await Task.Delay(200);
|
|
|
|
Assert.False(cancelled);
|
|
cts.Cancel();
|
|
}
|
|
}
|