Merge task branch for: Worker: UsageGate — Queue ab Schwelle pausieren
This commit is contained in:
@@ -8,6 +8,7 @@ using ClaudeDo.Worker.Planning;
|
||||
using ClaudeDo.Worker.Queue;
|
||||
using ClaudeDo.Worker.Runner;
|
||||
using ClaudeDo.Worker.Tests.Infrastructure;
|
||||
using ClaudeDo.Worker.Usage;
|
||||
using ClaudeDo.Worker.Worktrees;
|
||||
using ClaudeDo.Worker.Config;
|
||||
using Microsoft.Extensions.Logging.Abstractions;
|
||||
@@ -75,7 +76,8 @@ public sealed class AddSubtaskToolTests : IDisposable
|
||||
var picker = new ClaudeDo.Worker.Queue.QueuePicker(dbFactory);
|
||||
var runCancels = new RunCancellationRegistry();
|
||||
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, runCancels);
|
||||
var queue = new QueueService(dbFactory, runner, cfg, NullLogger<QueueService>.Instance, waker, picker, overrideSlot, state, runCancels);
|
||||
var queue = new QueueService(dbFactory, runner, cfg, NullLogger<QueueService>.Instance, waker, picker, overrideSlot, state, runCancels,
|
||||
new FakeUsageGate(), new UsageState(), broadcaster);
|
||||
var maintenance = new WorktreeMaintenanceService(dbFactory, git, NullLogger<WorktreeMaintenanceService>.Instance);
|
||||
var merge = new TaskMergeService(dbFactory, git, broadcaster, state, NullLogger<TaskMergeService>.Instance);
|
||||
var aggregator = new PlanningAggregator(dbFactory, git, NullLogger<PlanningAggregator>.Instance);
|
||||
|
||||
@@ -10,6 +10,7 @@ using ClaudeDo.Worker.Planning;
|
||||
using ClaudeDo.Worker.Queue;
|
||||
using ClaudeDo.Worker.Runner;
|
||||
using ClaudeDo.Worker.Tests.Infrastructure;
|
||||
using ClaudeDo.Worker.Usage;
|
||||
using ClaudeDo.Worker.Worktrees;
|
||||
using Microsoft.Extensions.Logging.Abstractions;
|
||||
using TaskStatus = ClaudeDo.Data.Models.TaskStatus;
|
||||
@@ -96,7 +97,8 @@ public sealed class BatchMcpToolsTests : IDisposable
|
||||
var runCancels = new RunCancellationRegistry();
|
||||
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, runCancels);
|
||||
return new QueueService(dbFactory, runner, cfg, NullLogger<QueueService>.Instance,
|
||||
new QueueWaker(), new QueuePicker(dbFactory), overrideSlot, state, runCancels);
|
||||
new QueueWaker(), new QueuePicker(dbFactory), overrideSlot, state, runCancels,
|
||||
new FakeUsageGate(), new UsageState(), broadcaster);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
|
||||
@@ -10,6 +10,7 @@ using ClaudeDo.Worker.Planning;
|
||||
using ClaudeDo.Worker.Queue;
|
||||
using ClaudeDo.Worker.Runner;
|
||||
using ClaudeDo.Worker.Tests.Infrastructure;
|
||||
using ClaudeDo.Worker.Usage;
|
||||
using ClaudeDo.Worker.Worktrees;
|
||||
using Microsoft.AspNetCore.SignalR;
|
||||
using Microsoft.Extensions.Logging.Abstractions;
|
||||
@@ -161,7 +162,8 @@ public sealed class ExternalMcpServiceTests : IDisposable
|
||||
var picker = new ClaudeDo.Worker.Queue.QueuePicker(dbFactory);
|
||||
var runCancels = new RunCancellationRegistry();
|
||||
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, runCancels);
|
||||
return new QueueService(dbFactory, runner, cfg, NullLogger<QueueService>.Instance, waker, picker, overrideSlot, state, runCancels);
|
||||
return new QueueService(dbFactory, runner, cfg, NullLogger<QueueService>.Instance, waker, picker, overrideSlot, state, runCancels,
|
||||
new FakeUsageGate(), new UsageState(), broadcaster);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
using ClaudeDo.Worker.Usage;
|
||||
using ClaudeDo.Worker.Usage.Interfaces;
|
||||
|
||||
namespace ClaudeDo.Worker.Tests.Infrastructure;
|
||||
|
||||
public sealed class FakeUsageGate : IUsageGate
|
||||
{
|
||||
public UsageGateDecision Decision { get; set; } = new(false, null);
|
||||
|
||||
public Task<UsageGateDecision> EvaluateAsync(CancellationToken ct = default) => Task.FromResult(Decision);
|
||||
}
|
||||
@@ -6,6 +6,7 @@ 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;
|
||||
|
||||
@@ -59,7 +60,8 @@ public sealed class QueueServiceSlotGuardTests : IDisposable
|
||||
_waker = new QueueWaker();
|
||||
var picker = new QueuePicker(dbFactory);
|
||||
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, built.RunCancels);
|
||||
var service = new QueueService(dbFactory, runner, _cfg, NullLogger<QueueService>.Instance, _waker, picker, overrideSlot, state, built.RunCancels);
|
||||
var service = new QueueService(dbFactory, runner, _cfg, NullLogger<QueueService>.Instance, _waker, picker, overrideSlot, state, built.RunCancels,
|
||||
new FakeUsageGate(), new UsageState(), broadcaster);
|
||||
return (service, fake);
|
||||
}
|
||||
|
||||
|
||||
@@ -7,6 +7,7 @@ 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;
|
||||
|
||||
@@ -44,12 +45,16 @@ public sealed class QueueServiceTests : IDisposable
|
||||
}
|
||||
|
||||
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)
|
||||
Func<string, string, IReadOnlyList<string>, Func<string, Task>, CancellationToken, Task<RunResult>>? handler = null,
|
||||
FakeUsageGate? usageGate = null)
|
||||
{
|
||||
var fake = new FakeClaudeProcess(handler);
|
||||
var broadcaster = new HubBroadcaster(new CapturingHubContext());
|
||||
_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();
|
||||
@@ -60,7 +65,9 @@ public sealed class QueueServiceTests : IDisposable
|
||||
_waker = new QueueWaker();
|
||||
var picker = new QueuePicker(dbFactory);
|
||||
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, built.RunCancels);
|
||||
var service = new QueueService(dbFactory, runner, _cfg, NullLogger<QueueService>.Instance, _waker, picker, overrideSlot, state, built.RunCancels);
|
||||
_usageGate = usageGate ?? new FakeUsageGate();
|
||||
var service = new QueueService(dbFactory, runner, _cfg, NullLogger<QueueService>.Instance, _waker, picker, overrideSlot, state, built.RunCancels,
|
||||
_usageGate, new UsageState(), broadcaster);
|
||||
return (service, fake);
|
||||
}
|
||||
|
||||
@@ -342,4 +349,107 @@ public sealed class QueueServiceTests : IDisposable
|
||||
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,105 @@
|
||||
using ClaudeDo.Data.Repositories;
|
||||
using ClaudeDo.Worker.Tests.Infrastructure;
|
||||
using ClaudeDo.Worker.Usage;
|
||||
using Microsoft.Extensions.Logging.Abstractions;
|
||||
|
||||
namespace ClaudeDo.Worker.Tests.Usage;
|
||||
|
||||
public sealed class UsageGateTests : IDisposable
|
||||
{
|
||||
private readonly DbFixture _db = new();
|
||||
|
||||
public void Dispose() => _db.Dispose();
|
||||
|
||||
private async Task SetThresholdsAsync(int fiveHourPct, int sevenDayPct)
|
||||
{
|
||||
using var ctx = _db.CreateContext();
|
||||
var repo = new AppSettingsRepository(ctx);
|
||||
var settings = await repo.GetAsync();
|
||||
settings.UsageGateFiveHourPct = fiveHourPct;
|
||||
settings.UsageGateSevenDayPct = sevenDayPct;
|
||||
await repo.UpdateAsync(settings);
|
||||
}
|
||||
|
||||
private UsageGate CreateGate(UsageState state) =>
|
||||
new(_db.CreateFactory(), state, NullLogger<UsageGate>.Instance);
|
||||
|
||||
private static UsageState StateWithSnapshot(double fiveHourPct, double sevenDayPct)
|
||||
{
|
||||
var state = new UsageState();
|
||||
state.ReportSuccess(new UsageSnapshot(
|
||||
new UsageBucket(fiveHourPct, null),
|
||||
new UsageBucket(sevenDayPct, null),
|
||||
Array.Empty<UsageLimitRow>(),
|
||||
DateTime.UtcNow));
|
||||
return state;
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task BothUnderThreshold_NotBlocked()
|
||||
{
|
||||
await SetThresholdsAsync(80, 90);
|
||||
var decision = await CreateGate(StateWithSnapshot(50, 60)).EvaluateAsync();
|
||||
Assert.False(decision.IsBlocked);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task FiveHourAtOrOverThreshold_Blocked()
|
||||
{
|
||||
await SetThresholdsAsync(80, 90);
|
||||
var decision = await CreateGate(StateWithSnapshot(85, 60)).EvaluateAsync();
|
||||
Assert.True(decision.IsBlocked);
|
||||
Assert.Contains("5h", decision.Reason);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task SevenDayAtOrOverThreshold_Blocked()
|
||||
{
|
||||
await SetThresholdsAsync(80, 90);
|
||||
var decision = await CreateGate(StateWithSnapshot(50, 95)).EvaluateAsync();
|
||||
Assert.True(decision.IsBlocked);
|
||||
Assert.Contains("7d", decision.Reason);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task ExactlyOnThreshold_Blocked()
|
||||
{
|
||||
await SetThresholdsAsync(80, 90);
|
||||
var decision = await CreateGate(StateWithSnapshot(80, 60)).EvaluateAsync();
|
||||
Assert.True(decision.IsBlocked);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task ThresholdZero_ThatBucketNeverGates()
|
||||
{
|
||||
await SetThresholdsAsync(0, 90);
|
||||
var decision = await CreateGate(StateWithSnapshot(100, 60)).EvaluateAsync();
|
||||
Assert.False(decision.IsBlocked);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task BothThresholdsZero_AlwaysFree()
|
||||
{
|
||||
await SetThresholdsAsync(0, 0);
|
||||
var decision = await CreateGate(StateWithSnapshot(100, 100)).EvaluateAsync();
|
||||
Assert.False(decision.IsBlocked);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task NoSnapshotYet_FailsOpen()
|
||||
{
|
||||
await SetThresholdsAsync(80, 90);
|
||||
var decision = await CreateGate(new UsageState()).EvaluateAsync();
|
||||
Assert.False(decision.IsBlocked);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task LastPollFailed_FailsOpen()
|
||||
{
|
||||
await SetThresholdsAsync(80, 90);
|
||||
var state = StateWithSnapshot(95, 95);
|
||||
state.ReportFailure("boom", DateTime.UtcNow);
|
||||
var decision = await CreateGate(state).EvaluateAsync();
|
||||
Assert.False(decision.IsBlocked);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user