126 lines
6.7 KiB
JavaScript
126 lines
6.7 KiB
JavaScript
/* ============================================================================
|
|
* 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?, restore? }
|
|
* → { ok, known:true } cid already verified — nothing to send
|
|
* (unless restore:true and the server
|
|
* lost its copy — plan 018)
|
|
* → { 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,
|
|
serverHasCid = async () => true, // (cid) → does the server still hold these bytes?
|
|
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);
|
|
const restore = body.restore === true && known && known.status === 'verified' && !(await serverHasCid(cid));
|
|
if (known && known.status === 'verified' && !restore) 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 };
|
|
}
|