// Server-side on-disk store for the social-post corpus. // // Layout, per channel: // transcripts/channels//posts/YYYY-MM.jsonl — one post per line // transcripts/channels//posts-archive — " " lines // // Month-sharded JSONL rather than a directory per post: an account with 50k // posts would otherwise create 50k directories. The archive file mirrors how // yt-dlp's `--download-archive` drives incremental sync today (sync() in // common/ytdlp/runYtdlp.ts stops at the first already-archived entry; the post // fetchers stop the same way), and uses the same " " line format // so common/lib/archive.ts parses it unchanged. import { writeJsonAtomic } from "./jsonFile-server"; import path from "node:path"; import { readdir, mkdir, readFile } from "node:fs/promises"; import { appendFile } from "node:fs/promises"; import { readArchive } from "./archive"; import { comparePostsNewestFirst, isPostAvailability, monthShardFromCreatedAt, parsePost, postSlug, type Post, type PostAvailability, type PostAvailabilityMap, type PostAvailabilityRecord, type PostPlatform, } from "./posts"; export const POSTS_DIRNAME = "posts"; export const POSTS_ARCHIVE_FILENAME = "posts-archive"; export function channelPostsDir(channelRoot: string): string { return path.join(channelRoot, POSTS_DIRNAME); } export function channelPostsArchivePath(channelRoot: string): string { return path.join(channelRoot, POSTS_ARCHIVE_FILENAME); } function shardPath(channelRoot: string, shard: string): string { return path.join(channelPostsDir(channelRoot), `${shard}.jsonl`); } // The set of post ids already archived for this channel, used by every fetcher // to stop paging once it reaches previously-synced content. export async function readSeenPostIds( channelRoot: string, ): Promise> { const archive = await readArchive(channelPostsArchivePath(channelRoot)); return archive.ids; } // Append ids to the archive. Idempotent at the caller's discretion — writePosts // filters against the seen set before calling, so duplicate lines only appear // if two fetches race, and readArchive dedupes on read anyway. async function appendArchiveIds( channelRoot: string, platform: PostPlatform, ids: ReadonlyArray, ): Promise { if (ids.length === 0) return; const lines = ids.map((id) => `${platform} ${id}\n`).join(""); await appendFile(channelPostsArchivePath(channelRoot), lines, "utf8"); } export type WritePostsResult = { written: number; skipped: number; // already present in the archive shards: string[]; // YYYY-MM shards touched }; // Append posts to their month shards and record their ids in the archive. // Posts whose id is already archived are skipped, so a re-run over overlapping // pages is a no-op rather than a source of duplicates. export async function writePosts( channelRoot: string, posts: ReadonlyArray, ): Promise { if (posts.length === 0) return { written: 0, skipped: 0, shards: [] }; await mkdir(channelPostsDir(channelRoot), { recursive: true }); const seen = await readSeenPostIds(channelRoot); // Group by shard so each file is opened once, and dedupe within the batch // itself (a fetcher may legitimately return the same post twice across a // cursor boundary). const byShard = new Map(); const accepted: Post[] = []; let skipped = 0; for (const post of posts) { if (seen.has(post.id)) { skipped++; continue; } seen.add(post.id); accepted.push(post); const shard = monthShardFromCreatedAt(post.createdAt); const bucket = byShard.get(shard); if (bucket) bucket.push(post); else byShard.set(shard, [post]); } if (accepted.length === 0) { return { written: 0, skipped, shards: [] }; } for (const [shard, shardPosts] of byShard) { const body = shardPosts.map((p) => JSON.stringify(p)).join("\n") + "\n"; await appendFile(shardPath(channelRoot, shard), body, "utf8"); } // Archive AFTER the shard writes: a crash between the two re-fetches the // posts (harmless — the next writePosts dedupes them), whereas archiving // first would lose them permanently. const byPlatform = new Map(); for (const post of accepted) { const bucket = byPlatform.get(post.platform); if (bucket) bucket.push(post.id); else byPlatform.set(post.platform, [post.id]); } for (const [platform, ids] of byPlatform) { await appendArchiveIds(channelRoot, platform, ids); } return { written: accepted.length, skipped, shards: [...byShard.keys()].sort(), }; } export type UpsertPostsResult = WritePostsResult & { // Archived posts whose stored record was superseded by a newer one. updated: number; // Archived posts that arrived again unchanged (or older than what is // stored). unchanged: number; }; // What makes two records of one post differ, for an update. function postContentKey(p: Post): string { return JSON.stringify([ p.text, p.author, p.links, p.media ?? null, p.forum?.editedAt ?? null, p.forum?.position ?? null, p.forum?.threadTitle ?? null, ]); } // Is `incoming` a newer record of the same post than `stored`? A later edit // time wins; an older one never replaces a newer one; with no edit times to // compare (or equal ones), different content is taken as the later read. export function supersedesPost(stored: Post, incoming: Post): boolean { if (postContentKey(stored) === postContentKey(incoming)) return false; const a = stored.forum?.editedAt; const b = incoming.forum?.editedAt; if (a && b && a !== b) return b > a; if (a && !b) return false; return true; } // writePosts, plus UPDATES: a post already archived whose incoming record is // newer (supersedesPost — an edit made since it was archived) is appended // again to its month shard. The JSONL stays append-only; every reader takes the // LAST record of an id (readAllPosts, the index build), so the newer one wins. // For sources that re-read posts they already hold (a forum page saved twice). export async function upsertPosts( channelRoot: string, posts: ReadonlyArray, ): Promise { // One record per id from the batch: the newest of them. const batch = new Map(); for (const p of posts) { const prev = batch.get(p.id); if (!prev || supersedesPost(prev, p)) batch.set(p.id, p); } const seen = await readSeenPostIds(channelRoot); const fresh = [...batch.values()].filter((p) => !seen.has(p.id)); const known = [...batch.values()].filter((p) => seen.has(p.id)); const written = await writePosts(channelRoot, fresh); let updated = 0; let unchanged = 0; const shards = new Set(written.shards); if (known.length > 0) { const stored = new Map(); for (const p of await readAllPosts(channelRoot)) stored.set(p.id, p); const byShard = new Map(); for (const p of known) { const old = stored.get(p.id); if (old && !supersedesPost(old, p)) { unchanged++; continue; } // Keep the shard the post already lives in: its createdAt does not // change with an edit. const shard = monthShardFromCreatedAt(old?.createdAt ?? p.createdAt); const bucket = byShard.get(shard); const rec = old ? { ...p, createdAt: old.createdAt, uploadDate: old.uploadDate } : p; if (bucket) bucket.push(rec); else byShard.set(shard, [rec]); updated++; } if (byShard.size > 0) await mkdir(channelPostsDir(channelRoot), { recursive: true }); for (const [shard, shardPosts] of byShard) { const body = shardPosts.map((p) => JSON.stringify(p)).join("\n") + "\n"; await appendFile(shardPath(channelRoot, shard), body, "utf8"); shards.add(shard); } } return { written: written.written, skipped: written.skipped, shards: [...shards].sort(), updated, unchanged, }; } // List the channel's month shards, oldest first. export async function listPostShards( channelRoot: string, ): Promise { let entries: string[]; try { entries = await readdir(channelPostsDir(channelRoot)); } catch { return []; } return entries .filter((e) => e.endsWith(".jsonl")) .map((e) => e.slice(0, -".jsonl".length)) .sort(); } // Read one month shard. Malformed lines are skipped rather than failing the // whole shard — an append that was interrupted mid-line must not make every // earlier post in that month unreadable. // // THE DIRECTORY IS THE CHANNEL. Each line carries `channelSlug` and `slug` // ("/") stamped at fetch time, and renameChannel moves the // directory without rewriting a line — nor would a re-fetch repair it: dedupe // is by id, so an unchanged post is never written again. Trusting the stamp // left every post of a renamed channel naming its OLD slug: the index keyed // the page by the new one while the record inside said the old, so postsCache // could not find a post on its own page, PostModal and threads threw, and MCP // links named a channel that no longer exists. Every reader comes through here // with `channels/` as the root, so both fields are re-derived from it. export async function readPostShard( channelRoot: string, shard: string, ): Promise { let text: string; try { text = await readFile(shardPath(channelRoot, shard), "utf8"); } catch { return []; } const channelSlug = path.basename(channelRoot); const posts: Post[] = []; for (const line of text.split("\n")) { const trimmed = line.trim(); if (!trimmed) continue; let raw: unknown; try { raw = JSON.parse(trimmed); } catch { continue; } const post = parsePost(raw); if (!post) continue; if (post.channelSlug !== channelSlug) { post.channelSlug = channelSlug; post.slug = postSlug(channelSlug, post.id); } posts.push(post); } return posts; } // Every archived post for a channel, newest first. Deduped by id: an // interrupted archive write can leave the same post in two shards. export async function readAllPosts(channelRoot: string): Promise { const shards = await listPostShards(channelRoot); const byId = new Map(); for (const shard of shards) { for (const post of await readPostShard(channelRoot, shard)) { byId.set(post.id, post); } } return [...byId.values()].sort(comparePostsNewestFirst); } // Cheap count for the dashboard/snapshot without materializing every post. export async function countPosts(channelRoot: string): Promise { const archive = await readArchive(channelPostsArchivePath(channelRoot)); return archive.ids.size; } // The newest archived post's createdAt, used as the incremental-fetch // watermark. Only the newest shard is read. export async function latestPostCreatedAt( channelRoot: string, ): Promise { const shards = await listPostShards(channelRoot); for (let i = shards.length - 1; i >= 0; i--) { const posts = await readPostShard(channelRoot, shards[i]); if (posts.length === 0) continue; let newest = posts[0].createdAt; for (const post of posts) { if (post.createdAt > newest) newest = post.createdAt; } return newest; } return null; } // The oldest archived post's createdAt: where the older-posts backfill starts. // Only the oldest non-empty shard is read (listPostShards sorts oldest first). export async function oldestPostCreatedAt( channelRoot: string, ): Promise { const shards = await listPostShards(channelRoot); for (const shard of shards) { const posts = await readPostShard(channelRoot, shard); if (posts.length === 0) continue; let oldest = posts[0].createdAt; for (const post of posts) { if (post.createdAt < oldest) oldest = post.createdAt; } return oldest; } return null; } // --------------------------------------------------------------------------- // Availability sidecar ("is this post still up?") // --------------------------------------------------------------------------- // // One file per channel rather than per post: posts live in month-sharded JSONL // with no per-post directory, so the per-video availability.json layout does // not transfer. The JSONL itself stays append-only and is never rewritten — // the sidecar is the source of truth and is merged into the served page tree // at index time, exactly as video availability is merged into a summary. export const POSTS_AVAILABILITY_FILENAME = "posts-availability.json"; export function postsAvailabilityPath(channelRoot: string): string { return path.join(channelRoot, POSTS_AVAILABILITY_FILENAME); } export async function readPostAvailability( channelRoot: string, ): Promise { try { const raw = await readFile(postsAvailabilityPath(channelRoot), "utf8"); const parsed = JSON.parse(raw) as unknown; if (!parsed || typeof parsed !== "object") return {}; const out: PostAvailabilityMap = {}; for (const [id, rec] of Object.entries(parsed as Record)) { const r = rec as Partial; if (!isPostAvailability(r?.availability)) continue; if (typeof r.checkedAt !== "string") continue; out[id] = { availability: r.availability, checkedAt: r.checkedAt, ...(Array.isArray(r.history) ? { history: r.history } : {}), }; } return out; } catch { return {}; } } export async function writePostAvailability( channelRoot: string, map: PostAvailabilityMap, ): Promise { const file = postsAvailabilityPath(channelRoot); await mkdir(path.dirname(file), { recursive: true }); await writeJsonAtomic(file, map); } // Fold new observations in, appending to history ONLY when the availability // actually changes. That keeps the file small and, more importantly, preserves // the moment a post was first seen deleted — which is the fact worth keeping. export function mergePostAvailability( existing: PostAvailabilityMap, observations: ReadonlyMap, checkedAt: string, ): { map: PostAvailabilityMap; newlyDeleted: string[]; changed: number } { const map: PostAvailabilityMap = { ...existing }; const newlyDeleted: string[] = []; let changed = 0; for (const [id, availability] of observations) { const prev = map[id]; if (prev?.availability === availability) { // Same verdict — just refresh when we last looked. map[id] = { ...prev, checkedAt }; continue; } changed++; if (availability === "deleted" && prev?.availability !== "deleted") { newlyDeleted.push(id); } map[id] = { availability, checkedAt, history: [ ...(prev?.history ?? []), ...(prev ? [{ availability: prev.availability, at: prev.checkedAt }] : []), ], }; } return { map, newlyDeleted, changed }; } // --------------------------------------------------------------------------- // Fetch state sidecar // --------------------------------------------------------------------------- // Small per-channel record of the last fetch attempt. Drives the social // channel's snapshot buckets (fetchFailed / needsCookies) without re-deriving // them from job logs. export type PostFetchState = { lastFetchedAt?: string; lastError?: string; // Set when the fetcher stopped because credentials are missing/expired — // the post-corpus analogue of the video pipeline's needs_auth outcome. needsCookies?: boolean; lastFetchedCount?: number; // The TIMELINE walk's resume point (gallery-dl's cursor). Only a normal fetch // reads it; the older-posts backfill keeps its own position in `older`. cursor?: string; // The older-posts backfill (a fetcher's `fetchOlder`), which walks the // account's history backwards through search windows. Kept apart from // `cursor` so a normal fetch never resumes backwards, and carried through // every write a normal fetch makes. older?: OlderBackfillState; }; // Where the older-posts backfill stands. The window is in UTC days, as X's // `since:`/`until:` search operators take them: posts created on or after // `since` and before `until`. export type OlderBackfillPosition = { since: string; // YYYY-MM-DD until: string; // YYYY-MM-DD // Resume INSIDE the window: only posts with an id at or below this one (the // oldest post the walk has read in this window). maxId?: string; // Windows in a row, ending with the last one walked, that held no posts. emptyWindows: number; }; export type OlderBackfillState = OlderBackfillPosition & { // The date the walk stops at (YYYY-MM-DD), when one was given. floor?: string; // The account's creation time, once a fetched post has carried it: the walk // never goes below it. accountCreatedAt?: string; complete?: boolean; completedAt?: string; // Why the walk ended, as a sentence ("reached the floor date 2020-01-01"). completeReason?: string; lastRunAt?: string; lastWritten?: number; lastError?: string; }; export const POST_FETCH_STATE_FILENAME = "posts-state.json"; export function postFetchStatePath(channelRoot: string): string { return path.join(channelRoot, POST_FETCH_STATE_FILENAME); } export async function readPostFetchState( channelRoot: string, ): Promise { try { const raw = await readFile(postFetchStatePath(channelRoot), "utf8"); const parsed = JSON.parse(raw) as unknown; if (!parsed || typeof parsed !== "object") return null; return parsed as PostFetchState; } catch { return null; } } export async function writePostFetchState( channelRoot: string, state: PostFetchState, ): Promise { const file = postFetchStatePath(channelRoot); await mkdir(path.dirname(file), { recursive: true }); await writeJsonAtomic(file, state); }