Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions apps/server/src/usage/UsageService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,9 @@
* `(size, mtime)`. A cold 30-day scan of ~1.4 GB lands around 2-3 seconds; warm
* scans only reparse files that changed, and a file that merely grew resumes
* from its cached parse position so only the appended bytes are read.
* SQLite readers query live databases each scan so WAL writes remain visible.
* OpenCode's SQLite reader queries the live database each scan so WAL writes
* remain visible. Antigravity databases are memoised in memory while the
* database and its WAL keep the same `(size, mtime, ctime)`.
*
* @module UsageService
*/
Expand Down Expand Up @@ -51,7 +53,7 @@ import { resolveCodexHomeLayout } from "../provider/Drivers/CodexHomeLayout.ts";
import { resolveAntigravityInstanceDirectories } from "../provider/antigravityAuthSupport.ts";
import { mergeProviderInstanceEnvironment } from "../provider/ProviderInstanceEnvironment.ts";
import { readOpenCodeUsage } from "./opencodeUsageReader.ts";
import { readAntigravityUsage } from "./antigravityUsageReader.ts";
import { makeAntigravityUsageCache, readAntigravityUsage } from "./antigravityUsageReader.ts";
import { readCursorAccountUsage } from "./cursorUsageReader.ts";
import { UsageAggregator } from "./usageAggregation.ts";
import { createOverrideRateTable, parseRateTable, type RateTable } from "./usagePricing.ts";
Expand Down Expand Up @@ -160,6 +162,7 @@ export const make = Effect.gen(function* () {
const platform = yield* HostProcessPlatform;

const fileCache: ScanCache = new Map();
const antigravityCache = makeAntigravityUsageCache();
const sourceCache = new Map<string, typeof CachedSource.Type>();
let cacheDirty = false;
const isWithinDirectory = (filePath: string, dir: string) => {
Expand Down Expand Up @@ -564,7 +567,7 @@ export const make = Effect.gen(function* () {
antigravityDirs.add(yield* fileSystem.realPath(dir).pipe(Effect.orElseSucceed(() => dir)));
}
const antigravity = yield* Effect.promise(() =>
readAntigravityUsage([...antigravityDirs], windowStartMs),
readAntigravityUsage([...antigravityDirs], windowStartMs, antigravityCache),
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
);
for (const dir of antigravityDirs) {
const exists = yield* fileSystem
Expand Down
36 changes: 34 additions & 2 deletions apps/server/src/usage/antigravityUsageReader.ts
Original file line number Diff line number Diff line change
Expand Up @@ -257,10 +257,28 @@ async function readDatabase(path: string, fallbackTimestamp: number): Promise<Us
}
}

/** Reads and merges aliases across every configured Antigravity store before date filtering. */
interface CachedDatabase {
readonly fingerprint: string;
readonly candidates: readonly UsageCandidate[];
}

/**
* Parsed databases keyed by canonical path. An entry is reused only while both
* the database and its `-wal` sidecar keep the same size, mtime and ctime: new
* rows land in the WAL without touching the main file until a checkpoint, and
* ctime moves on any content write even when mtime is restored.
*/
export const makeAntigravityUsageCache = () => new Map<string, CachedDatabase>();

/**
* Reads and merges aliases across every configured Antigravity store before date
* filtering. With a cache, databases unchanged since the previous read are not
* decoded again.
*/
export async function readAntigravityUsage(
conversationsDirectories: string | readonly string[],
sinceMs: number,
cache?: Map<string, CachedDatabase>,
) {
const roots =
typeof conversationsDirectories === "string"
Expand Down Expand Up @@ -351,7 +369,18 @@ export async function readAntigravityUsage(
if (visited.has(canonical)) continue;
visited.add(canonical);
const stat = await NodeFSP.stat(path);
const candidates = await readDatabase(path, stat.mtimeMs);
const wal = await NodeFSP.stat(`${path}-wal`).catch(() => null);
const fingerprint = [stat, wal]
.map((file) => (file ? `${file.size}:${file.mtimeMs}:${file.ctimeMs}` : "-"))
.join("/");
const cached = cache?.get(canonical);
let candidates: readonly UsageCandidate[];
if (cached?.fingerprint === fingerprint) {
candidates = cached.candidates;
} else {
candidates = await readDatabase(path, stat.mtimeMs);
cache?.set(canonical, { fingerprint, candidates });
}
const fileIndex = files.length;
files.push({ root, path, records: [] });
for (const [index, candidate] of candidates.entries()) {
Expand All @@ -365,6 +394,9 @@ export async function readAntigravityUsage(
}
};
for (const root of roots) await walk(root, root);
if (cache !== undefined) {
for (const key of cache.keys()) if (!visited.has(key)) cache.delete(key);
}
for (const [index, group] of groups.entries()) {
if (group.parent === index && group.record.timestampMs >= sinceMs) {
files[group.fileIndex]!.records.push(group.record);
Expand Down
101 changes: 100 additions & 1 deletion apps/server/src/usage/usageTranscriptReader.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ import { afterEach, assert, beforeEach, describe, it } from "@effect/vitest";
import { readTranscriptRecords } from "./usageTranscriptReader.ts";
import { readOpenCodeUsage } from "./opencodeUsageReader.ts";
import { readCursorAccountUsage } from "./cursorUsageReader.ts";
import { readAntigravityUsage } from "./antigravityUsageReader.ts";
import { makeAntigravityUsageCache, readAntigravityUsage } from "./antigravityUsageReader.ts";

let dir: string;

Expand Down Expand Up @@ -725,6 +725,105 @@ describe("SQLite usage readers", () => {
);
});

const antigravityGeneration = (responseId: string) =>
new Uint8Array(
protoBytes(1, [
...protoBytes(4, [
...protoNumber(2, 100),
...protoNumber(3, 40),
...protoText(11, responseId),
]),
...protoText(19, "Gemini 3 Pro"),
...protoBytes(9, protoBytes(4, protoNumber(1, 1780000000))),
]),
);

it("reuses an unchanged Antigravity database instead of decoding it again", async () => {
const path = NodePath.join(dir, "session-1.db");
const db = new NodeSqlite.DatabaseSync(path);
try {
db.exec("CREATE TABLE gen_metadata (idx INTEGER, data BLOB)");
db.prepare("INSERT INTO gen_metadata VALUES (?, ?)").run(0, antigravityGeneration("r-1"));
const cache = makeAntigravityUsageCache();
assert.deepStrictEqual((await readAntigravityUsage(dir, 0, cache)).errors, []);

// An exclusive lock makes a fresh read fail without touching the file, so
// only a cache hit can still return the earlier records.
db.exec("BEGIN EXCLUSIVE");
assert.strictEqual((await readAntigravityUsage(dir, 0)).errors.length, 1);
const cached = await readAntigravityUsage(dir, 0, cache);
assert.deepStrictEqual(cached.errors, []);
assert.deepStrictEqual(
cached.files.flatMap((file) => file.records).map((record) => record.dedupeKey),
["antigravity:11:r-1"],
);
db.exec("ROLLBACK");
} finally {
db.close();
}
});

it("rereads an Antigravity database rewritten with its size and mtime restored", async () => {
const path = NodePath.join(dir, "session-1.db");
const db = new NodeSqlite.DatabaseSync(path);
try {
db.exec("CREATE TABLE gen_metadata (idx INTEGER, data BLOB)");
db.prepare("INSERT INTO gen_metadata VALUES (?, ?)").run(0, antigravityGeneration("r-1"));
} finally {
db.close();
}
await NodeFSP.utimes(path, 1780000000, 1780000000);
const cache = makeAntigravityUsageCache();
assert.deepStrictEqual((await readAntigravityUsage(dir, 0, cache)).errors, []);

// ctime has the kernel's timestamp granularity, which can be a few
// milliseconds, so repeat the forged rewrite until it lands on a later tick
// than the cached read, as any real rewrite does.
const cached = await NodeFSP.stat(path);
do {
await NodeFSP.writeFile(path, Buffer.alloc(cached.size));
await NodeFSP.utimes(path, 1780000000, 1780000000);
} while ((await NodeFSP.stat(path)).ctimeMs === cached.ctimeMs);
const restored = await NodeFSP.stat(path);
assert.strictEqual(restored.size, cached.size);
assert.strictEqual(restored.mtimeMs, cached.mtimeMs);
const next = await readAntigravityUsage(dir, 0, cache);
assert.deepStrictEqual(next.errors, [path]);
assert.deepStrictEqual(
next.files.flatMap((file) => file.records),
[],
);
});

it("rereads an Antigravity database when only its WAL changed", async () => {
const path = NodePath.join(dir, "session-1.db");
const db = new NodeSqlite.DatabaseSync(path);
try {
db.exec(
"PRAGMA journal_mode = WAL; PRAGMA wal_autocheckpoint = 0; CREATE TABLE gen_metadata (idx INTEGER, data BLOB)",
);
const insert = db.prepare("INSERT INTO gen_metadata VALUES (?, ?)");
insert.run(0, antigravityGeneration("r-1"));
const cache = makeAntigravityUsageCache();
const first = await readAntigravityUsage(dir, 0, cache);
assert.strictEqual(first.files.flatMap((file) => file.records).length, 1);

const before = await NodeFSP.stat(path);
insert.run(1, antigravityGeneration("r-2"));
const after = await NodeFSP.stat(path);
assert.strictEqual(after.size, before.size);
assert.strictEqual(after.mtimeMs, before.mtimeMs);

const next = await readAntigravityUsage(dir, 0, cache);
assert.deepStrictEqual(
next.files.flatMap((file) => file.records).map((record) => record.dedupeKey),
["antigravity:11:r-1", "antigravity:11:r-2"],
);
} finally {
db.close();
}
});

it("reads Antigravity step-only stores and reports malformed databases", async () => {
const db = new NodeSqlite.DatabaseSync(NodePath.join(dir, "steps.db"));
try {
Expand Down
Loading