commit 3dc9b95c1387303ebce7294a5e4d1d78429e516a
parent fa48d29d2e54a8a6e14254152e4d580ed2b6f6d5
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Wed, 20 May 2026 20:54:20 -0400
Update buildIndex to stream and avoid OOM
Diffstat:
3 files changed, 202 insertions(+), 133 deletions(-)
diff --git a/common/controller/buildIndex.ts b/common/controller/buildIndex.ts
@@ -20,7 +20,9 @@ import {
stat,
writeFile,
} from "node:fs/promises";
-import type { Dirent } from "node:fs";
+import { createWriteStream } from "node:fs";
+import type { Dirent, WriteStream } from "node:fs";
+import { once } from "node:events";
import { open } from "lmdb";
import { parseVtt, type Cue } from "../lib/vtt";
import { parseWhisper } from "../lib/whisper";
@@ -78,8 +80,9 @@ import {
type SubTrack,
} from "../lib/videoStatus";
import { loadAvailability } from "../lib/availability-server";
+import { AVAILABILITY_FILENAME } from "../lib/availability";
-const SCHEMA_VERSION = 7;
+const SCHEMA_VERSION = 8;
type IndexKey = [string, string, string];
type ChannelKey = [string, string, string];
@@ -90,6 +93,8 @@ type MtimeRecord = {
metaMs: number;
transcriptMs: number | null;
subsMs: number | null;
+ availabilityMs: number | null;
+ isDeleted: boolean;
indexKey: IndexKey;
};
@@ -115,6 +120,7 @@ type LiveEntry = {
transcriptKind: IndexTranscript["kind"] | null;
subTracks: SubTrack[];
subsMs: number | null;
+ availabilityMs: number | null;
};
async function exists(p: string): Promise<boolean> {
@@ -210,6 +216,14 @@ async function scanSource(
// ignore
}
}
+ let availabilityMs: number | null = null;
+ try {
+ availabilityMs = (
+ await stat(path.join(fullVideoDir, AVAILABILITY_FILENAME))
+ ).mtimeMs;
+ } catch {
+ availabilityMs = null;
+ }
live.push({
channelSlug: ch.name,
handling: cfg.handling,
@@ -222,6 +236,7 @@ async function scanSource(
transcriptKind: picked?.kind ?? null,
subTracks,
subsMs,
+ availabilityMs,
});
}
}
@@ -358,7 +373,8 @@ export async function buildIndex({
} else if (
prev.metaMs !== s.metaMs ||
prev.transcriptMs !== s.transcriptMs ||
- (prev.subsMs ?? null) !== s.subsMs
+ (prev.subsMs ?? null) !== s.subsMs ||
+ (prev.availabilityMs ?? null) !== s.availabilityMs
) {
changed.push(s);
}
@@ -637,11 +653,19 @@ export async function buildIndex({
if (parsedSubs.length > 0) subs.put(indexKey, parsedSubs);
else subs.remove(indexKey);
+ let isDeleted = false;
+ if (s.availabilityMs !== null) {
+ const availability = await loadAvailability(videoFullDir);
+ isDeleted = availability?.availability === "deleted";
+ }
+
byChannel.put(indexToChannelKey(indexKey), 1);
mtimes.put(pk, {
metaMs: s.metaMs,
transcriptMs: s.transcriptMs,
subsMs: s.subsMs,
+ availabilityMs: s.availabilityMs,
+ isDeleted,
indexKey,
});
} catch (err) {
@@ -675,45 +699,153 @@ export async function buildIndex({
let pagesSkipped = 0;
let pagesDeleted = 0;
- type PendingEntry = { encoded: string; id: string };
-
- const writePage = async (
- channelSlug: string,
- idx: number,
- entries: PendingEntry[],
- ): Promise<void> => {
- const body = `[${entries.map((e) => e.encoded).join(",")}]`;
- const hash = createHash("sha1").update(body).digest("hex");
- const prev = pageHashes.get([channelSlug, idx]);
- const outPath = path.join(
- transcriptsOutDir,
- channelSlug,
- transcriptPageFileName(idx),
- );
- if (prev?.hash === hash && (await exists(outPath))) {
- pagesSkipped++;
- return;
- }
- const sizeBytes = Buffer.byteLength(body, "utf8");
- const tmp = `${outPath}.tmp-${process.pid}`;
- await writeFile(tmp, body);
- await rename(tmp, outPath);
- pageHashes.put([channelSlug, idx], {
- hash,
- entryCount: entries.length,
- sizeBytes,
- });
- pagesWritten++;
+ // Streaming page writer: writes each entry directly to the tmp output as it
+ // arrives instead of buffering the full page in memory. Hash is computed
+ // incrementally; on close, the hash is compared to the previously stored one
+ // — equal means we delete the tmp and skip, different means we atomically
+ // rename in place. Avoids the `[...].map(e => e.encoded).join(",")` pattern
+ // that briefly held the entire concatenated page body in V8, which OOMed
+ // when a single live_chat track produced a multi-hundred-MB encoded entry.
+ type PageWriterOpts = {
+ outDir: string;
+ fileName: (idx: number) => string;
+ maxPageBytes: number;
+ ensureDir: boolean;
+ getPrevHash: (idx: number) => string | undefined;
+ setHash: (idx: number, record: PageHashRecord) => void;
+ onLog: (msg: string) => void;
+ };
+ type PageWriterResult = {
+ pageCount: number;
+ pagesWritten: number;
+ pagesSkipped: number;
+ slugToPage: Record<string, number>;
+ };
+
+ const createPageWriter = (opts: PageWriterOpts) => {
+ let pageIdx = 0;
+ let stream: WriteStream | null = null;
+ let tmpPath = "";
+ let outPath = "";
+ let hash = createHash("sha1");
+ let payloadBytes = 0;
+ let totalBytes = 0;
+ let entryCount = 0;
+ let firstInPage = true;
+ let pagesWrittenLocal = 0;
+ let pagesSkippedLocal = 0;
+ let dirEnsured = !opts.ensureDir;
+ const slugToPage: Record<string, number> = {};
+
+ const writeChunk = async (chunk: string): Promise<void> => {
+ const s = stream;
+ if (!s) throw new Error("page writer: stream not open");
+ hash.update(chunk);
+ totalBytes += Buffer.byteLength(chunk, "utf8");
+ if (!s.write(chunk)) {
+ await once(s, "drain");
+ }
+ };
+
+ const openPage = async (): Promise<void> => {
+ if (!dirEnsured) {
+ await mkdir(opts.outDir, { recursive: true });
+ dirEnsured = true;
+ }
+ outPath = path.join(opts.outDir, opts.fileName(pageIdx));
+ tmpPath = `${outPath}.tmp-${process.pid}`;
+ stream = createWriteStream(tmpPath);
+ hash = createHash("sha1");
+ payloadBytes = 0;
+ totalBytes = 0;
+ entryCount = 0;
+ firstInPage = true;
+ await writeChunk("[");
+ };
+
+ const closePage = async (): Promise<void> => {
+ const s = stream;
+ if (!s) return;
+ await writeChunk("]");
+ await new Promise<void>((resolve, reject) => {
+ s.end((err?: NodeJS.ErrnoException | null) => {
+ if (err) reject(err);
+ else resolve();
+ });
+ });
+ const digest = hash.digest("hex");
+ const prev = opts.getPrevHash(pageIdx);
+ if (prev === digest && (await exists(outPath))) {
+ await rm(tmpPath, { force: true });
+ pagesSkippedLocal++;
+ } else {
+ await rename(tmpPath, outPath);
+ opts.setHash(pageIdx, {
+ hash: digest,
+ entryCount,
+ sizeBytes: totalBytes,
+ });
+ pagesWrittenLocal++;
+ }
+ stream = null;
+ };
+
+ const push = async (encoded: string, id: string): Promise<void> => {
+ const entryBytes = Buffer.byteLength(encoded, "utf8");
+ const commaBytes = firstInPage ? 0 : 1;
+ const delta = entryBytes + commaBytes;
+
+ if (stream && !firstInPage && payloadBytes + delta + 2 > opts.maxPageBytes) {
+ await closePage();
+ pageIdx++;
+ }
+ if (!stream) await openPage();
+
+ if (entryBytes + 2 > opts.maxPageBytes) {
+ opts.onLog(
+ ` ${id}: entry (${entryBytes} bytes) exceeds page budget; emitting solo page.`,
+ );
+ }
+
+ slugToPage[id] = pageIdx;
+ if (!firstInPage) await writeChunk(",");
+ await writeChunk(encoded);
+ firstInPage = false;
+ payloadBytes += delta;
+ entryCount++;
+ };
+
+ const finish = async (): Promise<PageWriterResult> => {
+ if (stream) {
+ await closePage();
+ pageIdx++;
+ }
+ return {
+ pageCount: pageIdx,
+ pagesWritten: pagesWrittenLocal,
+ pagesSkipped: pagesSkippedLocal,
+ slugToPage,
+ };
+ };
+
+ return { push, finish };
};
for (const channelSlug of Array.from(channelConfigs.keys()).sort()) {
const channelDir = path.join(transcriptsOutDir, channelSlug);
await mkdir(channelDir, { recursive: true });
- let pageIdx = 0;
- let buffer: PendingEntry[] = [];
- let payloadBytes = 0;
- const slugToPage: Record<string, number> = {};
+ const writer = createPageWriter({
+ outDir: channelDir,
+ fileName: transcriptPageFileName,
+ maxPageBytes: maxTranscriptPageBytes,
+ ensureDir: false,
+ getPrevHash: (idx) => pageHashes.get([channelSlug, idx])?.hash,
+ setHash: (idx, record) => {
+ pageHashes.put([channelSlug, idx], record);
+ },
+ onLog: log,
+ });
for (const { key } of byChannel.getRange({
start: [channelSlug],
@@ -727,34 +859,13 @@ export async function buildIndex({
const cueList = cues.get(indexKey);
const detail: TranscriptDetail = { ...summary, cues: cueList };
const encoded = JSON.stringify(detail);
- const entryBytes = Buffer.byteLength(encoded, "utf8");
- const commaBytes = buffer.length === 0 ? 0 : 1;
- const delta = entryBytes + commaBytes;
-
- if (buffer.length > 0 && payloadBytes + delta + 2 > maxTranscriptPageBytes) {
- await writePage(channelSlug, pageIdx, buffer);
- pageIdx++;
- buffer = [];
- payloadBytes = 0;
- }
-
- if (entryBytes + 2 > maxTranscriptPageBytes) {
- log(
- ` ${summary.slug}: entry (${entryBytes} bytes) exceeds page budget; emitting solo page.`,
- );
- }
-
- slugToPage[summary.id] = pageIdx;
- buffer.push({ encoded, id: summary.id });
- payloadBytes += delta;
+ await writer.push(encoded, summary.id);
}
- if (buffer.length > 0) {
- await writePage(channelSlug, pageIdx, buffer);
- pageIdx++;
- }
-
- const pageCount = pageIdx;
+ const { pageCount, pagesWritten: chPagesWritten, pagesSkipped: chPagesSkipped, slugToPage } =
+ await writer.finish();
+ pagesWritten += chPagesWritten;
+ pagesSkipped += chPagesSkipped;
const keep = new Set<string>(["manifest.json"]);
for (let i = 0; i < pageCount; i++) keep.add(transcriptPageFileName(i));
@@ -807,18 +918,15 @@ export async function buildIndex({
`Transcript pages: ${pagesWritten} written, ${pagesSkipped} unchanged, ${pagesDeleted} stale removed.`,
);
- // Build a lookup from indexKey → isDeleted by reading each video's
- // availability.json sidecar. Keyed off mtimes, which holds the canonical
- // (channelSlug, videoDir) → indexKey mapping; only "deleted" entries get
- // recorded so lookup defaults to "not deleted". Built once and reused by
- // both the sub-page loop and the cross-channel summaries loop.
+ // Build a lookup from indexKey → isDeleted from MtimeRecord, which caches the
+ // availability check done at mutation-processing time. Previously this loop
+ // read availability.json for every video on every build (~28k sequential
+ // disk reads), which stalled GC right before the heaviest phase.
const deletedByIndexKey = new Set<string>();
- for (const { key, value } of mtimes.getRange()) {
- const pk = key as PathKey;
- const ik = (value as MtimeRecord).indexKey;
- const videoDir = path.join(channelsDir, pk[0], "data", pk[1]);
- const availability = await loadAvailability(videoDir);
- if (availability?.availability === "deleted") {
+ for (const { value } of mtimes.getRange()) {
+ const rec = value as MtimeRecord;
+ if (rec.isDeleted) {
+ const ik = rec.indexKey;
deletedByIndexKey.add(`${ik[0]}\x00${ik[1]}\x00${ik[2]}`);
}
}
@@ -830,42 +938,25 @@ export async function buildIndex({
let subsTotalCount = 0;
let subsLiveChatTotalCount = 0;
- const writeSubsPage = async (
- channelSlug: string,
- idx: number,
- entries: PendingEntry[],
- ): Promise<void> => {
- const body = `[${entries.map((e) => e.encoded).join(",")}]`;
- const hash = createHash("sha1").update(body).digest("hex");
- const prev = subPageHashes.get([channelSlug, idx]);
- const outPath = path.join(subsOutDir, channelSlug, subsPageFileName(idx));
- if (prev?.hash === hash && (await exists(outPath))) {
- subsPagesSkipped++;
- return;
- }
- const sizeBytes = Buffer.byteLength(body, "utf8");
- const tmp = `${outPath}.tmp-${process.pid}`;
- await writeFile(tmp, body);
- await rename(tmp, outPath);
- subPageHashes.put([channelSlug, idx], {
- hash,
- entryCount: entries.length,
- sizeBytes,
- });
- subsPagesWritten++;
- };
-
for (const channelSlug of Array.from(channelConfigs.keys()).sort()) {
const cfg = channelConfigs.get(channelSlug)!;
const subsChannelDir = path.join(subsOutDir, channelSlug);
- let subPageIdx = 0;
- let subBuffer: PendingEntry[] = [];
- let subPayloadBytes = 0;
- const subSlugToPage: Record<string, number> = {};
const tracksInChannel = new Set<string>();
let videoCount = 0;
let liveChatCount = 0;
+ const subWriter = createPageWriter({
+ outDir: subsChannelDir,
+ fileName: subsPageFileName,
+ maxPageBytes: maxTranscriptPageBytes,
+ ensureDir: true,
+ getPrevHash: (idx) => subPageHashes.get([channelSlug, idx])?.hash,
+ setHash: (idx, record) => {
+ subPageHashes.put([channelSlug, idx], record);
+ },
+ onLog: log,
+ });
+
for (const { key } of byChannel.getRange({
start: [channelSlug],
end: [channelSlug, ""],
@@ -891,38 +982,14 @@ export async function buildIndex({
if (hasLiveChat) liveChatCount++;
const detail: SubsDetail = { ...display, tracks };
const encoded = JSON.stringify(detail);
- const entryBytes = Buffer.byteLength(encoded, "utf8");
- const commaBytes = subBuffer.length === 0 ? 0 : 1;
- const delta = entryBytes + commaBytes;
-
- if (
- subBuffer.length > 0 &&
- subPayloadBytes + delta + 2 > maxTranscriptPageBytes
- ) {
- await mkdir(subsChannelDir, { recursive: true });
- await writeSubsPage(channelSlug, subPageIdx, subBuffer);
- subPageIdx++;
- subBuffer = [];
- subPayloadBytes = 0;
- }
-
- if (entryBytes + 2 > maxTranscriptPageBytes) {
- log(
- ` ${summary.slug}: sub entry (${entryBytes} bytes) exceeds page budget; emitting solo page.`,
- );
- }
-
- subSlugToPage[summary.id] = subPageIdx;
- subBuffer.push({ encoded, id: summary.id });
- subPayloadBytes += delta;
+ await subWriter.push(encoded, summary.id);
videoCount++;
}
- if (subBuffer.length > 0) {
- await mkdir(subsChannelDir, { recursive: true });
- await writeSubsPage(channelSlug, subPageIdx, subBuffer);
- subPageIdx++;
- }
+ const { pageCount: subPageCount, pagesWritten: chSubsWritten, pagesSkipped: chSubsSkipped, slugToPage: subSlugToPage } =
+ await subWriter.finish();
+ subsPagesWritten += chSubsWritten;
+ subsPagesSkipped += chSubsSkipped;
if (videoCount === 0) {
// No subs for this channel — remove any stale dir + hashes.
@@ -936,7 +1003,6 @@ export async function buildIndex({
continue;
}
- const subPageCount = subPageIdx;
const subKeep = new Set<string>(["manifest.json"]);
for (let i = 0; i < subPageCount; i++) subKeep.add(subsPageFileName(i));
const subExisting = await readdir(subsChannelDir).catch(
diff --git a/export/CHANGELOG.md b/export/CHANGELOG.md
@@ -6,3 +6,6 @@
- Composable layered search. The single search bar is now a query builder: any number of layers can be combined with AND / OR / NOT and arbitrary nesting. Each layer targets a scope (transcripts, live chat, or title/channel), and matches from every contributing layer are surfaced in the result list with a per-layer colour swatch. Per-layer results are memoised in IndexedDB so editing a deeper leaf only re-runs that layer against its already-narrowed scope. Simple one-keyword search still looks like a single input — the builder collapses to compact mode when there's only one layer. The composite query serialises into a new `qt=` URL parameter; legacy `?q=&m=&re=` links auto-migrate to a one-layer tree.
- New `/changelog` page that renders the export's `CHANGELOG.md`. Linked from the right side of the sticky header.
- Per-heading copy-link buttons on the changelog page for permalinks to any section.
+
+### Fixed
+- `build:index` no longer runs out of memory on large datasets. Both the transcripts and the subs page writers now stream each entry directly to disk and hash it incrementally instead of materialising the joined page body in memory, which previously OOMed when a single live-chat track encoded to hundreds of MB. The post-processing "is this video deleted?" pass is now an in-memory LMDB scan (the availability check is cached in the mtime record at mutation time, schema bumped to 8 to invalidate the old cache) instead of ~28k sequential `availability.json` reads on every build. The `build:index` script also pre-sets `--max-old-space-size=8192` as a backstop.
diff --git a/export/package.json b/export/package.json
@@ -5,7 +5,7 @@
"type": "module",
"scripts": {
"dev": "next dev",
- "build:index": "tsx ../common/bin/build-index.ts",
+ "build:index": "NODE_OPTIONS=--max-old-space-size=8192 tsx ../common/bin/build-index.ts",
"prebuild": "pnpm run build:index",
"build": "next build",
"start": "serve out",