Restore an evicted server copy from an online holder
This commit is contained in:
@@ -160,7 +160,7 @@
|
|||||||
// Send one saved video to the server so it can hash + validate it and add its
|
// 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
|
// 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).
|
// ("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 dev = await ensureDevice();
|
||||||
const rec = await window.DeviceDB.getFile(videoId);
|
const rec = await window.DeviceDB.getFile(videoId);
|
||||||
const file = window.OPFS.getFileObject ? await window.OPFS.getFileObject(videoId) : null;
|
const file = window.OPFS.getFileObject ? await window.OPFS.getFileObject(videoId) : null;
|
||||||
@@ -168,7 +168,7 @@
|
|||||||
const h = { 'X-Device': dev.deviceId + '.' + dev.secret };
|
const h = { 'X-Device': dev.deviceId + '.' + dev.secret };
|
||||||
const t = await (await fetch('/api/p2p/intake', {
|
const t = await (await fetch('/api/p2p/intake', {
|
||||||
method: 'POST', headers: { ...h, 'Content-Type': 'application/json' },
|
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' }));
|
})).json().catch(() => ({ ok: false, error: 'server unreachable' }));
|
||||||
if (!t.ok || t.known) { if (t.known) changed(); return t; }
|
if (!t.ok || t.known) { if (t.known) changed(); return t; }
|
||||||
const r = await (await fetch(t.url, { method: 'PUT', headers: h, body: file }))
|
const r = await (await fetch(t.url, { method: 'PUT', headers: h, body: file }))
|
||||||
@@ -280,6 +280,19 @@
|
|||||||
timer = setTimeout(sync, 5000);
|
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) {
|
function start(h) {
|
||||||
hooks = { ...hooks, ...(h || {}) };
|
hooks = { ...hooks, ...(h || {}) };
|
||||||
const idle = window.requestIdleCallback || ((fn) => setTimeout(fn, 1));
|
const idle = window.requestIdleCallback || ((fn) => setTimeout(fn, 1));
|
||||||
|
|||||||
@@ -25,5 +25,5 @@ green, app boots with no JS errors, P2P on by default, offline boot works).
|
|||||||
| 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 | done | Let a device hand a file to the server for hashing and validation | 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 | done | Download a verified file from another device over WebRTC | STUN only |
|
| 017 | 017-peer-transfer-1faaa7 | Download a verified file from another device over WebRTC | done | Download a verified file from another device over WebRTC | STUN only |
|
||||||
| 018 | 018-server-rehydrate-from-peer-4fb8bd | Restore an evicted server copy from an online holder | in-progress | | |
|
| 018 | 018-server-rehydrate-from-peer-4fb8bd | Restore an evicted server copy from an online holder | done | Restore an evicted server copy from an online holder | |
|
||||||
| 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 | | |
|
||||||
|
|||||||
@@ -118,3 +118,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-hub.js`, `p2p-hub.test.js`, `p2p-intake.js`, `p2p-intake.test.js`, `p2p-client.js` byte-identical to the pre-tested versions; hub tests 4 pass, intake tests 5 pass; all 15 server test files 0 fail; `SERVER_OK`; frontend syntax ok; browser check in Chromium reproduced the expected JSON (first `/api/streams` -> `restoring: true`, holder uploaded, server adopted the exact cid). No leftover processes.
|
||||||
|
- Executor Findings (verbatim): All patches applied successfully without conflicts. All 4 plan steps executed: (1) both git patches applied, (2) import modified to include createRehydrator, (3) serverHasCid option added to registerIntakeRoutes, (4) p2pRehydrate instance created with correct parameters, (5) /api/streams catch block replaced with rehydration logic. Hub tests: 4 pass. Intake tests: 5 pass. All server tests pass (0 fail across all suites). Browser test confirms complete rehydration flow: device detects request, uploads file, server adopts it into cache. Expected JSON output matched exactly.
|
||||||
@@ -118,3 +118,28 @@ export async function holdersPayload({ videoId, cid, p2pDb, isOnline, staleDays,
|
|||||||
}
|
}
|
||||||
return { ok: true, staleDays, cids: out };
|
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;
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|||||||
@@ -85,3 +85,21 @@ test('holders payload: persistent rows, online flag, stale marker, share-off hid
|
|||||||
]);
|
]);
|
||||||
expect(JSON.stringify(p)).not.toContain('dev_'); // never leak device ids
|
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
|
||||||
|
});
|
||||||
|
|||||||
@@ -3,8 +3,10 @@
|
|||||||
* (docs/p2p-architecture.md flow 7). The ONLY path by which bytes that did
|
* (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.
|
* 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
|
* → { 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
|
* → { ok, ticket, url, expiresAt } PUT the bytes to url within 30 min
|
||||||
* PUT /api/p2p/intake/:ticket (device) raw file body
|
* PUT /api/p2p/intake/:ticket (device) raw file body
|
||||||
* → { ok, cid, adopted } admitted (+ adopted into the media cache)
|
* → { ok, cid, adopted } admitted (+ adopted into the media cache)
|
||||||
@@ -30,6 +32,7 @@ const clean = (v, max) => String(v || '').replace(/[\u0000-\u001f\u007f<>]+/g, '
|
|||||||
export function registerIntakeRoutes(app, deps) {
|
export function registerIntakeRoutes(app, deps) {
|
||||||
const {
|
const {
|
||||||
cfg, p2pDb, gate, requireDevice, validateMedia, admitFile, adopt = null,
|
cfg, p2pDb, gate, requireDevice, validateMedia, admitFile, adopt = null,
|
||||||
|
serverHasCid = async () => true, // (cid) → does the server still hold these bytes?
|
||||||
now = () => Date.now(), log = console,
|
now = () => Date.now(), log = console,
|
||||||
} = deps;
|
} = deps;
|
||||||
const tickets = new Map(); // ticket -> { deviceId, videoId, cid, size, expiresAt, busy }
|
const tickets = new Map(); // ticket -> { deviceId, videoId, cid, size, expiresAt, busy }
|
||||||
@@ -48,7 +51,8 @@ export function registerIntakeRoutes(app, deps) {
|
|||||||
if (!(size > 0) || size > cfg.intakeMaxBytes) return c.json({ ok: false, error: 'bad size' }, 413);
|
if (!(size > 0) || size > cfg.intakeMaxBytes) return c.json({ ok: false, error: 'bad size' }, 413);
|
||||||
if (cid) {
|
if (cid) {
|
||||||
const known = await p2pDb.getContent(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);
|
if (known && known.status === 'revoked') return c.json({ ok: false, error: 'revoked' }, 410);
|
||||||
}
|
}
|
||||||
sweep();
|
sweep();
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
// Device → server intake: hash, validate, (scan), admit (plan 016). Needs ffmpeg.
|
// Device → server intake: hash, validate, (scan), admit (plan 016). Needs ffmpeg.
|
||||||
import { test, expect, beforeAll } from 'bun:test';
|
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 { spawnSync } from 'node:child_process';
|
||||||
import { tmpdir } from 'node:os';
|
import { tmpdir } from 'node:os';
|
||||||
import { join } from 'node:path';
|
import { join } from 'node:path';
|
||||||
@@ -64,6 +64,27 @@ test('a valid file is hashed by the server, validated, admitted and its uploader
|
|||||||
expect(await (await req('/api/p2p/intake', { method: 'POST', body: { videoId: 'upAAAAAAAA1', cid: goodCid, size: good.length } })).json()).toEqual({ ok: true, known: true });
|
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 () => {
|
test('claimed cid mismatch, truncated media, and non-media are rejected and deleted', async () => {
|
||||||
const cases = [
|
const cases = [
|
||||||
[{ cid: 'f'.repeat(64), bytes: good }, 400, /hash mismatch/],
|
[{ cid: 'f'.repeat(64), bytes: good }, 400, /hash mismatch/],
|
||||||
|
|||||||
@@ -53,7 +53,7 @@ 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 { registerIntakeRoutes } from './p2p-intake.js';
|
||||||
import { createP2pHub, holdersPayload } from './p2p-hub.js';
|
import { createP2pHub, holdersPayload, createRehydrator } 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';
|
||||||
import QRCode from 'qrcode';
|
import QRCode from 'qrcode';
|
||||||
@@ -711,6 +711,16 @@ app.get('/api/streams', async (c) => {
|
|||||||
},
|
},
|
||||||
});
|
});
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
|
// The source failed. If a device holds a verified copy, ask it to send one
|
||||||
|
// to the server so the video comes back (P2P flow 8).
|
||||||
|
let restoring = false;
|
||||||
|
try { restoring = await p2pRehydrate(videoId); } catch { /* best effort */ }
|
||||||
|
if (restoring) {
|
||||||
|
return c.json({
|
||||||
|
ok: false, restoring: true,
|
||||||
|
error: 'This video is unavailable at the source — a device that has it is sending a copy to the server. Try again in a minute.',
|
||||||
|
}, 503);
|
||||||
|
}
|
||||||
return c.json({ ok: false, error: err.message }, 500);
|
return c.json({ ok: false, error: err.message }, 500);
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
@@ -1989,6 +1999,14 @@ const intake = registerIntakeRoutes(app, {
|
|||||||
validateMedia: (path, expected, opts) => validateMedia(path, expected,
|
validateMedia: (path, expected, opts) => validateMedia(path, expected,
|
||||||
{ ...opts, ffmpeg: FFMPEG, ffprobe: process.env.FFPROBE_PATH || 'ffprobe' }),
|
{ ...opts, ffmpeg: FFMPEG, ffprobe: process.env.FFPROBE_PATH || 'ffprobe' }),
|
||||||
adopt: (videoId, path, info) => media.adoptFile(videoId, path, info),
|
adopt: (videoId, path, info) => media.adoptFile(videoId, path, info),
|
||||||
|
serverHasCid: async (cid) => !!(await fileForCid(cid)),
|
||||||
|
});
|
||||||
|
|
||||||
|
// Source gone + server copy evicted → ask one online holder to send it back
|
||||||
|
// through intake (docs/p2p-architecture.md flow 8).
|
||||||
|
const p2pRehydrate = createRehydrator({
|
||||||
|
p2pDb, hub: p2pHub, enabled: () => P2P.enabled,
|
||||||
|
hasServerCopy: async (id) => !!(await media.getReady(id).catch(() => null)),
|
||||||
});
|
});
|
||||||
|
|
||||||
// ============================================================================
|
// ============================================================================
|
||||||
|
|||||||
Reference in New Issue
Block a user