Add P2P tables, config flags and db helpers
This commit is contained in:
236
server/p2p-db.js
Normal file
236
server/p2p-db.js
Normal file
@@ -0,0 +1,236 @@
|
||||
/* ============================================================================
|
||||
* p2p-db.js — tables + queries for peer-to-peer sharing
|
||||
* (docs/p2p-architecture.md). Shares the libsql client from db.js.
|
||||
*
|
||||
* p2p_content one row per verified file (cid = sha256 of the bytes); never
|
||||
* deleted, only revoked — the catalog grows over time
|
||||
* p2p_devices registered devices (secret stored as sha256)
|
||||
* p2p_holders which device holds which cid; PERSISTENT (no TTL) with
|
||||
* last_verified_at — the UI decides what is "stale"
|
||||
* video_views per-video per-day view counts (retention criteria)
|
||||
* All timestamps are ms epochs.
|
||||
* ========================================================================== */
|
||||
import { db } from './db.js';
|
||||
|
||||
export async function initP2pSchema() {
|
||||
await db.executeMultiple(`
|
||||
CREATE TABLE IF NOT EXISTS p2p_content (
|
||||
cid TEXT PRIMARY KEY,
|
||||
video_id TEXT NOT NULL,
|
||||
size INTEGER NOT NULL,
|
||||
height INTEGER NOT NULL DEFAULT 0,
|
||||
vcodec TEXT,
|
||||
acodec TEXT,
|
||||
duration REAL NOT NULL DEFAULT 0,
|
||||
meta TEXT NOT NULL DEFAULT '{}',
|
||||
origin TEXT NOT NULL, -- server | intake
|
||||
status TEXT NOT NULL DEFAULT 'verified', -- verified | revoked
|
||||
scan TEXT NOT NULL DEFAULT 'skipped', -- skipped | clean
|
||||
created_at INTEGER NOT NULL,
|
||||
verified_at INTEGER NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_p2p_content_video ON p2p_content (video_id, created_at DESC);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS p2p_devices (
|
||||
device_id TEXT PRIMARY KEY,
|
||||
secret_hash TEXT NOT NULL,
|
||||
fingerprint TEXT,
|
||||
profile TEXT,
|
||||
share INTEGER NOT NULL DEFAULT 1,
|
||||
created_at INTEGER NOT NULL,
|
||||
last_seen_at INTEGER NOT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS p2p_holders (
|
||||
cid TEXT NOT NULL,
|
||||
device_id TEXT NOT NULL,
|
||||
status TEXT NOT NULL DEFAULT 'active', -- active | removed
|
||||
trust TEXT NOT NULL DEFAULT 'reported', -- reported | challenged
|
||||
first_reported_at INTEGER NOT NULL,
|
||||
last_verified_at INTEGER NOT NULL,
|
||||
removed_at INTEGER,
|
||||
PRIMARY KEY (cid, device_id)
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_p2p_holders_device ON p2p_holders (device_id, status);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS video_views (
|
||||
video_id TEXT NOT NULL,
|
||||
day TEXT NOT NULL,
|
||||
n INTEGER NOT NULL DEFAULT 0,
|
||||
PRIMARY KEY (video_id, day)
|
||||
);
|
||||
`);
|
||||
// media_cache.sha256 — the cid of the current <id>.<gen>.mp4 (plan 009).
|
||||
try { await db.execute('ALTER TABLE media_cache ADD COLUMN sha256 TEXT'); }
|
||||
catch (e) { if (!/duplicate column/i.test(String(e.message))) throw e; }
|
||||
}
|
||||
|
||||
const rowsOf = (r) => r.rows.map((row) => {
|
||||
const o = {};
|
||||
r.columns.forEach((c, i) => { o[c] = row[i]; });
|
||||
return o;
|
||||
});
|
||||
|
||||
// ---- content ----------------------------------------------------------------
|
||||
|
||||
export async function upsertContent(c) {
|
||||
await db.execute({
|
||||
sql: `INSERT INTO p2p_content (cid, video_id, size, height, vcodec, acodec, duration, meta, origin, status, scan, created_at, verified_at)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 'verified', ?, ?, ?)
|
||||
ON CONFLICT(cid) DO UPDATE SET verified_at = excluded.verified_at, scan = excluded.scan`,
|
||||
args: [c.cid, c.videoId, c.size, c.height || 0, c.vcodec || null, c.acodec || null, c.duration || 0,
|
||||
JSON.stringify(c.meta || {}), c.origin, c.scan || 'skipped', c.now, c.now],
|
||||
});
|
||||
}
|
||||
|
||||
export async function getContent(cid) {
|
||||
const r = await db.execute({ sql: 'SELECT * FROM p2p_content WHERE cid = ?', args: [cid] });
|
||||
return rowsOf(r)[0] || null;
|
||||
}
|
||||
|
||||
export async function listContentForVideo(videoId) {
|
||||
const r = await db.execute({
|
||||
sql: "SELECT * FROM p2p_content WHERE video_id = ? AND status = 'verified' ORDER BY created_at DESC LIMIT 10",
|
||||
args: [videoId],
|
||||
});
|
||||
return rowsOf(r);
|
||||
}
|
||||
|
||||
export async function knownCids(cids) {
|
||||
if (!cids.length) return new Set();
|
||||
const out = new Set();
|
||||
for (let i = 0; i < cids.length; i += 200) {
|
||||
const part = cids.slice(i, i + 200);
|
||||
const r = await db.execute({
|
||||
sql: `SELECT cid FROM p2p_content WHERE status = 'verified' AND cid IN (${part.map(() => '?').join(',')})`,
|
||||
args: part,
|
||||
});
|
||||
for (const row of r.rows) out.add(row[0]);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
export async function revokeContent(cid) {
|
||||
await db.execute({ sql: "UPDATE p2p_content SET status = 'revoked' WHERE cid = ?", args: [cid] });
|
||||
}
|
||||
|
||||
// ---- devices ------------------------------------------------------------------
|
||||
|
||||
export async function createDevice({ deviceId, secretHash, fingerprint, profile, now }) {
|
||||
await db.execute({
|
||||
sql: `INSERT INTO p2p_devices (device_id, secret_hash, fingerprint, profile, share, created_at, last_seen_at)
|
||||
VALUES (?, ?, ?, ?, 1, ?, ?)`,
|
||||
args: [deviceId, secretHash, fingerprint || null, profile || null, now, now],
|
||||
});
|
||||
}
|
||||
|
||||
export async function getDevice(deviceId) {
|
||||
const r = await db.execute({ sql: 'SELECT * FROM p2p_devices WHERE device_id = ?', args: [deviceId] });
|
||||
return rowsOf(r)[0] || null;
|
||||
}
|
||||
|
||||
export async function touchDevice(deviceId, { now, share, profile } = {}) {
|
||||
await db.execute({
|
||||
sql: `UPDATE p2p_devices SET last_seen_at = ?,
|
||||
share = COALESCE(?, share), profile = COALESCE(?, profile)
|
||||
WHERE device_id = ?`,
|
||||
args: [now, share === undefined ? null : (share ? 1 : 0), profile || null, deviceId],
|
||||
});
|
||||
}
|
||||
|
||||
// ---- holders (persistent; never expired by time) ------------------------------
|
||||
|
||||
export async function upsertHolder({ cid, deviceId, trust = 'reported', now }) {
|
||||
await db.execute({
|
||||
sql: `INSERT INTO p2p_holders (cid, device_id, status, trust, first_reported_at, last_verified_at)
|
||||
VALUES (?, ?, 'active', ?, ?, ?)
|
||||
ON CONFLICT(cid, device_id) DO UPDATE SET
|
||||
status = 'active', removed_at = NULL, last_verified_at = excluded.last_verified_at,
|
||||
trust = CASE WHEN p2p_holders.trust = 'challenged' OR excluded.trust = 'challenged'
|
||||
THEN 'challenged' ELSE 'reported' END`,
|
||||
args: [cid, deviceId, trust, now, now],
|
||||
});
|
||||
}
|
||||
|
||||
export async function setHolderTrust({ cid, deviceId, trust, now }) {
|
||||
await db.execute({
|
||||
sql: 'UPDATE p2p_holders SET trust = ?, last_verified_at = ? WHERE cid = ? AND device_id = ?',
|
||||
args: [trust, now, cid, deviceId],
|
||||
});
|
||||
}
|
||||
|
||||
export async function removeHolder({ cid, deviceId, now }) {
|
||||
await db.execute({
|
||||
sql: "UPDATE p2p_holders SET status = 'removed', removed_at = ? WHERE cid = ? AND device_id = ? AND status = 'active'",
|
||||
args: [now, cid, deviceId],
|
||||
});
|
||||
}
|
||||
|
||||
// A full report: every active holding of this device NOT in `keep` is removed.
|
||||
export async function removeHoldersExcept({ deviceId, keep, now }) {
|
||||
const r = await db.execute({
|
||||
sql: "SELECT cid FROM p2p_holders WHERE device_id = ? AND status = 'active'",
|
||||
args: [deviceId],
|
||||
});
|
||||
const keepSet = new Set(keep);
|
||||
let removed = 0;
|
||||
for (const row of r.rows) {
|
||||
if (keepSet.has(row[0])) continue;
|
||||
await removeHolder({ cid: row[0], deviceId, now });
|
||||
removed++;
|
||||
}
|
||||
return removed;
|
||||
}
|
||||
|
||||
export async function activeHoldingsOf(deviceId) {
|
||||
const r = await db.execute({
|
||||
sql: "SELECT cid FROM p2p_holders WHERE device_id = ? AND status = 'active'",
|
||||
args: [deviceId],
|
||||
});
|
||||
return r.rows.map((row) => row[0]);
|
||||
}
|
||||
|
||||
// Holders of one cid, joined with the device's share flag. Newest check first.
|
||||
export async function listHolders(cid, limit = 50) {
|
||||
const r = await db.execute({
|
||||
sql: `SELECT h.device_id, h.trust, h.first_reported_at, h.last_verified_at, d.share
|
||||
FROM p2p_holders h JOIN p2p_devices d ON d.device_id = h.device_id
|
||||
WHERE h.cid = ? AND h.status = 'active'
|
||||
ORDER BY h.last_verified_at DESC LIMIT ?`,
|
||||
args: [cid, limit],
|
||||
});
|
||||
return rowsOf(r);
|
||||
}
|
||||
|
||||
// ---- views + stats --------------------------------------------------------------
|
||||
|
||||
export function dayKey(ms) {
|
||||
return new Date(ms).toISOString().slice(0, 10);
|
||||
}
|
||||
|
||||
export async function addView(videoId, now) {
|
||||
await db.execute({
|
||||
sql: `INSERT INTO video_views (video_id, day, n) VALUES (?, ?, 1)
|
||||
ON CONFLICT(video_id, day) DO UPDATE SET n = n + 1`,
|
||||
args: [videoId, dayKey(now)],
|
||||
});
|
||||
}
|
||||
|
||||
export async function viewsSince(videoId, sinceMs) {
|
||||
const r = await db.execute({
|
||||
sql: 'SELECT COALESCE(SUM(n), 0) FROM video_views WHERE video_id = ? AND day >= ?',
|
||||
args: [videoId, dayKey(sinceMs)],
|
||||
});
|
||||
return Number(r.rows[0][0]) || 0;
|
||||
}
|
||||
|
||||
export async function p2pStats() {
|
||||
const one = async (sql) => Number((await db.execute(sql)).rows[0][0]) || 0;
|
||||
return {
|
||||
content: await one("SELECT COUNT(*) FROM p2p_content WHERE status = 'verified'"),
|
||||
revoked: await one("SELECT COUNT(*) FROM p2p_content WHERE status = 'revoked'"),
|
||||
devices: await one('SELECT COUNT(*) FROM p2p_devices'),
|
||||
holders: await one("SELECT COUNT(*) FROM p2p_holders WHERE status = 'active'"),
|
||||
heldCids: await one("SELECT COUNT(DISTINCT cid) FROM p2p_holders WHERE status = 'active'"),
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user