Let a device hand a file to the server for hashing and validation
This commit is contained in:
@@ -27,6 +27,8 @@ services:
|
|||||||
# P2P_KEEP_MIN_VIEWS: "3" # server keeps copies with ≥ this many views…
|
# P2P_KEEP_MIN_VIEWS: "3" # server keeps copies with ≥ this many views…
|
||||||
# P2P_KEEP_DAYS: "30" # …in this many days
|
# P2P_KEEP_DAYS: "30" # …in this many days
|
||||||
# P2P_KEEP_RECENT_DAYS: "14" # …or played this recently
|
# P2P_KEEP_RECENT_DAYS: "14" # …or played this recently
|
||||||
|
# P2P_INTAKE_DIR: "/app/data/p2p-intake" # quarantine for device uploads (never served)
|
||||||
|
# P2P_INTAKE_MAX_BYTES: "3221225472" # 3 GiB
|
||||||
# Optional: force yt-dlp search instead of InnerTube API
|
# Optional: force yt-dlp search instead of InnerTube API
|
||||||
# SEARCH_INNERTUBE: "0" # force yt-dlp search
|
# SEARCH_INNERTUBE: "0" # force yt-dlp search
|
||||||
# Optional: override yt-dlp binary path if you mount a custom one
|
# Optional: override yt-dlp binary path if you mount a custom one
|
||||||
|
|||||||
@@ -7505,6 +7505,13 @@ async function renderSettings() {
|
|||||||
<span>This device<small id="p2pDeviceInfo">…</small></span>
|
<span>This device<small id="p2pDeviceInfo">…</small></span>
|
||||||
<span id="p2pHoldingCount" class="set-stat">…</span>
|
<span id="p2pHoldingCount" class="set-stat">…</span>
|
||||||
</div>
|
</div>
|
||||||
|
<div class="set-row hidden" id="p2pContributeRow">
|
||||||
|
<span>
|
||||||
|
Saved videos the server can't verify yet
|
||||||
|
<small>Uploading lets the server check them (hash + media validation) so other devices can get them. Uses your upload bandwidth.</small>
|
||||||
|
</span>
|
||||||
|
<button id="p2pContributeBtn" class="btn">Verify & share</button>
|
||||||
|
</div>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
<div class="set-group">
|
<div class="set-group">
|
||||||
@@ -7695,6 +7702,26 @@ async function renderSettings() {
|
|||||||
: P.isConnected() ? `Online · peer ${P.peer()}` : 'Registered · not connected';
|
: P.isConnected() ? `Online · peer ${P.peer()}` : 'Registered · not connected';
|
||||||
const files = window.DeviceDB ? await window.DeviceDB.listFiles() : [];
|
const files = window.DeviceDB ? await window.DeviceDB.listFiles() : [];
|
||||||
cnt.textContent = `${files.filter((f) => f.state === 'verified').length} verified / ${files.length} saved`;
|
cnt.textContent = `${files.filter((f) => f.state === 'verified').length} verified / ${files.length} saved`;
|
||||||
|
const unknown = window.P2PClient ? await window.P2PClient.unknownVideos() : [];
|
||||||
|
const row = $('p2pContributeRow');
|
||||||
|
const btn = $('p2pContributeBtn');
|
||||||
|
if (row && btn && unknown.length) {
|
||||||
|
row.classList.remove('hidden');
|
||||||
|
btn.textContent = `Verify & share ${unknown.length} saved video${unknown.length === 1 ? '' : 's'}`;
|
||||||
|
btn.onclick = async () => {
|
||||||
|
btn.disabled = true;
|
||||||
|
let ok = 0;
|
||||||
|
for (let i = 0; i < unknown.length; i++) {
|
||||||
|
const v = videoById(unknown[i]) || {};
|
||||||
|
btn.textContent = `Uploading ${i + 1} of ${unknown.length}…`;
|
||||||
|
const r = await window.P2PClient.contribute(unknown[i], { title: v.title || '', channel: v.channel || '' });
|
||||||
|
if (r && r.ok) ok++;
|
||||||
|
}
|
||||||
|
toast(`Verified ${ok} of ${unknown.length} saved video${unknown.length === 1 ? '' : 's'}`);
|
||||||
|
btn.disabled = false;
|
||||||
|
row.classList.add('hidden');
|
||||||
|
};
|
||||||
|
}
|
||||||
})();
|
})();
|
||||||
|
|
||||||
// ---- Cache management ----
|
// ---- Cache management ----
|
||||||
|
|||||||
@@ -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
|
// Bytes [offset, offset+length) of a saved video — answers the server's
|
||||||
// P2P range challenges (p2p-client.js). null when the file is missing.
|
// P2P range challenges (p2p-client.js). null when the file is missing.
|
||||||
async readRange(videoId, offset, length) {
|
async readRange(videoId, offset, length) {
|
||||||
|
|||||||
@@ -6,6 +6,9 @@
|
|||||||
* changed() a save/delete happened → re-report soon
|
* changed() a save/delete happened → re-report soon
|
||||||
* device() { deviceId, secret } or null
|
* device() { deviceId, secret } or null
|
||||||
* authHeaders() { 'X-Device': … } for other P2P calls
|
* 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()
|
* onMessage(type, fn) / signal(to, data) / peer()
|
||||||
* live /ws/p2p socket (plan 014): the
|
* live /ws/p2p socket (plan 014): the
|
||||||
* device is "online" while it is open
|
* device is "online" while it is open
|
||||||
@@ -33,6 +36,7 @@
|
|||||||
let again = false;
|
let again = false;
|
||||||
let timer = null;
|
let timer = null;
|
||||||
let lastSync = 0;
|
let lastSync = 0;
|
||||||
|
let lastUnknown = []; // cids the server did not recognise at the last report
|
||||||
|
|
||||||
function device() {
|
function device() {
|
||||||
try {
|
try {
|
||||||
@@ -130,6 +134,7 @@
|
|||||||
if (r.status === 401) { try { localStorage.removeItem(KEY); } catch { /* ignore */ } return null; }
|
if (r.status === 401) { try { localStorage.removeItem(KEY); } catch { /* ignore */ } return null; }
|
||||||
const j = await r.json();
|
const j = await r.json();
|
||||||
if (!j || !j.ok) return null;
|
if (!j || !j.ok) return null;
|
||||||
|
lastUnknown = Array.isArray(j.unknown) ? j.unknown : [];
|
||||||
const accepted = new Set(j.accepted || []);
|
const accepted = new Set(j.accepted || []);
|
||||||
for (const rec of recs) {
|
for (const rec of recs) {
|
||||||
if (rec.cid && accepted.has(rec.cid) && rec.state !== 'verified') { rec.state = 'verified'; await window.DeviceDB.putFile(rec); }
|
if (rec.cid && accepted.has(rec.cid) && rec.state !== 'verified') { rec.state = 'verified'; await window.DeviceDB.putFile(rec); }
|
||||||
@@ -151,6 +156,33 @@
|
|||||||
return j;
|
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 ------------------------------------
|
// ---- presence socket (/ws/p2p) — plan 014 ------------------------------------
|
||||||
// Open while sharing OR receiving is on. Credentials go in the first message,
|
// Open while sharing OR receiving is on. Credentials go in the first message,
|
||||||
// never the URL. Reconnects with backoff (5 s … 5 min).
|
// never the URL. Reconnects with backoff (5 s … 5 min).
|
||||||
@@ -258,7 +290,7 @@
|
|||||||
}
|
}
|
||||||
|
|
||||||
window.P2PClient = {
|
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),
|
onMessage, signal, peer: () => myPeer, isConnected: () => !!(ws && ws.readyState === 1 && myPeer),
|
||||||
};
|
};
|
||||||
}());
|
}());
|
||||||
|
|||||||
@@ -23,7 +23,7 @@ green, app boots with no JS errors, P2P on by default, offline boot works).
|
|||||||
| 013 | 013-device-identity-and-holdings-3ba493 | Register devices and report verified holdings to the server | done | Register devices and report verified holdings to the server | persistent holders, no TTL |
|
| 013 | 013-device-identity-and-holdings-3ba493 | Register devices and report verified holdings to the server | done | Register devices and report verified holdings to the server | persistent holders, no TTL |
|
||||||
| 014 | 014-p2p-presence-hub-ceced8 | Add the /ws/p2p presence and signalling hub and the holders endpoint | done | Add the /ws/p2p presence and signalling hub and the holders endpoint | stale flag, never hidden |
|
| 014 | 014-p2p-presence-hub-ceced8 | Add the /ws/p2p presence and signalling hub and the holders endpoint | done | Add the /ws/p2p presence and signalling hub and the holders endpoint | stale flag, never hidden |
|
||||||
| 015 | 015-availability-ui-and-settings-3b9397 | Show peer availability with stale markers and add Sharing settings | done | Show peer availability with stale markers and add Sharing settings | |
|
| 015 | 015-availability-ui-and-settings-3b9397 | Show peer availability with stale markers and add Sharing settings | done | Show peer availability with stale markers and add Sharing settings | |
|
||||||
| 016 | 016-intake-and-server-verification-cfe031 | Let a device hand a file to the server for hashing and validation | in-progress | | needs ffmpeg for tests |
|
| 016 | 016-intake-and-server-verification-cfe031 | Let a device hand a file to the server for hashing and validation | done | Let a device hand a file to the server for hashing and validation | needs ffmpeg for tests |
|
||||||
| 017 | 017-peer-transfer-1faaa7 | Download a verified file from another device over WebRTC | queued | | STUN only |
|
| 017 | 017-peer-transfer-1faaa7 | Download a verified file from another device over WebRTC | queued | | STUN only |
|
||||||
| 018 | 018-server-rehydrate-from-peer-4fb8bd | Restore an evicted server copy from an online holder | queued | | |
|
| 018 | 018-server-rehydrate-from-peer-4fb8bd | Restore an evicted server copy from an online holder | queued | | |
|
||||||
| 019 | 019-admin-p2p-panel-4dc623 | Add a P2P panel to the admin page | queued | | |
|
| 019 | 019-admin-p2p-panel-4dc623 | Add a P2P panel to the admin page | queued | | |
|
||||||
|
|||||||
@@ -148,3 +148,9 @@ Output ONLY the following, no other prose:
|
|||||||
3. `Findings:` — max 10 lines.
|
3. `Findings:` — max 10 lines.
|
||||||
|
|
||||||
Do not commit. Do not push. Do not touch files outside the Steps.
|
Do not commit. Do not push. Do not touch files outside the Steps.
|
||||||
|
|
||||||
|
## Execution log
|
||||||
|
|
||||||
|
- Executor: in-session Agent (haiku). Attempts: 1. Fix rounds: 0.
|
||||||
|
- Orchestrator re-ran Verification: `p2p-intake.js`, `p2p-intake.test.js`, `media-cache.js`, `media-cache.test.js`, `p2p-client.js`, `opfs.js` byte-identical to the pre-tested 016 versions; intake tests 4 pass; media-cache 27 pass; all 14 server test files 0 fail; `SERVER_OK`; frontend syntax ok; browser check in Chromium reproduced the expected JSON (unknown file -> contribute -> accepted, cid valid, second report accepted 1 / unknown 0). No leftover processes.
|
||||||
|
- Executor Findings (verbatim): All three patches applied without errors. server/package.json updated to run p2p-intake.test.js. server/server.js updated with validateMedia import and registerIntakeRoutes import and intake route initialization. frontend/app.js updated with HTML and JS for contribute button. docker-compose.yml updated with commented P2P_INTAKE_DIR and P2P_INTAKE_MAX_BYTES. All 4 expected intake tests passed, all 27 media-cache tests passed, full test suite passed with 0 failures across all test files. Server build successful. Frontend syntax check passed. Browser test output matches expected JSON exactly.
|
||||||
@@ -522,6 +522,42 @@ export function createMediaCache({
|
|||||||
}
|
}
|
||||||
|
|
||||||
// The "Broken" button: drop the copy and fetch it again at high priority.
|
// 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) {
|
async function redownload(id) {
|
||||||
if (!isMediaId(id)) throw new MediaSkip('invalid video id');
|
if (!isMediaId(id)) throw new MediaSkip('invalid video id');
|
||||||
const existing = jobs.get(id);
|
const existing = jobs.get(id);
|
||||||
@@ -676,5 +712,5 @@ export function createMediaCache({
|
|||||||
return n;
|
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 };
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -261,6 +261,26 @@ describe('media cache jobs', () => {
|
|||||||
expect(await dbmod.getMedia('retAAAAAAA1')).not.toBeNull();
|
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 () => {
|
test('free-disk guard skips caching', async () => {
|
||||||
await clearDb();
|
await clearDb();
|
||||||
const { cache, calls } = makeCache({ freeBytes: () => 1024 ** 3, minFreeBytes: 5 * 1024 ** 3 });
|
const { cache, calls } = makeCache({ freeBytes: () => 1024 ** 3, minFreeBytes: 5 * 1024 ** 3 });
|
||||||
|
|||||||
121
server/p2p-intake.js
Normal file
121
server/p2p-intake.js
Normal file
@@ -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 };
|
||||||
|
}
|
||||||
103
server/p2p-intake.test.js
Normal file
103
server/p2p-intake.test.js
Normal file
@@ -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/);
|
||||||
|
});
|
||||||
@@ -6,7 +6,7 @@
|
|||||||
"scripts": {
|
"scripts": {
|
||||||
"start": "bun server.js",
|
"start": "bun server.js",
|
||||||
"dev": "bun --hot server.js",
|
"dev": "bun --hot server.js",
|
||||||
"test": "bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./ytdlp-pool.test.js && bun test ./p2p-db.test.js && bun test ./p2p-admit.test.js && bun test ./p2p-retention.test.js && bun test ./p2p-routes.test.js && bun test ./p2p-hub.test.js"
|
"test": "bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./ytdlp-pool.test.js && bun test ./p2p-db.test.js && bun test ./p2p-admit.test.js && bun test ./p2p-retention.test.js && bun test ./p2p-routes.test.js && bun test ./p2p-hub.test.js && bun test --timeout 60000 ./p2p-intake.test.js"
|
||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@hono/node-server": "^1.14.0",
|
"@hono/node-server": "^1.14.0",
|
||||||
|
|||||||
@@ -42,7 +42,7 @@ import { createHash } from 'node:crypto';
|
|||||||
import { brotliCompressSync, constants as zlibConstants } from 'node:zlib';
|
import { brotliCompressSync, constants as zlibConstants } from 'node:zlib';
|
||||||
import { initDb, upsertUser, recordVideoAccess, getUserData, createProfile, getProfile, saveProfile, createSharedPlaylist, getSharedPlaylist, queueInboxPlaylist, listInbox, deleteInboxItem, countInbox,
|
import { initDb, upsertUser, recordVideoAccess, getUserData, createProfile, getProfile, saveProfile, createSharedPlaylist, getSharedPlaylist, queueInboxPlaylist, listInbox, deleteInboxItem, countInbox,
|
||||||
getMedia, upsertMedia, deleteMedia, listMedia, listMediaLru, touchMedia, mediaStats } from './db.js';
|
getMedia, upsertMedia, deleteMedia, listMedia, listMediaLru, touchMedia, mediaStats } from './db.js';
|
||||||
import { createMediaCache, HIGH, LOW } from './media-cache.js';
|
import { createMediaCache, HIGH, LOW, validateMedia } from './media-cache.js';
|
||||||
import * as notesDb from './db.js';
|
import * as notesDb from './db.js';
|
||||||
import { registerNoteRoutes, parseLrc, sanitizeLyrics } from './notes.js';
|
import { registerNoteRoutes, parseLrc, sanitizeLyrics } from './notes.js';
|
||||||
import { createRemoteHub } from './remote.js';
|
import { createRemoteHub } from './remote.js';
|
||||||
@@ -52,6 +52,7 @@ import { admitFile } from './p2p-admit.js';
|
|||||||
import { P2P } from './p2p-config.js';
|
import { P2P } from './p2p-config.js';
|
||||||
import * as p2pDb from './p2p-db.js';
|
import * as p2pDb from './p2p-db.js';
|
||||||
import { registerP2pRoutes } from './p2p-routes.js';
|
import { registerP2pRoutes } from './p2p-routes.js';
|
||||||
|
import { registerIntakeRoutes } from './p2p-intake.js';
|
||||||
import { createP2pHub, holdersPayload } from './p2p-hub.js';
|
import { createP2pHub, holdersPayload } from './p2p-hub.js';
|
||||||
import { sha256Range } from './hash.js';
|
import { sha256Range } from './hash.js';
|
||||||
import * as innertube from './innertube.js';
|
import * as innertube from './innertube.js';
|
||||||
@@ -1981,6 +1982,15 @@ app.get('/api/p2p/holders', p2p.gate, async (c) => {
|
|||||||
return c.json(payload, 200, { 'Cache-Control': 'no-store' });
|
return c.json(payload, 200, { 'Cache-Control': 'no-store' });
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// Device → server intake: the server hashes, validates (and scans when
|
||||||
|
// P2P_MALWARE_SCAN=1) before a cid is admitted. docs/p2p-architecture.md flow 7.
|
||||||
|
const intake = registerIntakeRoutes(app, {
|
||||||
|
cfg: P2P, p2pDb, gate: p2p.gate, requireDevice: p2p.requireDevice, admitFile,
|
||||||
|
validateMedia: (path, expected, opts) => validateMedia(path, expected,
|
||||||
|
{ ...opts, ffmpeg: FFMPEG, ffprobe: process.env.FFPROBE_PATH || 'ffprobe' }),
|
||||||
|
adopt: (videoId, path, info) => media.adoptFile(videoId, path, info),
|
||||||
|
});
|
||||||
|
|
||||||
// ============================================================================
|
// ============================================================================
|
||||||
// GET /sw.js — serve the service worker with BUILD_TAG injected
|
// GET /sw.js — serve the service worker with BUILD_TAG injected
|
||||||
//
|
//
|
||||||
|
|||||||
Reference in New Issue
Block a user