231 lines
13 KiB
Diff
231 lines
13 KiB
Diff
--- /dev/null
|
|
+++ b/server/p2p-intake.js
|
|
@@ -0,0 +1,121 @@
|
|
+/* ============================================================================
|
|
+ * p2p-intake.js — a device hands a file to the server for verification
|
|
+ * (docs/p2p-architecture.md flow 7). The ONLY path by which bytes that did
|
|
+ * not come from the server's own fetch can become verified content.
|
|
+ *
|
|
+ * POST /api/p2p/intake (device) { videoId, cid?, size, title?, channel? }
|
|
+ * → { ok, known:true } cid already verified — nothing to send
|
|
+ * → { ok, ticket, url, expiresAt } PUT the bytes to url within 30 min
|
|
+ * PUT /api/p2p/intake/:ticket (device) raw file body
|
|
+ * → { ok, cid, adopted } admitted (+ adopted into the media cache)
|
|
+ * → 4xx { ok:false, error } rejected; nothing kept
|
|
+ *
|
|
+ * Order is fixed: bytes land in P2P_INTAKE_DIR (never served) → the SERVER
|
|
+ * hashes them while writing → a claimed cid must match → validateMedia() →
|
|
+ * admitFile() (malware scan only when P2P_MALWARE_SCAN=1) → holder row for
|
|
+ * the uploader → optional adoption into the media cache. Any failure deletes
|
|
+ * the file.
|
|
+ * ========================================================================== */
|
|
+import { createHash, randomBytes } from 'node:crypto';
|
|
+import { createWriteStream, mkdirSync, unlinkSync } from 'node:fs';
|
|
+import { join } from 'node:path';
|
|
+
|
|
+const CID_RE = /^[0-9a-f]{64}$/;
|
|
+const ID_RE = /^[A-Za-z0-9_-]{6,64}$/;
|
|
+const TICKET_TTL = 30 * 60_000;
|
|
+const MAX_ACTIVE = 2;
|
|
+// Display text from the device — untrusted: control chars/brackets stripped, capped.
|
|
+const clean = (v, max) => String(v || '').replace(/[\u0000-\u001f\u007f<>]+/g, ' ').replace(/\s+/g, ' ').trim().slice(0, max);
|
|
+
|
|
+export function registerIntakeRoutes(app, deps) {
|
|
+ const {
|
|
+ cfg, p2pDb, gate, requireDevice, validateMedia, admitFile, adopt = null,
|
|
+ now = () => Date.now(), log = console,
|
|
+ } = deps;
|
|
+ const tickets = new Map(); // ticket -> { deviceId, videoId, cid, size, expiresAt, busy }
|
|
+ let active = 0;
|
|
+ mkdirSync(cfg.intakeDir, { recursive: true });
|
|
+
|
|
+ const sweep = () => { const t = now(); for (const [k, v] of tickets) if (!v.busy && v.expiresAt < t) tickets.delete(k); };
|
|
+
|
|
+ app.post('/api/p2p/intake', gate, requireDevice, async (c) => {
|
|
+ const d = c.get('device');
|
|
+ const body = await c.req.json().catch(() => ({}));
|
|
+ const videoId = String(body.videoId || '');
|
|
+ const cid = body.cid ? String(body.cid).toLowerCase() : null;
|
|
+ const size = Number(body.size);
|
|
+ if (!ID_RE.test(videoId) || (cid && !CID_RE.test(cid))) return c.json({ ok: false, error: 'bad videoId or cid' }, 400);
|
|
+ if (!(size > 0) || size > cfg.intakeMaxBytes) return c.json({ ok: false, error: 'bad size' }, 413);
|
|
+ if (cid) {
|
|
+ const known = await p2pDb.getContent(cid);
|
|
+ if (known && known.status === 'verified') return c.json({ ok: true, known: true });
|
|
+ if (known && known.status === 'revoked') return c.json({ ok: false, error: 'revoked' }, 410);
|
|
+ }
|
|
+ sweep();
|
|
+ for (const v of tickets.values()) if (v.deviceId === d.device_id) return c.json({ ok: false, error: 'an upload from this device is already open' }, 429);
|
|
+ const ticket = randomBytes(16).toString('hex');
|
|
+ const expiresAt = now() + TICKET_TTL;
|
|
+ const meta = { title: clean(body.title, 300), channel: clean(body.channel, 200) };
|
|
+ tickets.set(ticket, { deviceId: d.device_id, videoId, cid, size, meta, expiresAt, busy: false });
|
|
+ return c.json({ ok: true, ticket, url: `/api/p2p/intake/${ticket}`, expiresAt });
|
|
+ });
|
|
+
|
|
+ app.put('/api/p2p/intake/:ticket', gate, requireDevice, async (c) => {
|
|
+ const d = c.get('device');
|
|
+ const ticket = c.req.param('ticket');
|
|
+ const tk = tickets.get(ticket);
|
|
+ if (!tk || tk.deviceId !== d.device_id || tk.expiresAt < now()) return c.json({ ok: false, error: 'unknown or expired ticket' }, 404);
|
|
+ if (tk.busy) return c.json({ ok: false, error: 'upload already running' }, 409);
|
|
+ if (active >= MAX_ACTIVE) return c.json({ ok: false, error: 'server busy, retry later' }, 503);
|
|
+ tk.busy = true;
|
|
+ active++;
|
|
+ const path = join(cfg.intakeDir, ticket + '.mp4');
|
|
+ let keep = false;
|
|
+ try {
|
|
+ const h = createHash('sha256');
|
|
+ let got = 0;
|
|
+ const out = createWriteStream(path);
|
|
+ const reader = c.req.raw.body ? c.req.raw.body.getReader() : null;
|
|
+ if (!reader) throw Object.assign(new Error('empty body'), { status: 400 });
|
|
+ try {
|
|
+ for (;;) {
|
|
+ const { done, value } = await reader.read();
|
|
+ if (done) break;
|
|
+ got += value.byteLength;
|
|
+ if (got > tk.size) throw Object.assign(new Error('more bytes than announced'), { status: 413 });
|
|
+ h.update(value);
|
|
+ if (!out.write(value)) await new Promise((r) => out.once('drain', r));
|
|
+ }
|
|
+ } finally {
|
|
+ await new Promise((r) => out.end(r));
|
|
+ }
|
|
+ if (got !== tk.size) throw Object.assign(new Error(`size mismatch (${got} of ${tk.size})`), { status: 400 });
|
|
+ const cid = h.digest('hex');
|
|
+ if (tk.cid && cid !== tk.cid) throw Object.assign(new Error('content hash mismatch'), { status: 400 });
|
|
+ let probe;
|
|
+ try { probe = await validateMedia(path, 0, { codecs: ['h264', 'hevc'] }); }
|
|
+ catch (e) { throw Object.assign(new Error(e.message), { status: 422 }); }
|
|
+ const adm = await admitFile(
|
|
+ { path, cid, videoId: tk.videoId, size: got, height: probe.height, vcodec: probe.vcodec, acodec: probe.acodec, duration: probe.duration, meta: tk.meta, origin: 'intake' },
|
|
+ { cfg, upsertContent: p2pDb.upsertContent, log },
|
|
+ );
|
|
+ if (!adm.ok) throw Object.assign(new Error('not admitted: ' + adm.reason), { status: 422 });
|
|
+ await p2pDb.upsertHolder({ cid, deviceId: d.device_id, trust: 'challenged', now: now() });
|
|
+ let adopted = false;
|
|
+ if (adopt) {
|
|
+ try { adopted = !!(await adopt(tk.videoId, path, { sha256: cid, probe, meta: tk.meta })).adopted; keep = adopted; }
|
|
+ catch (e) { log.warn?.(`[p2p] adopt ${tk.videoId} failed: ${e.message}`); }
|
|
+ }
|
|
+ log.info?.(`[p2p] intake ${tk.videoId} ${cid.slice(0, 12)} admitted${adopted ? ' + adopted' : ''}`);
|
|
+ return c.json({ ok: true, cid, adopted });
|
|
+ } catch (e) {
|
|
+ return c.json({ ok: false, error: e.message }, e.status || 500);
|
|
+ } finally {
|
|
+ tickets.delete(ticket);
|
|
+ active--;
|
|
+ if (!keep) { try { unlinkSync(path); } catch { /* never written */ } }
|
|
+ }
|
|
+ });
|
|
+
|
|
+ return { openTickets: () => tickets.size, active: () => active };
|
|
+}
|
|
--- /dev/null
|
|
+++ b/server/p2p-intake.test.js
|
|
@@ -0,0 +1,103 @@
|
|
+// Device → server intake: hash, validate, (scan), admit (plan 016). Needs ffmpeg.
|
|
+import { test, expect, beforeAll } from 'bun:test';
|
|
+import { mkdtempSync, readFileSync, readdirSync } from 'node:fs';
|
|
+import { spawnSync } from 'node:child_process';
|
|
+import { tmpdir } from 'node:os';
|
|
+import { join } from 'node:path';
|
|
+import { createHash } from 'node:crypto';
|
|
+import { Hono } from 'hono';
|
|
+
|
|
+const root = mkdtempSync(join(tmpdir(), 'ytp-intake-'));
|
|
+process.env.DB_PATH = join(root, 'test.db');
|
|
+const dbmod = await import('./db.js');
|
|
+const p2pDb = await import('./p2p-db.js');
|
|
+const { registerP2pRoutes } = await import('./p2p-routes.js');
|
|
+const { registerIntakeRoutes } = await import('./p2p-intake.js');
|
|
+const { admitFile } = await import('./p2p-admit.js');
|
|
+const { validateMedia } = await import('./media-cache.js');
|
|
+
|
|
+const quiet = { warn() {}, info() {} };
|
|
+const intakeDir = join(root, 'intake');
|
|
+const cfg = { enabled: true, malwareScan: false, scanCmd: 'true', staleDays: 7, intakeDir, intakeMaxBytes: 50 * 1024 * 1024 };
|
|
+let app;
|
|
+let dev;
|
|
+let good;
|
|
+let goodCid;
|
|
+const adopted = [];
|
|
+
|
|
+const req = (path, { method = 'GET', body, raw } = {}) => app.request(path, {
|
|
+ method,
|
|
+ headers: { ...(raw ? {} : { 'Content-Type': 'application/json' }), 'X-Device': dev.deviceId + '.' + dev.secret },
|
|
+ body: raw || (body ? JSON.stringify(body) : undefined),
|
|
+});
|
|
+
|
|
+beforeAll(async () => {
|
|
+ await dbmod.initDb();
|
|
+ await p2pDb.initP2pSchema();
|
|
+ const f = join(root, 'good.mp4');
|
|
+ const r = spawnSync('ffmpeg', ['-v', 'error', '-y', '-f', 'lavfi', '-i', 'testsrc=size=320x180:rate=25:duration=4', '-f', 'lavfi', '-i', 'sine=frequency=440:duration=4',
|
|
+ '-c:v', 'libx264', '-preset', 'ultrafast', '-pix_fmt', 'yuv420p', '-c:a', 'aac', '-shortest', '-movflags', '+faststart', f]);
|
|
+ if (r.status !== 0) throw new Error('ffmpeg fixture failed: ' + r.stderr);
|
|
+ good = readFileSync(f);
|
|
+ goodCid = createHash('sha256').update(good).digest('hex');
|
|
+ app = new Hono();
|
|
+ const p2p = registerP2pRoutes(app, { cfg, p2pDb, fileForCid: async () => null, sha256Range: async () => '', log: quiet });
|
|
+ registerIntakeRoutes(app, {
|
|
+ cfg, p2pDb, gate: p2p.gate, requireDevice: p2p.requireDevice, validateMedia, admitFile, log: quiet,
|
|
+ adopt: async (videoId, path, info) => { adopted.push({ videoId, sha256: info.sha256, meta: info.meta }); return { adopted: false }; },
|
|
+ });
|
|
+ dev = await (await app.request('/api/p2p/device', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: '{}' })).json();
|
|
+});
|
|
+
|
|
+test('a valid file is hashed by the server, validated, admitted and its uploader becomes a holder', async () => {
|
|
+ const t = await (await req('/api/p2p/intake', { method: 'POST', body: { videoId: 'upAAAAAAAA1', cid: goodCid, size: good.length, title: 'Song <b>x</b>', channel: 'Choir' } })).json();
|
|
+ expect(t.ok).toBe(true);
|
|
+ const r = await req(t.url, { method: 'PUT', raw: good });
|
|
+ expect(r.status).toBe(200);
|
|
+ expect(await r.json()).toEqual({ ok: true, cid: goodCid, adopted: false });
|
|
+ expect((await p2pDb.getContent(goodCid))).toMatchObject({ origin: 'intake', status: 'verified', scan: 'skipped', video_id: 'upAAAAAAAA1' });
|
|
+ expect((await p2pDb.listHolders(goodCid))[0]).toMatchObject({ device_id: dev.deviceId, trust: 'challenged' });
|
|
+ expect(adopted).toEqual([{ videoId: 'upAAAAAAAA1', sha256: goodCid, meta: { title: 'Song b x /b', channel: 'Choir' } }]);
|
|
+ expect(JSON.parse((await p2pDb.getContent(goodCid)).meta).title).toBe('Song b x /b');
|
|
+ expect(readdirSync(intakeDir)).toEqual([]); // not adopted → deleted
|
|
+ // Already verified → nothing to upload next time.
|
|
+ expect(await (await req('/api/p2p/intake', { method: 'POST', body: { videoId: 'upAAAAAAAA1', cid: goodCid, size: good.length } })).json()).toEqual({ ok: true, known: true });
|
|
+});
|
|
+
|
|
+test('claimed cid mismatch, truncated media, and non-media are rejected and deleted', async () => {
|
|
+ const cases = [
|
|
+ [{ cid: 'f'.repeat(64), bytes: good }, 400, /hash mismatch/],
|
|
+ [{ bytes: good.subarray(0, Math.floor(good.length * 0.6)) }, 422, /validation/],
|
|
+ [{ bytes: Buffer.alloc(200 * 1024, 7) }, 422, /validation/],
|
|
+ ];
|
|
+ for (const [c, status, msg] of cases) {
|
|
+ const t = await (await req('/api/p2p/intake', { method: 'POST', body: { videoId: 'upAAAAAAAA2', cid: c.cid, size: c.bytes.length } })).json();
|
|
+ const r = await req(t.url, { method: 'PUT', raw: c.bytes });
|
|
+ expect(r.status).toBe(status);
|
|
+ expect((await r.json()).error).toMatch(msg);
|
|
+ }
|
|
+ expect(readdirSync(intakeDir)).toEqual([]);
|
|
+ expect((await p2pDb.listContentForVideo('upAAAAAAAA2')).length).toBe(0);
|
|
+});
|
|
+
|
|
+test('size limits and tickets are enforced', async () => {
|
|
+ expect((await req('/api/p2p/intake', { method: 'POST', body: { videoId: 'upAAAAAAAA3', size: cfg.intakeMaxBytes + 1 } })).status).toBe(413);
|
|
+ const t = await (await req('/api/p2p/intake', { method: 'POST', body: { videoId: 'upAAAAAAAA3', size: 10 } })).json();
|
|
+ expect((await req('/api/p2p/intake', { method: 'POST', body: { videoId: 'upAAAAAAAA4', size: 10 } })).status).toBe(429); // one open per device
|
|
+ const r = await req(t.url, { method: 'PUT', raw: Buffer.alloc(11) });
|
|
+ expect(r.status).toBe(413);
|
|
+ expect((await req(t.url, { method: 'PUT', raw: Buffer.alloc(10) })).status).toBe(404); // ticket consumed
|
|
+});
|
|
+
|
|
+test('malware scan on: an infected verdict is never admitted', async () => {
|
|
+ const cfg2 = { ...cfg, malwareScan: true, scanCmd: 'false' }; // exit 1 = infected
|
|
+ const app2 = new Hono();
|
|
+ const p2p = registerP2pRoutes(app2, { cfg: cfg2, p2pDb, fileForCid: async () => null, sha256Range: async () => '', log: quiet });
|
|
+ registerIntakeRoutes(app2, { cfg: cfg2, p2pDb, gate: p2p.gate, requireDevice: p2p.requireDevice, validateMedia, admitFile, log: quiet });
|
|
+ const h = { 'X-Device': dev.deviceId + '.' + dev.secret };
|
|
+ const t = await (await app2.request('/api/p2p/intake', { method: 'POST', headers: { ...h, 'Content-Type': 'application/json' }, body: JSON.stringify({ videoId: 'upAAAAAAAA5', size: good.length }) })).json();
|
|
+ await p2pDb.revokeContent(goodCid); // make sure it is not simply "already known"
|
|
+ const r = await app2.request(t.url, { method: 'PUT', headers: h, body: good });
|
|
+ expect(r.status).toBe(422);
|
|
+ expect((await r.json()).error).toMatch(/scan infected/);
|
|
+});
|