From ea56d286b30d340e439dbebce6e87701b836a1aa Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 30 Sep 2026 07:18:52 +0000 Subject: [PATCH] Hash every validated server copy and register it as verified content --- Dockerfile | 8 +++ plans/INDEX.md | 2 +- .../009-server-content-hash-186e7f.md | 6 ++ server/hash.js | 26 +++++++++ server/media-cache.js | 55 ++++++++++++++++++- server/media-cache.test.js | 33 +++++++++++ server/p2p-admit.js | 53 ++++++++++++++++++ server/p2p-admit.test.js | 52 ++++++++++++++++++ server/package.json | 2 +- server/server.js | 15 ++++- 10 files changed, 245 insertions(+), 7 deletions(-) rename plans/{active => done}/009-server-content-hash-186e7f.md (93%) create mode 100644 server/hash.js create mode 100644 server/p2p-admit.js create mode 100644 server/p2p-admit.test.js diff --git a/Dockerfile b/Dockerfile index 6c7b78d..4a64db6 100644 --- a/Dockerfile +++ b/Dockerfile @@ -24,6 +24,14 @@ RUN apt-get update -qq && \ ca-certificates && \ rm -rf /var/lib/apt/lists/* +# Optional malware scanner for P2P admission (P2P_MALWARE_SCAN=1). Off by +# default: build with --build-arg INSTALL_CLAMAV=1 to include it. +ARG INSTALL_CLAMAV=0 +RUN if [ "$INSTALL_CLAMAV" = "1" ]; then \ + apt-get update -qq && apt-get install -y --no-install-recommends clamav clamav-freshclam && \ + freshclam --quiet || true; rm -rf /var/lib/apt/lists/*; \ + fi + # ---- Install yt-dlp ---- RUN curl -fsSL \ https://github.com/yt-dlp/yt-dlp/releases/latest/download/yt-dlp \ diff --git a/plans/INDEX.md b/plans/INDEX.md index e3a3d72..e5e4ad8 100644 --- a/plans/INDEX.md +++ b/plans/INDEX.md @@ -16,7 +16,7 @@ green, app boots with no JS errors, P2P on by default, offline boot works). | 006 | 006-innertube-search-48066b | Answer searches from YouTube InnerTube directly with yt-dlp fallback | done | Answer searches from YouTube InnerTube directly with yt-dlp fallback | search 4–5 s → ~0.8 s | | 007 | 007-ytdlp-worker-045800 | Keep one long-lived yt-dlp worker process instead of spawning per call | done | Keep one long-lived yt-dlp worker process instead of spawning per call | ~1 s per yt-dlp call | | 008 | 008-p2p-schema-and-config-127966 | Add P2P tables, config flags and db helpers | done | Add P2P tables, config flags and db helpers | P2P ON, malware scan OFF by default | -| 009 | 009-server-content-hash-186e7f | Hash every validated server copy and register it as verified content | in-progress | | uses plans/patches/009-* | +| 009 | 009-server-content-hash-186e7f | Hash every validated server copy and register it as verified content | done | Hash every validated server copy and register it as verified content | uses plans/patches/009-* | | 010 | 010-views-and-retention-d0c6ca | Count views and evict server copies by retention criteria before LRU | queued | | | | 011 | 011-browser-sha256-e1793d | Add an incremental SHA-256 library for the browser and node tests | queued | | | | 012 | 012-device-file-registry-288d55 | Add the on-device IndexedDB file registry and hash saves while downloading | queued | | browser harness | diff --git a/plans/active/009-server-content-hash-186e7f.md b/plans/done/009-server-content-hash-186e7f.md similarity index 93% rename from plans/active/009-server-content-hash-186e7f.md rename to plans/done/009-server-content-hash-186e7f.md index 138a5f9..45a504f 100644 --- a/plans/active/009-server-content-hash-186e7f.md +++ b/plans/done/009-server-content-hash-186e7f.md @@ -271,3 +271,9 @@ test('scanFile maps exit codes', async () => { expect((await scanFile(file, '/nonexistent/scanner')).result).toBe('error'); }); ``` + +## Execution log + +- Executor: in-session Agent (haiku). Attempts: 1. Fix rounds: 0. +- Orchestrator re-ran Verification: `server/media-cache.js` and `media-cache.test.js` byte-identical to the pre-tested versions; `hash.js`, `p2p-admit.js`, `p2p-admit.test.js` identical; p2p-admit 5 pass; media-cache 25 pass; all 9 server test files 0 fail; `SERVER_OK`; the 3 server.js edits present. +- Executor Findings (verbatim): Plan 009 executed successfully. All steps completed: created server/hash.js, server/p2p-admit.js, server/p2p-admit.test.js; applied both media-cache patches without errors (25/25 tests pass); updated package.json test script; updated server.js imports to namespace import for p2pDb; added onReady callback to createMediaCache; added X-Content-SHA256 header to cachedDownloadResponse; added cid field to cachedStreamsPayload; updated Dockerfile with optional ClamAV. Server builds successfully; all 74 tests pass (9 test suites). diff --git a/server/hash.js b/server/hash.js new file mode 100644 index 0000000..eaac618 --- /dev/null +++ b/server/hash.js @@ -0,0 +1,26 @@ +/* hash.js — streaming SHA-256 of files on disk (never loads a whole video). + * The hex digest of a validated file is its P2P content id (cid). */ +import { createHash } from 'node:crypto'; +import { createReadStream } from 'node:fs'; + +export function sha256File(path) { + return new Promise((resolve, reject) => { + const h = createHash('sha256'); + createReadStream(path, { highWaterMark: 1024 * 1024 }) + .on('data', (d) => h.update(d)) + .on('error', reject) + .on('end', () => resolve(h.digest('hex'))); + }); +} + +// Hash of bytes [offset, offset+length) — used for holder range challenges. +export function sha256Range(path, offset, length) { + return new Promise((resolve, reject) => { + if (!(length > 0)) { resolve(createHash('sha256').digest('hex')); return; } + const h = createHash('sha256'); + createReadStream(path, { start: offset, end: offset + length - 1 }) + .on('data', (d) => h.update(d)) + .on('error', reject) + .on('end', () => resolve(h.digest('hex'))); + }); +} diff --git a/server/media-cache.js b/server/media-cache.js index 6d39552..c690dd3 100644 --- a/server/media-cache.js +++ b/server/media-cache.js @@ -24,6 +24,7 @@ import { statSync, statfsSync, unlinkSync, } from 'node:fs'; import { join } from 'node:path'; +import { sha256File } from './hash.js'; export const HIGH = 0; // explicit save / Broken re-download export const LOW = 1; // auto-cache after a play @@ -167,6 +168,9 @@ export function createMediaCache({ freeBytes = () => { const s = statfsSync(dir); return s.bavail * s.bsize; }, now = () => Date.now(), log = console, + hashFile = sha256File, // server-computed SHA-256 = the file's P2P content id + onReady = null, // ({ id, gen, path, sha256, size, … }) after a validated copy lands + backfillDelayMs = 30_000, // hash pre-existing copies this long after init() } = {}) { const tmpDir = join(dir, '.tmp'); const fileFor = (id, gen, kind = 'mp4') => join(dir, `${id}.${gen}.${kind}`); @@ -236,6 +240,20 @@ export function createMediaCache({ } } + // Tell the P2P layer a validated file with a server-computed hash exists. + // Never throws and never delays the caller. + function notifyReady(id, gen, sha256, probe, metaJson) { + if (!onReady || !sha256) return; + let meta = {}; + try { meta = JSON.parse(metaJson || '{}'); } catch { /* keep {} */ } + Promise.resolve() + .then(() => onReady({ + id, gen, path: fileFor(id, gen), sha256, size: probe.size, height: probe.height, + vcodec: probe.vcodec, acodec: probe.acodec, duration: probe.duration, meta, + })) + .catch((e) => log.warn?.(`[media] onReady ${id} failed: ${e.message}`)); + } + async function setStatus(job, status) { job.status = status; await db.upsertMedia(job.id, { status, updated_at: now() }); @@ -308,6 +326,7 @@ export function createMediaCache({ '-map', '0:a:0', '-c', 'copy', '-movflags', '+faststart', '-f', 'mp4', m4a]); if (ex.code !== 0 || sizeOf(m4a) < 1024) throw new Error('audio sidecar failed: ' + tail(ex.stderr)); + const sha256 = await hashFile(norm); const prev = await db.getMedia(id); const gen = ((prev && prev.gen) || 0) + 1; renameSync(norm, fileFor(id, gen)); @@ -318,13 +337,14 @@ export function createMediaCache({ size: probe.size + sizeOf(fileFor(id, gen, 'm4a')), height: probe.height, vcodec: probe.vcodec, acodec: probe.acodec, duration: probe.duration, optimized: 0, meta: JSON.stringify(metaFromInfo(info, probe.duration)), - attempts: 0, error: null, retry_at: 0, updated_at: t, last_access: t, + attempts: 0, error: null, retry_at: 0, updated_at: t, last_access: t, sha256, }; await db.upsertMedia(id, fields); removeFiles(id, gen); job.status = 'ready'; log.info?.(`[media] cached ${id} ${probe.height}p ${fmtMB(fields.size)}`); enqueueOptimize(id); + notifyReady(id, gen, sha256, probe, fields.meta); return { ...(prev || {}), video_id: id, ...fields }; } catch (err) { job.status = 'failed'; @@ -427,6 +447,7 @@ export function createMediaCache({ log.info?.(`[media] kept original ${id}: ${probe.vcodec} ${fmtMB(probe.size)} vs ${fmtMB(oldSize)}`); return; } + const sha256 = await hashFile(out); const gen = row.gen + 1; renameSync(out, fileFor(id, gen)); const oldAudio = fileFor(id, row.gen, 'm4a'); @@ -434,8 +455,9 @@ export function createMediaCache({ try { linkSync(oldAudio, newAudio); } catch { copyFileSync(oldAudio, newAudio); } await db.upsertMedia(id, { gen, size: probe.size + sizeOf(newAudio), height: probe.height, vcodec: probe.vcodec, - duration: probe.duration, optimized: OPT_REV, updated_at: now(), + duration: probe.duration, optimized: OPT_REV, updated_at: now(), sha256, }); + notifyReady(id, gen, sha256, probe, cur.meta); removeFiles(id, gen, OLD_GEN_GRACE_MS); log.info?.(`[media] optimized ${id}: ${fmtMB(oldSize)} → ${fmtMB(probe.size)} (${probe.vcodec})`); } catch (err) { @@ -622,7 +644,34 @@ export function createMediaCache({ } const s = await db.mediaStats(); log.info?.(`[media] ${s.count} cached (${fmtMB(s.bytes)}), ${resumed} job(s) resumed, dir ${dir}`); + if (backfillDelayMs >= 0) { + const t = setTimeout(() => backfillHashes().catch(() => {}), backfillDelayMs); + t.unref?.(); + } } - return { init, ensureCached, redownload, verify, getReady, filePath, touch, status, stats, isMediaId }; + // Copies cached before hashing existed get their sha256 in the background, + // one at a time, so their content ids can be admitted too. + async function backfillHashes() { + let n = 0; + for (const r of await db.listMedia()) { + if (r.status !== 'ready' || r.sha256) continue; + const p = fileFor(r.video_id, r.gen); + if (!existsSync(p)) continue; + try { + const sha256 = await hashFile(p); + const cur = await db.getMedia(r.video_id); + if (!cur || cur.status !== 'ready' || cur.gen !== r.gen) continue; + await db.upsertMedia(r.video_id, { sha256 }); + notifyReady(r.video_id, r.gen, sha256, + { size: sizeOf(p), height: r.height, vcodec: r.vcodec, acodec: r.acodec, duration: r.duration }, r.meta); + n++; + } catch (e) { + log.warn?.(`[media] hash backfill ${r.video_id}: ${e.message}`); + } + } + return n; + } + + return { init, ensureCached, redownload, verify, getReady, filePath, touch, status, stats, isMediaId, backfillHashes }; } diff --git a/server/media-cache.test.js b/server/media-cache.test.js index a2c22df..3bc7630 100644 --- a/server/media-cache.test.js +++ b/server/media-cache.test.js @@ -10,6 +10,8 @@ const root = mkdtempSync(join(tmpdir(), 'ytp-media-test-')); process.env.DB_PATH = join(root, 'test.db'); const dbmod = await import('./db.js'); const { createMediaCache, validateMedia, MediaSkip, HIGH, LOW, OPT_REV } = await import('./media-cache.js'); +const { initP2pSchema } = await import('./p2p-db.js'); +const { sha256File } = await import('./hash.js'); const fx = (name) => join(root, name); const ff = (...args) => { @@ -26,6 +28,7 @@ async function until(fn, ms = 20000) { beforeAll(async () => { await dbmod.initDb(); + await initP2pSchema(); // adds media_cache.sha256 // 10 s 640x360 H.264 + AAC — the shape of a good cached copy. ff('-f', 'lavfi', '-i', 'testsrc=size=640x360:rate=30:duration=10', '-f', 'lavfi', '-i', 'sine=frequency=440:duration=10', '-c:v', 'libx264', '-preset', 'ultrafast', '-pix_fmt', 'yuv420p', '-c:a', 'aac', '-shortest', fx('good.mp4')); @@ -102,6 +105,7 @@ function makeCache({ fixture = 'good.mp4', duration = 10, ...opts } = {}) { download: async (id, out) => { calls.push(id); await sleep(30); copyFileSync(fx(state.fixture), out); return out; }, transcode: { enabled: false }, freeBytes: () => 100 * 1024 ** 3, + backfillDelayMs: -1, ...opts, }); return { cache, dir, calls, state }; @@ -110,6 +114,35 @@ async function clearDb() { for (const r of await dbmod.listMedia()) await dbmod.deleteMedia(r.video_id); } +describe('content hashes (P2P cids)', () => { + test('a promoted copy stores the sha256 of its mp4 and reports it once', async () => { + await clearDb(); + const seen = []; + const { cache, dir } = makeCache({ onReady: (i) => { seen.push(i); } }); + await cache.init(); + const row = await cache.ensureCached('hashAAAAAA1', { priority: HIGH }); + const want = await sha256File(join(dir, `hashAAAAAA1.${row.gen}.mp4`)); + expect((await dbmod.getMedia('hashAAAAAA1')).sha256).toBe(want); + expect(await until(() => seen.length === 1)).toBe(true); + expect(seen[0]).toMatchObject({ id: 'hashAAAAAA1', gen: row.gen, sha256: want, vcodec: 'h264' }); + expect(seen[0].meta.title).toBe('T hashAAAAAA1'); + }); + + test('backfill hashes ready copies that have no sha256 yet', async () => { + await clearDb(); + const seen = []; + const { cache } = makeCache({ onReady: (i) => { seen.push(i); } }); + await cache.init(); + await cache.ensureCached('hashAAAAAA2', { priority: HIGH }); + await dbmod.upsertMedia('hashAAAAAA2', { sha256: null }); + seen.length = 0; + expect(await cache.backfillHashes()).toBe(1); + expect((await dbmod.getMedia('hashAAAAAA2')).sha256).toMatch(/^[0-9a-f]{64}$/); + expect(await until(() => seen.length === 1)).toBe(true); + expect(await cache.backfillHashes()).toBe(0); // nothing left + }); +}); + describe('media cache jobs', () => { test('fetches once, dedupes concurrent requests, then serves from disk', async () => { await clearDb(); diff --git a/server/p2p-admit.js b/server/p2p-admit.js new file mode 100644 index 0000000..51eb80e --- /dev/null +++ b/server/p2p-admit.js @@ -0,0 +1,53 @@ +/* ============================================================================ + * p2p-admit.js — the ONLY way a content id enters p2p_content + * (docs/p2p-architecture.md, "Security rules"). + * + * Callers must already have: the complete file on the server's own disk, its + * SHA-256 computed BY THE SERVER, and validateMedia() passed. This adds the + * optional malware scan (P2P_MALWARE_SCAN=1, off by default) and writes the + * row. The scan command gets the path as its last argument: exit 0 = clean, + * 1 = infected (rejected), anything else = scanner error (not admitted now). + * ========================================================================== */ +import { spawn } from 'node:child_process'; + +export function scanFile(path, cmd) { + const parts = String(cmd).split(/\s+/).filter(Boolean); + return new Promise((resolve) => { + let child; + try { child = spawn(parts[0], [...parts.slice(1), path], { stdio: ['ignore', 'pipe', 'pipe'] }); } + catch (e) { resolve({ result: 'error', detail: e.message }); return; } + let out = ''; + child.stdout.on('data', (d) => { out = (out + d).slice(-2000); }); + child.stderr.on('data', (d) => { out = (out + d).slice(-2000); }); + child.on('error', (e) => resolve({ result: 'error', detail: e.message })); + child.on('close', (code) => resolve( + code === 0 ? { result: 'clean' } : code === 1 ? { result: 'infected', detail: out.trim() } : { result: 'error', detail: out.trim() || 'exit ' + code }, + )); + }); +} + +const CID_RE = /^[0-9a-f]{64}$/; + +// info: { path, cid, videoId, size, height, vcodec, acodec, duration, meta, origin } +// deps: { cfg (P2P config), upsertContent, scan = scanFile, now = Date.now, log = console } +// → { ok: true, scan } | { ok: false, reason } +export async function admitFile(info, deps) { + const { cfg, upsertContent, scan = scanFile, now = Date.now, log = console } = deps; + if (!cfg.enabled) return { ok: false, reason: 'p2p disabled' }; + if (!CID_RE.test(String(info.cid || ''))) return { ok: false, reason: 'bad cid' }; + let scanResult = 'skipped'; + if (cfg.malwareScan) { + const r = await scan(info.path, cfg.scanCmd); + if (r.result !== 'clean') { + log.warn?.(`[p2p] ${info.videoId} ${info.cid.slice(0, 12)} not admitted: scan ${r.result} ${r.detail || ''}`); + return { ok: false, reason: 'scan ' + r.result }; + } + scanResult = 'clean'; + } + await upsertContent({ + cid: info.cid, videoId: info.videoId, size: info.size, height: info.height, vcodec: info.vcodec, + acodec: info.acodec, duration: info.duration, meta: info.meta || {}, origin: info.origin, + scan: scanResult, now: now(), + }); + return { ok: true, scan: scanResult }; +} diff --git a/server/p2p-admit.test.js b/server/p2p-admit.test.js new file mode 100644 index 0000000..dcc569b --- /dev/null +++ b/server/p2p-admit.test.js @@ -0,0 +1,52 @@ +import { test, expect } from 'bun:test'; +import { mkdtempSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { createHash } from 'node:crypto'; +import { admitFile, scanFile } from './p2p-admit.js'; +import { sha256File, sha256Range } from './hash.js'; + +const dir = mkdtempSync(join(tmpdir(), 'ytp-admit-')); +const file = join(dir, 'f.bin'); +const bytes = Buffer.from(Array.from({ length: 300000 }, (_, i) => i % 251)); +writeFileSync(file, bytes); +const CID = createHash('sha256').update(bytes).digest('hex'); +const quiet = { warn() {}, info() {} }; +const cfg = (o = {}) => ({ enabled: true, malwareScan: false, scanCmd: 'true', ...o }); +const info = { path: file, cid: CID, videoId: 'dQw4w9WgXcQ', size: bytes.length, origin: 'server' }; + +test('sha256File / sha256Range match node:crypto', async () => { + expect(await sha256File(file)).toBe(CID); + const want = createHash('sha256').update(bytes.subarray(1000, 1000 + 65536)).digest('hex'); + expect(await sha256Range(file, 1000, 65536)).toBe(want); +}); + +test('scan off by default: admitted with scan=skipped', async () => { + const rows = []; + const r = await admitFile(info, { cfg: cfg(), upsertContent: async (c) => rows.push(c), log: quiet }); + expect(r).toEqual({ ok: true, scan: 'skipped' }); + expect(rows[0]).toMatchObject({ cid: CID, videoId: 'dQw4w9WgXcQ', origin: 'server', scan: 'skipped' }); +}); + +test('scan on: clean admits, infected and scanner errors do not', async () => { + for (const [result, ok] of [['clean', true], ['infected', false], ['error', false]]) { + const rows = []; + const r = await admitFile(info, { cfg: cfg({ malwareScan: true }), upsertContent: async (c) => rows.push(c), scan: async () => ({ result }), log: quiet }); + expect(r.ok).toBe(ok); + expect(rows.length).toBe(ok ? 1 : 0); + } +}); + +test('disabled P2P or a malformed cid never admits', async () => { + const rows = []; + expect((await admitFile(info, { cfg: cfg({ enabled: false }), upsertContent: async (c) => rows.push(c) })).ok).toBe(false); + expect((await admitFile({ ...info, cid: 'XYZ' }, { cfg: cfg(), upsertContent: async (c) => rows.push(c) })).ok).toBe(false); + expect(rows.length).toBe(0); +}); + +test('scanFile maps exit codes', async () => { + expect((await scanFile(file, 'true')).result).toBe('clean'); + expect((await scanFile(file, 'false')).result).toBe('infected'); // exit 1 + expect((await scanFile(file, 'sh -c "exit 2" --')).result).toBe('error'); + expect((await scanFile(file, '/nonexistent/scanner')).result).toBe('error'); +}); diff --git a/server/package.json b/server/package.json index 0671a1d..cdae12f 100644 --- a/server/package.json +++ b/server/package.json @@ -6,7 +6,7 @@ "scripts": { "start": "bun server.js", "dev": "bun --hot server.js", - "test": "bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./ytdlp-pool.test.js && bun test ./p2p-db.test.js" + "test": "bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./ytdlp-pool.test.js && bun test ./p2p-db.test.js && bun test ./p2p-admit.test.js" }, "dependencies": { "@hono/node-server": "^1.14.0", diff --git a/server/server.js b/server/server.js index 266aa7f..64e1b74 100644 --- a/server/server.js +++ b/server/server.js @@ -48,7 +48,9 @@ import { registerNoteRoutes, parseLrc, sanitizeLyrics } from './notes.js'; import { createRemoteHub } from './remote.js'; import { createPartyHub } from './party.js'; import { registerUploadRoutes } from './uploads.js'; -import { initP2pSchema } from './p2p-db.js'; +import { admitFile } from './p2p-admit.js'; +import { P2P } from './p2p-config.js'; +import * as p2pDb from './p2p-db.js'; import * as innertube from './innertube.js'; import QRCode from 'qrcode'; import { createYtdlpPool } from './ytdlp-pool.js'; @@ -1080,6 +1082,12 @@ const media = createMediaCache({ threads: envNum('MEDIA_THREADS', 2), maxSeconds: envNum('MEDIA_OPT_MAX_SECONDS', 3600), }, + // P2P (docs/p2p-architecture.md): every validated copy's server-computed + // hash becomes verified content, after the optional malware scan. + onReady: (info) => admitFile( + { ...info, cid: info.sha256, videoId: info.id, origin: 'server' }, + { cfg: P2P, upsertContent: p2pDb.upsertContent }, + ), }); // A cached copy the compression lane turned into HEVC is only handed to @@ -1096,6 +1104,7 @@ function cachedStreamsPayload(videoId, row) { let meta = {}; try { meta = JSON.parse(row.meta || '{}'); } catch { /* corrupt — use defaults */ } const url = `/api/media/${videoId}?g=${row.gen}`; + // data.cid is additive and web-only (like serverCached) — the Tauri bridge ignores it. return { meta: { id: videoId, @@ -1116,6 +1125,7 @@ function cachedStreamsPayload(videoId, row) { codec: row.vcodec || 'h264', }], serverCached: true, + cid: row.sha256 || null, }; } @@ -1370,6 +1380,7 @@ function cachedDownloadResponse(videoId, fp, row) { 'Content-Disposition': `attachment; filename="${videoId}.mp4"`, 'Cache-Control': 'no-store', 'Access-Control-Allow-Origin': '*', + ...(row.sha256 ? { 'X-Content-SHA256': row.sha256, 'Access-Control-Expose-Headers': 'X-Content-SHA256' } : {}), }, }); } @@ -2022,7 +2033,7 @@ app.get('/*', indexHtml); // ============================================================================ async function main() { await initDb(); - await initP2pSchema(); + await p2pDb.initP2pSchema(); console.log(`[ytplayer] DB ready`); await media.init(); notes.startBackups();