From d4cd202460204bc698dd3ecae154e944e4b0496a Mon Sep 17 00:00:00 2001 From: mika kuns Date: Wed, 5 Aug 2026 12:37:50 +0200 Subject: [PATCH] feat(worker): expose usage/model-usage hub surface and persist run model Adds GetUsageSnapshot/GetModelUsage/GetTaskUsage to WorkerHub (backed by a shared UsageSnapshotBuilder), a UsageUpdated broadcast fired after every UsageMonitorService poll cycle, and records the resolved model on each task_runs row so per-model/per-task usage can be reported from history. --- src/ClaudeDo.Worker/CLAUDE.md | 5 +- src/ClaudeDo.Worker/Hub/HubBroadcaster.cs | 3 + src/ClaudeDo.Worker/Hub/WorkerHub.cs | 114 ++++++++++++- src/ClaudeDo.Worker/Program.cs | 3 +- src/ClaudeDo.Worker/Runner/TaskRunner.cs | 1 + .../Usage/UsageMonitorService.cs | 13 +- .../Usage/UsageSnapshotBuilder.cs | 63 +++++++ .../Hub/TaskUsageHubTests.cs | 161 ++++++++++++++++++ .../Runner/RunModelPersistenceTests.cs | 105 ++++++++++++ .../Usage/UsageMonitorServiceTests.cs | 49 +++++- .../Usage/UsageSnapshotBuilderTests.cs | 145 ++++++++++++++++ 11 files changed, 650 insertions(+), 12 deletions(-) create mode 100644 src/ClaudeDo.Worker/Usage/UsageSnapshotBuilder.cs create mode 100644 tests/ClaudeDo.Worker.Tests/Hub/TaskUsageHubTests.cs create mode 100644 tests/ClaudeDo.Worker.Tests/Runner/RunModelPersistenceTests.cs create mode 100644 tests/ClaudeDo.Worker.Tests/Usage/UsageSnapshotBuilderTests.cs diff --git a/src/ClaudeDo.Worker/CLAUDE.md b/src/ClaudeDo.Worker/CLAUDE.md index cdcc88ea..ddbe7d4d 100644 --- a/src/ClaudeDo.Worker/CLAUDE.md +++ b/src/ClaudeDo.Worker/CLAUDE.md @@ -21,7 +21,7 @@ Worker/ Report/ — ClaudeHistoryReader, WeekReportPromptBuilder, WeekReportService; interfaces in Report/Interfaces/ Prime/ — daily-prep ("Prime Claude"): PrimeScheduler (BackgroundService), PrimeRunner (runs the daily prep), DailyPrepPrompt (fixed prompt + CLI args + LogPath() helper), NextDueCalculator, PrimeScheduleSignal; interfaces in Prime/Interfaces/ (IPrimeRunner, IPrimeClock, IPrimeScheduleSignal, IPrimeBroadcaster) Online/ — optional Online Inbox sync: OnlineInboxConfig (config record), Dtos (RemoteList/RemoteTask/MirrorTask), IOnlineInboxApi, OnlineInboxApiClient (typed HttpClient, bearer auth, HTTPS guard), OnlineTokenStore (DPAPI refresh-token store, Windows-only), StaticTokenAuthProvider (default/test IOnlineAuthProvider), ZitadelAuthProvider (OIDC discovery + refresh-token flow), OnlineSyncService (BackgroundService: reconcile loop), OnlineBacklog (Idle-backlog filter/query); interface in Online/Interfaces/ (IOnlineAuthProvider) - Usage/ — OAuth usage monitor: UsageModels (UsageBucket/UsageLimitRow/UsageSnapshot), ClaudeOAuthUsageClient (reads the access token Claude Code keeps fresh at `~/.claude/.credentials.json`, calls `GET https://api.anthropic.com/api/oauth/usage`; defensive parsing — missing/null buckets → null, missing `limits` → empty list; never logs the token), UsageState (threadsafe singleton; a failed poll never overwrites the last good snapshot, only sets `LastError`), UsageMonitorService (BackgroundService, polls on `usage_poll_interval_seconds`, one poll at startup, logs a failure at most once per distinct error message), TranscriptUsageReader (aggregates Claude Code transcript token usage from `~/.claude/projects/**/*.jsonl` by date/model/scope (ClaudeDo vs Other), deduped by requestId, with a per-file length+mtime cache), UsageGate (reads `UsageState` + `AppSettings.UsageGateFiveHourPct`/`UsageGateSevenDayPct`, returns a `UsageGateDecision(IsBlocked, Reason)`; `Utilization` from `UsageBucket` is already a 0–100 percent, compared directly against the threshold with `>=`; threshold `0` = that bucket never gates; fail-open — no snapshot yet, a failed last poll, or a settings-read error all resolve to not-blocked); interfaces in Usage/Interfaces/ (IUsageClient, ITranscriptUsageReader, IUsageGate) + Usage/ — OAuth usage monitor: UsageModels (UsageBucket/UsageLimitRow/UsageSnapshot), ClaudeOAuthUsageClient (reads the access token Claude Code keeps fresh at `~/.claude/.credentials.json`, calls `GET https://api.anthropic.com/api/oauth/usage`; defensive parsing — missing/null buckets → null, missing `limits` → empty list; never logs the token), UsageState (threadsafe singleton; a failed poll never overwrites the last good snapshot, only sets `LastError`), UsageMonitorService (BackgroundService, polls on `usage_poll_interval_seconds`, one poll at startup, logs a failure at most once per distinct error message, broadcasts `HubBroadcaster.UsageUpdated` after every tick via `UsageSnapshotBuilder`), UsageSnapshotBuilder (builds the Hub-facing `UsageSnapshotDto` from `UsageState` + `IUsageGate` + `AppSettings` thresholds — the one place `WorkerHub.GetUsageSnapshot` and `UsageMonitorService` share the stale/threshold/gate logic), TranscriptUsageReader (aggregates Claude Code transcript token usage from `~/.claude/projects/**/*.jsonl` by date/model/scope (ClaudeDo vs Other), deduped by requestId, with a per-file length+mtime cache), UsageGate (reads `UsageState` + `AppSettings.UsageGateFiveHourPct`/`UsageGateSevenDayPct`, returns a `UsageGateDecision(IsBlocked, Reason)`; `Utilization` from `UsageBucket` is already a 0–100 percent, compared directly against the threshold with `>=`; threshold `0` = that bucket never gates; fail-open — no snapshot yet, a failed last poll, or a settings-read error all resolve to not-blocked); interfaces in Usage/Interfaces/ (IUsageClient, ITranscriptUsageReader, IUsageGate) ``` Interfaces (e.g. `IQueueWaker`, `IPrimeClock`, `ITaskStateService`) live in an `Interfaces/` subfolder within their area; the namespace stays the area namespace. @@ -181,8 +181,9 @@ Each CLI invocation is recorded in the `task_runs` table via `TaskRunRepository` - Agents/settings/lists: `GetAgents`, `RefreshAgents`, `RestoreDefaultAgents`, `GetAppSettings`, `UpdateAppSettings`, `UpdateList`, `UpdateListConfig`, `GetListConfig`, `UpdateTaskAgentSettings` - Reports/notes/prep: `GetWeekReport`, `GenerateWeekReport`, `GetDailyNotes`, `AddDailyNote`, `UpdateDailyNote`, `DeleteDailyNote`, `RunDailyPrepNow`, `ClearMyDay`, `GetLastPrepLog`, `ListPrimeSchedules`, `UpsertPrimeSchedule`, `DeletePrimeSchedule` - Diagnostics: `GetRecentLogs` (last 30 min of buffered log records, all levels, for the Log Visualizer overlay) +- Usage: `GetUsageSnapshot() -> UsageSnapshotDto` (built by `UsageSnapshotBuilder` from `UsageState` + `IUsageGate` + `AppSettings` gate thresholds; percentages/limits/`FetchedAtUtc` null and `IsStale=true` when no snapshot has landed yet; `IsStale` also trips on a failed last poll or a snapshot older than 3× `usage_poll_interval_seconds`), `GetModelUsage(from, to) -> IReadOnlyList` (thin wrapper over `ITranscriptUsageReader.ReadAsync`), `GetTaskUsage(from, to) -> IReadOnlyList` (top consumers from `task_runs` joined to task/list, grouped per task — `Runs`/summed `TokensIn`/`TokensOut` (null token columns count as 0, never dropped), `Model` from that task's most recent run — sorted by total tokens descending, capped at 100) -**HubBroadcaster** events: `TaskStarted`, `TaskFinished`, `TaskMessage`, `WorktreeUpdated`, `TaskUpdated`, `RunCreated`, `ListUpdated`, `WorkerLog`, `PrimeFired`, `PrepStarted`, `PrepLine`, `PrepFinished`, `PlanningMergeStarted`, `PlanningSubtaskMerged`, `PlanningMergeConflict`, `PlanningMergeAborted`, `PlanningCompleted`, `RefineStarted`, `RefineFinished` +**HubBroadcaster** events: `TaskStarted`, `TaskFinished`, `TaskMessage`, `WorktreeUpdated`, `TaskUpdated`, `RunCreated`, `ListUpdated`, `WorkerLog`, `PrimeFired`, `PrepStarted`, `PrepLine`, `PrepFinished`, `PlanningMergeStarted`, `PlanningSubtaskMerged`, `PlanningMergeConflict`, `PlanningMergeAborted`, `PlanningCompleted`, `RefineStarted`, `RefineFinished`, `UsageUpdated` (carries the same `UsageSnapshotDto` as `GetUsageSnapshot`; `UsageMonitorService` fires it after every poll cycle, success or failure, via the shared `UsageSnapshotBuilder`) `WorkerLog` carries two sources: the hand-curated business events (`_broadcaster.WorkerLog(...)` in TaskRunner/TaskMergeService/TaskResetService) **and** every Serilog **Warn/Error** event, re-broadcast by `BroadcastLogSink` (deduped within a 120 s per-message window; SignalR plumbing source-contexts filtered to avoid feedback loops). The sink also buffers **all** levels into `LogRingBuffer` for `GetRecentLogs`. diff --git a/src/ClaudeDo.Worker/Hub/HubBroadcaster.cs b/src/ClaudeDo.Worker/Hub/HubBroadcaster.cs index 0bbb8f88..49464696 100644 --- a/src/ClaudeDo.Worker/Hub/HubBroadcaster.cs +++ b/src/ClaudeDo.Worker/Hub/HubBroadcaster.cs @@ -38,6 +38,9 @@ public sealed class HubBroadcaster : IPrimeBroadcaster, IRefineBroadcaster public Task RunCreated(string taskId, int runNumber, bool isRetry) => _hub.Clients.All.SendAsync("RunCreated", taskId, runNumber, isRetry); + public Task UsageUpdated(UsageSnapshotDto snapshot) => + _hub.Clients.All.SendAsync("UsageUpdated", snapshot); + public Task WorkerLog(string message, WorkerLogLevel level, DateTime timestampUtc) => _hub.Clients.All.SendAsync("WorkerLog", message, level, timestampUtc); diff --git a/src/ClaudeDo.Worker/Hub/WorkerHub.cs b/src/ClaudeDo.Worker/Hub/WorkerHub.cs index 4010f656..9c6bf219 100644 --- a/src/ClaudeDo.Worker/Hub/WorkerHub.cs +++ b/src/ClaudeDo.Worker/Hub/WorkerHub.cs @@ -17,6 +17,8 @@ using ClaudeDo.Worker.Report; using ClaudeDo.Worker.Report.Interfaces; using ClaudeDo.Worker.Skills; using ClaudeDo.Worker.State; +using ClaudeDo.Worker.Usage; +using ClaudeDo.Worker.Usage.Interfaces; using ClaudeDo.Worker.Worktrees; using System.Text.Json; using TaskStatus = ClaudeDo.Data.Models.TaskStatus; @@ -103,6 +105,49 @@ public record OnlineInboxConfigInput( string Scopes, string RedirectUri); +public record UsageLimitDto( + string Kind, + string Group, + double Percent, + string Severity, + DateTimeOffset? ResetsAt, + string? ScopeModelDisplayName, + bool IsActive); + +public record UsageSnapshotDto( + double? FiveHourPercent, + DateTimeOffset? FiveHourResetsAt, + double? SevenDayPercent, + DateTimeOffset? SevenDayResetsAt, + IReadOnlyList Limits, + int FiveHourThresholdPct, + int SevenDayThresholdPct, + bool IsGateBlocked, + string? GateReason, + DateTime? FetchedAtUtc, + bool IsStale, + string? LastError); + +public record ModelUsageRowDto( + DateOnly Date, + string Model, + string Scope, + long InputTokens, + long OutputTokens, + long CacheReadTokens, + long CacheCreationTokens, + int Messages); + +public record TaskUsageRowDto( + string TaskId, + string TaskTitle, + string ListId, + string ListName, + string? Model, + int Runs, + long TokensIn, + long TokensOut); + public sealed class WorkerHub : Microsoft.AspNetCore.SignalR.Hub { private static readonly string Version = @@ -136,6 +181,8 @@ public sealed class WorkerHub : Microsoft.AspNetCore.SignalR.Hub private readonly IInteractiveLaunchSpecService? _interactiveLaunchSpec; private readonly WorktreeManager? _worktreeManager; private readonly Data.Git.GitService? _git; + private readonly UsageSnapshotBuilder? _usageSnapshotBuilder; + private readonly ITranscriptUsageReader? _usageReader; public WorkerHub( QueueService queue, @@ -165,7 +212,9 @@ public sealed class WorkerHub : Microsoft.AspNetCore.SignalR.Hub LogRingBuffer? logBuffer = null, IInteractiveLaunchSpecService? interactiveLaunchSpec = null, WorktreeManager? worktreeManager = null, - Data.Git.GitService? git = null) + Data.Git.GitService? git = null, + UsageSnapshotBuilder? usageSnapshotBuilder = null, + ITranscriptUsageReader? usageReader = null) { _queue = queue; _waker = waker; @@ -195,6 +244,8 @@ public sealed class WorkerHub : Microsoft.AspNetCore.SignalR.Hub _interactiveLaunchSpec = interactiveLaunchSpec; _worktreeManager = worktreeManager; _git = git; + _usageSnapshotBuilder = usageSnapshotBuilder; + _usageReader = usageReader; } // Persistence boundary for the session_skills JSON-array columns (task/list/global). @@ -978,4 +1029,65 @@ public sealed class WorkerHub : Microsoft.AspNetCore.SignalR.Hub _onlineTokenStore.Clear(); } #pragma warning restore CA1416 + + public Task GetUsageSnapshot() => HubGuard(() => + { + if (_usageSnapshotBuilder is null) + throw new InvalidOperationException("Usage snapshot builder is not configured."); + return _usageSnapshotBuilder.BuildAsync(Context.ConnectionAborted); + }); + + public Task> GetModelUsage(DateOnly from, DateOnly to) => HubGuard(async () => + { + if (_usageReader is null) + throw new InvalidOperationException("Transcript usage reader is not configured."); + var rows = await _usageReader.ReadAsync(from, to, Context.ConnectionAborted); + return (IReadOnlyList)rows + .Select(r => new ModelUsageRowDto( + r.Date, r.Model, r.Scope == UsageScope.ClaudeDo ? "claudedo" : "other", + r.InputTokens, r.OutputTokens, r.CacheReadTokens, r.CacheCreationTokens, r.Messages)) + .ToList(); + }); + + public async Task> GetTaskUsage(DateOnly from, DateOnly to) + { + var fromDt = from.ToDateTime(TimeOnly.MinValue); + var toDt = to.ToDateTime(TimeOnly.MaxValue); + + await using var ctx = await _dbFactory.CreateDbContextAsync(Context.ConnectionAborted); + var runs = await ctx.TaskRuns + .Where(r => r.StartedAt != null && r.StartedAt >= fromDt && r.StartedAt <= toDt) + .ToListAsync(Context.ConnectionAborted); + + if (runs.Count == 0) return Array.Empty(); + + var taskIds = runs.Select(r => r.TaskId).Distinct().ToList(); + var tasks = await ctx.Tasks.Where(t => taskIds.Contains(t.Id)).ToListAsync(Context.ConnectionAborted); + var taskById = tasks.ToDictionary(t => t.Id); + + var listIds = tasks.Select(t => t.ListId).Distinct().ToList(); + var lists = await ctx.Lists.Where(l => listIds.Contains(l.Id)).ToDictionaryAsync(l => l.Id, Context.ConnectionAborted); + + return runs + .Where(r => taskById.ContainsKey(r.TaskId)) + .GroupBy(r => r.TaskId) + .Select(g => + { + var task = taskById[g.Key]; + var listName = lists.TryGetValue(task.ListId, out var list) ? list.Name : ""; + var latestModel = g.OrderByDescending(r => r.StartedAt).First().Model; + return new TaskUsageRowDto( + task.Id, + task.Title, + task.ListId, + listName, + latestModel, + g.Count(), + g.Sum(r => (long)(r.TokensIn ?? 0)), + g.Sum(r => (long)(r.TokensOut ?? 0))); + }) + .OrderByDescending(r => r.TokensIn + r.TokensOut) + .Take(100) + .ToList(); + } } diff --git a/src/ClaudeDo.Worker/Program.cs b/src/ClaudeDo.Worker/Program.cs index 9b74097a..c18b3b70 100644 --- a/src/ClaudeDo.Worker/Program.cs +++ b/src/ClaudeDo.Worker/Program.cs @@ -206,8 +206,9 @@ builder.Services.AddHttpClient(client => { client.Timeout = TimeSpan.FromSeconds(5); }); -builder.Services.AddHostedService(); builder.Services.AddSingleton(); +builder.Services.AddSingleton(); +builder.Services.AddHostedService(); // Loopback-only bind. Firewall is irrelevant for 127.0.0.1. builder.WebHost.UseUrls($"http://127.0.0.1:{cfg.SignalRPort}"); diff --git a/src/ClaudeDo.Worker/Runner/TaskRunner.cs b/src/ClaudeDo.Worker/Runner/TaskRunner.cs index 9f32a51d..820b2457 100644 --- a/src/ClaudeDo.Worker/Runner/TaskRunner.cs +++ b/src/ClaudeDo.Worker/Runner/TaskRunner.cs @@ -318,6 +318,7 @@ public sealed class TaskRunner Prompt = prompt, LogPath = logPath, StartedAt = DateTime.UtcNow, + Model = config.Model, }; using (var context = _dbFactory.CreateDbContext()) diff --git a/src/ClaudeDo.Worker/Usage/UsageMonitorService.cs b/src/ClaudeDo.Worker/Usage/UsageMonitorService.cs index 43fe314b..8d013348 100644 --- a/src/ClaudeDo.Worker/Usage/UsageMonitorService.cs +++ b/src/ClaudeDo.Worker/Usage/UsageMonitorService.cs @@ -1,4 +1,5 @@ using ClaudeDo.Worker.Config; +using ClaudeDo.Worker.Hub; using ClaudeDo.Worker.Usage.Interfaces; namespace ClaudeDo.Worker.Usage; @@ -7,6 +8,8 @@ namespace ClaudeDo.Worker.Usage; /// Polls on and keeps /// current. Polls once immediately at startup. A failure is logged as a /// warning at most once per distinct error message, to avoid log spam on a persistent outage. +/// Broadcasts after every poll cycle, success or failure, +/// so the UI can reflect a stale/blocked state as soon as it happens. /// public sealed class UsageMonitorService : BackgroundService { @@ -14,15 +17,20 @@ public sealed class UsageMonitorService : BackgroundService private readonly UsageState _state; private readonly WorkerConfig _config; private readonly ILogger _logger; + private readonly UsageSnapshotBuilder _snapshotBuilder; + private readonly HubBroadcaster _broadcaster; private string? _lastLoggedError; public UsageMonitorService( - IUsageClient client, UsageState state, WorkerConfig config, ILogger logger) + IUsageClient client, UsageState state, WorkerConfig config, ILogger logger, + UsageSnapshotBuilder snapshotBuilder, HubBroadcaster broadcaster) { _client = client; _state = state; _config = config; _logger = logger; + _snapshotBuilder = snapshotBuilder; + _broadcaster = broadcaster; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) @@ -64,5 +72,8 @@ public sealed class UsageMonitorService : BackgroundService _lastLoggedError = ex.Message; } } + + var dto = await _snapshotBuilder.BuildAsync(ct); + await _broadcaster.UsageUpdated(dto); } } diff --git a/src/ClaudeDo.Worker/Usage/UsageSnapshotBuilder.cs b/src/ClaudeDo.Worker/Usage/UsageSnapshotBuilder.cs new file mode 100644 index 00000000..efd6535a --- /dev/null +++ b/src/ClaudeDo.Worker/Usage/UsageSnapshotBuilder.cs @@ -0,0 +1,63 @@ +using ClaudeDo.Data; +using ClaudeDo.Data.Repositories; +using ClaudeDo.Worker.Config; +using ClaudeDo.Worker.Hub; +using ClaudeDo.Worker.Usage.Interfaces; +using Microsoft.EntityFrameworkCore; + +namespace ClaudeDo.Worker.Usage; + +/// +/// Builds the Hub-facing usage snapshot DTO from UsageState + IUsageGate + AppSettings +/// thresholds. Shared by WorkerHub.GetUsageSnapshot and UsageMonitorService's post-poll +/// broadcast so both surfaces agree on staleness/threshold logic. +/// +public sealed class UsageSnapshotBuilder +{ + private readonly UsageState _state; + private readonly IUsageGate _gate; + private readonly IDbContextFactory _dbFactory; + private readonly WorkerConfig _cfg; + + public UsageSnapshotBuilder( + UsageState state, IUsageGate gate, IDbContextFactory dbFactory, WorkerConfig cfg) + { + _state = state; + _gate = gate; + _dbFactory = dbFactory; + _cfg = cfg; + } + + public async Task BuildAsync(CancellationToken ct = default) + { + var snapshot = _state.Snapshot; + var lastError = _state.LastError; + + Data.Models.AppSettingsEntity settings; + using (var context = _dbFactory.CreateDbContext()) + settings = await new AppSettingsRepository(context).GetAsync(ct); + + var decision = await _gate.EvaluateAsync(ct); + + var maxAge = TimeSpan.FromSeconds(_cfg.UsagePollIntervalSeconds * 3); + var isStale = snapshot is null || lastError is not null || (DateTime.UtcNow - snapshot.FetchedAtUtc) > maxAge; + + var limits = (snapshot?.Limits ?? Array.Empty()) + .Select(l => new UsageLimitDto(l.Kind, l.Group, l.Percent, l.Severity, l.ResetsAt, l.ScopeModelDisplayName, l.IsActive)) + .ToList(); + + return new UsageSnapshotDto( + snapshot?.FiveHour?.Utilization, + snapshot?.FiveHour?.ResetsAt, + snapshot?.SevenDay?.Utilization, + snapshot?.SevenDay?.ResetsAt, + limits, + settings.UsageGateFiveHourPct, + settings.UsageGateSevenDayPct, + decision.IsBlocked, + decision.Reason, + snapshot?.FetchedAtUtc, + isStale, + lastError); + } +} diff --git a/tests/ClaudeDo.Worker.Tests/Hub/TaskUsageHubTests.cs b/tests/ClaudeDo.Worker.Tests/Hub/TaskUsageHubTests.cs new file mode 100644 index 00000000..2fc4d7a7 --- /dev/null +++ b/tests/ClaudeDo.Worker.Tests/Hub/TaskUsageHubTests.cs @@ -0,0 +1,161 @@ +using ClaudeDo.Data.Models; +using ClaudeDo.Data.Repositories; +using ClaudeDo.Worker.Hub; +using ClaudeDo.Worker.Tests.Infrastructure; +using Xunit; +using TaskStatus = ClaudeDo.Data.Models.TaskStatus; + +namespace ClaudeDo.Worker.Tests.Hub; + +public sealed class TaskUsageHubTests : IDisposable +{ + private readonly DbFixture _db = new(); + + public void Dispose() => _db.Dispose(); + + private WorkerHub CreateHub() + { + var broadcaster = new HubBroadcaster(new CapturingHubContext()); + var hub = new WorkerHub( + null!, null!, null!, null!, broadcaster, _db.CreateFactory(), + null!, null!, null!, null!, null!, null!, null!, null!, null!, null!, null!, null!, null!, + null!, new ClaudeDo.Worker.Online.OnlineInboxConfig(), new ClaudeDo.Worker.Online.OnlineTokenStore(), + new ClaudeDo.Worker.Runner.PendingQuestionRegistry(), null!); + hub.Clients = new FakeHubCallerClients(new RecordingClientProxy()); + hub.Context = new FakeHubCallerContext(); + return hub; + } + + private async Task SeedListAsync(string name = "L") + { + using var ctx = _db.CreateContext(); + var listId = Guid.NewGuid().ToString(); + await new ListRepository(ctx).AddAsync(new ListEntity + { + Id = listId, Name = name, CreatedAt = DateTime.UtcNow, + }); + return listId; + } + + private async Task SeedTaskAsync(string listId, string title = "T") + { + using var ctx = _db.CreateContext(); + var taskId = Guid.NewGuid().ToString(); + await new TaskRepository(ctx).AddAsync(new TaskEntity + { + Id = taskId, ListId = listId, Title = title, Status = TaskStatus.Done, + CreatedAt = DateTime.UtcNow, CommitType = "feat", + }); + return taskId; + } + + private async Task SeedRunAsync( + string taskId, DateTime startedAt, int? tokensIn, int? tokensOut, string? model, int runNumber = 1) + { + using var ctx = _db.CreateContext(); + await new TaskRunRepository(ctx).AddAsync(new TaskRunEntity + { + Id = Guid.NewGuid().ToString(), + TaskId = taskId, + RunNumber = runNumber, + IsRetry = false, + Prompt = "p", + StartedAt = startedAt, + TokensIn = tokensIn, + TokensOut = tokensOut, + Model = model, + }); + } + + [Fact] + public async Task Groups_multiple_runs_for_same_task_and_sums_tokens() + { + var listId = await SeedListAsync(); + var taskId = await SeedTaskAsync(listId, "Task A"); + var day = new DateTime(2026, 8, 1); + await SeedRunAsync(taskId, day, 100, 50, "sonnet", 1); + await SeedRunAsync(taskId, day.AddHours(1), 200, 75, "sonnet", 2); + await SeedRunAsync(taskId, day.AddHours(2), 300, 25, "sonnet", 3); + + var hub = CreateHub(); + var rows = await hub.GetTaskUsage(DateOnly.FromDateTime(day), DateOnly.FromDateTime(day)); + + var row = Assert.Single(rows); + Assert.Equal(taskId, row.TaskId); + Assert.Equal(3, row.Runs); + Assert.Equal(600, row.TokensIn); + Assert.Equal(150, row.TokensOut); + } + + [Fact] + public async Task Sorts_descending_by_total_tokens() + { + var listId = await SeedListAsync(); + var smallTask = await SeedTaskAsync(listId, "Small"); + var bigTask = await SeedTaskAsync(listId, "Big"); + var day = new DateTime(2026, 8, 1); + await SeedRunAsync(smallTask, day, 10, 10, "haiku"); + await SeedRunAsync(bigTask, day, 1000, 1000, "opus"); + + var hub = CreateHub(); + var rows = await hub.GetTaskUsage(DateOnly.FromDateTime(day), DateOnly.FromDateTime(day)); + + Assert.Equal(2, rows.Count); + Assert.Equal(bigTask, rows[0].TaskId); + Assert.Equal(smallTask, rows[1].TaskId); + } + + [Fact] + public async Task Respects_date_range() + { + var listId = await SeedListAsync(); + var taskId = await SeedTaskAsync(listId); + var inRange = new DateTime(2026, 8, 5); + var outOfRange = new DateTime(2026, 8, 1); + await SeedRunAsync(taskId, inRange, 100, 100, "sonnet", 1); + await SeedRunAsync(taskId, outOfRange, 999, 999, "sonnet", 2); + + var hub = CreateHub(); + var rows = await hub.GetTaskUsage(DateOnly.FromDateTime(inRange), DateOnly.FromDateTime(inRange)); + + var row = Assert.Single(rows); + Assert.Equal(1, row.Runs); + Assert.Equal(100, row.TokensIn); + } + + [Fact] + public async Task Null_token_values_are_counted_as_zero_not_excluded() + { + var listId = await SeedListAsync(); + var taskId = await SeedTaskAsync(listId); + var day = new DateTime(2026, 8, 1); + await SeedRunAsync(taskId, day, null, null, null); + + var hub = CreateHub(); + var rows = await hub.GetTaskUsage(DateOnly.FromDateTime(day), DateOnly.FromDateTime(day)); + + var row = Assert.Single(rows); + Assert.Equal(1, row.Runs); + Assert.Equal(0, row.TokensIn); + Assert.Equal(0, row.TokensOut); + } + + [Fact] + public async Task Caps_at_100_rows() + { + var listId = await SeedListAsync(); + var day = new DateTime(2026, 8, 1); + for (var i = 0; i < 105; i++) + { + var taskId = await SeedTaskAsync(listId, $"Task {i}"); + await SeedRunAsync(taskId, day, i, 0, "sonnet"); + } + + var hub = CreateHub(); + var rows = await hub.GetTaskUsage(DateOnly.FromDateTime(day), DateOnly.FromDateTime(day)); + + Assert.Equal(100, rows.Count); + // Highest token counts (100..4) survive the cap; the lowest (0..3) are dropped. + Assert.All(rows, r => Assert.True(r.TokensIn >= 5)); + } +} diff --git a/tests/ClaudeDo.Worker.Tests/Runner/RunModelPersistenceTests.cs b/tests/ClaudeDo.Worker.Tests/Runner/RunModelPersistenceTests.cs new file mode 100644 index 00000000..33426fbb --- /dev/null +++ b/tests/ClaudeDo.Worker.Tests/Runner/RunModelPersistenceTests.cs @@ -0,0 +1,105 @@ +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.Runner; +using ClaudeDo.Worker.Tests.Infrastructure; +using Microsoft.Extensions.Logging.Abstractions; +using TaskStatus = ClaudeDo.Data.Models.TaskStatus; +using Xunit; + +namespace ClaudeDo.Worker.Tests.Runner; + +/// Verifies TaskRunner persists the resolved model (task -> list -> AppSettings.DefaultModel) +/// onto the task_runs row it creates, since the hub's usage-by-model report reads it back from there. +public sealed class RunModelPersistenceTests : IDisposable +{ + private readonly DbFixture _db = new(); + private readonly string _tempDir; + private readonly WorkerConfig _cfg; + + public RunModelPersistenceTests() + { + _tempDir = Path.Combine(Path.GetTempPath(), $"cd_runmodel_{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 TaskRunner BuildRunner() + { + var dbFactory = _db.CreateFactory(); + var state = TaskStateServiceBuilder.Build(dbFactory).State; + var wt = new WorktreeManager(new GitService(), dbFactory, _cfg, NullLogger.Instance); + var fake = new FakeClaudeProcess((_, _, _, _, _) => + Task.FromResult(new RunResult { ExitCode = 0, ResultMarkdown = "ok" })); + return new TaskRunner(fake, dbFactory, new HubBroadcaster(new CapturingHubContext()), wt, + new ClaudeArgsBuilder(), _cfg, NullLogger.Instance, state, new TaskRunTokenRegistry(), + new AttachmentStore(), new FakeSessionSkillSeeder()); + } + + private async Task SeedAsync(string? taskModel, string? listModel) + { + using var ctx = _db.CreateContext(); + ctx.Lists.Add(new ListEntity { Id = "l1", Name = "L", WorkingDir = null, CreatedAt = DateTime.UtcNow }); + if (listModel is not null) + ctx.ListConfigs.Add(new ListConfigEntity { ListId = "l1", Model = listModel }); + ctx.Tasks.Add(new TaskEntity + { + Id = "t1", ListId = "l1", Title = "Task", Status = TaskStatus.Idle, CreatedAt = DateTime.UtcNow, + Model = taskModel, + }); + await ctx.SaveChangesAsync(); + } + + private async Task RunAndGetPersistedModelAsync() + { + var runner = BuildRunner(); + using (var ctx = _db.CreateContext()) + await runner.RunAsync((await new TaskRepository(ctx).GetByIdAsync("t1"))!, "slot-1", CancellationToken.None); + + using var readCtx = _db.CreateContext(); + var run = await new TaskRunRepository(readCtx).GetLatestByTaskIdAsync("t1"); + return run?.Model; + } + + [Fact] + public async Task Task_level_model_override_is_persisted_on_the_run() + { + await SeedAsync(taskModel: "opus", listModel: "haiku"); + + var persisted = await RunAndGetPersistedModelAsync(); + + Assert.Equal("opus", persisted); + } + + [Fact] + public async Task List_default_model_is_persisted_when_task_has_no_override() + { + await SeedAsync(taskModel: null, listModel: "haiku"); + + var persisted = await RunAndGetPersistedModelAsync(); + + Assert.Equal("haiku", persisted); + } + + [Fact] + public async Task Global_default_model_is_persisted_when_task_and_list_have_no_override() + { + await SeedAsync(taskModel: null, listModel: null); + + string globalDefault; + using (var ctx = _db.CreateContext()) + { + var settings = await new AppSettingsRepository(ctx).GetAsync(); + globalDefault = settings.DefaultModel; + } + + var persisted = await RunAndGetPersistedModelAsync(); + + Assert.Equal(globalDefault, persisted); + } +} diff --git a/tests/ClaudeDo.Worker.Tests/Usage/UsageMonitorServiceTests.cs b/tests/ClaudeDo.Worker.Tests/Usage/UsageMonitorServiceTests.cs index ebd06965..48fb8a17 100644 --- a/tests/ClaudeDo.Worker.Tests/Usage/UsageMonitorServiceTests.cs +++ b/tests/ClaudeDo.Worker.Tests/Usage/UsageMonitorServiceTests.cs @@ -1,12 +1,18 @@ using ClaudeDo.Worker.Config; +using ClaudeDo.Worker.Hub; +using ClaudeDo.Worker.Tests.Infrastructure; using ClaudeDo.Worker.Usage; using ClaudeDo.Worker.Usage.Interfaces; using Microsoft.Extensions.Logging.Abstractions; namespace ClaudeDo.Worker.Tests.Usage; -public sealed class UsageMonitorServiceTests +public sealed class UsageMonitorServiceTests : IDisposable { + private readonly DbFixture _db = new(); + + public void Dispose() => _db.Dispose(); + private sealed class FakeClient : IUsageClient { public Queue> Results { get; } = new(); @@ -20,15 +26,31 @@ public sealed class UsageMonitorServiceTests } } + private sealed class FakeGate : IUsageGate + { + public Task EvaluateAsync(CancellationToken ct = default) => + Task.FromResult(new UsageGateDecision(false, null)); + } + private static UsageSnapshot MakeSnapshot() => new(new UsageBucket(1, null), null, [], DateTime.UtcNow); + private (UsageMonitorService Service, UsageState State, CapturingHubContext Hub) CreateService(FakeClient client, WorkerConfig? cfg = null) + { + var state = new UsageState(); + var config = cfg ?? new WorkerConfig(); + var builder = new UsageSnapshotBuilder(state, new FakeGate(), _db.CreateFactory(), config); + var hubContext = new CapturingHubContext(); + var broadcaster = new HubBroadcaster(hubContext); + var service = new UsageMonitorService(client, state, config, NullLogger.Instance, builder, broadcaster); + return (service, state, hubContext); + } + [Fact] public async Task TickAsync_Success_UpdatesState() { var client = new FakeClient(); client.Results.Enqueue(MakeSnapshot); - var state = new UsageState(); - var service = new UsageMonitorService(client, state, new WorkerConfig(), NullLogger.Instance); + var (service, state, _) = CreateService(client); await service.TickAsync(CancellationToken.None); @@ -41,8 +63,7 @@ public sealed class UsageMonitorServiceTests { var client = new FakeClient(); client.Results.Enqueue(() => throw new InvalidOperationException("network unreachable")); - var state = new UsageState(); - var service = new UsageMonitorService(client, state, new WorkerConfig(), NullLogger.Instance); + var (service, state, _) = CreateService(client); await service.TickAsync(CancellationToken.None); @@ -56,8 +77,7 @@ public sealed class UsageMonitorServiceTests var client = new FakeClient(); client.Results.Enqueue(MakeSnapshot); client.Results.Enqueue(() => throw new InvalidOperationException("down")); - var state = new UsageState(); - var service = new UsageMonitorService(client, state, new WorkerConfig(), NullLogger.Instance); + var (service, state, _) = CreateService(client); await service.TickAsync(CancellationToken.None); await service.TickAsync(CancellationToken.None); @@ -65,4 +85,19 @@ public sealed class UsageMonitorServiceTests Assert.NotNull(state.Snapshot); Assert.Equal("down", state.LastError); } + + [Fact] + public async Task TickAsync_BroadcastsUsageUpdated_OnSuccessAndFailure() + { + var client = new FakeClient(); + client.Results.Enqueue(MakeSnapshot); + client.Results.Enqueue(() => throw new InvalidOperationException("down")); + var (service, _, hubContext) = CreateService(client); + + await service.TickAsync(CancellationToken.None); + await service.TickAsync(CancellationToken.None); + + var calls = hubContext.Proxy.Calls.Where(c => c.Method == "UsageUpdated").ToList(); + Assert.Equal(2, calls.Count); + } } diff --git a/tests/ClaudeDo.Worker.Tests/Usage/UsageSnapshotBuilderTests.cs b/tests/ClaudeDo.Worker.Tests/Usage/UsageSnapshotBuilderTests.cs new file mode 100644 index 00000000..9b6aaf1e --- /dev/null +++ b/tests/ClaudeDo.Worker.Tests/Usage/UsageSnapshotBuilderTests.cs @@ -0,0 +1,145 @@ +using ClaudeDo.Data.Repositories; +using ClaudeDo.Worker.Config; +using ClaudeDo.Worker.Tests.Infrastructure; +using ClaudeDo.Worker.Usage; +using ClaudeDo.Worker.Usage.Interfaces; + +namespace ClaudeDo.Worker.Tests.Usage; + +public sealed class UsageSnapshotBuilderTests : IDisposable +{ + private readonly DbFixture _db = new(); + + public void Dispose() => _db.Dispose(); + + private sealed class FakeGate : IUsageGate + { + private readonly UsageGateDecision _decision; + public FakeGate(UsageGateDecision decision) => _decision = decision; + public Task EvaluateAsync(CancellationToken ct = default) => Task.FromResult(_decision); + } + + 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 UsageSnapshotBuilder CreateBuilder(UsageState state, UsageGateDecision decision, WorkerConfig? cfg = null) => + new(state, new FakeGate(decision), _db.CreateFactory(), cfg ?? new WorkerConfig()); + + [Fact] + public async Task Maps_buckets_limits_thresholds_and_gate_state() + { + await SetThresholdsAsync(70, 85); + + var fiveHourReset = DateTimeOffset.UtcNow.AddHours(2); + var sevenDayReset = DateTimeOffset.UtcNow.AddDays(3); + var limit = new UsageLimitRow("session", "default", 42.5, "warn", fiveHourReset, "sonnet", true); + var snapshot = new UsageSnapshot( + new UsageBucket(55.5, fiveHourReset), + new UsageBucket(88.8, sevenDayReset), + new[] { limit }, + DateTime.UtcNow); + + var state = new UsageState(); + state.ReportSuccess(snapshot); + + var builder = CreateBuilder(state, new UsageGateDecision(true, "5h-Limit 90% >= 70%")); + + var dto = await builder.BuildAsync(); + + Assert.Equal(55.5, dto.FiveHourPercent); + Assert.Equal(fiveHourReset, dto.FiveHourResetsAt); + Assert.Equal(88.8, dto.SevenDayPercent); + Assert.Equal(sevenDayReset, dto.SevenDayResetsAt); + Assert.Equal(70, dto.FiveHourThresholdPct); + Assert.Equal(85, dto.SevenDayThresholdPct); + Assert.True(dto.IsGateBlocked); + Assert.Equal("5h-Limit 90% >= 70%", dto.GateReason); + Assert.Equal(snapshot.FetchedAtUtc, dto.FetchedAtUtc); + Assert.False(dto.IsStale); + Assert.Null(dto.LastError); + + var mappedLimit = Assert.Single(dto.Limits); + Assert.Equal(limit.Kind, mappedLimit.Kind); + Assert.Equal(limit.Group, mappedLimit.Group); + Assert.Equal(limit.Percent, mappedLimit.Percent); + Assert.Equal(limit.Severity, mappedLimit.Severity); + Assert.Equal(limit.ResetsAt, mappedLimit.ResetsAt); + Assert.Equal(limit.ScopeModelDisplayName, mappedLimit.ScopeModelDisplayName); + Assert.Equal(limit.IsActive, mappedLimit.IsActive); + } + + [Fact] + public async Task No_snapshot_yields_null_percents_stale_and_not_blocked() + { + await SetThresholdsAsync(80, 90); + var state = new UsageState(); + var builder = CreateBuilder(state, new UsageGateDecision(false, null)); + + var dto = await builder.BuildAsync(); + + Assert.Null(dto.FiveHourPercent); + Assert.Null(dto.SevenDayPercent); + Assert.Empty(dto.Limits); + Assert.Null(dto.FetchedAtUtc); + Assert.True(dto.IsStale); + Assert.False(dto.IsGateBlocked); + } + + [Fact] + public async Task Snapshot_older_than_4x_poll_interval_is_stale() + { + await SetThresholdsAsync(80, 90); + var cfg = new WorkerConfig { UsagePollIntervalSeconds = 60 }; + var state = new UsageState(); + state.ReportSuccess(new UsageSnapshot( + new UsageBucket(10, null), null, Array.Empty(), + DateTime.UtcNow.AddSeconds(-4 * 60))); + + var builder = CreateBuilder(state, new UsageGateDecision(false, null), cfg); + + var dto = await builder.BuildAsync(); + + Assert.True(dto.IsStale); + } + + [Fact] + public async Task Fresh_snapshot_is_not_stale() + { + await SetThresholdsAsync(80, 90); + var cfg = new WorkerConfig { UsagePollIntervalSeconds = 60 }; + var state = new UsageState(); + state.ReportSuccess(new UsageSnapshot( + new UsageBucket(10, null), null, Array.Empty(), + DateTime.UtcNow)); + + var builder = CreateBuilder(state, new UsageGateDecision(false, null), cfg); + + var dto = await builder.BuildAsync(); + + Assert.False(dto.IsStale); + } + + [Fact] + public async Task Failed_last_poll_marks_stale_even_with_fresh_snapshot() + { + await SetThresholdsAsync(80, 90); + var state = new UsageState(); + state.ReportSuccess(new UsageSnapshot( + new UsageBucket(10, null), null, Array.Empty(), DateTime.UtcNow)); + state.ReportFailure("boom", DateTime.UtcNow); + + var builder = CreateBuilder(state, new UsageGateDecision(false, null)); + + var dto = await builder.BuildAsync(); + + Assert.True(dto.IsStale); + Assert.Equal("boom", dto.LastError); + } +}