feat: add StatsCache for shared per-model usage polling
This commit is contained in:
@@ -0,0 +1,83 @@
|
||||
import { fetchStats, type UsageStats } from "./stats";
|
||||
import { type LlamaSwapConfig } from "./util";
|
||||
|
||||
const POLL_MS = 5000;
|
||||
const ACTIVITY_THROTTLE_MS = 2000;
|
||||
|
||||
export type FetchFn = (cfg: LlamaSwapConfig, modelId: string) => Promise<UsageStats | null>;
|
||||
|
||||
export class StatsCache {
|
||||
private keys = new Set<string>();
|
||||
private values = new Map<string, UsageStats | undefined>();
|
||||
private listeners = new Set<() => void>();
|
||||
private timer?: ReturnType<typeof setInterval>;
|
||||
private throttleTimer?: ReturnType<typeof setTimeout>;
|
||||
private lastRefresh = 0;
|
||||
|
||||
constructor(
|
||||
private cfg: LlamaSwapConfig,
|
||||
private fetchFn: FetchFn = fetchStats,
|
||||
private pollMs = POLL_MS,
|
||||
private throttleMs = ACTIVITY_THROTTLE_MS,
|
||||
) {}
|
||||
|
||||
setConfig(cfg: LlamaSwapConfig): void {
|
||||
this.cfg = cfg;
|
||||
this.values.clear();
|
||||
}
|
||||
|
||||
register(key: string): void {
|
||||
if (this.keys.has(key)) return;
|
||||
this.keys.add(key);
|
||||
if (!this.timer) {
|
||||
void this.refresh();
|
||||
this.timer = setInterval(() => void this.refresh(), this.pollMs);
|
||||
}
|
||||
}
|
||||
|
||||
unregister(key: string): void {
|
||||
this.keys.delete(key);
|
||||
this.values.delete(key);
|
||||
if (this.keys.size === 0 && this.timer) {
|
||||
clearInterval(this.timer);
|
||||
this.timer = undefined;
|
||||
}
|
||||
}
|
||||
|
||||
get(key: string): UsageStats | undefined {
|
||||
return this.values.get(key);
|
||||
}
|
||||
|
||||
scheduleRefresh(): void {
|
||||
const now = Date.now();
|
||||
const wait = Math.max(0, this.throttleMs - (now - this.lastRefresh));
|
||||
if (this.throttleTimer) clearTimeout(this.throttleTimer);
|
||||
this.throttleTimer = setTimeout(() => {
|
||||
this.throttleTimer = undefined;
|
||||
void this.refresh();
|
||||
}, wait);
|
||||
}
|
||||
|
||||
async refresh(): Promise<void> {
|
||||
this.lastRefresh = Date.now();
|
||||
for (const key of this.keys) {
|
||||
try {
|
||||
const stats = await this.fetchFn(this.cfg, key);
|
||||
if (stats) this.values.set(key, stats);
|
||||
} catch {
|
||||
}
|
||||
}
|
||||
this.emit();
|
||||
}
|
||||
|
||||
onChange(listener: () => void): () => void {
|
||||
this.listeners.add(listener);
|
||||
return () => {
|
||||
this.listeners.delete(listener);
|
||||
};
|
||||
}
|
||||
|
||||
private emit(): void {
|
||||
for (const listener of this.listeners) listener();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,85 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { test } from "node:test";
|
||||
import { StatsCache, type FetchFn } from "../src/lib/stats-cache";
|
||||
import { type LlamaSwapConfig } from "../src/lib/util";
|
||||
|
||||
const cfg: LlamaSwapConfig = { baseUrl: "http://x" };
|
||||
|
||||
function stubFetch(result: unknown, counter: { count: number }): FetchFn {
|
||||
return async () => {
|
||||
counter.count++;
|
||||
return result as never;
|
||||
};
|
||||
}
|
||||
|
||||
async function flush(): Promise<void> {
|
||||
await new Promise((r) => setTimeout(r, 0));
|
||||
}
|
||||
|
||||
test("register polls once immediately and caches the value", async () => {
|
||||
const calls = { count: 0 };
|
||||
const cache = new StatsCache(cfg, stubFetch({ totalRequests: 3, totalInputTokens: 1, totalOutputTokens: 2, genP95: 4 }, calls), 1000, 30);
|
||||
cache.register("all");
|
||||
await flush();
|
||||
assert.equal(calls.count, 1);
|
||||
assert.equal(cache.get("all")!.totalRequests, 3);
|
||||
cache.unregister("all");
|
||||
});
|
||||
|
||||
test("multiple keys are each polled", async () => {
|
||||
const keys: string[] = [];
|
||||
const cache = new StatsCache(
|
||||
cfg,
|
||||
async (_c, key) => {
|
||||
keys.push(key);
|
||||
return { totalRequests: 1, totalInputTokens: 1, totalOutputTokens: 1, genP95: 1 };
|
||||
},
|
||||
1000,
|
||||
30,
|
||||
);
|
||||
cache.register("a");
|
||||
cache.register("b");
|
||||
await flush();
|
||||
assert.deepEqual(keys.sort(), ["a", "b"]);
|
||||
cache.unregister("a");
|
||||
cache.unregister("b");
|
||||
});
|
||||
|
||||
test("unregister stops the interval and clears the value", async () => {
|
||||
const calls = { count: 0 };
|
||||
const cache = new StatsCache(cfg, stubFetch({ totalRequests: 1, totalInputTokens: 1, totalOutputTokens: 1, genP95: 1 }, calls), 10, 30);
|
||||
cache.register("a");
|
||||
await flush();
|
||||
assert.equal(calls.count, 1);
|
||||
cache.unregister("a");
|
||||
await new Promise((r) => setTimeout(r, 40));
|
||||
assert.equal(cache.get("a"), undefined);
|
||||
assert.equal(calls.count, 1);
|
||||
});
|
||||
|
||||
test("scheduleRefresh throttles to one poll per window", async () => {
|
||||
const calls = { count: 0 };
|
||||
const cache = new StatsCache(cfg, stubFetch({ totalRequests: 1, totalInputTokens: 1, totalOutputTokens: 1, genP95: 1 }, calls), 10000, 30);
|
||||
cache.register("a");
|
||||
await flush();
|
||||
const before = calls.count;
|
||||
cache.scheduleRefresh();
|
||||
cache.scheduleRefresh();
|
||||
cache.scheduleRefresh();
|
||||
await flush();
|
||||
assert.equal(calls.count, before);
|
||||
await new Promise((r) => setTimeout(r, 60));
|
||||
assert.equal(calls.count, before + 1);
|
||||
cache.unregister("a");
|
||||
});
|
||||
|
||||
test("fetch failure keeps the last-known value", async () => {
|
||||
const fail: FetchFn = async () => {
|
||||
throw new Error("boom");
|
||||
};
|
||||
const cache = new StatsCache(cfg, fail, 10000, 30);
|
||||
cache.register("a");
|
||||
await flush();
|
||||
assert.equal(cache.get("a"), undefined);
|
||||
cache.unregister("a");
|
||||
});
|
||||
Reference in New Issue
Block a user