Queue executable speed and peer-to-peer sharing plans with tested patches and harnesses

This commit is contained in:
Claude
2026-09-29 19:03:12 +00:00
parent a0ffc1b493
commit c17bd4c9ac
54 changed files with 6369 additions and 11 deletions

View File

@@ -0,0 +1,63 @@
--- a/server/media-cache.test.js
+++ b/server/media-cache.test.js
@@ -10,6 +10,8 @@
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 @@
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 @@
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 @@
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();

View File

@@ -0,0 +1,120 @@
--- a/server/media-cache.js
+++ b/server/media-cache.js
@@ -24,6 +24,7 @@
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 @@
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 @@
}
}
+ // 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 @@
'-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 @@
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 @@
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 @@
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 @@
}
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?.();
+ }
+ }
+
+ // 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 };
+ return { init, ensureCached, redownload, verify, getReady, filePath, touch, status, stats, isMediaId, backfillHashes };
}

View File

@@ -0,0 +1,44 @@
--- a/server/media-cache.js
+++ b/server/media-cache.js
@@ -226,7 +226,10 @@
let { bytes } = await db.mediaStats();
if (bytes + est > maxBytes) {
const cutoff = now() - EVICT_PROTECT_MS;
- for (const r of await db.listMediaLru()) {
+ // Retention order (plan 010) when the server provides it: copies that
+ // are neither top nor recent go first. Otherwise plain LRU.
+ const order = db.listMediaEvictionOrder ? await db.listMediaEvictionOrder() : await db.listMediaLru();
+ for (const r of order) {
if (bytes + est <= maxBytes) break;
if (r.video_id === id || r.last_access > cutoff || jobs.has(r.video_id) || optActive === r.video_id) continue;
await evict(r.video_id, 'budget');
--- a/server/media-cache.test.js
+++ b/server/media-cache.test.js
@@ -240,6 +240,27 @@
expect(await cache.getReady('lruAAAAAAA3')).not.toBeNull();
});
+ test('eviction follows listMediaEvictionOrder when the db provides it', async () => {
+ await clearDb();
+ const size = readFileSync(fx('good.mp4')).length;
+ const db = {
+ getMedia: dbmod.getMedia, upsertMedia: dbmod.upsertMedia, deleteMedia: dbmod.deleteMedia,
+ listMedia: dbmod.listMedia, listMediaLru: dbmod.listMediaLru, touchMedia: dbmod.touchMedia,
+ mediaStats: dbmod.mediaStats,
+ // Retention says the MORE recently played #2 is the one to drop (e.g. no views).
+ listMediaEvictionOrder: async () => (await dbmod.listMediaLru()).reverse(),
+ };
+ const { cache } = makeCache({ db, maxBytes: size * 2 + 60 * 250 * 1024 + 50_000, duration: 10 });
+ await cache.init();
+ await cache.ensureCached('retAAAAAAA1', { priority: HIGH });
+ await cache.ensureCached('retAAAAAAA2', { priority: HIGH });
+ await dbmod.upsertMedia('retAAAAAAA1', { last_access: 1000 });
+ await dbmod.upsertMedia('retAAAAAAA2', { last_access: 2000 });
+ await cache.ensureCached('retAAAAAAA3', { priority: HIGH });
+ expect(await dbmod.getMedia('retAAAAAAA2')).toBeNull();
+ expect(await dbmod.getMedia('retAAAAAAA1')).not.toBeNull();
+ });
+
test('free-disk guard skips caching', async () => {
await clearDb();
const { cache, calls } = makeCache({ freeBytes: () => 1024 ** 3, minFreeBytes: 5 * 1024 ** 3 });

View File

@@ -0,0 +1,75 @@
--- a/frontend/opfs-worker.js
+++ b/frontend/opfs-worker.js
@@ -13,12 +13,20 @@
* Out messages:
* { type: 'unsupported' } → caller falls back to main thread
* { type: 'progress', received } → bytes written so far
- * { type: 'done', ext } → file stored as <videoId>.<ext>
+ * { type: 'done', ext, sha256, expectedSha, size }
+ * → file stored as <videoId>.<ext>;
+ * sha256 = hash of the stored bytes
+ * (P2P content id), expectedSha = the
+ * server's X-Content-SHA256 or null
* { type: 'error', error } → failed; .part cleaned up
* ========================================================================== */
'use strict';
+// Incremental SHA-256 (frontend/sha256.js): the file is hashed while it is
+// written, so the device knows its content id without re-reading the file.
+try { importScripts('/sha256.js'); } catch { /* hashing unavailable — save still works */ }
+
async function getVideosDir() {
const root = await navigator.storage.getDirectory();
return root.getDirectoryHandle('videos', { create: true });
@@ -60,6 +68,7 @@
dir = await getVideosDir();
const partHandle = await dir.getFileHandle(partName, { create: true });
const access = await partHandle.createSyncAccessHandle();
+ const hasher = self.Sha256 ? self.Sha256.create() : null;
let offset = 0;
try {
const reader = res.body.getReader();
@@ -67,6 +76,7 @@
const { done, value } = await reader.read();
if (done) break;
access.write(value, { at: offset });
+ if (hasher) hasher.update(value);
offset += value.byteLength;
self.postMessage({ type: 'progress', received: offset });
}
@@ -81,6 +91,14 @@
if (expected > 0 && offset !== expected) {
throw new Error(`download cut short (${offset} of ${expected} bytes)`);
}
+ // The server names the hash of what it sent (media cache copies). A
+ // mismatch means the bytes were damaged on the way — never keep them.
+ const sha256 = hasher ? hasher.hex() : null;
+ const sent = (res.headers.get('x-content-sha256') || '').trim().toLowerCase();
+ const expectedSha = /^[0-9a-f]{64}$/.test(sent) ? sent : null;
+ if (sha256 && expectedSha && sha256 !== expectedSha) {
+ throw new Error('integrity check failed (content hash mismatch)');
+ }
// Finalize: .part → permanent name. Prefer the native rename, but treat
// ANY move() failure as "unavailable" and fall back to a chunked copy —
@@ -111,7 +129,7 @@
await dir.removeEntry(partName);
}
- self.postMessage({ type: 'done', ext });
+ self.postMessage({ type: 'done', ext, sha256, expectedSha, size: offset });
} catch (err) {
// Never leave a corrupt partial behind
try { if (dir && partName) await dir.removeEntry(partName); } catch { /* gone */ }
--- a/frontend/opfs.js
+++ b/frontend/opfs.js
@@ -136,7 +136,7 @@
};
worker.onmessage = (e) => {
const m = e.data || {};
- if (m.type === 'done') finish({ ok: true });
+ if (m.type === 'done') finish({ ok: true, sha256: m.sha256 || null, expectedSha: m.expectedSha || null, size: m.size || 0 });
else if (m.type === 'unsupported') finish({ ok: false, fallback: true });
else if (m.type === 'error') finish({ ok: false, error: m.error });
// 'progress' messages are informational; ignored here

View File

@@ -0,0 +1,22 @@
--- a/frontend/opfs.js
+++ b/frontend/opfs.js
@@ -108,6 +108,19 @@
}
},
+ // Bytes [offset, offset+length) of a saved video — answers the server's
+ // P2P range challenges (p2p-client.js). null when the file is missing.
+ async readRange(videoId, offset, length) {
+ try {
+ const found = await findHandle(videoId);
+ if (!found) return null;
+ const file = await found[0].getFile();
+ return new Uint8Array(await file.slice(offset, offset + length).arrayBuffer());
+ } catch {
+ return null;
+ }
+ },
+
revokeUrl(url) {
if (url && _blobUrls.has(url)) {
URL.revokeObjectURL(url);

View File

@@ -0,0 +1,105 @@
--- a/frontend/p2p-client.js
+++ b/frontend/p2p-client.js
@@ -6,6 +6,9 @@
* changed() a save/delete happened → re-report soon
* device() { deviceId, secret } or null
* authHeaders() { 'X-Device': … } for other P2P calls
+ * onMessage(type, fn) / signal(to, data) / peer()
+ * live /ws/p2p socket (plan 014): the
+ * device is "online" while it is open
*
* What it does, in order, each sync:
* 1. registers the device once (localStorage ytpDevice)
@@ -148,6 +151,74 @@
return j;
}
+ // ---- presence socket (/ws/p2p) — plan 014 ------------------------------------
+ // Open while sharing OR receiving is on. Credentials go in the first message,
+ // never the URL. Reconnects with backoff (5 s … 5 min).
+ let ws = null;
+ let myPeer = null;
+ let wsRetry = 0;
+ let wsTimer = null;
+ const listeners = new Map(); // type -> Set<fn>
+
+ const wantSocket = () => {
+ const s = hooks.getSettings();
+ return s.p2pShare !== false || s.p2pReceive !== false;
+ };
+
+ function onMessage(type, fn) {
+ if (!listeners.has(type)) listeners.set(type, new Set());
+ listeners.get(type).add(fn);
+ return () => listeners.get(type).delete(fn);
+ }
+ function emit(m) {
+ for (const fn of listeners.get(m.type) || []) { try { fn(m); } catch { /* a listener's bug is not ours */ } }
+ }
+
+ function connect() {
+ if (ws || !wantSocket()) return;
+ const d = device();
+ if (!d) return;
+ const proto = location.protocol === 'https:' ? 'wss:' : 'ws:';
+ let sock;
+ try { sock = new WebSocket(`${proto}//${location.host}/ws/p2p`); } catch { return; }
+ ws = sock;
+ let ping = null;
+ sock.onopen = () => {
+ sock.send(JSON.stringify({ type: 'auth', device: d.deviceId, secret: d.secret }));
+ // The server drops sockets idle for 120 s (server.js websocketHandler).
+ ping = setInterval(() => { try { sock.send('{"type":"ping"}'); } catch { /* closing */ } }, 50_000);
+ };
+ sock.onmessage = (e) => {
+ let m;
+ try { m = JSON.parse(e.data); } catch { return; }
+ if (m.type === 'hello') { myPeer = m.peer; wsRetry = 0; }
+ emit(m);
+ };
+ sock.onclose = () => {
+ clearInterval(ping);
+ if (ws !== sock) return;
+ ws = null;
+ myPeer = null;
+ if (!wantSocket()) return;
+ clearTimeout(wsTimer);
+ wsTimer = setTimeout(connect, Math.min(300_000, 5000 * 2 ** wsRetry++));
+ };
+ }
+
+ function disconnect() {
+ clearTimeout(wsTimer);
+ const sock = ws;
+ ws = null;
+ myPeer = null;
+ if (sock) { try { sock.close(); } catch { /* gone */ } }
+ }
+
+ function signal(to, data) {
+ if (!ws || ws.readyState !== 1) return false;
+ ws.send(JSON.stringify({ type: 'signal', to, data }));
+ return true;
+ }
+
async function syncOnce() {
if (!window.OPFS || !window.OPFS.isSupported() || !window.DeviceDB || !window.Sha256) return null;
if (navigator.onLine === false) return null;
@@ -158,6 +229,7 @@
await hashPending(recs);
const res = await report(recs, dev);
lastSync = Date.now();
+ if (wantSocket()) connect(); else disconnect();
return res;
}
@@ -185,5 +257,8 @@
});
}
- window.P2PClient = { start, changed, sync, device, authHeaders, config };
+ window.P2PClient = {
+ start, changed, sync, device, authHeaders, config,
+ onMessage, signal, peer: () => myPeer, isConnected: () => !!(ws && ws.readyState === 1 && myPeer),
+ };
}());

View File

@@ -0,0 +1,213 @@
--- /dev/null
+++ b/server/p2p-hub.js
@@ -0,0 +1,120 @@
+/* ============================================================================
+ * p2p-hub.js — live presence + WebRTC signalling for P2P devices, and the
+ * availability payload (docs/p2p-architecture.md flows 4–5).
+ *
+ * /ws/p2p (credentials never go in the URL — proxies log URLs)
+ * you → server {type:'auth', device, secret} first message, within 10 s
+ * server → you {type:'hello', peer} after auth
+ * you → peer {type:'signal', to:<peer>, data} relayed as
+ * server → peer {type:'signal', from:<your peer>, data}
+ * server → you {type:'error', error}
+ * Server-initiated messages (plan 018) go through hub.send(deviceId, msg).
+ *
+ * "Online" lives only in memory: a server restart forgets it, a closed
+ * socket drops it. Holder ROWS are persistent (p2p_holders) — the payload
+ * reports both: `online` (now) and `lastVerifiedAt` / `stale` (history).
+ * Peers are addressed by an opaque id (peerIdOf), never by device id.
+ * ========================================================================== */
+import { createHash, timingSafeEqual } from 'node:crypto';
+
+const sha = (s) => createHash('sha256').update(String(s)).digest('hex');
+const same = (a, b) => { const x = Buffer.from(String(a)), y = Buffer.from(String(b)); return x.length === y.length && timingSafeEqual(x, y); };
+export const peerIdOf = (deviceId) => sha('peer:' + deviceId).slice(0, 12);
+
+const MAX_MSG = 64 * 1024;
+const BUDGET_PER_MIN = 300;
+
+export function createP2pHub({ getDevice, enabled = () => true, now = () => Date.now(), log = console } = {}) {
+ const online = new Map(); // deviceId -> ws
+ const byPeer = new Map(); // peer -> deviceId (online only)
+
+ const send = (ws, m) => { try { ws.send(JSON.stringify(m)); } catch { /* gone */ } };
+
+ function upgrade(req, server) {
+ if (!enabled()) return new Response('p2p disabled', { status: 404 });
+ const ok = server.upgrade(req, { data: { hub: 'p2p', authed: false, budget: [] } });
+ return ok ? undefined : new Response('upgrade failed', { status: 400 });
+ }
+
+ function open(ws) {
+ const t = setTimeout(() => { if (!ws.data.authed) { try { ws.close(4401, 'auth timeout'); } catch { /* gone */ } } }, 10_000);
+ t.unref?.();
+ }
+
+ async function auth(ws, m) {
+ const deviceId = String(m.device || '');
+ const secret = String(m.secret || '');
+ const d = /^dev_[0-9a-f]{16}$/.test(deviceId) ? await getDevice(deviceId).catch(() => null) : null;
+ if (!d || !same(d.secret_hash, sha(secret))) { try { ws.close(4401, 'unknown device'); } catch { /* gone */ } return; }
+ ws.data.authed = true;
+ ws.data.deviceId = deviceId;
+ ws.data.peer = peerIdOf(deviceId);
+ const prev = online.get(deviceId);
+ if (prev && prev !== ws) { prev.data.deviceId = null; try { prev.close(4409, 'replaced'); } catch { /* gone */ } }
+ online.set(deviceId, ws);
+ byPeer.set(ws.data.peer, deviceId);
+ send(ws, { type: 'hello', peer: ws.data.peer });
+ }
+
+ async function message(ws, raw) {
+ const t = now();
+ ws.data.budget = ws.data.budget.filter((x) => t - x < 60_000);
+ if (ws.data.budget.length >= BUDGET_PER_MIN) { send(ws, { type: 'error', error: 'slow down' }); return; }
+ ws.data.budget.push(t);
+ const s = typeof raw === 'string' ? raw : Buffer.from(raw).toString('utf8');
+ if (s.length > MAX_MSG) return;
+ let m;
+ try { m = JSON.parse(s); } catch { return; }
+ if (!ws.data.authed) { if (m && m.type === 'auth' && !ws.data.authing) { ws.data.authing = true; await auth(ws, m); } return; }
+ if (m && m.type === 'signal' && typeof m.to === 'string') {
+ const target = online.get(byPeer.get(m.to));
+ if (!target || target === ws) { send(ws, { type: 'error', error: 'peer offline', to: m.to }); return; }
+ send(target, { type: 'signal', from: ws.data.peer, data: m.data });
+ }
+ }
+
+ function close(ws) {
+ const { deviceId, peer } = ws.data || {};
+ if (deviceId && online.get(deviceId) === ws) {
+ online.delete(deviceId);
+ byPeer.delete(peer);
+ }
+ }
+
+ return {
+ upgrade,
+ websocket: { open, message, close },
+ isOnline: (deviceId) => online.has(deviceId),
+ send: (deviceId, msg) => { const ws = online.get(deviceId); if (!ws) return false; send(ws, msg); return true; },
+ onlineCount: () => online.size,
+ };
+}
+
+// GET /api/p2p/holders payload. Only devices that share are listed; holders
+// are never dropped for age — `stale` just flags an old last check.
+export async function holdersPayload({ videoId, cid, p2pDb, isOnline, staleDays, serverHas, now = Date.now() }) {
+ const contents = cid
+ ? [await p2pDb.getContent(cid)].filter((c) => c && c.status === 'verified')
+ : await p2pDb.listContentForVideo(videoId);
+ const staleMs = staleDays * 86400_000;
+ const out = [];
+ for (const c of contents) {
+ const holders = (await p2pDb.listHolders(c.cid, 50))
+ .filter((h) => Number(h.share) === 1)
+ .map((h) => ({
+ peer: peerIdOf(h.device_id),
+ online: isOnline(h.device_id),
+ lastVerifiedAt: Number(h.last_verified_at),
+ stale: now - Number(h.last_verified_at) > staleMs,
+ trust: h.trust,
+ }))
+ .sort((a, b) => (b.online - a.online) || (b.lastVerifiedAt - a.lastVerifiedAt));
+ out.push({
+ cid: c.cid, videoId: c.video_id, size: Number(c.size), height: Number(c.height), vcodec: c.vcodec,
+ serverHas: !!(await serverHas(c.cid)),
+ counts: { holders: holders.length, online: holders.filter((h) => h.online).length, fresh: holders.filter((h) => !h.stale).length },
+ holders,
+ });
+ }
+ return { ok: true, staleDays, cids: out };
+}
--- /dev/null
+++ b/server/p2p-hub.test.js
@@ -0,0 +1,87 @@
+// Presence/signalling hub and the holders payload (plan 014).
+import { test, expect, beforeAll } from 'bun:test';
+import { mkdtempSync } from 'node:fs';
+import { tmpdir } from 'node:os';
+import { join } from 'node:path';
+import { createHash } from 'node:crypto';
+
+const root = mkdtempSync(join(tmpdir(), 'ytp-p2p-hub-'));
+process.env.DB_PATH = join(root, 'test.db');
+const dbmod = await import('./db.js');
+const p2pDb = await import('./p2p-db.js');
+const { createP2pHub, holdersPayload, peerIdOf } = await import('./p2p-hub.js');
+
+const sha = (s) => createHash('sha256').update(s).digest('hex');
+const SECRET_A = '1'.repeat(64);
+const SECRET_B = '2'.repeat(64);
+const CID = 'a'.repeat(64);
+const DAY = 86400_000;
+const NOW = Date.UTC(2026, 8, 29, 12);
+
+function fakeWs(data) {
+ return { data, sent: [], closed: null, send(s) { this.sent.push(JSON.parse(s)); }, close(code) { this.closed = code; } };
+}
+
+beforeAll(async () => {
+ await dbmod.initDb();
+ await p2pDb.initP2pSchema();
+ await p2pDb.createDevice({ deviceId: 'dev_000000000000000a', secretHash: sha(SECRET_A), now: NOW });
+ await p2pDb.createDevice({ deviceId: 'dev_000000000000000b', secretHash: sha(SECRET_B), now: NOW });
+ await p2pDb.createDevice({ deviceId: 'dev_000000000000000c', secretHash: sha('x'), now: NOW });
+ await p2pDb.touchDevice('dev_000000000000000c', { now: NOW, share: false });
+ await p2pDb.upsertContent({ cid: CID, videoId: 'dQw4w9WgXcQ', size: 1234, height: 720, vcodec: 'h264', origin: 'server', now: NOW });
+ await p2pDb.upsertHolder({ cid: CID, deviceId: 'dev_000000000000000a', now: NOW - 1 * DAY });
+ await p2pDb.upsertHolder({ cid: CID, deviceId: 'dev_000000000000000b', now: NOW - 30 * DAY }); // stale, still listed
+ await p2pDb.upsertHolder({ cid: CID, deviceId: 'dev_000000000000000c', now: NOW }); // share off → hidden
+});
+
+test('auth on open, hello with opaque peer id, bad secret is closed', async () => {
+ const hub = createP2pHub({ getDevice: p2pDb.getDevice });
+ const a = fakeWs({ budget: [] });
+ hub.websocket.open(a);
+ await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: 'x', data: {} })); // ignored before auth
+ expect(a.sent.length).toBe(0);
+ await hub.websocket.message(a, JSON.stringify({ type: 'auth', device: 'dev_000000000000000a', secret: SECRET_A }));
+ expect(a.sent[0]).toEqual({ type: 'hello', peer: peerIdOf('dev_000000000000000a') });
+ expect(hub.isOnline('dev_000000000000000a')).toBe(true);
+ const bad = fakeWs({ budget: [] });
+ hub.websocket.open(bad);
+ await hub.websocket.message(bad, JSON.stringify({ type: 'auth', device: 'dev_000000000000000b', secret: SECRET_A }));
+ expect(bad.closed).toBe(4401);
+ expect(hub.isOnline('dev_000000000000000b')).toBe(false);
+});
+
+test('signals are relayed between online peers only; close drops presence', async () => {
+ const hub = createP2pHub({ getDevice: p2pDb.getDevice });
+ const a = fakeWs({ budget: [] });
+ const b = fakeWs({ budget: [] });
+ hub.websocket.open(a);
+ hub.websocket.open(b);
+ await hub.websocket.message(a, JSON.stringify({ type: 'auth', device: 'dev_000000000000000a', secret: SECRET_A }));
+ await hub.websocket.message(b, JSON.stringify({ type: 'auth', device: 'dev_000000000000000b', secret: SECRET_B }));
+ await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: b.data.peer, data: { sdp: 'x' } }));
+ expect(b.sent.at(-1)).toEqual({ type: 'signal', from: a.data.peer, data: { sdp: 'x' } });
+ hub.websocket.close(b);
+ expect(hub.isOnline('dev_000000000000000b')).toBe(false);
+ await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: b.data.peer, data: {} }));
+ expect(a.sent.at(-1)).toMatchObject({ type: 'error', error: 'peer offline' });
+ expect(hub.send('dev_000000000000000a', { type: 'ping' })).toBe(true);
+ expect(hub.send('dev_000000000000000b', { type: 'ping' })).toBe(false);
+});
+
+test('holders payload: persistent rows, online flag, stale marker, share-off hidden', async () => {
+ const online = new Set(['dev_000000000000000b']);
+ const p = await holdersPayload({
+ videoId: 'dQw4w9WgXcQ', p2pDb, isOnline: (d) => online.has(d), staleDays: 7,
+ serverHas: async () => false, now: NOW,
+ });
+ expect(p.staleDays).toBe(7);
+ expect(p.cids.length).toBe(1);
+ const c = p.cids[0];
+ expect(c).toMatchObject({ cid: CID, size: 1234, height: 720, serverHas: false, counts: { holders: 2, online: 1, fresh: 1 } });
+ expect(c.holders.map((h) => [h.peer, h.online, h.stale])).toEqual([
+ [peerIdOf('dev_000000000000000b'), true, true], // online first even though stale
+ [peerIdOf('dev_000000000000000a'), false, false],
+ ]);
+ expect(JSON.stringify(p)).not.toContain('dev_'); // never leak device ids
+});

View File

@@ -0,0 +1,107 @@
--- /dev/null
+++ b/frontend/p2p-core.js
@@ -0,0 +1,55 @@
+/* ============================================================================
+ * p2p-core.js — pure helpers for peer-to-peer UI (window.P2PCore / node)
+ *
+ * formatAvailability(payload, now) turns GET /api/p2p/holders into the line
+ * under the now-playing title, e.g.
+ * "📡 On 3 devices · 1 online now · last checked 2d ago"
+ * Holders are persistent (docs/p2p-architecture.md): they are never hidden
+ * for age; when every holder is older than the server's staleDays the line
+ * says so and is marked stale.
+ * ========================================================================== */
+(function (root) {
+ 'use strict';
+
+ function ago(ms, now) {
+ const s = Math.max(0, Math.round((now - ms) / 1000));
+ if (s < 90) return 'just now';
+ const m = Math.round(s / 60);
+ if (m < 90) return m + 'm ago';
+ const h = Math.round(m / 60);
+ if (h < 48) return h + 'h ago';
+ return Math.round(h / 24) + 'd ago';
+ }
+
+ function formatAvailability(payload, now) {
+ if (!payload || !payload.ok || !Array.isArray(payload.cids)) return null;
+ const holders = [];
+ for (const c of payload.cids) for (const h of c.holders || []) holders.push(h);
+ if (!holders.length) return null;
+ const online = holders.filter((h) => h.online).length;
+ const fresh = holders.filter((h) => !h.stale).length;
+ const newest = Math.max(...holders.map((h) => Number(h.lastVerifiedAt) || 0));
+ const n = holders.length;
+ const allStale = fresh === 0;
+ const text = `📡 On ${n} device${n === 1 ? '' : 's'} · ${online ? online + ' online now' : 'none online'}`
+ + ` · last checked ${ago(newest, now)}` + (allStale ? ' (not checked recently)' : '');
+ const title = holders.map((h) =>
+ `${h.peer} · ${h.online ? 'online' : 'offline'} · checked ${ago(Number(h.lastVerifiedAt) || 0, now)}`
+ + `${h.stale ? ' (stale)' : ''}${h.trust === 'challenged' ? ' · spot-checked' : ''}`).join('\n');
+ return { text, title, stale: allStale, online, holders: n };
+ }
+
+ // Holders worth trying for a download: online first, then freshest check.
+ function pickPeers(payload, cid) {
+ if (!payload || !Array.isArray(payload.cids)) return [];
+ const c = payload.cids.find((x) => x.cid === cid) || payload.cids[0];
+ if (!c) return [];
+ return (c.holders || []).filter((h) => h.online)
+ .sort((a, b) => (b.trust === 'challenged') - (a.trust === 'challenged') || b.lastVerifiedAt - a.lastVerifiedAt)
+ .map((h) => ({ peer: h.peer, cid: c.cid, size: c.size }));
+ }
+
+ const P2PCore = { ago, formatAvailability, pickPeers };
+ if (typeof module !== 'undefined' && module.exports) module.exports = P2PCore;
+ else root.P2PCore = P2PCore;
+})(typeof globalThis !== 'undefined' ? globalThis : this);
--- /dev/null
+++ b/frontend/p2p-core.test.js
@@ -0,0 +1,46 @@
+'use strict';
+const { test } = require('node:test');
+const assert = require('node:assert');
+const C = require('./p2p-core');
+
+const NOW = Date.UTC(2026, 8, 29, 12);
+const H = 3600_000;
+const payload = (holders, extra = {}) => ({ ok: true, staleDays: 7, cids: [{ cid: 'a'.repeat(64), size: 10, holders, ...extra }] });
+
+test('ago buckets', () => {
+ assert.strictEqual(C.ago(NOW - 30_000, NOW), 'just now');
+ assert.strictEqual(C.ago(NOW - 20 * 60_000, NOW), '20m ago');
+ assert.strictEqual(C.ago(NOW - 5 * H, NOW), '5h ago');
+ assert.strictEqual(C.ago(NOW - 72 * H, NOW), '3d ago');
+});
+
+test('no holders → nothing to show', () => {
+ assert.strictEqual(C.formatAvailability(payload([]), NOW), null);
+ assert.strictEqual(C.formatAvailability(null, NOW), null);
+});
+
+test('mixed holders: counts, newest check, not stale', () => {
+ const f = C.formatAvailability(payload([
+ { peer: 'p1', online: true, lastVerifiedAt: NOW - 30 * 24 * H, stale: true, trust: 'reported' },
+ { peer: 'p2', online: false, lastVerifiedAt: NOW - 48 * H, stale: false, trust: 'challenged' },
+ ]), NOW);
+ assert.strictEqual(f.text, '📡 On 2 devices · 1 online now · last checked 2d ago');
+ assert.strictEqual(f.stale, false);
+ assert.match(f.title, /p1 · online · checked 30d ago \(stale\)/);
+ assert.match(f.title, /p2 · offline · checked 2d ago · spot-checked/);
+});
+
+test('every holder stale → still listed, flagged', () => {
+ const f = C.formatAvailability(payload([{ peer: 'p1', online: false, lastVerifiedAt: NOW - 20 * 24 * H, stale: true }]), NOW);
+ assert.strictEqual(f.text, '📡 On 1 device · none online · last checked 20d ago (not checked recently)');
+ assert.strictEqual(f.stale, true);
+});
+
+test('pickPeers: online only, spot-checked first', () => {
+ const p = payload([
+ { peer: 'off', online: false, lastVerifiedAt: NOW, trust: 'challenged' },
+ { peer: 'on1', online: true, lastVerifiedAt: NOW, trust: 'reported' },
+ { peer: 'on2', online: true, lastVerifiedAt: NOW - H, trust: 'challenged' },
+ ]);
+ assert.deepStrictEqual(C.pickPeers(p).map((x) => x.peer), ['on2', 'on1']);
+});

View File

@@ -0,0 +1,91 @@
--- a/frontend/p2p-client.js
+++ b/frontend/p2p-client.js
@@ -6,6 +6,9 @@
* changed() a save/delete happened → re-report soon
* device() { deviceId, secret } or null
* authHeaders() { 'X-Device': … } for other P2P calls
+ * contribute(videoId) / unknownVideos()
+ * hand a saved file the server doesn't
+ * know yet to /api/p2p/intake (plan 016)
* onMessage(type, fn) / signal(to, data) / peer()
* live /ws/p2p socket (plan 014): the
* device is "online" while it is open
@@ -33,6 +36,7 @@
let again = false;
let timer = null;
let lastSync = 0;
+ let lastUnknown = []; // cids the server did not recognise at the last report
function device() {
try {
@@ -130,6 +134,7 @@
if (r.status === 401) { try { localStorage.removeItem(KEY); } catch { /* ignore */ } return null; }
const j = await r.json();
if (!j || !j.ok) return null;
+ lastUnknown = Array.isArray(j.unknown) ? j.unknown : [];
const accepted = new Set(j.accepted || []);
for (const rec of recs) {
if (rec.cid && accepted.has(rec.cid) && rec.state !== 'verified') { rec.state = 'verified'; await window.DeviceDB.putFile(rec); }
@@ -151,6 +156,33 @@
return j;
}
+ // ---- intake (plan 016) ---------------------------------------------------------
+ // Send one saved video to the server so it can hash + validate it and add its
+ // cid to the catalog. Uploads the whole file: only on the user's request
+ // ("Verify & share") or when the server asks for a copy (plan 018).
+ async function contribute(videoId, { cid = null, title = '', channel = '' } = {}) {
+ const dev = await ensureDevice();
+ const rec = await window.DeviceDB.getFile(videoId);
+ const file = window.OPFS.getFileObject ? await window.OPFS.getFileObject(videoId) : null;
+ if (!file) return { ok: false, error: 'not saved on this device' };
+ const h = { 'X-Device': dev.deviceId + '.' + dev.secret };
+ const t = await (await fetch('/api/p2p/intake', {
+ method: 'POST', headers: { ...h, 'Content-Type': 'application/json' },
+ body: JSON.stringify({ videoId, cid: cid || (rec && rec.cid) || undefined, size: file.size, title, channel }),
+ })).json().catch(() => ({ ok: false, error: 'server unreachable' }));
+ if (!t.ok || t.known) { if (t.known) changed(); return t; }
+ const r = await (await fetch(t.url, { method: 'PUT', headers: h, body: file }))
+ .json().catch(() => ({ ok: false, error: 'upload failed' }));
+ if (r.ok) changed();
+ return r;
+ }
+
+ async function unknownVideos() {
+ if (!lastUnknown.length || !window.DeviceDB) return [];
+ const recs = await window.DeviceDB.listFiles();
+ return recs.filter((r) => r.cid && lastUnknown.includes(r.cid)).map((r) => r.videoId);
+ }
+
// ---- presence socket (/ws/p2p) — plan 014 ------------------------------------
// Open while sharing OR receiving is on. Credentials go in the first message,
// never the URL. Reconnects with backoff (5 s … 5 min).
@@ -258,7 +290,7 @@
}
window.P2PClient = {
- start, changed, sync, device, authHeaders, config,
+ start, changed, sync, device, authHeaders, config, contribute, unknownVideos,
onMessage, signal, peer: () => myPeer, isConnected: () => !!(ws && ws.readyState === 1 && myPeer),
};
}());
--- a/frontend/opfs.js
+++ b/frontend/opfs.js
@@ -108,6 +108,17 @@
}
},
+ // The saved File itself (disk-backed, not read into memory) — used as an
+ // upload body by P2P intake. null when not saved.
+ async getFileObject(videoId) {
+ try {
+ const found = await findHandle(videoId);
+ return found ? await found[0].getFile() : null;
+ } catch {
+ return null;
+ }
+ },
+
// Bytes [offset, offset+length) of a saved video — answers the server's
// P2P range challenges (p2p-client.js). null when the file is missing.
async readRange(videoId, offset, length) {

View File

@@ -0,0 +1,81 @@
--- a/server/media-cache.js
+++ b/server/media-cache.js
@@ -522,6 +522,42 @@
}
// The "Broken" button: drop the copy and fetch it again at high priority.
+ // Promote a file that came from somewhere else (P2P intake / rehydrate from
+ // a device) as this video's server copy. The caller has already hashed it
+ // and run validateMedia(). The bytes are kept EXACTLY (no remux — the cid
+ // must stay true); only the audio sidecar is derived. Never replaces an
+ // existing ready copy or races a running fetch.
+ async function adoptFile(id, src, { sha256, probe, meta = {} }) {
+ if (!isMediaId(id)) throw new Error('adopt: bad video id');
+ const cur = await db.getMedia(id);
+ if (cur && cur.status === 'ready' && readyFilesExist(cur)) return { adopted: false, reason: 'already cached' };
+ if (jobs.has(id)) return { adopted: false, reason: 'fetch in progress' };
+ await makeRoom(id, probe.size);
+ const prefix = `${id}-${now()}-adopt`;
+ const m4a = join(tmpDir, `${prefix}.m4a`);
+ try {
+ const ex = await run(ffmpeg, ['-v', 'error', '-nostdin', '-y', '-i', src,
+ '-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 gen = ((cur && cur.gen) || 0) + 1;
+ try { renameSync(src, fileFor(id, gen)); } catch { copyFileSync(src, fileFor(id, gen)); unlinkSync(src); }
+ renameSync(m4a, fileFor(id, gen, 'm4a'));
+ const t = now();
+ await db.upsertMedia(id, {
+ status: 'ready', gen, size: probe.size + sizeOf(fileFor(id, gen, 'm4a')),
+ height: probe.height, vcodec: probe.vcodec, acodec: probe.acodec, duration: probe.duration,
+ optimized: OPT_REV, // re-encoding would change the bytes other devices hold
+ meta: JSON.stringify(meta), priority: HIGH, auto: 0, attempts: 0, error: null, retry_at: 0,
+ created_at: (cur && cur.created_at) || t, updated_at: t, last_access: t, sha256,
+ });
+ removeFiles(id, gen);
+ log.info?.(`[media] adopted ${id} ${probe.height}p ${fmtMB(probe.size)} from P2P`);
+ return { adopted: true, gen };
+ } finally {
+ sweepTmp(prefix);
+ }
+ }
+
async function redownload(id) {
if (!isMediaId(id)) throw new MediaSkip('invalid video id');
const existing = jobs.get(id);
@@ -676,5 +712,5 @@
return n;
}
- return { init, ensureCached, redownload, verify, getReady, filePath, touch, status, stats, isMediaId, backfillHashes };
+ return { init, ensureCached, redownload, verify, getReady, filePath, touch, status, stats, isMediaId, backfillHashes, adoptFile };
}
--- a/server/media-cache.test.js
+++ b/server/media-cache.test.js
@@ -261,6 +261,26 @@
expect(await dbmod.getMedia('retAAAAAAA1')).not.toBeNull();
});
+ test('adoptFile promotes a validated outside file byte-for-byte, once', async () => {
+ await clearDb();
+ const { cache, dir, calls } = makeCache();
+ await cache.init();
+ const src = join(root, 'adopt-src.mp4');
+ copyFileSync(fx('good-fs.mp4'), src);
+ const bytes = readFileSync(src);
+ const probe = await validateMedia(src, 0);
+ const r = await cache.adoptFile('adoptAAAAA1', src, { sha256: 'e'.repeat(64), probe, meta: { title: 'from a device' } });
+ expect(r.adopted).toBe(true);
+ const row = await cache.getReady('adoptAAAAA1');
+ expect(row.sha256).toBe('e'.repeat(64));
+ expect(row.optimized).toBe(OPT_REV);
+ expect(readFileSync(join(dir, `adoptAAAAA1.${row.gen}.mp4`)).equals(bytes)).toBe(true);
+ expect(existsSync(join(dir, `adoptAAAAA1.${row.gen}.m4a`))).toBe(true);
+ expect(calls.length).toBe(0); // nothing fetched from the source
+ copyFileSync(fx('good-fs.mp4'), src);
+ expect((await cache.adoptFile('adoptAAAAA1', src, { sha256: 'f'.repeat(64), probe })).adopted).toBe(false);
+ });
+
test('free-disk guard skips caching', async () => {
await clearDb();
const { cache, calls } = makeCache({ freeBytes: () => 1024 ** 3, minFreeBytes: 5 * 1024 ** 3 });

View File

@@ -0,0 +1,230 @@
--- /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/);
+});

View File

@@ -0,0 +1,280 @@
--- /dev/null
+++ b/frontend/p2p-transfer.js
@@ -0,0 +1,184 @@
+/* ============================================================================
+ * p2p-transfer.js — download a verified file from another device over a
+ * WebRTC data channel (docs/p2p-architecture.md flow 6). window.P2PTransfer.
+ *
+ * start({ canShare }) serve requests from other devices (called once)
+ * download({ videoId, cid, size, peers, onProgress }) → { ok, error? }
+ *
+ * Signalling rides P2PClient.signal()/onMessage('signal') over /ws/p2p:
+ * {k:'offer', xid, cid, sdp} · {k:'answer', xid, sdp} · {k:'ice', xid, cand}
+ * {k:'deny', xid, reason}
+ * Data channel "file" (ordered): requester sends {t:'get'}; sender answers
+ * {t:'meta', size}, 64 KiB binary frames (back-pressured on bufferedAmount),
+ * then {t:'end'}. The receiver writes through p2p-recv-worker.js, which only
+ * commits the file when its SHA-256 equals the cid. STUN only — devices
+ * behind strict NATs won't connect (same as watch-party voice); the next
+ * holder is tried.
+ * ========================================================================== */
+(function () {
+ 'use strict';
+
+ const ICE = [{ urls: 'stun:stun.l.google.com:19302' }, { urls: 'stun:stun1.l.google.com:19302' }];
+ const FRAME = 64 * 1024;
+ const HIGH_WATER = 4 * 1024 * 1024;
+ const OPEN_TIMEOUT = 20_000;
+ const IDLE_TIMEOUT = 30_000;
+ const MAX_SERVING = 1;
+
+ let iceServers = ICE;
+ let canShare = () => true;
+ const sessions = new Map(); // xid -> { pc, onSignal }
+ let serving = 0;
+
+ const newXid = () => Math.random().toString(36).slice(2) + Date.now().toString(36);
+
+ function onSignal(m) {
+ const d = m && m.data;
+ if (!d || !d.xid) return;
+ const s = sessions.get(d.xid);
+ if (s) { s.onSignal(d, m.from); return; }
+ if (d.k === 'offer') serve(d, m.from).catch(() => {});
+ }
+
+ // ---- sender -------------------------------------------------------------------
+ async function serve(offer, from) {
+ const deny = (reason) => window.P2PClient.signal(from, { k: 'deny', xid: offer.xid, reason });
+ if (!canShare()) return deny('sharing off');
+ if (serving >= MAX_SERVING) return deny('busy');
+ const rec = await window.DeviceDB.getByCid(offer.cid);
+ const file = rec && window.OPFS.getFileObject ? await window.OPFS.getFileObject(rec.videoId) : null;
+ if (!rec || !file || file.size !== rec.size) return deny('not here');
+ serving++;
+ const pc = new RTCPeerConnection({ iceServers });
+ let done = false;
+ const finish = () => {
+ if (done) return;
+ done = true;
+ serving--;
+ sessions.delete(offer.xid);
+ try { pc.close(); } catch { /* closed */ }
+ };
+ sessions.set(offer.xid, {
+ pc,
+ onSignal: (d) => { if (d.k === 'ice' && d.cand) pc.addIceCandidate(d.cand).catch(() => {}); },
+ });
+ pc.onicecandidate = (e) => { if (e.candidate) window.P2PClient.signal(from, { k: 'ice', xid: offer.xid, cand: e.candidate.toJSON() }); };
+ pc.onconnectionstatechange = () => { if (['failed', 'closed', 'disconnected'].includes(pc.connectionState)) finish(); };
+ setTimeout(() => { if (pc.connectionState !== 'connected') finish(); }, OPEN_TIMEOUT);
+ pc.ondatachannel = (e) => {
+ const dc = e.channel;
+ dc.bufferedAmountLowThreshold = 1024 * 1024;
+ dc.onmessage = async (ev) => {
+ let msg;
+ try { msg = JSON.parse(ev.data); } catch { return; }
+ if (msg.t !== 'get') return;
+ try {
+ dc.send(JSON.stringify({ t: 'meta', size: file.size }));
+ for (let pos = 0; pos < file.size && !done; pos += 4 * FRAME) {
+ const buf = await file.slice(pos, pos + 4 * FRAME).arrayBuffer();
+ for (let i = 0; i < buf.byteLength; i += FRAME) {
+ if (dc.bufferedAmount > HIGH_WATER) {
+ await new Promise((r) => { dc.onbufferedamountlow = () => { dc.onbufferedamountlow = null; r(); }; });
+ }
+ dc.send(buf.slice(i, i + FRAME));
+ }
+ }
+ dc.send(JSON.stringify({ t: 'end' }));
+ } catch { finish(); }
+ };
+ dc.onclose = finish;
+ };
+ await pc.setRemoteDescription({ type: 'offer', sdp: offer.sdp });
+ const answer = await pc.createAnswer();
+ await pc.setLocalDescription(answer);
+ window.P2PClient.signal(from, { k: 'answer', xid: offer.xid, sdp: answer.sdp });
+ }
+
+ // ---- receiver -----------------------------------------------------------------
+ function tryPeer({ videoId, cid, size, peer, onProgress }) {
+ return new Promise((resolve) => {
+ const xid = newXid();
+ const pc = new RTCPeerConnection({ iceServers });
+ const worker = new Worker('/p2p-recv-worker.js');
+ let settled = false;
+ let idle = null;
+ let expected = size;
+ const end = (res) => {
+ if (settled) return;
+ settled = true;
+ clearTimeout(idle);
+ clearTimeout(openTimer);
+ sessions.delete(xid);
+ try { pc.close(); } catch { /* closed */ }
+ if (!res.ok) worker.postMessage({ op: 'abort' });
+ setTimeout(() => worker.terminate(), res.ok ? 0 : 2000);
+ resolve(res);
+ };
+ const bump = () => { clearTimeout(idle); idle = setTimeout(() => end({ ok: false, error: 'peer went quiet' }), IDLE_TIMEOUT); };
+ const openTimer = setTimeout(() => end({ ok: false, error: 'could not connect' }), OPEN_TIMEOUT);
+
+ worker.onmessage = (e) => {
+ const m = e.data || {};
+ if (m.op === 'progress' && onProgress) onProgress(m.received, expected);
+ if (m.op === 'done') end(m.ok ? { ok: true, sha256: m.sha256, size: m.size } : { ok: false, error: m.error });
+ };
+ worker.postMessage({ op: 'open', videoId });
+
+ sessions.set(xid, {
+ pc,
+ onSignal: (d) => {
+ if (d.k === 'answer') pc.setRemoteDescription({ type: 'answer', sdp: d.sdp }).catch(() => end({ ok: false, error: 'bad answer' }));
+ else if (d.k === 'ice' && d.cand) pc.addIceCandidate(d.cand).catch(() => {});
+ else if (d.k === 'deny') end({ ok: false, error: 'peer declined: ' + d.reason });
+ },
+ });
+ pc.onicecandidate = (e) => { if (e.candidate) window.P2PClient.signal(peer, { k: 'ice', xid, cand: e.candidate.toJSON() }); };
+ const dc = pc.createDataChannel('file', { ordered: true });
+ dc.binaryType = 'arraybuffer';
+ dc.onopen = () => { clearTimeout(openTimer); bump(); dc.send(JSON.stringify({ t: 'get' })); };
+ dc.onmessage = (e) => {
+ bump();
+ if (typeof e.data === 'string') {
+ let m;
+ try { m = JSON.parse(e.data); } catch { return; }
+ if (m.t === 'meta') {
+ if (m.size !== size) end({ ok: false, error: 'peer has a different file size' });
+ expected = m.size;
+ } else if (m.t === 'end') {
+ clearTimeout(idle);
+ worker.postMessage({ op: 'finish', cid, size });
+ }
+ return;
+ }
+ worker.postMessage({ op: 'chunk', buf: e.data }, [e.data]);
+ };
+ dc.onclose = () => { if (!settled) setTimeout(() => end({ ok: false, error: 'channel closed' }), 5000); };
+ (async () => {
+ try {
+ const offer = await pc.createOffer();
+ await pc.setLocalDescription(offer);
+ if (!window.P2PClient.signal(peer, { k: 'offer', xid, cid, sdp: offer.sdp })) end({ ok: false, error: 'not connected to the P2P hub' });
+ } catch (err) { end({ ok: false, error: err.message }); }
+ })();
+ });
+ }
+
+ // Try each online holder in turn until one delivers a verified file.
+ async function download({ videoId, cid, size, peers, onProgress }) {
+ if (!window.P2PClient || !window.P2PClient.isConnected()) return { ok: false, error: 'not connected to the P2P hub' };
+ let last = { ok: false, error: 'no online device has this video' };
+ for (const p of peers || []) {
+ last = await tryPeer({ videoId, cid, size, peer: p.peer, onProgress });
+ if (last.ok) return last;
+ }
+ return last;
+ }
+
+ function start(opts = {}) {
+ if (opts.canShare) canShare = opts.canShare;
+ if (opts.iceServers) iceServers = opts.iceServers;
+ window.P2PClient.onMessage('signal', onSignal);
+ }
+
+ window.P2PTransfer = { start, download };
+}());
--- /dev/null
+++ b/frontend/p2p-recv-worker.js
@@ -0,0 +1,90 @@
+/* ============================================================================
+ * p2p-recv-worker.js — writes a file arriving from another device into OPFS
+ * while hashing it; commits it ONLY when the SHA-256 equals the expected
+ * content id (docs/p2p-architecture.md flow 6).
+ *
+ * In: { op:'open', videoId } start videos/<videoId>.p2p.part
+ * { op:'chunk', buf } ArrayBuffer (transferred)
+ * { op:'finish', cid, size } verify + rename to <videoId>.mp4
+ * { op:'abort' } drop the partial file
+ * Out: { op:'opened' } | { op:'progress', received } |
+ * { op:'done', ok:true, sha256, size } | { op:'done', ok:false, error }
+ * The ".part" suffix keeps OPFS.listVideos() from ever listing a partial file.
+ * ========================================================================== */
+'use strict';
+importScripts('/sha256.js');
+
+let dir = null;
+let handle = null;
+let access = null;
+let partName = '';
+let finalName = '';
+let hasher = null;
+let offset = 0;
+let lastProgress = 0;
+
+async function cleanup() {
+ try { if (access) access.close(); } catch { /* closed */ }
+ access = null;
+ try { if (dir && partName) await dir.removeEntry(partName); } catch { /* gone */ }
+}
+
+// Messages are handled strictly one after another: an async handler would
+// otherwise let 'chunk' or 'abort' run while 'open' is still awaiting OPFS.
+let chain = Promise.resolve();
+self.onmessage = (e) => { chain = chain.then(() => onOp(e.data || {})); };
+
+async function onOp(m) {
+ try {
+ if (m.op === 'open') {
+ const root = await navigator.storage.getDirectory();
+ dir = await root.getDirectoryHandle('videos', { create: true });
+ partName = `${m.videoId}.p2p.part`;
+ finalName = `${m.videoId}.mp4`;
+ handle = await dir.getFileHandle(partName, { create: true });
+ access = await handle.createSyncAccessHandle();
+ access.truncate(0);
+ hasher = self.Sha256.create();
+ offset = 0;
+ self.postMessage({ op: 'opened' });
+ } else if (m.op === 'chunk') {
+ const u8 = new Uint8Array(m.buf);
+ access.write(u8, { at: offset });
+ hasher.update(u8);
+ offset += u8.byteLength;
+ if (offset - lastProgress > 1024 * 1024) { lastProgress = offset; self.postMessage({ op: 'progress', received: offset }); }
+ } else if (m.op === 'finish') {
+ access.truncate(offset);
+ access.flush();
+ access.close();
+ access = null;
+ const sha256 = hasher.hex();
+ if (offset !== m.size) throw new Error(`size mismatch (${offset} of ${m.size})`);
+ if (sha256 !== m.cid) throw new Error('content hash mismatch');
+ try { await dir.removeEntry(finalName); } catch { /* no previous copy */ }
+ let renamed = false;
+ if (typeof handle.move === 'function') { try { await handle.move(finalName); renamed = true; } catch { /* copy below */ } }
+ if (!renamed) {
+ const out = await (await dir.getFileHandle(finalName, { create: true })).createSyncAccessHandle();
+ try {
+ const file = await handle.getFile();
+ for (let pos = 0; pos < file.size; pos += 8 * 1024 * 1024) {
+ const buf = new Uint8Array(await file.slice(pos, pos + 8 * 1024 * 1024).arrayBuffer());
+ out.write(buf, { at: pos });
+ }
+ out.truncate(file.size);
+ out.flush();
+ } finally { out.close(); }
+ await dir.removeEntry(partName);
+ }
+ partName = '';
+ self.postMessage({ op: 'done', ok: true, sha256, size: offset });
+ } else if (m.op === 'abort') {
+ await cleanup();
+ self.postMessage({ op: 'done', ok: false, error: 'aborted' });
+ }
+ } catch (err) {
+ await cleanup();
+ self.postMessage({ op: 'done', ok: false, error: err && err.message ? err.message : String(err) });
+ }
+}

View File

@@ -0,0 +1,40 @@
--- a/frontend/p2p-client.js
+++ b/frontend/p2p-client.js
@@ -160,7 +160,7 @@
// Send one saved video to the server so it can hash + validate it and add its
// cid to the catalog. Uploads the whole file: only on the user's request
// ("Verify & share") or when the server asks for a copy (plan 018).
- async function contribute(videoId, { cid = null, title = '', channel = '' } = {}) {
+ async function contribute(videoId, { cid = null, title = '', channel = '', restore = false } = {}) {
const dev = await ensureDevice();
const rec = await window.DeviceDB.getFile(videoId);
const file = window.OPFS.getFileObject ? await window.OPFS.getFileObject(videoId) : null;
@@ -168,7 +168,7 @@
const h = { 'X-Device': dev.deviceId + '.' + dev.secret };
const t = await (await fetch('/api/p2p/intake', {
method: 'POST', headers: { ...h, 'Content-Type': 'application/json' },
- body: JSON.stringify({ videoId, cid: cid || (rec && rec.cid) || undefined, size: file.size, title, channel }),
+ body: JSON.stringify({ videoId, cid: cid || (rec && rec.cid) || undefined, size: file.size, title, channel, restore }),
})).json().catch(() => ({ ok: false, error: 'server unreachable' }));
if (!t.ok || t.known) { if (t.known) changed(); return t; }
const r = await (await fetch(t.url, { method: 'PUT', headers: h, body: file }))
@@ -280,6 +280,19 @@
timer = setTimeout(sync, 5000);
}
+ // Plan 018: the server lost its copy and the source is gone — it asks one
+ // holder to send the file back. Only while sharing is on, one at a time,
+ // and only for the exact cid this device holds.
+ let restoring = false;
+ onMessage('upload-request', async (m) => {
+ if (restoring || hooks.getSettings().p2pShare === false) return;
+ const rec = window.DeviceDB ? await window.DeviceDB.getFile(m.videoId) : null;
+ if (!rec || rec.cid !== m.cid) return;
+ restoring = true;
+ try { await contribute(m.videoId, { cid: m.cid, restore: true }); } catch { /* next request retries */ }
+ finally { restoring = false; }
+ });
+
function start(h) {
hooks = { ...hooks, ...(h || {}) };
const idle = window.requestIdleCallback || ((fn) => setTimeout(fn, 1));

View File

@@ -0,0 +1,125 @@
--- a/server/p2p-hub.js
+++ b/server/p2p-hub.js
@@ -118,3 +118,28 @@
}
return { ok: true, staleDays, cids: out };
}
+
+// Flow 8 (plan 018): the source is gone and the server evicted its copy, so
+// ask ONE online holder that shares to upload it through intake
+// ({type:'upload-request', videoId, cid}). At most one request per video per
+// 10 minutes; returns true when a device was asked (or recently was).
+export function createRehydrator({ p2pDb, hub, hasServerCopy, enabled = () => true, now = () => Date.now() }) {
+ const asked = new Map(); // videoId -> ms
+ return async function rehydrate(videoId) {
+ if (!enabled()) return false;
+ if (await hasServerCopy(videoId)) return false;
+ const last = asked.get(videoId);
+ if (last && now() - last < 10 * 60_000) return true;
+ for (const c of await p2pDb.listContentForVideo(videoId)) {
+ for (const h of await p2pDb.listHolders(c.cid, 50)) {
+ if (Number(h.share) !== 1 || !hub.isOnline(h.device_id)) continue;
+ if (hub.send(h.device_id, { type: 'upload-request', videoId, cid: c.cid })) {
+ asked.set(videoId, now());
+ if (asked.size > 5000) asked.delete(asked.keys().next().value);
+ return true;
+ }
+ }
+ }
+ return false;
+ };
+}
--- a/server/p2p-hub.test.js
+++ b/server/p2p-hub.test.js
@@ -85,3 +85,21 @@
]);
expect(JSON.stringify(p)).not.toContain('dev_'); // never leak device ids
});
+
+test('rehydrator asks one online sharing holder, once per 10 min, only when the server has no copy', async () => {
+ const { createRehydrator } = await import('./p2p-hub.js');
+ const sent = [];
+ let t = NOW;
+ let serverHas = false;
+ const hub = { isOnline: (d) => d === 'dev_000000000000000a' || d === 'dev_000000000000000c', send: (d, m) => { sent.push([d, m]); return true; } };
+ const rehydrate = createRehydrator({ p2pDb, hub, hasServerCopy: async () => serverHas, now: () => t });
+ expect(await rehydrate('dQw4w9WgXcQ')).toBe(true);
+ // c is online but has sharing off; a is online and shares.
+ expect(sent).toEqual([['dev_000000000000000a', { type: 'upload-request', videoId: 'dQw4w9WgXcQ', cid: CID }]]);
+ expect(await rehydrate('dQw4w9WgXcQ')).toBe(true); // recently asked — no second message
+ expect(sent.length).toBe(1);
+ t += 11 * 60_000;
+ serverHas = true;
+ expect(await rehydrate('dQw4w9WgXcQ')).toBe(false); // server copy is back
+ expect(await rehydrate('unknownVid1')).toBe(false); // nobody holds it
+});
--- a/server/p2p-intake.js
+++ b/server/p2p-intake.js
@@ -3,8 +3,10 @@
* (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? }
+ * 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)
@@ -30,6 +32,7 @@
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 }
@@ -48,7 +51,8 @@
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 });
+ 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();
--- a/server/p2p-intake.test.js
+++ b/server/p2p-intake.test.js
@@ -1,6 +1,6 @@
// 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 { mkdtempSync, readFileSync, readdirSync, unlinkSync } from 'node:fs';
import { spawnSync } from 'node:child_process';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
@@ -64,6 +64,27 @@
expect(await (await req('/api/p2p/intake', { method: 'POST', body: { videoId: 'upAAAAAAAA1', cid: goodCid, size: good.length } })).json()).toEqual({ ok: true, known: true });
});
+test('restore: a known cid is uploaded again only when the server lost its copy', async () => {
+ let serverHas = true;
+ const app3 = new Hono();
+ const p2p = registerP2pRoutes(app3, { cfg, p2pDb, fileForCid: async () => null, sha256Range: async () => '', log: quiet });
+ const got = [];
+ registerIntakeRoutes(app3, {
+ cfg, p2pDb, gate: p2p.gate, requireDevice: p2p.requireDevice, validateMedia, admitFile, log: quiet,
+ serverHasCid: async () => serverHas,
+ adopt: async (videoId, path, info) => { got.push(info.sha256); unlinkSync(path); return { adopted: true }; }, // a real adopt moves the file
+ });
+ const h = { 'X-Device': dev.deviceId + '.' + dev.secret, 'Content-Type': 'application/json' };
+ const open = async () => (await app3.request('/api/p2p/intake', { method: 'POST', headers: h, body: JSON.stringify({ videoId: 'upAAAAAAAA1', cid: goodCid, size: good.length, restore: true }) })).json();
+ expect(await open()).toEqual({ ok: true, known: true }); // server still has it
+ serverHas = false;
+ const t = await open();
+ expect(t.ticket).toMatch(/^[0-9a-f]{32}$/);
+ const r = await (await app3.request(t.url, { method: 'PUT', headers: { 'X-Device': h['X-Device'] }, body: good })).json();
+ expect(r).toEqual({ ok: true, cid: goodCid, adopted: true });
+ expect(got).toEqual([goodCid]);
+});
+
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/],