feat(usage): cache the TokenTracker export with single-flight refresh
This commit is contained in:
@@ -0,0 +1,98 @@
|
||||
using ClaudeDo.Worker.Usage.TokenTracker.Interfaces;
|
||||
|
||||
namespace ClaudeDo.Worker.Usage.TokenTracker;
|
||||
|
||||
/// <summary>
|
||||
/// Owns the cached export. One fetch covers <see cref="WindowDays"/> days because
|
||||
/// <c>buildSessionAnalytics</c> walks the whole history regardless of <c>--from</c>/<c>--to</c>
|
||||
/// (v0.88.4 returns the same rows either way), so narrowing the request would buy nothing while
|
||||
/// making every range switch pay the cost again.
|
||||
/// </summary>
|
||||
public sealed class TokenTrackerService
|
||||
{
|
||||
public const int WindowDays = 90;
|
||||
|
||||
private readonly ITokenTrackerClient _client;
|
||||
private readonly Func<DateTime> _utcNow;
|
||||
private readonly SemaphoreSlim _refreshLock = new(1, 1);
|
||||
private TokenTrackerProbe? _probe;
|
||||
|
||||
public TokenTrackerService(ITokenTrackerClient client, TokenTrackerState state, Func<DateTime>? utcNow = null)
|
||||
{
|
||||
_client = client;
|
||||
State = state;
|
||||
_utcNow = utcNow ?? (() => DateTime.UtcNow);
|
||||
}
|
||||
|
||||
public TokenTrackerState State { get; }
|
||||
|
||||
public async Task<TokenTrackerProbe> ProbeAsync(bool force = false, CancellationToken ct = default)
|
||||
{
|
||||
if (!force && _probe is { } cached) return cached;
|
||||
|
||||
var probe = await _client.ProbeAsync(ct);
|
||||
_probe = probe;
|
||||
return probe;
|
||||
}
|
||||
|
||||
/// <summary>Fetches only when the cached export is missing or older than
|
||||
/// <paramref name="maxAge"/>. Safe to call on every modal open.</summary>
|
||||
public Task EnsureFreshAsync(TimeSpan maxAge, CancellationToken ct = default) =>
|
||||
State.IsOlderThan(maxAge, _utcNow()) ? RefreshAsync(ct) : Task.CompletedTask;
|
||||
|
||||
public async Task RefreshAsync(CancellationToken ct = default)
|
||||
{
|
||||
// A refresh already in flight is the refresh this caller wanted — waiting for it and
|
||||
// returning is correct, and it stops a burst of range switches from spawning processes.
|
||||
if (!await _refreshLock.WaitAsync(0, ct))
|
||||
{
|
||||
await _refreshLock.WaitAsync(ct);
|
||||
_refreshLock.Release();
|
||||
return;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
var now = _utcNow();
|
||||
var to = DateOnly.FromDateTime(now.ToLocalTime());
|
||||
var from = to.AddDays(-WindowDays);
|
||||
|
||||
var run = await _client.ExportAsync(from, to, ct);
|
||||
if (!run.Ok)
|
||||
{
|
||||
State.ReportFailure(run.Error ?? "TokenTracker export failed.", now);
|
||||
return;
|
||||
}
|
||||
|
||||
var export = TokenTrackerExportParser.Parse(run.StdOut, now);
|
||||
if (export is null)
|
||||
{
|
||||
State.ReportFailure("TokenTracker returned output we could not parse.", now);
|
||||
return;
|
||||
}
|
||||
|
||||
if (export.FormatVersion != TokenTrackerExportParser.SupportedFormatVersion)
|
||||
{
|
||||
State.ReportFailure(
|
||||
$"Unexpected TokenTracker export format v{export.FormatVersion} " +
|
||||
$"(supported: v{TokenTrackerExportParser.SupportedFormatVersion}).",
|
||||
now);
|
||||
return;
|
||||
}
|
||||
|
||||
State.ReportSuccess(export);
|
||||
}
|
||||
finally
|
||||
{
|
||||
_refreshLock.Release();
|
||||
}
|
||||
}
|
||||
|
||||
public async Task<TokenTrackerRunResult> InstallAsync(
|
||||
IProgress<string>? output = null, CancellationToken ct = default)
|
||||
{
|
||||
var result = await _client.InstallAsync(output, ct);
|
||||
await ProbeAsync(force: true, ct);
|
||||
return result;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,47 @@
|
||||
using ClaudeDo.Worker.Usage.TokenTracker;
|
||||
using ClaudeDo.Worker.Usage.TokenTracker.Interfaces;
|
||||
|
||||
namespace ClaudeDo.Worker.Tests.Usage.TokenTracker;
|
||||
|
||||
public sealed class FakeTokenTrackerClient : ITokenTrackerClient
|
||||
{
|
||||
private readonly SemaphoreSlim _gate = new(0);
|
||||
|
||||
public TokenTrackerProbe Probe { get; set; } = new(true, "1.2.3", true, "v22.1.0", null);
|
||||
public TokenTrackerRunResult ExportResult { get; set; } =
|
||||
new(true, TokenTrackerFixtures.SessionsJson, null);
|
||||
public TokenTrackerRunResult InstallResult { get; set; } = new(true, "added 1 package", null);
|
||||
|
||||
public int ExportCalls;
|
||||
public int ProbeCalls;
|
||||
public int InstallCalls;
|
||||
public DateOnly? LastFrom;
|
||||
public DateOnly? LastTo;
|
||||
|
||||
/// <summary>When true, <see cref="ExportAsync"/> blocks until <see cref="Release"/> is called.</summary>
|
||||
public bool BlockExport { get; set; }
|
||||
|
||||
public void Release() => _gate.Release();
|
||||
|
||||
public Task<TokenTrackerProbe> ProbeAsync(CancellationToken ct = default)
|
||||
{
|
||||
Interlocked.Increment(ref ProbeCalls);
|
||||
return Task.FromResult(Probe);
|
||||
}
|
||||
|
||||
public async Task<TokenTrackerRunResult> ExportAsync(DateOnly from, DateOnly to, CancellationToken ct = default)
|
||||
{
|
||||
Interlocked.Increment(ref ExportCalls);
|
||||
LastFrom = from;
|
||||
LastTo = to;
|
||||
if (BlockExport) await _gate.WaitAsync(ct);
|
||||
return ExportResult;
|
||||
}
|
||||
|
||||
public Task<TokenTrackerRunResult> InstallAsync(IProgress<string>? output = null, CancellationToken ct = default)
|
||||
{
|
||||
Interlocked.Increment(ref InstallCalls);
|
||||
output?.Report("added 1 package");
|
||||
return Task.FromResult(InstallResult);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,138 @@
|
||||
using ClaudeDo.Worker.Usage.TokenTracker;
|
||||
|
||||
namespace ClaudeDo.Worker.Tests.Usage.TokenTracker;
|
||||
|
||||
public sealed class TokenTrackerServiceTests
|
||||
{
|
||||
private static readonly DateTime Now = new(2026, 8, 24, 9, 0, 0, DateTimeKind.Utc);
|
||||
|
||||
private static (TokenTrackerService Service, FakeTokenTrackerClient Client) Build(DateTime? now = null)
|
||||
{
|
||||
var client = new FakeTokenTrackerClient();
|
||||
var clock = now ?? Now;
|
||||
return (new TokenTrackerService(client, new TokenTrackerState(), () => clock), client);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task RefreshAsync_StoresParsedExport()
|
||||
{
|
||||
var (service, client) = Build();
|
||||
|
||||
await service.RefreshAsync();
|
||||
|
||||
Assert.Equal(1, client.ExportCalls);
|
||||
Assert.Equal(4, service.State.Export!.Sessions.Count);
|
||||
Assert.Null(service.State.LastError);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task RefreshAsync_RequestsTheConfiguredWindow()
|
||||
{
|
||||
var (service, client) = Build();
|
||||
|
||||
await service.RefreshAsync();
|
||||
|
||||
Assert.Equal(DateOnly.FromDateTime(Now.ToLocalTime()), client.LastTo);
|
||||
Assert.Equal(
|
||||
DateOnly.FromDateTime(Now.ToLocalTime()).AddDays(-TokenTrackerService.WindowDays),
|
||||
client.LastFrom);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task RefreshAsync_Failure_KeepsLastGoodExport()
|
||||
{
|
||||
var (service, client) = Build();
|
||||
await service.RefreshAsync();
|
||||
|
||||
client.ExportResult = new TokenTrackerRunResult(false, "", "tokentracker exited with code 1");
|
||||
await service.RefreshAsync();
|
||||
|
||||
Assert.Equal(4, service.State.Export!.Sessions.Count);
|
||||
Assert.Equal("tokentracker exited with code 1", service.State.LastError);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task RefreshAsync_UnparseableOutput_IsAnError()
|
||||
{
|
||||
var (service, client) = Build();
|
||||
client.ExportResult = new TokenTrackerRunResult(true, "not json", null);
|
||||
|
||||
await service.RefreshAsync();
|
||||
|
||||
Assert.Null(service.State.Export);
|
||||
Assert.NotNull(service.State.LastError);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task RefreshAsync_UnsupportedFormatVersion_IsAnErrorAndKeepsNoData()
|
||||
{
|
||||
var (service, client) = Build();
|
||||
client.ExportResult = new TokenTrackerRunResult(true, TokenTrackerFixtures.UnsupportedVersionJson, null);
|
||||
|
||||
await service.RefreshAsync();
|
||||
|
||||
Assert.Null(service.State.Export);
|
||||
Assert.Contains("99", service.State.LastError);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task EnsureFreshAsync_FreshExport_DoesNotFetchAgain()
|
||||
{
|
||||
var (service, client) = Build();
|
||||
await service.RefreshAsync();
|
||||
|
||||
await service.EnsureFreshAsync(TimeSpan.FromMinutes(15));
|
||||
|
||||
Assert.Equal(1, client.ExportCalls);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task EnsureFreshAsync_NoExportYet_Fetches()
|
||||
{
|
||||
var (service, client) = Build();
|
||||
|
||||
await service.EnsureFreshAsync(TimeSpan.FromMinutes(15));
|
||||
|
||||
Assert.Equal(1, client.ExportCalls);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task ConcurrentRefresh_RunsTheCliOnce()
|
||||
{
|
||||
var (service, client) = Build();
|
||||
client.BlockExport = true;
|
||||
|
||||
var first = service.RefreshAsync();
|
||||
var second = service.RefreshAsync();
|
||||
client.Release();
|
||||
await Task.WhenAll(first, second);
|
||||
|
||||
Assert.Equal(1, client.ExportCalls);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task ProbeAsync_CachesUntilForced()
|
||||
{
|
||||
var (service, client) = Build();
|
||||
|
||||
await service.ProbeAsync();
|
||||
await service.ProbeAsync();
|
||||
Assert.Equal(1, client.ProbeCalls);
|
||||
|
||||
await service.ProbeAsync(force: true);
|
||||
Assert.Equal(2, client.ProbeCalls);
|
||||
}
|
||||
|
||||
[Fact]
|
||||
public async Task InstallAsync_ReprobesAfterwards()
|
||||
{
|
||||
var (service, client) = Build();
|
||||
await service.ProbeAsync();
|
||||
|
||||
var result = await service.InstallAsync();
|
||||
|
||||
Assert.True(result.Ok);
|
||||
Assert.Equal(1, client.InstallCalls);
|
||||
Assert.Equal(2, client.ProbeCalls);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user