diff --git a/src/ClaudeDo.Worker/CLAUDE.md b/src/ClaudeDo.Worker/CLAUDE.md index 7f4c091d..69d5a07e 100644 --- a/src/ClaudeDo.Worker/CLAUDE.md +++ b/src/ClaudeDo.Worker/CLAUDE.md @@ -19,6 +19,7 @@ Worker/ Hub/ — WorkerHub, HubBroadcaster Logging/ — LogRingBuffer (30-min in-memory log window) + BroadcastLogSink (Serilog sink → footer + overlay) Report/ — ClaudeHistoryReader, WeekReportPromptBuilder, WeekReportService; interfaces in Report/Interfaces/ + Usage/ — TranscriptUsageReader: aggregates Claude Code transcript token usage (~/.claude/projects/**/*.jsonl) by date/model/scope (ClaudeDo vs Other), deduped by requestId, with a per-file (length+mtime) cache; interface in Usage/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) ``` diff --git a/src/ClaudeDo.Worker/Program.cs b/src/ClaudeDo.Worker/Program.cs index 2cb4f463..0cf996a7 100644 --- a/src/ClaudeDo.Worker/Program.cs +++ b/src/ClaudeDo.Worker/Program.cs @@ -19,6 +19,8 @@ using ClaudeDo.Worker.Refine; using ClaudeDo.Worker.Report; using ClaudeDo.Worker.Report.Interfaces; using ClaudeDo.Worker.Skills; +using ClaudeDo.Worker.Usage; +using ClaudeDo.Worker.Usage.Interfaces; using ClaudeDo.Worker.Worktrees; using Microsoft.EntityFrameworkCore; using Serilog; @@ -122,6 +124,9 @@ builder.Services.AddSingleton(_ => Environment.GetFolderPath(Environment.SpecialFolder.UserProfile), ".claude", "projects"))); builder.Services.AddSingleton(); +// Usage +builder.Services.AddSingleton(); + // Prime Claude builder.Services.AddSingleton(); builder.Services.AddSingleton(); diff --git a/src/ClaudeDo.Worker/Usage/Interfaces/ITranscriptUsageReader.cs b/src/ClaudeDo.Worker/Usage/Interfaces/ITranscriptUsageReader.cs new file mode 100644 index 00000000..2a297020 --- /dev/null +++ b/src/ClaudeDo.Worker/Usage/Interfaces/ITranscriptUsageReader.cs @@ -0,0 +1,7 @@ +namespace ClaudeDo.Worker.Usage.Interfaces; + +public interface ITranscriptUsageReader +{ + Task> ReadAsync( + DateOnly start, DateOnly end, CancellationToken ct = default); +} diff --git a/src/ClaudeDo.Worker/Usage/TranscriptUsageReader.cs b/src/ClaudeDo.Worker/Usage/TranscriptUsageReader.cs new file mode 100644 index 00000000..693c89a8 --- /dev/null +++ b/src/ClaudeDo.Worker/Usage/TranscriptUsageReader.cs @@ -0,0 +1,163 @@ +using System.Collections.Concurrent; +using System.Text.Json; +using ClaudeDo.Data; +using ClaudeDo.Worker.Config; +using ClaudeDo.Worker.Usage.Interfaces; + +namespace ClaudeDo.Worker.Usage; + +public sealed class TranscriptUsageReader : ITranscriptUsageReader +{ + private readonly string _projectsRoot; + private readonly string _centralRoot; + private readonly string _sandboxRoot; + private readonly ConcurrentDictionary _cache = new(); + + public TranscriptUsageReader(WorkerConfig cfg, string? projectsRoot = null) + { + _projectsRoot = projectsRoot ?? Paths.Expand("~/.claude/projects"); + _centralRoot = NormalizePath(cfg.CentralWorktreeRoot); + _sandboxRoot = NormalizePath(cfg.SandboxRoot); + } + + public Task> ReadAsync( + DateOnly start, DateOnly end, CancellationToken ct = default) + { + var seenKeys = new HashSet(); + var buckets = new Dictionary<(DateOnly Date, string Model, UsageScope Scope), Accumulator>(); + + if (Directory.Exists(_projectsRoot)) + { + foreach (var file in Directory.EnumerateFiles(_projectsRoot, "*.jsonl", SearchOption.AllDirectories)) + { + ct.ThrowIfCancellationRequested(); + + foreach (var record in GetOrReadFile(file)) + { + if (record.Date < start || record.Date > end) continue; + if (!seenKeys.Add(record.DedupeKey)) continue; + + var key = (record.Date, record.Model, record.Scope); + if (!buckets.TryGetValue(key, out var acc)) + { + acc = new Accumulator(); + buckets[key] = acc; + } + acc.Input += record.InputTokens; + acc.Output += record.OutputTokens; + acc.CacheRead += record.CacheReadTokens; + acc.CacheCreation += record.CacheCreationTokens; + acc.Messages++; + } + } + } + + var rows = buckets + .Select(kv => new UsageAggregateRow( + kv.Key.Date, kv.Key.Model, kv.Key.Scope, + kv.Value.Input, kv.Value.Output, kv.Value.CacheRead, kv.Value.CacheCreation, kv.Value.Messages)) + .OrderBy(r => r.Date).ThenBy(r => r.Model).ThenBy(r => r.Scope) + .ToList(); + + return Task.FromResult>(rows); + } + + private List GetOrReadFile(string file) + { + var info = new FileInfo(file); + if (_cache.TryGetValue(file, out var cached) && + cached.Length == info.Length && cached.LastWriteUtc == info.LastWriteTimeUtc) + { + return cached.Records; + } + + var records = ReadFile(file); + _cache[file] = new FileCacheEntry(info.Length, info.LastWriteTimeUtc, records); + return records; + } + + private List ReadFile(string file) + { + var records = new List(); + + foreach (var line in File.ReadLines(file)) + { + if (string.IsNullOrWhiteSpace(line)) continue; + + JsonDocument doc; + try { doc = JsonDocument.Parse(line); } + catch (JsonException) { continue; } + + using (doc) + { + var root = doc.RootElement; + if (root.ValueKind != JsonValueKind.Object) continue; + if (!root.TryGetProperty("type", out var typeEl) || typeEl.GetString() != "assistant") continue; + if (!root.TryGetProperty("timestamp", out var tsEl) || + !DateTimeOffset.TryParse(tsEl.GetString(), out var ts)) continue; + if (!root.TryGetProperty("message", out var msg) || msg.ValueKind != JsonValueKind.Object) continue; + if (!msg.TryGetProperty("model", out var modelEl) || modelEl.ValueKind != JsonValueKind.String) continue; + + var date = DateOnly.FromDateTime(ts.LocalDateTime); + var model = modelEl.GetString()!; + + long input = 0, output = 0, cacheRead = 0, cacheCreation = 0; + if (msg.TryGetProperty("usage", out var usage) && usage.ValueKind == JsonValueKind.Object) + { + input = GetLong(usage, "input_tokens"); + output = GetLong(usage, "output_tokens"); + cacheRead = GetLong(usage, "cache_read_input_tokens"); + cacheCreation = GetLong(usage, "cache_creation_input_tokens"); + } + + var cwd = TryGetString(root, "cwd") ?? ""; + var dedupeKey = TryGetString(root, "requestId") + ?? TryGetString(msg, "id") + ?? Guid.NewGuid().ToString(); + + records.Add(new UsageMessageRecord( + date, model, ResolveScope(cwd), input, output, cacheRead, cacheCreation, dedupeKey)); + } + } + + return records; + } + + private UsageScope ResolveScope(string cwd) + { + var norm = NormalizePath(cwd); + if (norm.Length == 0) return UsageScope.Other; + + // Sibling-strategy worktrees are placed next to whatever repo they belong to + // (no single root path), but always under a literal ".claudedo-worktrees" segment. + if (norm.Split('\\').Any(seg => seg == ".claudedo-worktrees")) return UsageScope.ClaudeDo; + if (IsUnderRoot(norm, _centralRoot)) return UsageScope.ClaudeDo; + if (IsUnderRoot(norm, _sandboxRoot)) return UsageScope.ClaudeDo; + return UsageScope.Other; + } + + private static bool IsUnderRoot(string normPath, string normRoot) => + normRoot.Length > 0 && (normPath == normRoot || normPath.StartsWith(normRoot + "\\", StringComparison.Ordinal)); + + private static string? TryGetString(JsonElement obj, string prop) => + obj.TryGetProperty(prop, out var el) && el.ValueKind == JsonValueKind.String ? el.GetString() : null; + + private static long GetLong(JsonElement obj, string prop) => + obj.TryGetProperty(prop, out var el) && el.TryGetInt64(out var v) ? v : 0; + + private static string NormalizePath(string p) => + (p ?? "").Replace('/', '\\').TrimEnd('\\').ToLowerInvariant(); + + private sealed class Accumulator + { + public long Input, Output, CacheRead, CacheCreation; + public int Messages; + } + + private sealed record FileCacheEntry(long Length, DateTime LastWriteUtc, List Records); + + private sealed record UsageMessageRecord( + DateOnly Date, string Model, UsageScope Scope, + long InputTokens, long OutputTokens, long CacheReadTokens, long CacheCreationTokens, + string DedupeKey); +} diff --git a/src/ClaudeDo.Worker/Usage/UsageModels.cs b/src/ClaudeDo.Worker/Usage/UsageModels.cs new file mode 100644 index 00000000..974d568c --- /dev/null +++ b/src/ClaudeDo.Worker/Usage/UsageModels.cs @@ -0,0 +1,17 @@ +namespace ClaudeDo.Worker.Usage; + +public enum UsageScope +{ + ClaudeDo, + Other, +} + +public sealed record UsageAggregateRow( + DateOnly Date, + string Model, + UsageScope Scope, + long InputTokens, + long OutputTokens, + long CacheReadTokens, + long CacheCreationTokens, + int Messages); diff --git a/tests/ClaudeDo.Worker.Tests/Usage/TranscriptUsageReaderTests.cs b/tests/ClaudeDo.Worker.Tests/Usage/TranscriptUsageReaderTests.cs new file mode 100644 index 00000000..2587fb5a --- /dev/null +++ b/tests/ClaudeDo.Worker.Tests/Usage/TranscriptUsageReaderTests.cs @@ -0,0 +1,202 @@ +using System.Text.Json; +using ClaudeDo.Worker.Config; +using ClaudeDo.Worker.Usage; + +namespace ClaudeDo.Worker.Tests.Usage; + +public class TranscriptUsageReaderTests : IDisposable +{ + private readonly string _root; + private readonly WorkerConfig _cfg; + + public TranscriptUsageReaderTests() + { + _root = Path.Combine(Path.GetTempPath(), $"tur_{Guid.NewGuid():N}"); + Directory.CreateDirectory(_root); + _cfg = new WorkerConfig + { + CentralWorktreeRoot = Path.Combine(_root, "central-worktrees"), + SandboxRoot = Path.Combine(_root, "sandbox"), + }; + } + + public void Dispose() { try { Directory.Delete(_root, true); } catch { } } + + private string WriteSession(string projectDir, string file, params string[] lines) + { + var dir = Path.Combine(_root, "projects", projectDir); + Directory.CreateDirectory(dir); + var path = Path.Combine(dir, file); + File.WriteAllLines(path, lines); + return path; + } + + private static string AssistantLine( + string cwd, string ts, string model, long input, long output, long cacheRead, long cacheCreation, + string? requestId = null, string? messageId = null) => + JsonSerializer.Serialize(new Dictionary + { + ["type"] = "assistant", + ["cwd"] = cwd, + ["timestamp"] = ts, + ["requestId"] = requestId, + ["message"] = new Dictionary + { + ["id"] = messageId, + ["model"] = model, + ["usage"] = new Dictionary + { + ["input_tokens"] = input, + ["output_tokens"] = output, + ["cache_read_input_tokens"] = cacheRead, + ["cache_creation_input_tokens"] = cacheCreation, + }, + }, + }); + + private TranscriptUsageReader MakeReader() => + new(_cfg, Path.Combine(_root, "projects")); + + [Fact] + public async Task Aggregates_By_Date_And_Model_Including_Cache_Tokens() + { + WriteSession("proj", "s.jsonl", + AssistantLine(@"C:\Dev\App", "2026-06-01T08:00:00Z", "claude-sonnet-5", 10, 20, 3, 1), + AssistantLine(@"C:\Dev\App", "2026-06-01T09:00:00Z", "claude-sonnet-5", 5, 6, 1, 0), + AssistantLine(@"C:\Dev\App", "2026-06-02T08:00:00Z", "claude-opus-5", 100, 200, 30, 10)); + + var reader = MakeReader(); + var result = await reader.ReadAsync(new DateOnly(2026, 6, 1), new DateOnly(2026, 6, 3)); + + Assert.Equal(2, result.Count); + var day1 = Assert.Single(result, r => r.Model == "claude-sonnet-5"); + Assert.Equal(new DateOnly(2026, 6, 1), day1.Date); + Assert.Equal(15, day1.InputTokens); + Assert.Equal(26, day1.OutputTokens); + Assert.Equal(4, day1.CacheReadTokens); + Assert.Equal(1, day1.CacheCreationTokens); + Assert.Equal(2, day1.Messages); + + var day2 = Assert.Single(result, r => r.Model == "claude-opus-5"); + Assert.Equal(new DateOnly(2026, 6, 2), day2.Date); + Assert.Equal(100, day2.InputTokens); + Assert.Equal(1, day2.Messages); + } + + [Fact] + public async Task Same_RequestId_Across_Files_Counts_Once() + { + WriteSession("proj-a", "s1.jsonl", + AssistantLine(@"C:\Dev\App", "2026-06-01T08:00:00Z", "claude-sonnet-5", 10, 20, 0, 0, requestId: "req-1")); + WriteSession("proj-b", "s2.jsonl", + AssistantLine(@"C:\Dev\App", "2026-06-01T08:00:00Z", "claude-sonnet-5", 10, 20, 0, 0, requestId: "req-1")); + + var reader = MakeReader(); + var result = await reader.ReadAsync(new DateOnly(2026, 6, 1), new DateOnly(2026, 6, 3)); + + var row = Assert.Single(result); + Assert.Equal(1, row.Messages); + Assert.Equal(10, row.InputTokens); + } + + [Fact] + public async Task Scope_Split_Covers_Sibling_Central_And_Sandbox_Roots() + { + var siblingCwd = Path.Combine(_root, "some-repo", ".claudedo-worktrees", "list-slug", "task-id"); + var centralCwd = Path.Combine(_cfg.CentralWorktreeRoot, "list-slug", "task-id"); + var sandboxCwd = Path.Combine(_cfg.SandboxRoot, "task-id"); + var otherCwd = Path.Combine(_root, "some-repo"); + + WriteSession("proj", "s.jsonl", + AssistantLine(siblingCwd, "2026-06-01T08:00:00Z", "claude-sonnet-5", 1, 1, 0, 0, requestId: "sibling"), + AssistantLine(centralCwd, "2026-06-01T08:00:00Z", "claude-sonnet-5", 1, 1, 0, 0, requestId: "central"), + AssistantLine(sandboxCwd, "2026-06-01T08:00:00Z", "claude-sonnet-5", 1, 1, 0, 0, requestId: "sandbox"), + AssistantLine(otherCwd, "2026-06-01T08:00:00Z", "claude-sonnet-5", 1, 1, 0, 0, requestId: "other")); + + var reader = MakeReader(); + var result = await reader.ReadAsync(new DateOnly(2026, 6, 1), new DateOnly(2026, 6, 3)); + + Assert.Equal(2, result.Count); + var claudeDo = Assert.Single(result, r => r.Scope == UsageScope.ClaudeDo); + Assert.Equal(3, claudeDo.Messages); + var other = Assert.Single(result, r => r.Scope == UsageScope.Other); + Assert.Equal(1, other.Messages); + } + + [Fact] + public async Task Date_Filter_Excludes_Lines_Outside_Window() + { + WriteSession("proj", "s.jsonl", + AssistantLine(@"C:\Dev\App", "2026-05-01T08:00:00Z", "claude-sonnet-5", 1, 1, 0, 0), + AssistantLine(@"C:\Dev\App", "2026-06-02T08:00:00Z", "claude-sonnet-5", 5, 5, 0, 0), + AssistantLine(@"C:\Dev\App", "2026-07-01T08:00:00Z", "claude-sonnet-5", 1, 1, 0, 0)); + + var reader = MakeReader(); + var result = await reader.ReadAsync(new DateOnly(2026, 6, 1), new DateOnly(2026, 6, 3)); + + var row = Assert.Single(result); + Assert.Equal(5, row.InputTokens); + Assert.Equal(1, row.Messages); + } + + [Fact] + public async Task Malformed_Line_Does_Not_Abort_The_Run() + { + WriteSession("proj", "s.jsonl", + "this is not json", + AssistantLine(@"C:\Dev\App", "2026-06-01T08:00:00Z", "claude-sonnet-5", 5, 5, 0, 0)); + + var reader = MakeReader(); + var result = await reader.ReadAsync(new DateOnly(2026, 6, 1), new DateOnly(2026, 6, 3)); + + var row = Assert.Single(result); + Assert.Equal(1, row.Messages); + } + + [Fact] + public async Task Cache_Skips_Unchanged_File_And_Picks_Up_Appended_Lines() + { + var path = WriteSession("proj", "s.jsonl", + AssistantLine(@"C:\Dev\App", "2026-06-01T08:00:00Z", "claude-sonnet-5", 5, 5, 0, 0, requestId: "req-a")); + + var reader = MakeReader(); + var first = Assert.Single(await reader.ReadAsync(new DateOnly(2026, 6, 1), new DateOnly(2026, 6, 3))); + Assert.Equal(1, first.Messages); + + // Overwrite with different content but the SAME length and mtime: if the reader honored + // the cache it must still return the ORIGINAL aggregate, proving it did not re-read the file. + var originalBytes = File.ReadAllBytes(path); + var originalWriteUtc = File.GetLastWriteTimeUtc(path); + var tamperedLine = AssistantLine(@"C:\Dev\App", "2026-06-01T08:00:00Z", "claude-sonnet-5", 9, 9, 0, 0, requestId: "req-9"); + // Pad/truncate to the EXACT same byte length as the original file so the (length, mtime) + // cache key still matches — the reader must then serve the cached (stale) aggregate. + var tamperedPadded = tamperedLine.Length + 1 <= originalBytes.Length + ? tamperedLine.PadRight(originalBytes.Length - 1) + "\n" + : tamperedLine[..(originalBytes.Length - 1)] + "\n"; + var tamperedBytes = System.Text.Encoding.UTF8.GetBytes(tamperedPadded); + Assert.Equal(originalBytes.Length, tamperedBytes.Length); + File.WriteAllBytes(path, tamperedBytes); + File.SetLastWriteTimeUtc(path, originalWriteUtc); + + var stale = Assert.Single(await reader.ReadAsync(new DateOnly(2026, 6, 1), new DateOnly(2026, 6, 3))); + Assert.Equal(5, stale.InputTokens); + + // Now really append a new line: length/mtime change, so the file must be re-read. + File.AppendAllLines(path, new[] + { + AssistantLine(@"C:\Dev\App", "2026-06-01T09:00:00Z", "claude-sonnet-5", 3, 3, 0, 0, requestId: "req-b"), + }); + + var updated = Assert.Single(await reader.ReadAsync(new DateOnly(2026, 6, 1), new DateOnly(2026, 6, 3))); + Assert.Equal(2, updated.Messages); + } + + [Fact] + public async Task Missing_ProjectsRoot_Returns_Empty_Without_Throwing() + { + var reader = new TranscriptUsageReader(_cfg, Path.Combine(_root, "does-not-exist")); + var result = await reader.ReadAsync(new DateOnly(2026, 6, 1), new DateOnly(2026, 6, 3)); + + Assert.Empty(result); + } +}