diff --git a/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.ts b/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.ts index 845fc74b9..951765967 100644 --- a/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.ts +++ b/apps/bulldozer-js/src/databases/low-level/implementations/lmdb.ts @@ -1,7 +1,7 @@ import { encodeBase64 } from "@hexclave/shared/dist/utils/bytes"; import { wait } from "@hexclave/shared/dist/utils/promises"; -import { traceSpan } from "../../../otel.js"; import * as lmdb from "lmdb"; +import { traceSpan } from "../../../otel.js"; import { DatabaseSeq } from "../../index.js"; import { LowLevelDatabase, LowLevelDatabaseDebugEntry, LowLevelKvDump, LowLevelKvStore } from "../index.js"; @@ -51,6 +51,33 @@ function validateValue(name: string, value: ArrayBuffer) { if (value.byteLength > 2_000_000_000) throw new Error(`KV store ${name} must be <= 2GB`); } +type LmdbActivityStats = { + puts: number, + transactions: number, + waitUntilAvailableResolves: number, + waitUntilDurableResolves: number, + waitUntilAvailableResolveTotalMs: number, + waitUntilDurableResolveTotalMs: number, +}; + +function emptyActivityStats(): LmdbActivityStats { + return { + puts: 0, + transactions: 0, + waitUntilAvailableResolves: 0, + waitUntilDurableResolves: 0, + waitUntilAvailableResolveTotalMs: 0, + waitUntilDurableResolveTotalMs: 0, + }; +} + +function hasActivity(stats: LmdbActivityStats): boolean { + return stats.puts > 0 + || stats.transactions > 0 + || stats.waitUntilAvailableResolves > 0 + || stats.waitUntilDurableResolves > 0; +} + export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: string, simulateReadMissDelayMs?: number }): LowLevelDatabase { const dbId = options.dbId ?? "default"; const simulateReadMissDelayMs = options.simulateReadMissDelayMs ?? 0; @@ -62,6 +89,33 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri const debugEntriesByStoreId = new Map<`${"store" | "dump"}-${string}`, () => Promise>(); const seqToAvailability = new Map>(); const seqToDurability = new Map>(); + let activityStats = emptyActivityStats(); + let activityWindowStartedAt = performance.now(); + const activityInterval = setInterval(() => { + if (!hasActivity(activityStats)) return; + const now = performance.now(); + const elapsedMs = now - activityWindowStartedAt; + const elapsedSeconds = elapsedMs / 1000; + console.debug("bulldozer-js low-level lmdb activity", { + dbId, + elapsedMs, + putsPerSecond: activityStats.puts / elapsedSeconds, + transactionsPerSecond: activityStats.transactions / elapsedSeconds, + waitUntilAvailableResolvesPerSecond: activityStats.waitUntilAvailableResolves / elapsedSeconds, + waitUntilDurableResolvesPerSecond: activityStats.waitUntilDurableResolves / elapsedSeconds, + averageSeqToAvailabilityResolveMs: activityStats.waitUntilAvailableResolves === 0 ? 0 : activityStats.waitUntilAvailableResolveTotalMs / activityStats.waitUntilAvailableResolves, + averageSeqToDurabilityResolveMs: activityStats.waitUntilDurableResolves === 0 ? 0 : activityStats.waitUntilDurableResolveTotalMs / activityStats.waitUntilDurableResolves, + mapSizes: { + seqToAvailability: seqToAvailability.size, + seqToDurability: seqToDurability.size, + debugEntriesByStoreId: debugEntriesByStoreId.size, + }, + currentVersion, + }); + activityStats = emptyActivityStats(); + activityWindowStartedAt = now; + }, 5_000); + activityInterval.unref(); const initialSeq = [dbId, initialSeqId] as unknown as LmdbSeq; const toSeq = (seqId: string) => [dbId, seqId] as unknown as LmdbSeq; @@ -79,14 +133,20 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri return seqToDurability.get(seqId) ?? Promise.resolve(); }; const rememberAvailability = (seqId: string, promise: Promise) => { + const insertedAt = performance.now(); const availability = traceSpan({ description: "bulldozer-js.low-level.lmdb.availability", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => await promise.then(() => { + activityStats.waitUntilAvailableResolveTotalMs += performance.now() - insertedAt; + activityStats.waitUntilAvailableResolves++; seqToAvailability.delete(seqId); })); availability.catch(() => {}); seqToAvailability.set(seqId, availability); }; const rememberDurability = (seqId: string, promise: Promise) => { + const insertedAt = performance.now(); const durability = traceSpan({ description: "bulldozer-js.low-level.lmdb.durability", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => await promise.then(async () => await root.flushed).then(() => { + activityStats.waitUntilDurableResolveTotalMs += performance.now() - insertedAt; + activityStats.waitUntilDurableResolves++; seqToDurability.delete(seqId); })); durability.catch(() => {}); @@ -101,6 +161,7 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri const version = nextVersion(); const seqId = nextSeqId(); const promise = traceSpan({ description: "bulldozer-js.low-level.lmdb.commit", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => await waitUntilAvailable(requiresSeq).then(async () => await root.transaction(() => { + activityStats.transactions++; return (async () => { await action(version); await meta.put("seq", version); @@ -131,11 +192,15 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri return key.buffer; }; const putWithVersion = async (db: BinaryDatabase, key: Buffer, value: Buffer, version: number) => { + activityStats.puts++; await db.put(key, value, version); }; const waitUntilAvailable = async (seq: DatabaseSeq) => { await getAvailabilityPromise(getSeqId(seq)); }; + const waitUntilDurable = async (seq: DatabaseSeq) => { + await getDurabilityPromise(getSeqId(seq)); + }; const waitUntilAllAvailable = async () => { await Promise.all(seqToAvailability.values()); }; @@ -197,6 +262,7 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri const seqId = nextSeqId(); const keys = values.map((_, index) => dumpKeyForVersion(version, index)); const promise = traceSpan({ description: "bulldozer-js.low-level.lmdb.insertAll.commit", attributes }, async () => await waitUntilAvailable(insertOptions?.requiresSeq ?? initialSeq).then(async () => await root.transaction(() => { + activityStats.transactions++; return (async () => { for (let i = 0; i < values.length; i++) { await putWithVersion(db, bufferFromArrayBuffer(keys[i]), bufferFromArrayBuffer(values[i]), version); @@ -268,7 +334,7 @@ export function declareLmdbLowLevelDatabase(options: { path: string, dbId?: stri await traceSpan({ description: "bulldozer-js.low-level.lmdb.waitUntilAvailable", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => await waitUntilAvailable(seq)); }, async waitUntilDurable(seq) { - await traceSpan({ description: "bulldozer-js.low-level.lmdb.waitUntilDurable", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => await getDurabilityPromise(getSeqId(seq))); + await traceSpan({ description: "bulldozer-js.low-level.lmdb.waitUntilDurable", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => await waitUntilDurable(seq)); }, async waitUntilReplicated(seq) { await traceSpan({ description: "bulldozer-js.low-level.lmdb.waitUntilReplicated", attributes: { "bulldozer.low_level.backend": "lmdb" } }, async () => {