feat(worker): Toggle "Continue on session limit reset" für Usage-Limit-Abbrüche

Klassifiziert einen echten Usage-Limit-Abbruch als eigene FailureReason
"usage_limit" (TaskRunner.ClassifyFailureReason: nur bei terminal_reason
"api_error" plus einem Limit-Muster im gerenderten Fehlertext, nicht an
Status==Failed allein). Neuer Toggle AutoContinueOnUsageLimit (app_settings,
Default aus) unter Settings → General → "Usage limit stop":

- UsageLimitAutoContinueCoordinator feuert pro Task genau einmal ContinueTask
  über OverrideSlotService, sobald das 5h-Fenster (UsageState.Snapshot.FiveHour
  .ResetsAt) tatsächlich zurückgesetzt ist; ein persistenter Marker
  (TaskEntity.UsageLimitAutoContinuedAt) verhindert einen zweiten Anlauf bei
  einem erneuten Limit-Treffer.
- QueueService schedult zusätzlich einen exakten Wake-Timer auf den
  Reset-Zeitpunkt, statt nur auf den 30s-Backstop zu warten.
- Fail-open durchgängig: kein Snapshot/keine Reset-Zeit → kein Timer, kein
  Continue, kein Throw. Toggle aus ändert das heutige Verhalten nicht.

Migration AddUsageLimitAutoContinue fügt beide Spalten hinzu; die von
`dotnet ef migrations add` mitgescaffoldete leere UpdateData auf app_settings
(columns/values: []) erzeugte ungültiges SQL ("near WHERE") und wurde entfernt
— TaskNumberMigrationTests deckte das über den vollen Migrate()-Pfad auf.
This commit is contained in:
mika kuns
2026-08-21 18:43:02 +02:00
parent 290dd1b614
commit 07dd75700d
30 changed files with 1505 additions and 19 deletions
@@ -77,8 +77,10 @@ public sealed class AddSubtaskToolTests : IDisposable
var picker = new ClaudeDo.Worker.Queue.QueuePicker(dbFactory);
var runCancels = new RunCancellationRegistry(NullLogger<RunCancellationRegistry>.Instance);
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, runCancels);
var usageState = new UsageState();
var queue = new QueueService(dbFactory, runner, cfg, NullLogger<QueueService>.Instance, waker, picker, overrideSlot, state, runCancels,
new FakeUsageGate(), new UsageState(), broadcaster);
new FakeUsageGate(), usageState, broadcaster,
new UsageLimitAutoContinueCoordinator(dbFactory, usageState, overrideSlot, broadcaster, NullLogger<UsageLimitAutoContinueCoordinator>.Instance));
var maintenance = new WorktreeMaintenanceService(dbFactory, git, NullLogger<WorktreeMaintenanceService>.Instance);
var merge = new TaskMergeService(dbFactory, git, broadcaster, state, new VerifyCommandRunner(), NullLogger<TaskMergeService>.Instance);
var aggregator = new PlanningAggregator(dbFactory, git, NullLogger<PlanningAggregator>.Instance);
+3 -1
View File
@@ -111,9 +111,11 @@ public sealed class BatchMcpToolsTests : IDisposable
NullLogger<TaskRunner>.Instance, state, new TaskRunTokenRegistry(), new AttachmentStore(), new FakeSessionSkillSeeder(), new FakeTranscriptUsageReader());
var runCancels = new RunCancellationRegistry(NullLogger<RunCancellationRegistry>.Instance);
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, runCancels);
var usageState = new UsageState();
return new QueueService(dbFactory, runner, cfg, NullLogger<QueueService>.Instance,
new QueueWaker(), new QueuePicker(dbFactory), overrideSlot, state, runCancels,
new FakeUsageGate(), new UsageState(), broadcaster);
new FakeUsageGate(), usageState, broadcaster,
new UsageLimitAutoContinueCoordinator(dbFactory, usageState, overrideSlot, broadcaster, NullLogger<UsageLimitAutoContinueCoordinator>.Instance));
}
[Fact]
@@ -178,8 +178,10 @@ public sealed class ExternalMcpServiceTests : IDisposable
var picker = new ClaudeDo.Worker.Queue.QueuePicker(dbFactory);
var runCancels = new RunCancellationRegistry(NullLogger<RunCancellationRegistry>.Instance);
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, runCancels);
var usageState = new UsageState();
return new QueueService(dbFactory, runner, cfg, NullLogger<QueueService>.Instance, waker, picker, overrideSlot, state, runCancels,
new FakeUsageGate(), new UsageState(), broadcaster);
new FakeUsageGate(), usageState, broadcaster,
new UsageLimitAutoContinueCoordinator(dbFactory, usageState, overrideSlot, broadcaster, NullLogger<UsageLimitAutoContinueCoordinator>.Instance));
}
[Fact]
@@ -83,9 +83,11 @@ public sealed class LifecycleMcpToolsTests : IDisposable
NullLogger<TaskRunner>.Instance, state, new TaskRunTokenRegistry(), new AttachmentStore(), new FakeSessionSkillSeeder(), new FakeTranscriptUsageReader());
var runCancels = new RunCancellationRegistry(NullLogger<RunCancellationRegistry>.Instance);
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, runCancels);
var usageState = new UsageState();
return new QueueService(dbFactory, runner, cfg, NullLogger<QueueService>.Instance,
new QueueWaker(), new QueuePicker(dbFactory), overrideSlot, state, runCancels,
new FakeUsageGate(), new UsageState(), broadcaster);
new FakeUsageGate(), usageState, broadcaster,
new UsageLimitAutoContinueCoordinator(dbFactory, usageState, overrideSlot, broadcaster, NullLogger<UsageLimitAutoContinueCoordinator>.Instance));
}
private async Task<TaskEntity> SeedTaskAsync(
@@ -60,8 +60,10 @@ public sealed class QueueStateMcpToolsTests : IDisposable
new FakeSessionSkillSeeder(), new FakeTranscriptUsageReader());
var picker = new QueuePicker(dbFactory);
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, built.RunCancels);
var resolvedUsageState = usageState ?? new UsageState();
var queue = new QueueService(dbFactory, runner, _cfg, NullLogger<QueueService>.Instance, new QueueWaker(), picker,
overrideSlot, built.State, built.RunCancels, new FakeUsageGate(), usageState ?? new UsageState(), broadcaster);
overrideSlot, built.State, built.RunCancels, new FakeUsageGate(), resolvedUsageState, broadcaster,
new UsageLimitAutoContinueCoordinator(dbFactory, resolvedUsageState, overrideSlot, broadcaster, NullLogger<UsageLimitAutoContinueCoordinator>.Instance));
return (queue, new QueueStateMcpTools(queue, dbFactory));
}
@@ -0,0 +1,207 @@
using ClaudeDo.Data;
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.Queue;
/// Covers the "Continue on session limit reset" toggle's coordinator in isolation from
/// QueueService's own timer plumbing (see QueueService.ScheduleResetWake for that part).
public sealed class UsageLimitAutoContinueCoordinatorTests : IDisposable
{
private readonly DbFixture _db = new();
private readonly string _tempDir;
private readonly WorkerConfig _cfg;
public UsageLimitAutoContinueCoordinatorTests()
{
_tempDir = Path.Combine(Path.GetTempPath(), $"cd_usagelimit_{Guid.NewGuid():N}");
Directory.CreateDirectory(_tempDir);
_cfg = new WorkerConfig { SandboxRoot = _tempDir, LogRoot = _tempDir };
}
public void Dispose() { _db.Dispose(); try { Directory.Delete(_tempDir, true); } catch { } }
private (UsageLimitAutoContinueCoordinator Coordinator, UsageState UsageState, FakeClaudeProcess Process) BuildCoordinator()
{
var dbFactory = _db.CreateFactory();
var state = TaskStateServiceBuilder.Build(dbFactory).State;
var wt = new WorktreeManager(new ClaudeDo.Data.Git.GitService(), dbFactory, _cfg, NullLogger<WorktreeManager>.Instance);
var fake = new FakeClaudeProcess();
var runner = new TaskRunner(fake, dbFactory, new HubBroadcaster(new CapturingHubContext()), wt,
new ClaudeArgsBuilder(), _cfg, NullLogger<TaskRunner>.Instance, state, new TaskRunTokenRegistry(),
new AttachmentStore(), new FakeSessionSkillSeeder(), new FakeTranscriptUsageReader());
var runCancels = new RunCancellationRegistry(NullLogger<RunCancellationRegistry>.Instance);
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, runCancels);
var usageState = new UsageState();
var coordinator = new UsageLimitAutoContinueCoordinator(
dbFactory, usageState, overrideSlot,
new HubBroadcaster(new CapturingHubContext()), NullLogger<UsageLimitAutoContinueCoordinator>.Instance);
return (coordinator, usageState, fake);
}
private async Task<string> SeedUsageLimitTaskAsync(bool withSessionId = true)
{
var listId = Guid.NewGuid().ToString();
var taskId = Guid.NewGuid().ToString();
using (var ctx = _db.CreateContext())
{
ctx.Lists.Add(new ListEntity { Id = listId, Name = "L", CreatedAt = DateTime.UtcNow });
ctx.Tasks.Add(new TaskEntity
{
Id = taskId, ListId = listId, Title = "T", Status = TaskStatus.Failed,
FailureReason = "usage_limit", CreatedAt = DateTime.UtcNow, FinishedAt = DateTime.UtcNow,
});
if (withSessionId)
{
ctx.TaskRuns.Add(new TaskRunEntity
{
Id = Guid.NewGuid().ToString(), TaskId = taskId, RunNumber = 1, IsRetry = false, Prompt = "p",
LogPath = Path.Combine(_tempDir, "log.ndjson"), StartedAt = DateTime.UtcNow,
FinishedAt = DateTime.UtcNow, SessionId = "sess-1", ExitCode = 1,
});
}
await ctx.SaveChangesAsync();
}
return taskId;
}
private async Task EnableToggleAsync()
{
using var ctx = _db.CreateContext();
var repo = new AppSettingsRepository(ctx);
var settings = await repo.GetAsync();
settings.AutoContinueOnUsageLimit = true;
await repo.UpdateAsync(settings);
}
private async Task<TaskEntity?> PollUntilCalledAsync(string taskId, FakeClaudeProcess fake)
{
var deadline = DateTime.UtcNow.AddSeconds(10);
while (DateTime.UtcNow < deadline && fake.CallCount == 0)
await Task.Delay(25);
await Task.Delay(150); // let the fire-and-forget continuation fully settle
using var ctx = _db.CreateContext();
return await new TaskRepository(ctx).GetByIdAsync(taskId);
}
[Fact]
public async Task Toggle_Off_Never_Schedules_Or_Continues()
{
var (coordinator, usageState, fake) = BuildCoordinator();
var taskId = await SeedUsageLimitTaskAsync();
usageState.ReportSuccess(new UsageSnapshot(
new UsageBucket(85, DateTimeOffset.UtcNow.AddSeconds(-1)), null, Array.Empty<UsageLimitRow>(), DateTime.UtcNow));
var wakeAt = await coordinator.GetScheduledWakeAtAsync(default);
await coordinator.RunAsync(default);
await Task.Delay(150);
Assert.Null(wakeAt);
Assert.Equal(0, fake.CallCount);
using var verify = _db.CreateContext();
var task = await new TaskRepository(verify).GetByIdAsync(taskId);
Assert.Null(task!.UsageLimitAutoContinuedAt);
}
[Fact]
public async Task Toggle_On_No_Snapshot_Yet_Schedules_Nothing_And_Does_Not_Crash()
{
var (coordinator, _, fake) = BuildCoordinator();
await SeedUsageLimitTaskAsync();
await EnableToggleAsync();
var wakeAt = await coordinator.GetScheduledWakeAtAsync(default);
await coordinator.RunAsync(default); // must not throw despite no usage snapshot
Assert.Null(wakeAt);
Assert.Equal(0, fake.CallCount);
}
[Fact]
public async Task Toggle_On_Reset_In_Future_Schedules_Wake_But_Does_Not_Continue_Yet()
{
var (coordinator, usageState, fake) = BuildCoordinator();
await SeedUsageLimitTaskAsync();
await EnableToggleAsync();
var resetsAt = DateTimeOffset.UtcNow.AddMinutes(5);
usageState.ReportSuccess(new UsageSnapshot(
new UsageBucket(85, resetsAt), null, Array.Empty<UsageLimitRow>(), DateTime.UtcNow));
var wakeAt = await coordinator.GetScheduledWakeAtAsync(default);
await coordinator.RunAsync(default);
await Task.Delay(150);
Assert.Equal(resetsAt, wakeAt);
Assert.Equal(0, fake.CallCount);
}
[Fact]
public async Task Toggle_On_Reset_Already_Passed_Continues_Exactly_Once()
{
var (coordinator, usageState, fake) = BuildCoordinator();
var taskId = await SeedUsageLimitTaskAsync();
await EnableToggleAsync();
usageState.ReportSuccess(new UsageSnapshot(
new UsageBucket(85, DateTimeOffset.UtcNow.AddSeconds(-1)), null, Array.Empty<UsageLimitRow>(), DateTime.UtcNow));
await coordinator.RunAsync(default);
var task = await PollUntilCalledAsync(taskId, fake);
Assert.Equal(1, fake.CallCount);
Assert.NotNull(task!.UsageLimitAutoContinuedAt);
// A second tick (e.g. the next backstop) must not fire it again.
await coordinator.RunAsync(default);
await Task.Delay(150);
Assert.Equal(1, fake.CallCount);
}
[Fact]
public async Task Missing_Reset_Time_Never_Continues_And_Never_Throws()
{
var (coordinator, usageState, fake) = BuildCoordinator();
await SeedUsageLimitTaskAsync();
await EnableToggleAsync();
// A snapshot exists, but this bucket has no reset time — fail open.
usageState.ReportSuccess(new UsageSnapshot(
new UsageBucket(85, null), null, Array.Empty<UsageLimitRow>(), DateTime.UtcNow));
var wakeAt = await coordinator.GetScheduledWakeAtAsync(default);
await coordinator.RunAsync(default);
await Task.Delay(150);
Assert.Null(wakeAt);
Assert.Equal(0, fake.CallCount);
}
[Fact]
public async Task Non_UsageLimit_Failure_Is_Never_Touched()
{
var (coordinator, usageState, fake) = BuildCoordinator();
var listId = Guid.NewGuid().ToString();
var taskId = Guid.NewGuid().ToString();
using (var ctx = _db.CreateContext())
{
ctx.Lists.Add(new ListEntity { Id = listId, Name = "L", CreatedAt = DateTime.UtcNow });
ctx.Tasks.Add(new TaskEntity { Id = taskId, ListId = listId, Title = "T", Status = TaskStatus.Failed,
FailureReason = "error", CreatedAt = DateTime.UtcNow });
await ctx.SaveChangesAsync();
}
await EnableToggleAsync();
usageState.ReportSuccess(new UsageSnapshot(
new UsageBucket(85, DateTimeOffset.UtcNow.AddSeconds(-1)), null, Array.Empty<UsageLimitRow>(), DateTime.UtcNow));
await coordinator.RunAsync(default);
await Task.Delay(150);
Assert.Equal(0, fake.CallCount);
}
}
@@ -99,6 +99,29 @@ public class FailureDiagnosisTests
{
Assert.Equal(expected, TaskRunner.ClassifyFailureReason(terminalReason));
}
[Theory]
[InlineData("You've hit your session limit resets 1pm (Europe/Berlin)")]
[InlineData("Claude usage limit reached, resets at 3pm")]
[InlineData("rate limited, please retry later")]
public void ClassifyFailureReason_ApiError_With_Limit_Text_Returns_UsageLimit(string messageText)
{
Assert.Equal("usage_limit", TaskRunner.ClassifyFailureReason("api_error", messageText));
}
[Fact]
public void ClassifyFailureReason_ApiError_With_Generic_Text_Stays_Error()
{
Assert.Equal("error", TaskRunner.ClassifyFailureReason("api_error", "Internal server error, please retry"));
}
[Theory]
[InlineData("max_turns")]
[InlineData("timeout")]
public void ClassifyFailureReason_Ignores_Limit_Text_For_NonApiError_Reasons(string terminalReason)
{
Assert.Equal(terminalReason, TaskRunner.ClassifyFailureReason(terminalReason, "usage limit reached"));
}
}
public sealed class FailureDiagnosisEndToEndTests : IDisposable
@@ -282,4 +305,32 @@ public sealed class FailureDiagnosisEndToEndTests : IDisposable
Assert.Equal("error", task.FailureReason);
Assert.NotEqual("max_turns", task.FailureReason);
}
[Fact]
public async Task Usage_Limit_Run_Sets_FailureReason_UsageLimit()
{
var dbFactory = _db.CreateFactory();
using (var ctx = _db.CreateContext())
{
ctx.Lists.Add(new ListEntity { Id = "l1", Name = "L", WorkingDir = null, CreatedAt = DateTime.UtcNow });
ctx.Tasks.Add(new TaskEntity { Id = "t1", ListId = "l1", Title = "T",
Status = TaskStatus.Running, CreatedAt = DateTime.UtcNow });
await ctx.SaveChangesAsync();
}
var fake = new FakeClaudeProcess((_, _, _, _, _) => Task.FromResult(new RunResult
{
ExitCode = 1,
TerminalReason = "api_error",
ResultMarkdown = "You've hit your session limit resets 1pm (Europe/Berlin)",
}));
var runner = MakeRunner(dbFactory, fake);
using (var ctx = _db.CreateContext())
await runner.RunAsync((await new TaskRepository(ctx).GetByIdAsync("t1"))!, "slot-1", default, alreadyClaimed: true);
using var verify = _db.CreateContext();
var task = await new TaskRepository(verify).GetByIdAsync("t1");
Assert.Equal(TaskStatus.Failed, task!.Status);
Assert.Equal("usage_limit", task.FailureReason);
}
}
@@ -65,8 +65,10 @@ public sealed class QueueServiceSlotFailureTests : IDisposable
new FakeSessionSkillSeeder(), new FakeTranscriptUsageReader());
var waker = new QueueWaker();
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, built.RunCancels);
var usageState = new UsageState();
var service = new QueueService(dbFactory, runner, _cfg, NullLogger<QueueService>.Instance, waker, picker,
overrideSlot, built.State, built.RunCancels, new FakeUsageGate(), new UsageState(), broadcaster);
overrideSlot, built.State, built.RunCancels, new FakeUsageGate(), usageState, broadcaster,
new UsageLimitAutoContinueCoordinator(dbFactory, usageState, overrideSlot, broadcaster, NullLogger<UsageLimitAutoContinueCoordinator>.Instance));
return (service, built.Hub, waker);
}
@@ -60,8 +60,10 @@ 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 usageState = new UsageState();
var service = new QueueService(dbFactory, runner, _cfg, NullLogger<QueueService>.Instance, _waker, picker, overrideSlot, state, built.RunCancels,
new FakeUsageGate(), new UsageState(), broadcaster);
new FakeUsageGate(), usageState, broadcaster,
new UsageLimitAutoContinueCoordinator(dbFactory, usageState, overrideSlot, broadcaster, NullLogger<UsageLimitAutoContinueCoordinator>.Instance));
return (service, fake);
}
@@ -69,8 +69,10 @@ public sealed class QueueServiceTests : IDisposable
var overrideSlot = new OverrideSlotService(dbFactory, runner, NullLogger<OverrideSlotService>.Instance, built.RunCancels);
_usageGate = usageGate ?? new FakeUsageGate();
_runCancels = built.RunCancels;
var resolvedUsageState = usageState ?? new UsageState();
var service = new QueueService(dbFactory, runner, _cfg, NullLogger<QueueService>.Instance, _waker, picker, overrideSlot, state, built.RunCancels,
_usageGate, usageState ?? new UsageState(), broadcaster);
_usageGate, resolvedUsageState, broadcaster,
new UsageLimitAutoContinueCoordinator(dbFactory, resolvedUsageState, overrideSlot, broadcaster, NullLogger<UsageLimitAutoContinueCoordinator>.Instance));
return (service, fake);
}