Files
NewTube/server/db.mjs
T
bruno 665a0f0ebd
CI / build-and-test (push) Successful in 14m43s
feat(providers): phases 7.3/7.4/7.6/8.1 — provenance, health, NDJSON, contrat unique
7.3: capturedAt/source au registre + 6 adaptateurs + module provenance.ts + ?debug=1 (search-transport.mjs). 7.4: ProviderHealthService + badge source degradee. 7.6: squelettes par provider + snapshots progressifs + transport NDJSON /api/search. 8.1: ProviderAdapter unifie (search enveloppe + channelContent/channelMeta/capabilities) via getProviderAdapter + test de contrat offline.
2026-09-30 07:57:11 -04:00

2171 lines
99 KiB
JavaScript

import fs from 'node:fs';
import path from 'node:path';
import { randomBytes, randomUUID } from 'node:crypto';
import Database from 'better-sqlite3';
const root = process.cwd();
const overrideDbFile = process.env.NEWTUBE_DB_FILE && String(process.env.NEWTUBE_DB_FILE).trim().length
? path.resolve(String(process.env.NEWTUBE_DB_FILE))
: null;
const dbDir = overrideDbFile ? path.dirname(overrideDbFile) : path.join(root, 'db');
const dbFile = overrideDbFile || path.join(dbDir, 'newtube.db');
// Try multiple schema locations to survive when /app/db is a mounted volume
const schemaCandidates = [
path.join(root, 'db', 'schema.sql'), // normal repo path (may be hidden by a volume)
path.join(root, 'db-schema', 'schema.sql'), // immutable path bundled in image
];
const schemaFile = schemaCandidates.find(p => fs.existsSync(p));
if (!fs.existsSync(dbDir)) {
fs.mkdirSync(dbDir, { recursive: true });
}
// Create DB and enable FKs
/**
* Garde-fou d'isolation (Phase 5).
*
* Tous les tests de `server/tests/` sont nommés `*.test.mjs` et DOIVENT
* travailler sur une base temporaire via `NEWTUBE_DB_FILE`. Sans ce garde-fou,
* une faute de frappe sur le nom de la variable (`NEW_TUBE_DB_PATH` au lieu de
* `NEWTUBE_DB_FILE`) fait ouvrir la base de developpement en lecture-ecriture :
* les fixtures du test atterrissent alors chez l'utilisateur, silencieusement.
* C'est deja arrive (lignes de cache factices dans `db/newtube.db`).
*
* On refuse donc explicitement, avec un message qui donne la marche a suivre,
* plutot que de laisser la corruption se produire.
*/
if (!overrideDbFile && process.argv.some((a) => /\.test\.mjs$/.test(a))) {
throw new Error(
`[db] Refus d'ouvrir la base de developpement (${dbFile}) depuis un test. `
+ 'Definissez process.env.NEWTUBE_DB_FILE sur un chemin temporaire AVANT '
+ "l'import de db.mjs (le nom exact est NEWTUBE_DB_FILE, sans \"NEW_TUBE_\").",
);
}
const db = new Database(dbFile);
db.pragma('foreign_keys = ON');
/**
* Chemin reel du fichier ouvert.
*
* Expose pour que les tests puissent AFFIRMER leur isolation au lieu de la
* supposer : un test qui se trompe de nom de variable d'environnement
* (`NEW_TUBE_DB_PATH` au lieu de `NEWTUBE_DB_FILE`) ouvre silencieusement la
* base de dev et y ecrit ses fixtures. Cela a deja pollue `db/newtube.db` avec
* des lignes de cache factices. Un `ok(getDbFile().includes('test'))` en tete de
* fichier echoue bruyamment au lieu de corrompre la base de l'utilisateur.
*/
export function getDbFile() {
return dbFile;
}
// Run schema if present (first boot)
if (schemaFile && fs.existsSync(schemaFile)) {
try {
const ddl = fs.readFileSync(schemaFile, 'utf8');
if (ddl && ddl.trim().length) {
db.exec(ddl);
console.log(`[db] Applied schema from ${schemaFile}`);
}
} catch (e) {
console.warn(`[db] Failed to apply schema from ${schemaFile}:`, e?.message || e);
}
} else {
console.warn('[db] No schema.sql found in expected locations:', schemaCandidates.join(', '));
}
// Lightweight idempotent migration runner.
// Applies SQL files from db/migrations/ in filename order, once each, tracked
// in the migrations table. Safe to run on every boot.
// When /app/db is shadowed by a Docker volume (which hides the image's migrations
// folder), fall back to the immutable /app/db-migrations copy bundled in the image.
(function runMigrations() {
try {
db.exec(`CREATE TABLE IF NOT EXISTS migrations (
name TEXT PRIMARY KEY,
applied_at TEXT NOT NULL
);`);
const migrationsDirs = [
path.join(root, 'db', 'migrations'),
path.join(root, 'db-migrations'),
];
const migrationsDir = migrationsDirs.find(p => fs.existsSync(p));
if (!migrationsDir) return;
const files = fs.readdirSync(migrationsDir)
.filter(f => f.endsWith('.sql'))
.sort();
const applied = new Set(
db.prepare('SELECT name FROM migrations').all().map(r => r.name)
);
for (const file of files) {
if (applied.has(file)) continue;
const sql = fs.readFileSync(path.join(migrationsDir, file), 'utf8');
if (!sql || !sql.trim()) continue;
// Migration files manage their own transactions (some contain BEGIN/COMMIT),
// so execute them as-is instead of wrapping in another transaction.
let failed = false;
try {
db.exec(sql);
} catch (e) {
failed = true;
// Idempotent migrations (IF NOT EXISTS / duplicate columns) can fail on
// databases already up to date — tolerate and continue.
console.warn(`[db] Migration ${file} reported an error:`, e?.message || e);
} finally {
// Une migration `BEGIN … COMMIT` qui échoue AU MILIEU (ex. un ALTER sur une
// colonne déjà présente) laisse la transaction ouverte dans better-sqlite3.
// Toutes les migrations SUIVANTES s'y exécutent alors « avec succès » puis
// sont annulées à la fermeture du process : leurs tables disparaissent
// sans trace alors que le journal affiche « applied ». On referme donc la
// transaction résiduelle avant de poursuivre.
if (db.inTransaction) {
try { db.exec('ROLLBACK'); } catch {}
}
}
// Record as applied even on tolerated errors so we don't re-run/re-warn each boot.
try {
db.prepare('INSERT INTO migrations (name, applied_at) VALUES (?, ?)').run(file, new Date().toISOString());
} catch {}
if (!failed) console.log(`[db] Migration applied: ${file}`);
}
} catch (e) {
console.warn('[db] Migration runner failed:', e?.message || e);
}
})();
// Lightweight schema upgrades for existing databases (SQLite is permissive)
(function ensurePlaylistSchemaUpgrades() {
try {
const colsPlaylists = db.prepare(`PRAGMA table_info(playlists)`).all();
const have = new Set(colsPlaylists.map(c => c.name));
if (!have.has('description')) db.exec(`ALTER TABLE playlists ADD COLUMN description TEXT`);
if (!have.has('thumbnail')) db.exec(`ALTER TABLE playlists ADD COLUMN thumbnail TEXT`);
if (!have.has('is_private')) db.exec(`ALTER TABLE playlists ADD COLUMN is_private INTEGER NOT NULL DEFAULT 1`);
} catch {}
try {
const colsItems = db.prepare(`PRAGMA table_info(playlist_items)`).all();
const have2 = new Set(colsItems.map(c => c.name));
if (!have2.has('thumbnail')) db.exec(`ALTER TABLE playlist_items ADD COLUMN thumbnail TEXT`);
} catch {}
// Phase 7.7 : bannière + description de chaîne. Ajoutés en JS plutôt que dans
// un fichier de migration car `db/schema.sql` est réappliqué à CHAQUE boot :
// un `ALTER TABLE … ADD COLUMN` en migration échouerait (colonne déjà
// présente) sur une base neuve, et ce fichier échouant laissait la
// transaction ouverte, annulant silencieusement les migrations suivantes.
try {
const colsChannels = db.prepare(`PRAGMA table_info(channels)`).all();
const haveCh = new Set(colsChannels.map(c => c.name));
if (!haveCh.has('banner_url')) db.exec(`ALTER TABLE channels ADD COLUMN banner_url TEXT`);
if (!haveCh.has('description')) db.exec(`ALTER TABLE channels ADD COLUMN description TEXT`);
} catch {}
try {
db.exec(`CREATE TABLE IF NOT EXISTS playlist_metrics (
id TEXT PRIMARY KEY,
user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
playlist_id TEXT NOT NULL REFERENCES playlists(id) ON DELETE CASCADE,
action TEXT NOT NULL,
meta_json TEXT,
created_at TEXT NOT NULL
);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_playlist_metrics_user_time ON playlist_metrics(user_id, created_at DESC);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_playlist_metrics_playlist_time ON playlist_metrics(playlist_id, created_at DESC);`);
} catch {}
})();
// (duplicate schema exec removed)
(function ensurePreferencesSchemaUpgrades() {
try {
const colsPrefs = db.prepare(`PRAGMA table_info(user_preferences)`).all();
const have = new Set(colsPrefs.map(c => c.name));
if (!have.has('default_providers')) db.exec(`ALTER TABLE user_preferences ADD COLUMN default_providers TEXT`);
if (!have.has('download_languages')) db.exec(`ALTER TABLE user_preferences ADD COLUMN download_languages TEXT`);
} catch {}
})();
// Helpers
export function nowIso() {
return new Date().toISOString();
}
export function getUserByUsername(username) {
return db.prepare('SELECT * FROM users WHERE username = ?').get(username);
}
export function getUserById(id) {
return db.prepare('SELECT * FROM users WHERE id = ?').get(id);
}
export function getUserByEmail(email) {
const clean = String(email || '').trim().toLowerCase();
if (!clean) return null;
return db.prepare('SELECT * FROM users WHERE LOWER(email) = ?').get(clean) || null;
}
/** Retrouve l'utilisateur lié à un compte OAuth externe (ex. Google sub). */
export function getUserByOAuth(provider, externalUserId) {
try {
ensureOAuthTables();
const row = db.prepare(`SELECT u.* FROM oauth_connections oc
JOIN users u ON u.id = oc.user_id
WHERE oc.provider = ? AND oc.external_user_id = ?`).get(provider, String(externalUserId || ''));
return row || null;
} catch { return null; }
}
export function insertUser({ id, username, email, passwordHash }) {
const ts = nowIso();
db.prepare(`INSERT INTO users (id, username, email, password_hash, is_active, created_at, updated_at)
VALUES (@id, @username, @email, @passwordHash, 1, @ts, @ts)`).run({ id, username, email, passwordHash, ts });
// Default preferences row (default_providers NULL = no multi-provider preference yet)
db.prepare(`INSERT INTO user_preferences (user_id, language, default_provider, default_providers, download_languages, theme, video_quality, region, version, updated_at)
VALUES (@id, 'en', 'youtube', NULL, @dl, 'system', 'auto', 'US', 1, @ts)`).run({ id, ts, dl: JSON.stringify(DEFAULT_DOWNLOAD_LANGUAGES) });
}// Sanitize a defaultProviders list: keep only known provider ids.
const KNOWN_PROVIDER_IDS = ['yt', 'dm', 'tw', 'pt', 'od', 'ru'];
// Langues de sous-titres/transcripts proposées dans les préférences.
// Défaut : français + anglais uniquement.
export const DEFAULT_DOWNLOAD_LANGUAGES = ['fr', 'en'];
export const SUPPORTED_DOWNLOAD_LANGUAGES = [
'fr', 'en', 'es', 'de', 'it', 'pt', 'nl', 'pl',
'ru', 'ar', 'zh', 'ja', 'hi', 'ko',
];
function sanitizeDownloadLanguages(value) {
let arr = value;
// Accept JSON string (from DB) or array (from API patch)
if (typeof arr === 'string') {
try { arr = JSON.parse(arr); } catch { arr = []; }
}
if (!Array.isArray(arr)) arr = [];
const cleaned = arr
.map(v => String(v || '').trim().toLowerCase().replace(/_/g, '-').split('-')[0])
.filter(v => /^[a-z]{2,3}$/.test(v) && SUPPORTED_DOWNLOAD_LANGUAGES.includes(v));
// Deduplicate, preserve order
return Array.from(new Set(cleaned));
}
function parseDownloadLanguagesForApi(stored) {
// NULL (ancienne ligne) = défaut fr+en ; '[]' explicite = aucune langue.
if (stored == null) return [...DEFAULT_DOWNLOAD_LANGUAGES];
try {
const parsed = JSON.parse(stored);
if (!Array.isArray(parsed)) return [...DEFAULT_DOWNLOAD_LANGUAGES];
return sanitizeDownloadLanguages(parsed);
} catch {
return [...DEFAULT_DOWNLOAD_LANGUAGES];
}
}
function sanitizeDefaultProviders(value) {
let arr = value;
// Accept JSON string (from DB) or array (from API patch)
if (typeof arr === 'string') {
try { arr = JSON.parse(arr); } catch { arr = []; }
}
if (!Array.isArray(arr)) arr = [];
const cleaned = arr
.map(v => String(v || '').trim().toLowerCase())
.filter(v => KNOWN_PROVIDER_IDS.includes(v));
// Deduplicate, preserve order
return Array.from(new Set(cleaned));
}
export function upsertPreferences(userId, patch) {
// Normalize the multi-provider selection if provided
let defaultProviders = undefined;
if (patch.defaultProviders !== undefined) {
defaultProviders = sanitizeDefaultProviders(patch.defaultProviders);
}
// Normalize the download/subtitle languages if provided (fr+en par défaut)
let downloadLanguages = undefined;
if (patch.downloadLanguages !== undefined) {
downloadLanguages = sanitizeDownloadLanguages(patch.downloadLanguages);
}
const current = db.prepare('SELECT * FROM user_preferences WHERE user_id = ?').get(userId);
if (!current) {
const merged = {
language: patch.language ?? 'en',
default_provider: patch.defaultProvider ?? 'youtube',
default_providers: defaultProviders !== undefined
? (defaultProviders.length ? JSON.stringify(defaultProviders) : null)
: null,
download_languages: downloadLanguages !== undefined
? JSON.stringify(downloadLanguages)
: JSON.stringify(DEFAULT_DOWNLOAD_LANGUAGES),
theme: patch.theme ?? 'system',
video_quality: patch.videoQuality ?? 'auto',
region: patch.region ?? 'US',
version: 1,
};
db.prepare(`INSERT INTO user_preferences (user_id, language, default_provider, default_providers, download_languages, theme, video_quality, region, version, updated_at)
VALUES (@userId, @language, @default_provider, @default_providers, @download_languages, @theme, @video_quality, @region, @version, @updated_at)`)
.run({ userId, ...merged, updated_at: nowIso() });
} else {
const storedList = current.default_providers;
const nextList = defaultProviders !== undefined
? (defaultProviders.length ? JSON.stringify(defaultProviders) : null)
: storedList;
const storedDl = Object.prototype.hasOwnProperty.call(current, 'download_languages')
? current.download_languages
: undefined;
const nextDl = downloadLanguages !== undefined
? JSON.stringify(downloadLanguages)
: (storedDl !== undefined ? storedDl : JSON.stringify(DEFAULT_DOWNLOAD_LANGUAGES));
const merged = {
language: patch.language ?? current.language,
default_provider: patch.defaultProvider ?? current.default_provider,
default_providers: nextList,
download_languages: nextDl,
theme: patch.theme ?? current.theme,
video_quality: patch.videoQuality ?? current.video_quality,
region: patch.region ?? current.region,
version: (current.version ?? 1) + 1,
updated_at: nowIso(),
};
db.prepare(`UPDATE user_preferences
SET language=@language, default_provider=@default_provider, default_providers=@default_providers,
download_languages=@download_languages,
theme=@theme, video_quality=@video_quality, region=@region, version=@version, updated_at=@updated_at
WHERE user_id=@userId`)
.run({ userId, ...merged });
}
}
export function getPreferences(userId) {
return db.prepare(`SELECT language, default_provider AS defaultProvider, default_providers AS defaultProviders,
download_languages AS downloadLanguages,
theme, video_quality AS videoQuality, region, version, updated_at
FROM user_preferences WHERE user_id = ?`).get(userId);
}
// Return preferences with defaultProviders parsed as an array (API shape)
export function getPreferencesForApi(userId) {
const prefs = getPreferences(userId) || {};
let list = null;
if (prefs.defaultProviders) {
try { list = JSON.parse(prefs.defaultProviders); } catch { list = null; }
}
return {
...prefs,
defaultProviders: Array.isArray(list) ? list : null,
downloadLanguages: parseDownloadLanguagesForApi(prefs.downloadLanguages),
};
}
export function insertSession({ id, userId, refreshTokenHash, isRemember, userAgent, deviceInfo, ip, expiresAt }) {
const ts = nowIso();
db.prepare(`INSERT INTO sessions (id, user_id, refresh_token_hash, user_agent, device_info, ip_address, is_remember, created_at, last_seen_at, expires_at)
VALUES (@id, @userId, @refreshTokenHash, @userAgent, @deviceInfo, @ip, @isRemember, @ts, @ts, @expiresAt)`)
.run({ id, userId, refreshTokenHash, userAgent, deviceInfo, ip, isRemember: isRemember ? 1 : 0, ts, expiresAt });
}
export function getSessionById(id) {
return db.prepare('SELECT * FROM sessions WHERE id = ?').get(id);
}
export function updateSessionToken(id, refreshTokenHash, expiresAt) {
const ts = nowIso();
db.prepare('UPDATE sessions SET refresh_token_hash = ?, last_seen_at = ?, expires_at = ?, revoked_at = NULL WHERE id = ?')
.run(refreshTokenHash, ts, expiresAt, id);
}
export function revokeSession(id) {
const ts = nowIso();
db.prepare('UPDATE sessions SET revoked_at = ? WHERE id = ?').run(ts, id);
}
export function revokeAllUserSessions(userId) {
const ts = nowIso();
db.prepare('UPDATE sessions SET revoked_at = ? WHERE user_id = ?').run(ts, userId);
}
export function listUserSessions(userId) {
return db.prepare('SELECT id, user_agent AS userAgent, device_info AS deviceInfo, ip_address AS ip, is_remember AS isRemember, created_at AS createdAt, last_seen_at AS lastSeenAt, expires_at AS expiresAt, revoked_at AS revokedAt FROM sessions WHERE user_id = ? ORDER BY created_at DESC').all(userId);
}
export function setUserLastLogin(userId) {
const ts = nowIso();
db.prepare('UPDATE users SET last_login_at = ?, updated_at = ? WHERE id = ?').run(ts, ts, userId);
}
export function insertLoginAudit({ userId, username, ip, userAgent, success, reason }) {
const id = cryptoRandomId();
const ts = nowIso();
db.prepare(`INSERT INTO login_audit (id, user_id, username, ip_address, user_agent, success, reason, created_at)
VALUES (@id, @userId, @username, @ip, @userAgent, @success, @reason, @ts)`)
.run({ id, userId, username, ip, userAgent, success: success ? 1 : 0, reason: reason || null, ts });
}
export function cryptoRandomId() {
// simple URL-safe base64 16 bytes id
return randomBytes(16).toString('base64url');
}
// -------------------- Telemetry (minimal product events) --------------------
export function insertTelemetryEvent({ userId, event, meta }) {
const id = cryptoRandomId();
const ts = nowIso();
try {
db.prepare(`INSERT INTO telemetry_events (id, user_id, event, meta_json, created_at)
VALUES (@id, @userId, @event, @metaJson, @ts)`)
.run({
id,
userId: userId || null,
event: String(event || '').slice(0, 120),
metaJson: meta ? JSON.stringify(meta) : null,
ts,
});
return { id, created_at: ts };
} catch {
// Telemetry must never break the request path
return null;
}
}
export function listTelemetryEvents({ event, limit = 100, since } = {}) {
const capped = Math.min(1000, Math.max(1, Number(limit || 100)));
const clauses = [];
const params = {};
if (event) { clauses.push('event = @event'); params.event = String(event); }
if (since) { clauses.push('created_at >= @since'); params.since = String(since); }
const where = clauses.length ? `WHERE ${clauses.join(' AND ')}` : '';
return db.prepare(`SELECT id, user_id AS userId, event, meta_json AS metaJson, created_at AS createdAt
FROM telemetry_events ${where}
ORDER BY created_at DESC LIMIT @limit`)
.all({ ...params, limit: capped });
}
export function countTelemetryEvents({ event, since } = {}) {
const clauses = [];
const params = {};
if (event) { clauses.push('event = @event'); params.event = String(event); }
if (since) { clauses.push('created_at >= @since'); params.since = String(since); }
const where = clauses.length ? `WHERE ${clauses.join(' AND ')}` : '';
const row = db.prepare(`SELECT COUNT(*) AS n FROM telemetry_events ${where}`).get(params);
return row ? row.n : 0;
}
export function cryptoRandomUUID() {
return randomUUID();
}
export default db;
// -------------------- Search History --------------------
// Ids courts (front : yt, dm, tw, pt, od, ru) <-> noms longs (API/Watch :
// youtube, dailymotion, twitch, peertube, odysee, rumble). Les filtres
// d'historique acceptent les deux formes ; le stockage garde la forme reçue.
const HISTORY_PROVIDER_ALIASES = {
youtube: ['youtube', 'yt'],
dailymotion: ['dailymotion', 'dm'],
twitch: ['twitch', 'tw'],
peertube: ['peertube', 'pt'],
odysee: ['odysee', 'od'],
rumble: ['rumble', 'ru'],
};
const HISTORY_SHORT_TO_LONG = { yt: 'youtube', dm: 'dailymotion', tw: 'twitch', pt: 'peertube', od: 'odysee', ru: 'rumble' };
/** Normalise un provider (court ou long, toute casse) vers le nom long, ou null. */
export function normalizeHistoryProvider(raw) {
const v = String(raw || '').trim().toLowerCase();
if (!v) return null;
if (HISTORY_PROVIDER_ALIASES[v]) return v;
return HISTORY_SHORT_TO_LONG[v] || null;
}
/** Échappe les jokers LIKE (`%`, `_`, `\`) d'une saisie utilisateur. */
export function escapeLikePattern(raw) {
return String(raw ?? '').replace(/\\/g, '\\\\').replace(/%/g, '\\%').replace(/_/g, '\\_');
}
// -------------------- Phase 5 : catalogue `videos` --------------------
/**
* `videos` est la **seule** table qui porte les metadonnees d'une video.
* `watch_history` et `playlist_items` les denormalisent (phase 5.4) : cette
* duplication est le bug qui fait qu'un like sans visionnage anterieur
* s'affiche sans titre.
*
* Provider : la cle est toujours le **nom long** (`normalizeHistoryProvider`).
* `video_tags.provider` est en revanche stocke tel quel par `likeVideo` — donc
* court ou long selon ce qu'envoie le front. C'est exactement pour cela que
* `listLikedVideos` normalise les deux cotes de la jointure.
*/
function ensureVideosTable() {
try {
db.exec(`CREATE TABLE IF NOT EXISTS videos (
provider TEXT NOT NULL, video_id TEXT NOT NULL,
title TEXT, thumbnail TEXT, duration_seconds INTEGER, views INTEGER,
published_at TEXT, url TEXT, kind TEXT,
channel_external_id TEXT, channel_name TEXT, channel_avatar_url TEXT,
width INTEGER, height INTEGER,
raw_json TEXT, captured_at TEXT NOT NULL, created_at TEXT,
PRIMARY KEY (provider, video_id)
);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_videos_captured ON videos(captured_at DESC);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_videos_channel ON videos(provider, channel_external_id);`);
} catch {}
}
ensureVideosTable();
/** Expression SQL normalisant une colonne provider vers le nom long. */
function providerLongExpr(col) {
const whens = Object.entries(HISTORY_SHORT_TO_LONG).map(([s, l]) => `WHEN '${s}' THEN '${l}'`).join(' ');
return `(CASE LOWER(${col}) ${whens} ELSE LOWER(${col}) END)`;
}
const positiveInt = (v) => (typeof v === 'number' && Number.isFinite(v) && v > 0 ? Math.round(v) : null);
const nonEmptyStr = (v) => (typeof v === 'string' && v.trim().length > 0 ? v.trim() : null);
/**
* Ecrit (ou met a jour) une ligne `videos`.
*
* **Best-effort par contrat** : une erreur ici ne doit jamais faire echouer
* l'ecriture fonctionnelle appelante (historique, like, playlist). D'ou le
* try/catch qui englobe tout et retourne `false`.
*
* Deux invariants :
* - `captured_at` est rafraichi a chaque re-observation (c'est l'indicateur de
* fraicheur de la phase 5.1) ;
* - les colonnes absentes ne sont **pas** ecrasees (`COALESCE`) : un appel
* minimal (un like ne fournit que titre + vignette) ne doit pas effacer les
* metadonnees vues lors d'une session de lecture complete.
*
* @returns {boolean} true si la ligne a ete ecrite
*/
export function upsertVideoRow(dto) {
try {
const provider = normalizeHistoryProvider(dto?.provider)
|| String(dto?.provider || '').trim().toLowerCase();
const videoId = String(dto?.videoId ?? dto?.video_id ?? '').trim();
if (!provider || !videoId) return false;
ensureVideosTable();
const now = nowIso();
// `raw_json` est un garde-fou de debug : tronque pour ne pas faire exploser
// la base avec le payload d'une recherche.
let rawJson = null;
if (dto?.raw && typeof dto.raw === 'object') {
try { rawJson = JSON.stringify(dto.raw).slice(0, 4096); } catch { rawJson = null; }
}
db.prepare(
`INSERT INTO videos (provider, video_id, title, thumbnail, duration_seconds, views,
published_at, url, kind, channel_external_id, channel_name,
channel_avatar_url, width, height, raw_json, captured_at, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(provider, video_id) DO UPDATE SET
title = COALESCE(excluded.title, videos.title),
thumbnail = COALESCE(excluded.thumbnail, videos.thumbnail),
duration_seconds = COALESCE(excluded.duration_seconds, videos.duration_seconds),
views = COALESCE(excluded.views, videos.views),
published_at = COALESCE(excluded.published_at, videos.published_at),
url = COALESCE(excluded.url, videos.url),
kind = COALESCE(excluded.kind, videos.kind),
channel_external_id = COALESCE(excluded.channel_external_id, videos.channel_external_id),
channel_name = COALESCE(excluded.channel_name, videos.channel_name),
channel_avatar_url = COALESCE(excluded.channel_avatar_url, videos.channel_avatar_url),
width = COALESCE(excluded.width, videos.width),
height = COALESCE(excluded.height, videos.height),
raw_json = COALESCE(excluded.raw_json, videos.raw_json),
captured_at = excluded.captured_at`,
).run(
provider, videoId,
nonEmptyStr(dto?.title), nonEmptyStr(dto?.thumbnail),
positiveInt(dto?.durationSeconds ?? dto?.duration), positiveInt(dto?.views),
nonEmptyStr(dto?.publishedAt), nonEmptyStr(dto?.url), nonEmptyStr(dto?.kind),
nonEmptyStr(dto?.channelExternalId ?? dto?.channelId), nonEmptyStr(dto?.channelName),
nonEmptyStr(dto?.channelAvatarUrl),
positiveInt(dto?.width), positiveInt(dto?.height),
rawJson, now, now,
);
return true;
} catch {
return false; // best-effort : voir le contrat ci-dessus
}
}
/** Ligne `videos` pour un couple (provider, videoId), provider court ou long. */
export function getVideoRow(provider, videoId) {
try {
const p = normalizeHistoryProvider(provider) || String(provider || '').trim().toLowerCase();
ensureVideosTable();
return db.prepare(`SELECT * FROM videos WHERE provider = ? AND video_id = ?`).get(p, String(videoId || '')) || null;
} catch { return null; }
}
export function countVideos() {
try { ensureVideosTable(); return db.prepare(`SELECT COUNT(1) AS n FROM videos`).get()?.n || 0; } catch { return 0; }
}
/**
* Phase 5.3 — backfill depuis les tables denormalisees, **rejouable** (la
* migration SQL le fait une fois au deploiement ; ceci permet de rattraper des
* lignes ecrites entre le deploiement et l'arrivee du nouveau code).
* `INSERT OR IGNORE` : n'ecrase jamais une ligne deja observee, qui est plus
* riche que la source de backfill.
*/
export function backfillVideosFromLegacy() {
const report = { watchHistory: 0, playlistItems: 0 };
try {
ensureVideosTable();
// `provider` est selectionne puis NORMALISE : `playlist_items` et
// `watch_history` conservent la forme recue (courte ou longue) selon l'appelant,
// alors que `videos` est cle par le nom long. Sans cette normalisation, un
// couple ('yt', 'id') et ('youtube', 'id') creeraient deux lignes et les
// lectures normalisees enVerifieraient toujours une des deux.
report.watchHistory = db.prepare(
`INSERT OR IGNORE INTO videos (provider, video_id, title, thumbnail, captured_at, created_at)
SELECT ${providerLongExpr('provider')}, video_id, title, thumbnail, last_watched_at, last_watched_at
FROM watch_history WHERE video_id IS NOT NULL AND title IS NOT NULL`,
).run().changes || 0;
report.playlistItems = db.prepare(
`INSERT OR IGNORE INTO videos (provider, video_id, title, thumbnail, captured_at, created_at)
SELECT ${providerLongExpr('provider')}, video_id, title, thumbnail, added_at, added_at
FROM playlist_items WHERE video_id IS NOT NULL AND title IS NOT NULL`,
).run().changes || 0;
} catch {}
return report;
}
export function insertSearchHistory({ userId, query, filters }) {
const id = cryptoRandomId();
const created_at = nowIso();
const filters_json = filters ? JSON.stringify(filters) : null;
db.prepare(`INSERT INTO search_history (id, user_id, query, filters_json, created_at)
VALUES (?, ?, ?, ?, ?)`)
.run(id, userId, query, filters_json, created_at);
return { id, created_at };
}
/** Variante Takeout : date d'origine conservée (ISO valide exigée). */
export function insertSearchHistoryAt({ userId, query, createdAt }) {
const id = cryptoRandomId();
const d = new Date(String(createdAt || ''));
const created_at = !Number.isNaN(d.getTime()) ? d.toISOString() : nowIso();
db.prepare(`INSERT INTO search_history (id, user_id, query, filters_json, created_at)
VALUES (?, ?, ?, ?, ?)`)
.run(id, userId, query, null, created_at);
return { id, created_at };
}
export function listSearchHistory({ userId, limit = 50, before, q, provider }) {
const lim = Math.max(1, Math.min(200, Number(limit) || 50));
const conds = ['user_id = ?'];
const params = [userId];
if (before) { conds.push('created_at < ?'); params.push(String(before)); }
if (typeof q === 'string' && q.trim().length > 0) {
const like = `%${escapeLikePattern(q.trim())}%`;
conds.push(`(query LIKE ? ESCAPE '\\' OR COALESCE(filters_json,'') LIKE ? ESCAPE '\\')`);
params.push(like, like);
}
const normProvider = normalizeHistoryProvider(provider);
if (normProvider) {
// filters_json stocke indifféremment ids courts (["yt","dm"]) et/ou nom
// long ("provider":"youtube") : on matche les deux formes (forme courte
// entre guillemets pour éviter les faux positifs de sous-chaînes).
const [long, short] = HISTORY_PROVIDER_ALIASES[normProvider];
conds.push(`(COALESCE(filters_json,'') LIKE ? ESCAPE '\\' OR COALESCE(filters_json,'') LIKE ? ESCAPE '\\')`);
params.push(`%${escapeLikePattern(long)}%`, `%\"${short}\"%`);
}
const where = `WHERE ${conds.join(' AND ')}`;
return db.prepare(`SELECT * FROM search_history ${where} ORDER BY created_at DESC LIMIT ?`).all(...params, lim);
}
export function deleteSearchHistoryById(userId, id) {
db.prepare(`DELETE FROM search_history WHERE id = ? AND user_id = ?`).run(id, userId);
}
export function deleteAllSearchHistory(userId) {
db.prepare(`DELETE FROM search_history WHERE user_id = ?`).run(userId);
}
// -------------------- Watch History --------------------
export function upsertWatchHistory({ userId, provider, videoId, title, thumbnail, watchedAt, progressSeconds = 0, durationSeconds = 0, lastPositionSeconds }) {
const now = nowIso();
const watched_at = watchedAt || now;
// Stockage canonique (nom long) : les filtres matchent court + long.
const normProvider = normalizeHistoryProvider(provider) || String(provider || '').trim().toLowerCase();
// Insert or update on unique (user_id, provider, video_id)
db.prepare(`INSERT INTO watch_history (id, user_id, provider, video_id, title, thumbnail, watched_at, progress_seconds, duration_seconds, last_position_seconds, last_watched_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(user_id, provider, video_id) DO UPDATE SET
title=COALESCE(excluded.title, title),
thumbnail=COALESCE(excluded.thumbnail, watch_history.thumbnail),
progress_seconds=MAX(excluded.progress_seconds, watch_history.progress_seconds),
duration_seconds=MAX(excluded.duration_seconds, watch_history.duration_seconds),
last_position_seconds=COALESCE(excluded.last_position_seconds, watch_history.last_position_seconds),
last_watched_at=excluded.last_watched_at`).run(
cryptoRandomId(), userId, normProvider, videoId, title || null, thumbnail || null, watched_at, progressSeconds, durationSeconds, (typeof lastPositionSeconds === 'number' ? lastPositionSeconds : null), now
);
// Phase 5.2 : alimente la table `videos` (best-effort, ne doit jamais faire
// echouer l'ecriture de l'historique).
upsertVideoRow({ provider: normProvider, videoId, title, thumbnail, durationSeconds });
// Return the row id
const row = db.prepare(`SELECT * FROM watch_history WHERE user_id = ? AND provider = ? AND video_id = ?`).get(userId, normProvider, videoId);
return row;
}
export function updateWatchHistoryById(id, { progressSeconds, lastPositionSeconds }) {
const now = nowIso();
const row = db.prepare(`SELECT * FROM watch_history WHERE id = ?`).get(id);
if (!row) return null;
const nextProgress = (typeof progressSeconds === 'number') ? Math.max(progressSeconds, row.progress_seconds || 0) : row.progress_seconds;
const nextLastPos = (typeof lastPositionSeconds === 'number') ? lastPositionSeconds : row.last_position_seconds;
db.prepare(`UPDATE watch_history SET progress_seconds = ?, last_position_seconds = ?, last_watched_at = ? WHERE id = ?`)
.run(nextProgress, nextLastPos, now, id);
return db.prepare(`SELECT * FROM watch_history WHERE id = ?`).get(id);
}
export function listWatchHistory({ userId, limit = 50, before, q, provider }) {
const lim = Math.max(1, Math.min(200, Number(limit) || 50));
const conds = ['user_id = ?'];
const params = [userId];
if (before) { conds.push('watched_at < ?'); params.push(String(before)); }
const normProvider = normalizeHistoryProvider(provider);
if (normProvider) {
// Lignes anciennes : forme courte possible -> matche les deux formes.
const [long, short] = HISTORY_PROVIDER_ALIASES[normProvider];
conds.push(`(provider = ? OR provider = ?)`);
params.push(long, short);
}
if (typeof q === 'string' && q.trim().length > 0) {
const like = `%${escapeLikePattern(q.trim())}%`;
conds.push(`(COALESCE(title,'') LIKE ? ESCAPE '\\' OR provider LIKE ? ESCAPE '\\' OR video_id LIKE ? ESCAPE '\\')`);
params.push(like, like, like);
}
const where = `WHERE ${conds.join(' AND ')}`;
return db.prepare(`SELECT * FROM watch_history ${where} ORDER BY watched_at DESC LIMIT ?`).all(...params, lim);
}
export function deleteWatchHistoryById(userId, id) {
db.prepare(`DELETE FROM watch_history WHERE id = ? AND user_id = ?`).run(id, userId);
}
export function deleteAllWatchHistory(userId) {
db.prepare(`DELETE FROM watch_history WHERE user_id = ?`).run(userId);
}
// -------------------- Transcript History --------------------
// Conserve les transcripts générés pour les rejouer sans régénération.
function ensureTranscriptHistoryTable() {
try {
db.exec(`CREATE TABLE IF NOT EXISTS transcript_history (
id TEXT PRIMARY KEY,
user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
provider TEXT NOT NULL,
video_id TEXT NOT NULL,
title TEXT,
thumbnail TEXT,
lang TEXT NOT NULL,
languages_json TEXT,
lines_json TEXT NOT NULL DEFAULT '[]',
line_count INTEGER NOT NULL DEFAULT 0,
char_count INTEGER NOT NULL DEFAULT 0,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);`);
db.exec(`CREATE UNIQUE INDEX IF NOT EXISTS uq_transcript_history_user_video_lang
ON transcript_history(user_id, provider, video_id, lang);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_transcript_history_user_time
ON transcript_history(user_id, updated_at DESC);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_transcript_history_user_provider
ON transcript_history(user_id, provider);`);
} catch {}
}
ensureTranscriptHistoryTable();
function transcriptRowToApi(row) {
if (!row) return null;
let lines = [];
let languages = [];
try { lines = JSON.parse(row.lines_json || '[]'); if (!Array.isArray(lines)) lines = []; } catch { lines = []; }
try { languages = JSON.parse(row.languages_json || '[]'); if (!Array.isArray(languages)) languages = []; } catch { languages = []; }
const preview = lines.slice(0, 3).map((l) => String(l?.text || '')).join(' ').slice(0, 280);
return {
id: row.id,
provider: row.provider,
video_id: row.video_id,
title: row.title,
thumbnail: row.thumbnail,
lang: row.lang,
languages,
line_count: row.line_count,
char_count: row.char_count,
preview,
created_at: row.created_at,
updated_at: row.updated_at,
lines,
};
}
export function upsertTranscriptHistory({ userId, provider, videoId, title, thumbnail, lang, languages, lines }) {
ensureTranscriptHistoryTable();
if (!userId || !provider || !videoId || !lang) return null;
const cleanLines = Array.isArray(lines) ? lines.slice(0, 2000).map((l) => ({
t: Number(l?.t) || 0,
dur: Number(l?.dur) || 0,
text: String(l?.text || '').slice(0, 2000),
})) : [];
if (!cleanLines.length) return null;
const now = nowIso();
const normLang = String(lang).trim().slice(0, 12).toLowerCase();
const normProvider = String(provider).trim().slice(0, 32).toLowerCase();
const normVideoId = String(videoId).trim().slice(0, 256);
const languages_json = JSON.stringify(Array.isArray(languages) ? languages.slice(0, 32) : []);
const lines_json = JSON.stringify(cleanLines);
const line_count = cleanLines.length;
const char_count = cleanLines.reduce((n, l) => n + String(l.text || '').length, 0);
const existing = db.prepare(
`SELECT * FROM transcript_history WHERE user_id = ? AND provider = ? AND video_id = ? AND lang = ?`
).get(userId, normProvider, normVideoId, normLang);
if (existing) {
db.prepare(`UPDATE transcript_history SET title = COALESCE(?, title), thumbnail = COALESCE(?, thumbnail),
languages_json = ?, lines_json = ?, line_count = ?, char_count = ?, updated_at = ? WHERE id = ?`)
.run(title || null, thumbnail || null, languages_json, lines_json, line_count, char_count, now, existing.id);
return transcriptRowToApi(db.prepare(`SELECT * FROM transcript_history WHERE id = ?`).get(existing.id));
}
const id = cryptoRandomId();
db.prepare(`INSERT INTO transcript_history (id, user_id, provider, video_id, title, thumbnail, lang, languages_json, lines_json, line_count, char_count, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`)
.run(id, userId, normProvider, normVideoId, title || null, thumbnail || null, normLang, languages_json, lines_json, line_count, char_count, now, now);
return transcriptRowToApi(db.prepare(`SELECT * FROM transcript_history WHERE id = ?`).get(id));
}
export function getTranscriptHistoryItem({ userId, provider, videoId, lang, id }) {
ensureTranscriptHistoryTable();
if (id) {
return transcriptRowToApi(db.prepare(`SELECT * FROM transcript_history WHERE id = ? AND user_id = ?`).get(id, userId));
}
if (!provider || !videoId) return null;
const normProvider = String(provider).trim().toLowerCase();
const normVideoId = String(videoId).trim();
if (lang) {
const normLang = String(lang).trim().slice(0, 12).toLowerCase();
return transcriptRowToApi(db.prepare(
`SELECT * FROM transcript_history WHERE user_id = ? AND provider = ? AND video_id = ? AND lang = ?`
).get(userId, normProvider, normVideoId, normLang));
}
return transcriptRowToApi(db.prepare(
`SELECT * FROM transcript_history WHERE user_id = ? AND provider = ? AND video_id = ? ORDER BY updated_at DESC LIMIT 1`
).get(userId, normProvider, normVideoId));
}
export function listTranscriptHistory({ userId, limit = 50, before, q, provider, lang }) {
ensureTranscriptHistoryTable();
const lim = Math.max(1, Math.min(200, Number(limit) || 50));
const conds = ['user_id = ?'];
const params = [userId];
if (before) { conds.push('updated_at < ?'); params.push(String(before)); }
const normProvider = normalizeHistoryProvider(provider);
if (normProvider) {
const [long, short] = HISTORY_PROVIDER_ALIASES[normProvider];
conds.push(`(provider = ? OR provider = ?)`);
params.push(long, short);
}
if (lang && String(lang).trim()) {
// `fr` matche `fr` et `fr-ca` ; casse indifférente (stockage minuscule).
const primary = String(lang).trim().split('-')[0].toLowerCase();
conds.push(`(lang = ? OR lang LIKE ? ESCAPE '\\')`);
params.push(primary, `${escapeLikePattern(primary)}-%`);
}
if (typeof q === 'string' && q.trim().length > 0) {
const like = `%${escapeLikePattern(q.trim())}%`;
conds.push(`(COALESCE(title,'') LIKE ? ESCAPE '\\' OR video_id LIKE ? ESCAPE '\\' OR lang LIKE ? ESCAPE '\\' OR COALESCE(lines_json,'') LIKE ? ESCAPE '\\')`);
params.push(like, like, like, like);
}
const where = `WHERE ${conds.join(' AND ')}`;
const rows = db.prepare(`SELECT * FROM transcript_history ${where} ORDER BY updated_at DESC LIMIT ?`).all(...params, lim);
return rows.map(transcriptRowToApi);
}
export function deleteTranscriptHistoryById(userId, id) {
ensureTranscriptHistoryTable();
db.prepare(`DELETE FROM transcript_history WHERE id = ? AND user_id = ?`).run(id, userId);
}
export function deleteAllTranscriptHistory(userId, provider) {
ensureTranscriptHistoryTable();
if (provider && String(provider).trim()) {
const norm = normalizeHistoryProvider(provider) || String(provider).trim().toLowerCase();
db.prepare(`DELETE FROM transcript_history WHERE user_id = ? AND provider = ?`).run(userId, norm);
return;
}
db.prepare(`DELETE FROM transcript_history WHERE user_id = ?`).run(userId);
}
// -------------------- Likes (via tags/video_tags) --------------------
function ensureTag(userId, name) {
const tag = db.prepare(`SELECT id FROM tags WHERE user_id = ? AND name = ?`).get(userId, name);
if (tag && tag.id) return tag.id;
const id = cryptoRandomId();
db.prepare(`INSERT INTO tags (id, user_id, name, created_at) VALUES (?, ?, ?, ?)`)
.run(id, userId, name, nowIso());
return id;
}
export function likeVideo({ userId, provider, videoId, title, thumbnail }) {
const tagId = ensureTag(userId, 'like');
// Upsert-like behavior; ignore if exists
db.prepare(`INSERT OR IGNORE INTO video_tags (user_id, provider, video_id, tag_id, created_at) VALUES (?, ?, ?, ?, ?)`)
.run(userId, provider, videoId, tagId, nowIso());
// Phase 5.2 : on alimente `videos` meme quand aucun titre/vignette n'est fourni.
// C'est le cas qui produisait le trou des likes : la ligne `video_tags`
// existait mais rien d'autre ne portait les metadonnees.
upsertVideoRow({ provider, videoId, title, thumbnail });
// Also update the watch_history table with the title and thumbnail
if (title || thumbnail) {
upsertWatchHistory({ userId, provider, videoId, title, thumbnail });
}
return { provider, video_id: videoId };
}
export function unlikeVideo({ userId, provider, videoId }) {
const tag = db.prepare(`SELECT id FROM tags WHERE user_id = ? AND name = ?`).get(userId, 'like');
if (!tag) return { removed: false };
const info = db.prepare(`DELETE FROM video_tags WHERE user_id = ? AND provider = ? AND video_id = ? AND tag_id = ?`)
.run(userId, provider, videoId, tag.id);
return { removed: (info.changes || 0) > 0 };
}
export function isVideoLiked({ userId, provider, videoId }) {
const tag = db.prepare(`SELECT id FROM tags WHERE user_id = ? AND name = ?`).get(userId, 'like');
if (!tag) return false;
const row = db.prepare(`SELECT 1 FROM video_tags WHERE user_id = ? AND provider = ? AND video_id = ? AND tag_id = ?`)
.get(userId, provider, videoId, tag.id);
return Boolean(row);
}
export function listLikedVideos({ userId, limit = 100, q }) {
try {
console.log(`[listLikedVideos] Récupération des vidéos aimées pour l'utilisateur ${userId}`);
// Vérifier que la table tags existe
const tableExists = db.prepare(`
SELECT name FROM sqlite_master
WHERE type='table' AND name='tags'
`).get();
if (!tableExists) {
console.error('[listLikedVideos] La table tags n\'existe pas');
return [];
}
// Récupérer le tag 'like' de l'utilisateur
const tag = db.prepare(`
SELECT id FROM tags
WHERE user_id = ? AND name = ?
`).get(userId, 'like');
if (!tag) {
console.log(`[listLikedVideos] Aucun tag 'like' trouvé pour l'utilisateur ${userId}`);
return [];
}
console.log(`[listLikedVideos] Tag ID: ${tag.id}`);
// Vérifier que la table video_tags existe
const videoTagsExists = db.prepare(`
SELECT name FROM sqlite_master
WHERE type='table' AND name='video_tags'
`).get();
if (!videoTagsExists) {
console.error('[listLikedVideos] La table video_tags n\'existe pas');
return [];
}
// Phase 5.3 : les metadonnees viennent de `videos` (catalogue partage), avec
// repli `watch_history` puis `playlist_items` pour les lignes anterieures a
// la migration.
//
// Les deux cotes sont normalises : `video_tags.provider` est stocke tel quel
// par `likeVideo` (court ou long selon le front) alors que `videos` est
// cle par le nom long. Sans cette normalisation, un like emis avec `yt`
// ne trouvait aucune ligne `youtube` — le bug d'origine.
const hasQ = typeof q === 'string' && q.trim().length > 0;
const like = `%${(q || '').trim()}%`;
const vtProvider = providerLongExpr('vt.provider');
const vProvider = providerLongExpr('v.provider');
const whProvider = providerLongExpr('wh.provider');
// Repli playlist en sous-requete correlee et NON en jointure : `playlist_items`
// n'est unique que par (playlist_id, provider, video_id), donc une video
// presente dans 3 playlists y apparaitrait 3 fois dans le resultat.
const pliTitle = `(SELECT pi.title FROM playlist_items pi
WHERE ${providerLongExpr('pi.provider')} = ${vtProvider}
AND pi.video_id = vt.video_id AND pi.title IS NOT NULL
ORDER BY pi.added_at ASC LIMIT 1)`;
const pliThumb = `(SELECT pi.thumbnail FROM playlist_items pi
WHERE ${providerLongExpr('pi.provider')} = ${vtProvider}
AND pi.video_id = vt.video_id AND pi.thumbnail IS NOT NULL
ORDER BY pi.added_at ASC LIMIT 1)`;
const titleExpr = `COALESCE(v.title, wh.title, ${pliTitle}, '')`;
const thumbExpr = `COALESCE(v.thumbnail, wh.thumbnail, ${pliThumb}, '')`;
const base = `
SELECT
vt.provider,
vt.video_id,
vt.created_at,
${titleExpr} AS title,
${thumbExpr} AS thumbnail,
wh.last_watched_at AS last_watched_at,
v.captured_at AS captured_at
FROM video_tags vt
LEFT JOIN videos v
ON v.provider = ${vProvider} AND v.video_id = vt.video_id
LEFT JOIN watch_history wh
ON wh.user_id = vt.user_id
AND ${whProvider} = ${vtProvider}
AND wh.video_id = vt.video_id
WHERE vt.user_id = ? AND vt.tag_id = ?
`;
const orderLimit = `
ORDER BY vt.created_at DESC
LIMIT ?
`;
const query = hasQ
? `${base} AND (${titleExpr} LIKE ? OR vt.provider LIKE ? OR vt.video_id LIKE ?)
${orderLimit}`
: `${base} ${orderLimit}`;
// La requete est devenue longue (normalisation provider x 3, sous-requetes
// de repli) : la logger en entier poluait la sortie a chaque appel.
console.log(`[listLikedVideos] ${hasQ ? 'recherche' : 'liste'} pour ${userId} (tag ${tag.id})`);
const rows = hasQ
? db.prepare(query).all(userId, tag.id, like, like, like, limit)
: db.prepare(query).all(userId, tag.id, limit);
console.log(`[listLikedVideos] ${rows.length} vidéos trouvées`);
return rows;
} catch (error) {
console.error('[listLikedVideos] Erreur:', error.message);
console.error(error.stack);
throw error; // Renvoyer l'erreur pour qu'elle soit gérée par le routeur
}
}
// -------------------- Playlists --------------------
// Metrics are observability-only: they must never break the core operation
// (e.g. a playlist DELETE must stay 204 even if the metrics insert fails).
export function recordPlaylistMetric({ userId, playlistId, action, meta }) {
try {
const id = cryptoRandomId();
const created_at = nowIso();
const meta_json = meta ? JSON.stringify(meta) : null;
db.prepare(`INSERT INTO playlist_metrics (id, user_id, playlist_id, action, meta_json, created_at)
VALUES (?, ?, ?, ?, ?, ?)`)
.run(id, userId, playlistId, action, meta_json, created_at);
return { id, created_at };
} catch {
return null;
}
}
export function createPlaylist({ userId, title, description, thumbnail, isPrivate = true }) {
const id = cryptoRandomUUID();
const now = nowIso();
const name = String(title || '').trim();
if (!name) throw new Error('title_required');
db.prepare(`INSERT INTO playlists (id, user_id, name, description, thumbnail, is_private, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`)
.run(id, userId, name, description || null, thumbnail || null, isPrivate ? 1 : 0, now, now);
recordPlaylistMetric({ userId, playlistId: id, action: 'create', meta: { title: name } });
return db.prepare(`SELECT id, user_id AS userId, name AS title, description, thumbnail, is_private AS isPrivate, created_at AS createdAt, updated_at AS updatedAt,
(SELECT COUNT(1) FROM playlist_items WHERE playlist_id = playlists.id) AS itemsCount
FROM playlists WHERE id = ?`).get(id);
}
export function listPlaylists({ userId, limit = 50, offset = 0, q }) {
const like = `%${(q || '').trim()}%`;
const hasQ = typeof q === 'string' && q.trim().length > 0;
const base = `SELECT id, user_id AS userId, name AS title, description, thumbnail, is_private AS isPrivate, created_at AS createdAt, updated_at AS updatedAt,
(SELECT COUNT(1) FROM playlist_items WHERE playlist_id = playlists.id) AS itemsCount
FROM playlists WHERE user_id = ?`;
const sql = hasQ ? `${base} AND (name LIKE ? OR COALESCE(description,'') LIKE ?)
ORDER BY updated_at DESC LIMIT ? OFFSET ?` :
`${base} ORDER BY updated_at DESC LIMIT ? OFFSET ?`;
return hasQ
? db.prepare(sql).all(userId, like, like, limit, offset)
: db.prepare(sql).all(userId, limit, offset);
}
// List public playlists (visible to everyone)
export function listPublicPlaylists({ limit = 50, offset = 0, q }) {
const like = `%${(q || '').trim()}%`;
const hasQ = typeof q === 'string' && q.trim().length > 0;
const base = `SELECT id, user_id AS userId, name AS title, description, thumbnail, is_private AS isPrivate, created_at AS createdAt, updated_at AS updatedAt,
(SELECT COUNT(1) FROM playlist_items WHERE playlist_id = playlists.id) AS itemsCount
FROM playlists WHERE is_private = 0`;
const sql = hasQ ? `${base} AND (name LIKE ? OR COALESCE(description,'') LIKE ?)
ORDER BY updated_at DESC LIMIT ? OFFSET ?` :
`${base} ORDER BY updated_at DESC LIMIT ? OFFSET ?`;
return hasQ
? db.prepare(sql).all(like, like, limit, offset)
: db.prepare(sql).all(limit, offset);
}
export function getPlaylistRaw(id) {
return db.prepare(`SELECT id, user_id AS userId, name AS title, description, thumbnail, is_private AS isPrivate, created_at AS createdAt, updated_at AS updatedAt
FROM playlists WHERE id = ?`).get(id);
}
// Get playlist with items if allowed for a given viewer (owner or public)
export function getPlaylistWithItemsIfAllowed({ viewerUserId, id, limit = 200, offset = 0 }) {
const pl = db.prepare(`SELECT id, user_id AS userId, name AS title, description, thumbnail, is_private AS isPrivate, created_at AS createdAt, updated_at AS updatedAt
FROM playlists WHERE id = ?`).get(id);
if (!pl) return null;
const isOwner = viewerUserId && pl.userId === viewerUserId;
if (!isOwner && Number(pl.isPrivate) === 1) return 'forbidden';
const items = db.prepare(`SELECT id, playlist_id AS playlistId, provider, video_id AS videoId, title, thumbnail, added_at AS addedAt, position
FROM playlist_items WHERE playlist_id = ?
ORDER BY position ASC LIMIT ? OFFSET ?`).all(id, limit, offset);
return { ...pl, items };
}
export function updatePlaylist({ userId, id, patch }) {
const cur = db.prepare(`SELECT * FROM playlists WHERE id = ?`).get(id);
if (!cur) return null;
if (cur.user_id !== userId) return 'forbidden';
const next = {
name: patch.title != null ? String(patch.title).trim() : cur.name,
description: patch.description != null ? String(patch.description).trim() : cur.description,
thumbnail: patch.thumbnail != null ? String(patch.thumbnail).trim() : cur.thumbnail,
is_private: typeof patch.isPrivate === 'boolean' ? (patch.isPrivate ? 1 : 0) : cur.is_private,
updated_at: nowIso(),
};
db.prepare(`UPDATE playlists SET name=@name, description=@description, thumbnail=@thumbnail, is_private=@is_private, updated_at=@updated_at WHERE id = ?`)
.run(id, next);
recordPlaylistMetric({ userId, playlistId: id, action: 'update', meta: { title: next.name } });
return db.prepare(`SELECT id, user_id AS userId, name AS title, description, thumbnail, is_private AS isPrivate, created_at AS createdAt, updated_at AS updatedAt
FROM playlists WHERE id = ?`).get(id);
}
export function deletePlaylist({ userId, id }) {
const cur = db.prepare(`SELECT * FROM playlists WHERE id = ?`).get(id);
if (!cur) return { removed: false };
if (cur.user_id !== userId) return 'forbidden';
// Record BEFORE the row disappears: playlist_metrics.playlist_id references
// playlists(id), so inserting after the DELETE violates the FK constraint.
try { recordPlaylistMetric({ userId, playlistId: id, action: 'delete' }); } catch {}
// Explicit child cleanup (works even on old DBs whose FK lacks ON DELETE CASCADE).
try { db.prepare(`DELETE FROM playlist_items WHERE playlist_id = ?`).run(id); } catch {}
try { db.prepare(`DELETE FROM playlist_metrics WHERE playlist_id = ?`).run(id); } catch {}
const info = db.prepare(`DELETE FROM playlists WHERE id = ?`).run(id);
return { removed: (info.changes || 0) > 0 };
}
export function listPlaylistItems({ playlistId, limit = 200, offset = 0 }) {
return db.prepare(`SELECT id, playlist_id AS playlistId, provider, video_id AS videoId, title, thumbnail, added_at AS addedAt, position
FROM playlist_items WHERE playlist_id = ?
ORDER BY position ASC LIMIT ? OFFSET ?`).all(playlistId, limit, offset);
}
export function addPlaylistVideo({ userId, playlistId, provider, videoId, title, thumbnail }) {
const pl = db.prepare(`SELECT * FROM playlists WHERE id = ?`).get(playlistId);
if (!pl) return 'not_found';
if (pl.user_id !== userId) return 'forbidden';
const now = nowIso();
const posRow = db.prepare(`SELECT COALESCE(MAX(position), 0) + 1 AS nextPos FROM playlist_items WHERE playlist_id = ?`).get(playlistId);
const position = Number(posRow?.nextPos || 1);
const id = cryptoRandomId();
db.prepare(`INSERT OR IGNORE INTO playlist_items (id, playlist_id, provider, video_id, title, thumbnail, added_at, position)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)`).run(id, playlistId, provider, videoId, title || null, thumbnail || null, now, position);
// Phase 5.2 : un ajout a une playlist est une observation de la video.
upsertVideoRow({ provider, videoId, title, thumbnail });
recordPlaylistMetric({ userId, playlistId, action: 'add_video', meta: { provider, videoId } });
// Return the row (if it existed we need to fetch by key)
const row = db.prepare(`SELECT id, playlist_id AS playlistId, provider, video_id AS videoId, title, thumbnail, added_at AS addedAt, position
FROM playlist_items WHERE playlist_id = ? AND provider = ? AND video_id = ?`).get(playlistId, provider, videoId);
return row;
}
export function removePlaylistVideo({ userId, playlistId, provider, videoId }) {
const pl = db.prepare(`SELECT * FROM playlists WHERE id = ?`).get(playlistId);
if (!pl) return 'not_found';
if (pl.user_id !== userId) return 'forbidden';
const info = db.prepare(`DELETE FROM playlist_items WHERE playlist_id = ? AND provider = ? AND video_id = ?`).run(playlistId, provider, videoId);
recordPlaylistMetric({ userId, playlistId, action: 'remove_video', meta: { provider, videoId } });
return { removed: (info.changes || 0) > 0 };
}
export function reorderPlaylistVideos({ userId, playlistId, order }) {
const pl = db.prepare(`SELECT * FROM playlists WHERE id = ?`).get(playlistId);
if (!pl) return 'not_found';
if (pl.user_id !== userId) return 'forbidden';
if (!Array.isArray(order) || order.length === 0) return { changed: 0 };
const tx = db.transaction((rows) => {
let i = 1, changed = 0;
for (const it of rows) {
// Support either item ids or provider+videoId pairs
if (it && typeof it === 'string') {
const info = db.prepare(`UPDATE playlist_items SET position = ? WHERE id = ? AND playlist_id = ?`).run(i, it, playlistId);
changed += info.changes || 0;
} else if (it && (it.id || (it.provider && it.videoId))) {
const info = it.id
? db.prepare(`UPDATE playlist_items SET position = ? WHERE id = ? AND playlist_id = ?`).run(i, it.id, playlistId)
: db.prepare(`UPDATE playlist_items SET position = ? WHERE provider = ? AND video_id = ? AND playlist_id = ?`).run(i, it.provider, it.videoId, playlistId);
changed += info.changes || 0;
}
i++;
}
return changed;
});
const changed = tx(order);
recordPlaylistMetric({ userId, playlistId, action: 'reorder', meta: { count: order.length } });
return { changed };
}
// -------------------- Download jobs --------------------
export function insertDownloadJob({ id, userId, provider, videoId, title, formatId, audioOnly, url }) {
const now = Date.now();
db.prepare(`INSERT INTO download_jobs (id, user_id, provider, video_id, title, state, progress, format_id, audio_only, url, created_at, updated_at)
VALUES (@id, @userId, @provider, @videoId, @title, 'queued', 0, @formatId, @audioOnly, @url, @now, @now)`)
.run({ id, userId, provider, videoId, title: title || null, formatId: formatId || null, audioOnly: audioOnly ? 1 : 0, url: url || null, now });
return getDownloadJob(id);
}
export function getDownloadJob(id) {
return db.prepare(`SELECT id, user_id AS userId, provider, video_id AS videoId, title, state, progress, format_id AS formatId,
audio_only AS audioOnly, url, file_name AS fileName, file_ext AS fileExt, file_size AS fileSize,
file_path AS filePath, error, created_at AS createdAt, updated_at AS updatedAt, completed_at AS completedAt
FROM download_jobs WHERE id = ?`).get(id);
}
export function listDownloadJobs({ userId, limit = 50, offset = 0, state }) {
const hasState = typeof state === 'string' && state.trim().length > 0;
// When userId is falsy, list jobs across all users (server-side maintenance).
const hasUser = userId != null && String(userId).length > 0;
const cols = `id, user_id AS userId, provider, video_id AS videoId, title, state, progress, format_id AS formatId,
audio_only AS audioOnly, url, file_name AS fileName, file_ext AS fileExt, file_size AS fileSize,
file_path AS filePath, error, created_at AS createdAt, updated_at AS updatedAt, completed_at AS completedAt`;
const where = [];
const params = [];
if (hasUser) { where.push('user_id = ?'); params.push(userId); }
if (hasState) { where.push('state = ?'); params.push(state.trim()); }
const whereSql = where.length ? ` WHERE ${where.join(' AND ')}` : '';
params.push(limit, offset);
return db.prepare(`SELECT ${cols} FROM download_jobs${whereSql} ORDER BY created_at DESC LIMIT ? OFFSET ?`).all(...params);
}
export function updateDownloadJob(id, patch = {}) {
const cur = getDownloadJob(id);
if (!cur) return null;
const next = {
state: patch.state != null ? String(patch.state) : cur.state,
progress: (typeof patch.progress === 'number' && patch.progress > (cur.progress || 0)) ? patch.progress : (cur.progress || 0),
title: patch.title != null ? patch.title : cur.title,
fileName: patch.fileName != null ? patch.fileName : cur.fileName,
fileExt: patch.fileExt != null ? patch.fileExt : cur.fileExt,
fileSize: patch.fileSize != null ? patch.fileSize : cur.fileSize,
filePath: patch.filePath != null ? patch.filePath : cur.filePath,
error: patch.error != null ? patch.error : cur.error,
completedAt: patch.completedAt != null ? patch.completedAt : cur.completedAt,
updatedAt: Date.now(),
};
db.prepare(`UPDATE download_jobs SET state=@state, progress=@progress, title=@title, file_name=@fileName, file_ext=@fileExt,
file_size=@fileSize, file_path=@filePath, error=@error, completed_at=@completedAt, updated_at=@updatedAt
WHERE id = @id`).run({ ...next, id });
return getDownloadJob(id);
}
export function deleteDownloadJob({ userId, id }) {
const info = db.prepare(`DELETE FROM download_jobs WHERE id = ? AND user_id = ?`).run(id, userId);
return (info.changes || 0) > 0;
}
// Mark jobs that were queued/running/merging when the API stopped as 'interrupted'
// so the user can retry them from the UI instead of polling forever.
export function resetActiveDownloadJobs() {
try {
const info = db.prepare(`UPDATE download_jobs
SET state = 'interrupted', error = 'api_restarted', updated_at = ?
WHERE state IN ('queued', 'running', 'merging')`).run(Date.now());
if ((info.changes || 0) > 0) console.log(`[db] Marked ${info.changes} active download job(s) as interrupted`);
} catch (e) {
console.warn('[db] resetActiveDownloadJobs failed:', e?.message || e);
}
}
// Count active (queued/running/merging) jobs for a user — used for concurrency quota
export function countActiveDownloadJobs(userId) {
const row = db.prepare(`SELECT COUNT(1) AS n FROM download_jobs WHERE user_id = ? AND state IN ('queued', 'running', 'merging')`).get(userId);
return row?.n || 0;
}
// Sum of completed job sizes (bytes) within the retention window — used for storage quota
export function sumCompletedDownloadBytes(userId, sinceMs) {
const row = sinceMs
? db.prepare(`SELECT COALESCE(SUM(file_size), 0) AS total FROM download_jobs WHERE user_id = ? AND state = 'completed' AND completed_at >= ?`).get(userId, sinceMs)
: db.prepare(`SELECT COALESCE(SUM(file_size), 0) AS total FROM download_jobs WHERE user_id = ? AND state = 'completed'`).get(userId);
return row?.total || 0;
}
// -------------------- Channels & Subscriptions --------------------
const CHANNEL_TTL_MS = Number(process.env.CHANNEL_TTL_MS || (6 * 60 * 60 * 1000));
export function getChannelByProviderExternalId(provider, externalId) {
return db.prepare(`SELECT * FROM channels WHERE provider = ? AND external_id = ?`).get(provider, externalId);
}
export function getChannelById(id) {
return db.prepare(`SELECT * FROM channels WHERE id = ?`).get(id);
}
export function channelRowToMeta(row) {
if (!row) return null;
return {
id: row.id,
provider: row.provider,
externalId: row.external_id,
title: row.title || null,
handle: row.handle || null,
avatarUrl: row.avatar_url || null,
bannerUrl: row.banner_url || null,
description: row.description || null,
url: row.url || null,
subsCount: typeof row.subs_count === 'number' ? row.subs_count : undefined,
verified: row.verified == null ? undefined : Boolean(row.verified),
lastRefreshedAt: row.last_refreshed_at || null,
};
}
/** Double défense : l'assainissement des URLs vit dans channel-registry, mais une
* bannière stockée en base ne doit jamais pouvoir survivre à un futur appelant
* qui oublierait `safeMeta`. Règle locale : http(s) ou rien. */
function sanitizeBannerUrl(value) {
return typeof value === 'string' && /^https?:\/\//i.test(value.trim()) ? value.trim() : null;
}
export function upsertChannelRow(meta) {
if (!meta || !meta.provider || !meta.externalId) {
throw new Error('invalid_channel_meta');
}
const now = Date.now();
const payload = {
provider: meta.provider,
external_id: meta.externalId,
title: meta.title || null,
handle: meta.handle || null,
avatar_url: meta.avatarUrl || null,
banner_url: sanitizeBannerUrl(meta.bannerUrl),
description: typeof meta.description === 'string' && meta.description.trim() ? meta.description.trim() : null,
url: meta.url || null,
subs_count: typeof meta.subsCount === 'number' ? meta.subsCount : null,
verified: meta.verified === undefined ? null : (meta.verified ? 1 : 0),
last_refreshed_at: meta.lastRefreshedAt ? Number(meta.lastRefreshedAt) : now,
};
db.prepare(`INSERT INTO channels (provider, external_id, title, handle, avatar_url, banner_url, description, url, subs_count, verified, last_refreshed_at)
VALUES (@provider, @external_id, @title, @handle, @avatar_url, @banner_url, @description, @url, @subs_count, @verified, @last_refreshed_at)
ON CONFLICT(provider, external_id) DO UPDATE SET
title=excluded.title,
handle=excluded.handle,
avatar_url=excluded.avatar_url,
banner_url=excluded.banner_url,
description=excluded.description,
url=excluded.url,
subs_count=excluded.subs_count,
verified=COALESCE(excluded.verified, channels.verified),
last_refreshed_at=excluded.last_refreshed_at`).run(payload);
const row = getChannelByProviderExternalId(meta.provider, meta.externalId);
return channelRowToMeta(row);
}
export async function ensureChannelFresh(provider, externalId, fetcher, opts = {}) {
if (typeof fetcher !== 'function') throw new Error('channel_fetcher_required');
const forceRefresh = Boolean(opts?.force);
const row = getChannelByProviderExternalId(provider, externalId);
if (!forceRefresh && row && row.last_refreshed_at && (Date.now() - row.last_refreshed_at) < CHANNEL_TTL_MS) {
return channelRowToMeta(row);
}
try {
const maybePromise = fetcher();
const data = maybePromise instanceof Promise ? await maybePromise : maybePromise;
const consolidated = {
provider,
externalId,
title: data?.title,
handle: data?.handle,
avatarUrl: data?.avatarUrl,
bannerUrl: data?.bannerUrl,
description: data?.description,
url: data?.url,
subsCount: data?.subsCount,
verified: data?.verified,
lastRefreshedAt: Date.now(),
};
return upsertChannelRow(consolidated);
} catch {
const fallback = {
provider,
externalId,
lastRefreshedAt: Date.now(),
};
return upsertChannelRow(fallback);
}
}
export function listSubscriptionsByUser(userId) {
const rows = db.prepare(`
SELECT s.id AS subscriptionId,
s.created_at AS createdAt,
c.id AS channelId,
c.provider,
c.external_id AS externalId,
c.title,
c.handle,
c.avatar_url AS avatarUrl,
c.url,
c.subs_count AS subsCount,
c.verified,
c.last_refreshed_at AS lastRefreshedAt
FROM subscriptions s
JOIN channels c ON c.id = s.channel_id
WHERE s.user_id = ?
ORDER BY COALESCE(c.title, c.external_id) COLLATE NOCASE ASC
`).all(userId);
return rows.map(subscriptionRowToDto);
}
export function subscribeChannel({ userId, provider, externalId, channelId }) {
if (!channelId) {
const channelRow = getChannelByProviderExternalId(provider, externalId);
if (!channelRow) throw new Error('channel_missing');
channelId = channelRow.id;
}
const ts = Date.now();
db.prepare(`INSERT OR IGNORE INTO subscriptions (user_id, channel_id, created_at)
VALUES (?, ?, ?)`).run(userId, channelId, ts);
const row = db.prepare(`
SELECT s.id AS subscriptionId,
s.created_at AS createdAt,
c.id AS channelId,
c.provider,
c.external_id AS externalId,
c.title,
c.handle,
c.avatar_url AS avatarUrl,
c.url,
c.subs_count AS subsCount,
c.verified,
c.last_refreshed_at AS lastRefreshedAt
FROM subscriptions s
JOIN channels c ON c.id = s.channel_id
WHERE s.user_id = ? AND s.channel_id = ?
`).get(userId, channelId);
return subscriptionRowToDto(row);
}
export function unsubscribeChannel({ userId, subscriptionId }) {
return db.prepare(`DELETE FROM subscriptions WHERE id = ? AND user_id = ?`).run(subscriptionId, userId);
}
export function isSubscribed({ userId, provider, externalId }) {
return db.prepare(`
SELECT 1
FROM subscriptions s
JOIN channels c ON c.id = s.channel_id
WHERE s.user_id = ? AND c.provider = ? AND c.external_id = ?
LIMIT 1
`).get(userId, provider, externalId) != null;
}
// -------------------- Subscription groups (façon PocketTube) --------------------
function subscriptionGroupRowToDto(row) {
if (!row) return null;
return {
id: row.id,
name: row.name,
color: row.color || '#ef4444',
icon: row.icon || '📁',
channelCount: typeof row.channelCount === 'number' ? row.channelCount : Number(row.channelCount || 0),
createdAt: row.created_at,
updatedAt: row.updated_at,
};
}
export function listSubscriptionGroups(userId) {
const rows = db.prepare(`
SELECT g.id, g.user_id, g.name, g.color, g.icon, g.created_at, g.updated_at,
COUNT(m.subscription_id) AS channelCount
FROM subscription_groups g
LEFT JOIN subscription_group_members m ON m.group_id = g.id
WHERE g.user_id = ?
GROUP BY g.id
ORDER BY g.created_at ASC
`).all(userId);
return rows.map(subscriptionGroupRowToDto);
}
export function createSubscriptionGroup({ userId, name, color, icon }) {
const cleanName = String(name || '').trim().slice(0, 80);
if (!cleanName) throw new Error('group_name_required');
const id = cryptoRandomId();
const now = Date.now();
try {
db.prepare(`INSERT INTO subscription_groups (id, user_id, name, color, icon, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?)`)
.run(id, userId, cleanName, String(color || '#ef4444').slice(0, 16), String(icon || '📁').slice(0, 16), now, now);
} catch (e) {
if (String(e?.message || '').includes('UNIQUE')) throw new Error('group_name_taken');
throw e;
}
return subscriptionGroupRowToDto(
db.prepare(`SELECT *, 0 AS channelCount FROM subscription_groups WHERE id = ?`).get(id)
);
}
export function updateSubscriptionGroup({ userId, groupId, name, color, icon }) {
const existing = db.prepare(`SELECT * FROM subscription_groups WHERE id = ? AND user_id = ?`).get(groupId, userId);
if (!existing) return null;
const next = {
name: name === undefined ? existing.name : String(name).trim().slice(0, 80),
color: color === undefined ? existing.color : String(color).slice(0, 16),
icon: icon === undefined ? existing.icon : String(icon).slice(0, 16),
};
if (!next.name) throw new Error('group_name_required');
try {
db.prepare(`UPDATE subscription_groups SET name = ?, color = ?, icon = ?, updated_at = ? WHERE id = ? AND user_id = ?`)
.run(next.name, next.color, next.icon, Date.now(), groupId, userId);
} catch (e) {
if (String(e?.message || '').includes('UNIQUE')) throw new Error('group_name_taken');
throw e;
}
const row = db.prepare(`
SELECT g.*, COUNT(m.subscription_id) AS channelCount
FROM subscription_groups g
LEFT JOIN subscription_group_members m ON m.group_id = g.id
WHERE g.id = ? AND g.user_id = ?
GROUP BY g.id
`).get(groupId, userId);
return subscriptionGroupRowToDto(row);
}
export function deleteSubscriptionGroup({ userId, groupId }) {
return db.prepare(`DELETE FROM subscription_groups WHERE id = ? AND user_id = ?`).run(groupId, userId);
}
/** Remplace le contenu d'un groupe (ids d'abonnements appartenant à l'utilisateur). */
export function setSubscriptionGroupMembers({ userId, groupId, subscriptionIds }) {
const group = db.prepare(`SELECT * FROM subscription_groups WHERE id = ? AND user_id = ?`).get(groupId, userId);
if (!group) return null;
const ids = [...new Set((subscriptionIds || []).map(Number).filter(Number.isFinite))];
const owned = ids.length
? db.prepare(`SELECT id FROM subscriptions WHERE user_id = ? AND id IN (${ids.map(() => '?').join(',')})`).all(userId, ...ids).map(r => r.id)
: [];
const tx = db.transaction(() => {
db.prepare(`DELETE FROM subscription_group_members WHERE group_id = ?`).run(groupId);
const now = Date.now();
const stmt = db.prepare(`INSERT OR IGNORE INTO subscription_group_members (group_id, subscription_id, added_at) VALUES (?, ?, ?)`);
for (const sid of owned) stmt.run(groupId, sid, now);
db.prepare(`UPDATE subscription_groups SET updated_at = ? WHERE id = ?`).run(now, groupId);
});
tx();
return updateSubscriptionGroup({ userId, groupId });
}
/** Remplace les groupes d'un abonnement (relation plusieurs-à-plusieurs). */
export function setSubscriptionGroups({ userId, subscriptionId, groupIds }) {
const sub = db.prepare(`SELECT id FROM subscriptions WHERE id = ? AND user_id = ?`).get(subscriptionId, userId);
if (!sub) return null;
const ids = [...new Set((groupIds || []).map(String))];
const owned = ids.length
? db.prepare(`SELECT id FROM subscription_groups WHERE user_id = ? AND id IN (${ids.map(() => '?').join(',')})`).all(userId, ...ids).map(r => r.id)
: [];
const tx = db.transaction(() => {
db.prepare(`DELETE FROM subscription_group_members WHERE subscription_id = ?`).run(subscriptionId);
const now = Date.now();
const stmt = db.prepare(`INSERT OR IGNORE INTO subscription_group_members (group_id, subscription_id, added_at) VALUES (?, ?, ?)`);
for (const gid of owned) stmt.run(gid, subscriptionId, now);
});
tx();
return owned;
}
/** Carte subscriptionId -> groupId[] pour tout l'utilisateur (une seule requête). */
export function listSubscriptionGroupMembersByUser(userId) {
const rows = db.prepare(`
SELECT m.subscription_id AS subscriptionId, m.group_id AS groupId
FROM subscription_group_members m
JOIN subscriptions s ON s.id = m.subscription_id AND s.user_id = ?
JOIN subscription_groups g ON g.id = m.group_id AND g.user_id = ?
`).all(userId, userId);
const map = {};
for (const r of rows) {
if (!map[r.subscriptionId]) map[r.subscriptionId] = [];
map[r.subscriptionId].push(r.groupId);
}
return map;
}
function subscriptionRowToDto(row) {
if (!row) return null;
const channel = channelRowToMeta({
id: row.channelId,
provider: row.provider,
external_id: row.externalId,
title: row.title,
handle: row.handle,
avatar_url: row.avatarUrl,
url: row.url,
subs_count: row.subsCount,
verified: row.verified,
last_refreshed_at: row.lastRefreshedAt,
});
return {
subscriptionId: row.subscriptionId,
createdAt: row.createdAt,
channel,
};
}
// -------------------- Step 17 : cache persistant recherche YouTube --------------------
function ensureYoutubeCacheTables() {
try {
db.exec(`CREATE TABLE IF NOT EXISTS youtube_search_cache (
q_hash TEXT PRIMARY KEY, q TEXT NOT NULL, payload_json TEXT NOT NULL,
source TEXT NOT NULL DEFAULT 'scrape', created_at INTEGER NOT NULL, expires_at INTEGER NOT NULL
);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_yt_cache_exp ON youtube_search_cache(expires_at);`);
db.exec(`CREATE TABLE IF NOT EXISTS youtube_metrics (
day TEXT PRIMARY KEY, scrape_calls INTEGER NOT NULL DEFAULT 0,
api_calls INTEGER NOT NULL DEFAULT 0, quota_units INTEGER NOT NULL DEFAULT 0, updated_at TEXT NOT NULL
);`);
} catch {}
}
ensureYoutubeCacheTables();
/**
* Task 4.2 — migration de `youtube_search_cache` vers `search_cache`.
*
* La clé est identique (`hashSearchKey` + préfixe `yt|`) : on copie donc les
* lignes telles quelles, sans retransformation. La table d'origine **n'est pas
* supprimée** : elle reste lue en repli pendant une version, ce qui rend la
* migration réversible (`DROP TABLE search_cache` suffit à revenir en arrière).
*
* `INSERT OR IGNORE` : ne jamais écraser une entrée plus fraîche déjà présente
* dans `search_cache` (une base peut avoir servi les deux tables).
*/
export function migrateYoutubeCacheToSearchCache() {
try {
ensureYoutubeCacheTables();
ensureSearchCacheTables();
const hasOld = db.prepare(`SELECT COUNT(1) AS n FROM youtube_search_cache`).get()?.n || 0;
if (hasOld === 0) return { migrated: 0, remaining: 0, legacyReadable: false };
const res = db.prepare(
`INSERT OR IGNORE INTO search_cache (cache_key, provider, q, payload_json, item_count, source, hit_count, created_at, expires_at)
SELECT q_hash, 'yt', q, payload_json,
CASE WHEN json_valid(payload_json) THEN json_array_length(payload_json) ELSE 0 END,
source, 0, created_at, expires_at
FROM youtube_search_cache WHERE expires_at > ?`,
).run(Date.now());
const migrated = res.changes || 0;
const remaining = db.prepare(`SELECT COUNT(1) AS n FROM youtube_search_cache WHERE expires_at > ?`).get(Date.now())?.n || 0;
return { migrated, remaining, legacyReadable: true };
} catch { return { migrated: 0, remaining: 0, legacyReadable: false, error: true }; }
}
/**
* Lecture YouTube. Essaie `search_cache` (table générique, phase 4) puis
* **replie** sur `youtube_search_cache` : une base upgradeée depuis une version
* antérieure peut n'avoir que des lignes dans l'ancienne table tant que la
* migration n'a pas tourné.
*/
export function getCachedYoutubeSearch(qHash) {
try {
const generic = getCachedSearch('yt', qHash);
if (generic) return { items: generic.items, source: generic.source };
} catch {}
try {
ensureYoutubeCacheTables();
const row = db.prepare(`SELECT payload_json AS payload, source, expires_at AS exp FROM youtube_search_cache WHERE q_hash = ?`).get(qHash);
if (!row) return null;
if (Date.now() >= Number(row.exp || 0)) {
try { db.prepare(`DELETE FROM youtube_search_cache WHERE q_hash = ?`).run(qHash); } catch {}
return null;
}
try { return { items: JSON.parse(String(row.payload || '[]')), source: row.source }; } catch { return null; }
} catch { return null; }
}
/** Écrit YouTube : table générique en primaire, ancienne table en repli. */
export function setCachedYoutubeSearch(qHash, q, items, source, ttlMs) {
const written = setCachedSearch('yt', qHash, q, items, source, ttlMs);
if (written) return;
try {
ensureYoutubeCacheTables();
const now = Date.now();
db.prepare(`INSERT INTO youtube_search_cache (q_hash, q, payload_json, source, created_at, expires_at)
VALUES (?, ?, ?, ?, ?, ?)
ON CONFLICT(q_hash) DO UPDATE SET q=excluded.q, payload_json=excluded.payload_json,
source=excluded.source, created_at=excluded.created_at, expires_at=excluded.expires_at`)
.run(qHash, String(q || '').slice(0, 300), JSON.stringify(items || []), source, now, now + ttlMs);
// Cap : garde 2000 entrées les plus fraîches
try { db.exec(`DELETE FROM youtube_search_cache WHERE q_hash NOT IN (SELECT q_hash FROM youtube_search_cache ORDER BY expires_at DESC LIMIT 2000)`); } catch {}
} catch {}
}
export function pruneYoutubeCache() {
// Phase 4.5 : la purge ne délaisse plus l'ancienne table derrière.
return pruneSearchCache();
}
export function incYoutubeMetrics({ scrapeCalls = 0, apiCalls = 0, quotaUnits = 0 } = {}) {
try {
ensureYoutubeCacheTables();
const day = new Date().toISOString().slice(0, 10);
db.prepare(`INSERT INTO youtube_metrics (day, scrape_calls, api_calls, quota_units, updated_at)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT(day) DO UPDATE SET scrape_calls = scrape_calls + ?, api_calls = api_calls + ?,
quota_units = quota_units + ?, updated_at = excluded.updated_at`)
.run(day, scrapeCalls, apiCalls, quotaUnits, new Date().toISOString(), scrapeCalls, apiCalls, quotaUnits);
} catch {}
}
export function getYoutubeMetricsToday() {
try {
ensureYoutubeCacheTables();
const day = new Date().toISOString().slice(0, 10);
return db.prepare(`SELECT * FROM youtube_metrics WHERE day = ?`).get(day) || { day, scrape_calls: 0, api_calls: 0, quota_units: 0 };
} catch { return { scrape_calls: 0, api_calls: 0, quota_units: 0 }; }
}
export function countYoutubeCacheRows() {
try { return db.prepare(`SELECT COUNT(1) AS n FROM youtube_search_cache`).get()?.n || 0; } catch { return 0; }
}
// -------------------- Phase 4.1 : cache de recherche générique --------------------
/**
* Cache de recherche **générique**, multi-fournisseur.
*
* Remplace `youtube_search_cache` (task 4.2) sans le supprimer : les anciennes
* fonctions YouTube restent des surcouches de compatibilité et lisent encore
* l'ancienne table en repli.
*
* Choix de conception (cf. `docs/plan-phases-catalogue-classification.md`) :
* - la clé reprend **exactement** le format YouTube `hashSearchKey`
* (`<provider>|<sha256(q|perPage|page|sort|filtersCacheKey)>`) : une ligne
* existante peut être migrée telle quelle, et la migration est réversible ;
* - `hit_count` est incrémenté **à chaque lecture** (y compris les lectures
* expirées) pour mesurer le ratio cache/total via `/api/providers/metrics` ;
* - un résultat **vide n'est jamais persisté** : une page vide transitoire
* (continuation expirée, raté réseau partiel) ne doit pas empoisonner le
* cache pendant 5-30 min.
*/
const SEARCH_CACHE_MAX_ROWS = 2000;
/** TTL par fournisseur. YT garde 30 min (quota Data API), les autres 5 min. */
export function searchCacheTtlMs(provider) {
const p = String(provider || '').toLowerCase();
if (p === 'yt') {
const override = Number(process.env.SEARCH_CACHE_TTL_MS_YT);
if (Number.isFinite(override) && override > 0) return override;
return 30 * 60 * 1000;
}
const perProvider = Number(process.env[`SEARCH_CACHE_TTL_MS_${p.toUpperCase()}`]);
if (Number.isFinite(perProvider) && perProvider > 0) return perProvider;
const general = Number(process.env.SEARCH_CACHE_TTL_MS_DEFAULT);
if (Number.isFinite(general) && general > 0) return general;
return 5 * 60 * 1000;
}
function ensureSearchCacheTables() {
try {
db.exec(`CREATE TABLE IF NOT EXISTS search_cache (
cache_key TEXT NOT NULL, provider TEXT NOT NULL, q TEXT NOT NULL DEFAULT '',
payload_json TEXT NOT NULL, item_count INTEGER NOT NULL DEFAULT 0,
source TEXT NOT NULL DEFAULT 'api', hit_count INTEGER NOT NULL DEFAULT 0,
created_at INTEGER NOT NULL, expires_at INTEGER NOT NULL,
PRIMARY KEY (cache_key, provider)
);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_search_cache_exp ON search_cache(expires_at);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_search_cache_prov ON search_cache(provider, expires_at);`);
} catch {}
}
ensureSearchCacheTables();
/**
* Lecture cache + bump du `hit_count`.
* @returns {{items: any[], source: string, hitCount: number}|null}
*/
export function getCachedSearch(provider, cacheKey) {
const p = String(provider || '').toLowerCase();
try {
ensureSearchCacheTables();
const row = db.prepare(
`SELECT payload_json AS payload, source, expires_at AS exp, hit_count AS hits, created_at AS createdAt
FROM search_cache WHERE cache_key = ? AND provider = ?`,
).get(String(cacheKey), p);
if (!row) return null;
const hits = Number(row.hits || 0) + 1;
try { db.prepare(`UPDATE search_cache SET hit_count = ? WHERE cache_key = ? AND provider = ?`).run(hits, String(cacheKey), p); } catch {}
if (Date.now() >= Number(row.exp || 0)) {
// Expiration paresseuse : la ligne est comptabilisée puis purgée.
try { db.prepare(`DELETE FROM search_cache WHERE cache_key = ? AND provider = ?`).run(String(cacheKey), p); } catch {}
return null;
}
try {
const items = JSON.parse(String(row.payload || '[]'));
if (!Array.isArray(items)) return null;
// `createdAt` = moment où l'amont a réellement été interrogé. C'est la
// SEULE date de capture honnête pour un hit de cache : `Date.now()`
// mentirait en annonçant « à l'instant » une réponse vieille de 4 min.
return { items, source: row.source, hitCount: hits, createdAt: Number(row.createdAt || 0) || null };
} catch { return null; }
} catch { return null; }
}
/** Écriture. Un tableau vide est ignoré (voir en-tête de section). */
export function setCachedSearch(provider, cacheKey, q, items, source, ttlMs) {
const p = String(provider || '').toLowerCase();
if (!Array.isArray(items) || items.length === 0) return false;
try {
ensureSearchCacheTables();
const now = Date.now();
// Un TTL explicitement fourni est respecte tel quel (y compris 0 ou
// negatif = « deja expire », utile pour vider un cache en test) ; seule
// l'absence de valeur (undefined/NaN) retombe sur le TTL du provider.
const ttl = Number.isFinite(ttlMs) ? ttlMs : searchCacheTtlMs(p);
db.prepare(
`INSERT INTO search_cache (cache_key, provider, q, payload_json, item_count, source, hit_count, created_at, expires_at)
VALUES (?, ?, ?, ?, ?, ?, 0, ?, ?)
ON CONFLICT(cache_key, provider) DO UPDATE SET q=excluded.q, payload_json=excluded.payload_json,
item_count=excluded.item_count, source=excluded.source, created_at=excluded.created_at,
expires_at=excluded.expires_at`,
) .run(String(cacheKey), p, String(q || '').slice(0, 300), JSON.stringify(items), items.length, String(source || 'api'), now, now + ttl);
// Uniquement le plafond : la purge des lignes expirees est un travail de
// fond (intervalle 10 min, cf. 4.5) et ne doit pas coutir un DELETE par
// ecriture.
enforceSearchCacheCap(p, SEARCH_CACHE_MAX_ROWS);
return true;
} catch { return false; }
}
/**
* Phase 3.12 — persistance des `pageToken` de contenu de chaîne.
*
* L'API YouTube pagine avec des jetons opaques, pas des numéros de page : pour
* atteindre la page N il faut le jeton produit en chargeant la page N-1. Ces
* jetons vivaient dans une Map **process-local**, ce qui avait deux conséquences
* réelles :
* 1. un redémarrage ou un passage derrière un load-balancer repart de la page 1 ;
* 2. au-delà de 500 clés, la Map faisait `clear()` — donc la pagination était
* réinitialisée SANS AUCUN SIGNE, silencieusement.
*
* On garde la Map comme cache L1 (rapide, et évite un accès SQLite par lecture)
* et on déporte la persistance dans `search_cache` (L2), ce qui couvre le
* redémarrage et le scale horizontal. Un éviction du L1 ne perd plus rien.
*
*.Namespace de provider distinct (`yt_tokens`) : ces lignes ne sont pas des
* résultats de recherche, et les mélanger aux lignes `yt` fausserait le plafond
* par fournisseur et les statistiques de cache.
*/
const PAGE_TOKEN_PROVIDER = 'yt_tokens';
const PAGE_TOKEN_TTL_MS = 5 * 60e3;
export function getCachedPageTokens(cacheKey) {
try {
ensureSearchCacheTables();
const row = db.prepare(
`SELECT payload_json AS payload, expires_at AS exp FROM search_cache WHERE cache_key = ? AND provider = ?`,
).get(String(cacheKey), PAGE_TOKEN_PROVIDER);
if (!row) return null;
if (Number(row.exp || 0) <= Date.now()) {
// Expiré : on ne le sert pas, et on le retire pour ne pas le relire.
try { db.prepare(`DELETE FROM search_cache WHERE cache_key = ? AND provider = ?`).run(String(cacheKey), PAGE_TOKEN_PROVIDER); } catch {}
return null;
}
const tokens = JSON.parse(row.payload);
return Array.isArray(tokens) ? tokens : null;
} catch { return null; }
}
export function setCachedPageTokens(cacheKey, tokens, ttlMs = PAGE_TOKEN_TTL_MS) {
if (!Array.isArray(tokens) || tokens.length === 0) return false;
try {
ensureSearchCacheTables();
const now = Date.now();
const ttl = Number.isFinite(ttlMs) ? ttlMs : PAGE_TOKEN_TTL_MS;
// Normalisation ICI, à la frontière de persistance, et pas chez l'appelant :
// `JSON.stringify` transforme un trou (`undefined`) en `null`, et un `filter`
// décalerait les index — le jeton de la page 2 se retrouverait à l'index 0
// et la page 3 renverrait celui de la page 4. Une pagination silencieusement
// décalée est indétectable en production ; on garantit donc l'indexation à
// l'écriture, pour tout appelant présent ou futur. `''` reste falsy, donc
// la lecture la traite comme un jeton absent.
const dense = Array.from({ length: tokens.length }, (_, i) => {
const t = tokens[i];
return typeof t === 'string' && t.length > 0 ? t : '';
});
db.prepare(
`INSERT INTO search_cache (cache_key, provider, q, payload_json, item_count, source, hit_count, created_at, expires_at)
VALUES (?, ?, '', ?, ?, 'page_token', 0, ?, ?)
ON CONFLICT(cache_key, provider) DO UPDATE SET payload_json=excluded.payload_json,
item_count=excluded.item_count, created_at=excluded.created_at, expires_at=excluded.expires_at`,
).run(String(cacheKey), PAGE_TOKEN_PROVIDER, JSON.stringify(dense), dense.length, now, now + ttl);
enforceSearchCacheCap(PAGE_TOKEN_PROVIDER, 500);
return true;
} catch { return false; }
}
/** Vide les jetons persistés (tests, et invalidation manuelle). */
export function clearCachedPageTokens() {
try {
ensureSearchCacheTables();
db.prepare(`DELETE FROM search_cache WHERE provider = ?`).run(PAGE_TOKEN_PROVIDER);
return true;
} catch { return false; }
}
/**
* Plafond par fournisseur : purge par `expires_at DESC`, donc les plus
* fraîches survivent (identique au plafond historique de `youtube_search_cache`).
*/
function enforceSearchCacheCap(provider, cap = SEARCH_CACHE_MAX_ROWS) {
try {
db.prepare(
`DELETE FROM search_cache WHERE provider = ? AND cache_key NOT IN (
SELECT cache_key FROM search_cache WHERE provider = ? ORDER BY expires_at DESC LIMIT ?)`,
).run(String(provider), String(provider), Math.max(1, cap));
} catch {}
}
/**
* Purge de fond : supprime les lignes expirees, puis applique le plafond.
* Appelee par l'intervalle de 10 min au boot (task 4.5).
*/
export function pruneSearchCache({ cap = SEARCH_CACHE_MAX_ROWS } = {}) {
let purged = 0;
try {
ensureSearchCacheTables();
purged = db.prepare(`DELETE FROM search_cache WHERE expires_at <= ?`).run(Date.now()).changes || 0;
} catch {}
try {
for (const r of db.prepare(`SELECT DISTINCT provider FROM search_cache`).all()) {
enforceSearchCacheCap(r.provider, cap);
}
} catch {}
return purged;
}
export function countSearchCacheRows(provider) {
try {
ensureSearchCacheTables();
const p = provider ? String(provider).toLowerCase() : null;
const row = p
? db.prepare(`SELECT COUNT(1) AS n FROM search_cache WHERE provider = ?`).get(p)
: db.prepare(`SELECT COUNT(1) AS n FROM search_cache`).get();
return row?.n || 0;
} catch { return 0; }
}
/** Ratio hit/total par fournisseur — alimenta `/api/providers/metrics`. */
export function searchCacheStats() {
try {
ensureSearchCacheTables();
return db.prepare(
`SELECT provider, COUNT(1) AS entries, COALESCE(SUM(hit_count), 0) AS hits,
COALESCE(SUM(item_count), 0) AS items
FROM search_cache GROUP BY provider`,
).all().map((r) => ({ ...r, hitRate: r.hits + r.entries > 0 ? r.hits / (r.hits + r.entries) : 0 }));
} catch { return []; }
}
// -------------------- Phase 4.3 : metriques fournisseurs --------------------
/**
* Compteurs par fenetre horaire et par fournisseur. Generalise `youtube_metrics`
* (qui reste en place pour le quota Data API, granularite jour).
*
* Volontairement **fin** : une ligne par (heure, provider), purgee au-dela de
* `PROVIDER_METRICS_RETENTION_H` (defaut 24 h). La bascule automatique (4.4) lit
* la fenetre d'1 h ; le cache lit `hit_count` de `search_cache`.
*/
function ensureProviderMetricsTable() {
try {
db.exec(`CREATE TABLE IF NOT EXISTS provider_metrics (
hour TEXT NOT NULL, provider TEXT NOT NULL,
calls INTEGER NOT NULL DEFAULT 0, ok_calls INTEGER NOT NULL DEFAULT 0,
errors INTEGER NOT NULL DEFAULT 0, fallback_calls INTEGER NOT NULL DEFAULT 0,
total_latency_ms INTEGER NOT NULL DEFAULT 0,
last_error TEXT, updated_at TEXT NOT NULL,
PRIMARY KEY (hour, provider)
);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_provider_metrics_hour ON provider_metrics(hour);`);
} catch {}
}
ensureProviderMetricsTable();
/** Incrémente les compteurs d'un appel fournisseur terminé. */
export function incProviderMetrics(provider, { ok = true, latencyMs = 0, fallback = false, error = null } = {}) {
try {
ensureProviderMetricsTable();
const p = String(provider || 'unknown').toLowerCase();
const now = new Date();
const hour = now.toISOString().slice(0, 13); // AAAA-MM-JJTHH
const lat = Number.isFinite(latencyMs) && latencyMs > 0 ? Math.min(Math.round(latencyMs), 600000) : 0;
db.prepare(
`INSERT INTO provider_metrics (hour, provider, calls, ok_calls, errors, fallback_calls, total_latency_ms, last_error, updated_at)
VALUES (?, ?, 1, ?, ?, ?, ?, ?, ?)
ON CONFLICT(hour, provider) DO UPDATE SET
calls = calls + 1,
ok_calls = ok_calls + excluded.ok_calls,
errors = errors + excluded.errors,
fallback_calls = fallback_calls + excluded.fallback_calls,
total_latency_ms = total_latency_ms + excluded.total_latency_ms,
last_error = COALESCE(excluded.last_error, provider_metrics.last_error),
updated_at = excluded.updated_at`,
).run(hour, p, ok ? 1 : 0, ok ? 0 : 1, fallback ? 1 : 0, lat, error ? String(error).slice(0, 500) : null, now.toISOString());
} catch {}
}
/** Instantané par fournisseur sur la fenetre demandee (defaut 1 h). */
export function providerMetricsSnapshot({ hours = 1 } = {}) {
try {
ensureProviderMetricsTable();
const since = new Date(Date.now() - Math.max(1, hours) * 3600e3).toISOString().slice(0, 13);
const rows = db.prepare(
`SELECT provider, COALESCE(SUM(calls),0) AS calls, COALESCE(SUM(ok_calls),0) AS okCalls,
COALESCE(SUM(errors),0) AS errors, COALESCE(SUM(fallback_calls),0) AS fallbacks,
COALESCE(SUM(total_latency_ms),0) AS totalLatencyMs
FROM provider_metrics WHERE hour >= ? GROUP BY provider`,
).all(since);
return rows.map((r) => ({
provider: r.provider,
calls: r.calls,
ok: r.okCalls,
errors: r.errors,
fallbacks: r.fallbacks,
avgLatencyMs: r.calls > 0 ? Math.round(r.totalLatencyMs / r.calls) : 0,
errorRate: r.calls > 0 ? r.errors / r.calls : 0,
}));
} catch { return []; }
}
export function purgeProviderMetrics() {
try {
const hours = Math.max(1, Number(process.env.PROVIDER_METRICS_RETENTION_H || 24));
const before = new Date(Date.now() - hours * 3600e3).toISOString().slice(0, 13);
return db.prepare(`DELETE FROM provider_metrics WHERE hour < ?`).run(before).changes || 0;
} catch { return 0; }
}
// -------------------- OAuth : connexions Google / Twitch (import favoris/abos) --------------------
function ensureOAuthTables() {
try {
db.exec(`CREATE TABLE IF NOT EXISTS oauth_connections (
user_id TEXT NOT NULL REFERENCES users(id) ON DELETE CASCADE,
provider TEXT NOT NULL,
external_user_id TEXT,
display_name TEXT,
avatar_url TEXT,
access_token TEXT,
refresh_token TEXT,
expires_at INTEGER,
scopes TEXT,
created_at INTEGER NOT NULL,
updated_at INTEGER NOT NULL,
PRIMARY KEY (user_id, provider)
);`);
db.exec(`CREATE INDEX IF NOT EXISTS idx_oauth_connections_user ON oauth_connections(user_id);`);
// Chaîne YouTube sélectionnée (multi-chaînes / marque) : ajout idempotent.
try {
const cols = db.prepare(`PRAGMA table_info(oauth_connections)`).all().map((c) => c.name);
if (!cols.includes('yt_channel_id')) db.exec(`ALTER TABLE oauth_connections ADD COLUMN yt_channel_id TEXT`);
if (!cols.includes('yt_page_id')) db.exec(`ALTER TABLE oauth_connections ADD COLUMN yt_page_id TEXT`);
} catch {}
} catch {}
}
ensureOAuthTables();
function oauthRowToDto(row, withTokens = false) {
if (!row) return null;
const dto = {
provider: row.provider,
externalUserId: row.external_user_id || null,
displayName: row.display_name || null,
avatarUrl: row.avatar_url || null,
scopes: row.scopes || null,
expiresAt: row.expires_at ?? null,
connectedAt: row.created_at ?? null,
updatedAt: row.updated_at ?? null,
};
if (withTokens) {
dto.accessToken = row.access_token || null;
dto.refreshToken = row.refresh_token || null;
}
dto.ytChannelId = row.yt_channel_id || null;
dto.ytPageId = row.yt_page_id || null;
return dto;
}
export function setOAuthChannel(userId, provider, { channelId, pageId }) {
ensureOAuthTables();
db.prepare(`UPDATE oauth_connections SET yt_channel_id = ?, yt_page_id = ?, updated_at = ?
WHERE user_id = ? AND provider = ?`)
.run(channelId || null, pageId || null, Date.now(), userId, provider);
return getOAuthConnection(userId, provider);
}
export function upsertOAuthConnection({ userId, provider, externalUserId, displayName, avatarUrl, accessToken, refreshToken, expiresAt, scopes }) {
ensureOAuthTables();
const now = Date.now();
const existing = db.prepare(`SELECT * FROM oauth_connections WHERE user_id = ? AND provider = ?`).get(userId, provider);
if (existing) {
db.prepare(`UPDATE oauth_connections SET external_user_id = ?, display_name = ?, avatar_url = ?,
access_token = COALESCE(?, access_token), refresh_token = COALESCE(?, refresh_token),
expires_at = COALESCE(?, expires_at), scopes = COALESCE(?, scopes), updated_at = ?
WHERE user_id = ? AND provider = ?`)
.run(
externalUserId ?? existing.external_user_id, displayName ?? existing.display_name,
avatarUrl ?? existing.avatar_url, accessToken ?? null, refreshToken ?? null,
expiresAt ?? null, scopes ?? null, now, userId, provider,
);
} else {
db.prepare(`INSERT INTO oauth_connections (user_id, provider, external_user_id, display_name, avatar_url,
access_token, refresh_token, expires_at, scopes, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`)
.run(userId, provider, externalUserId || null, displayName || null, avatarUrl || null,
accessToken || null, refreshToken || null, expiresAt ?? null, scopes || null, now, now);
}
return oauthRowToDto(db.prepare(`SELECT * FROM oauth_connections WHERE user_id = ? AND provider = ?`).get(userId, provider));
}
export function listOAuthConnections(userId) {
try {
ensureOAuthTables();
return db.prepare(`SELECT * FROM oauth_connections WHERE user_id = ?`).all(userId).map((r) => oauthRowToDto(r));
} catch { return []; }
}
export function getOAuthConnection(userId, provider, withTokens = false) {
try {
ensureOAuthTables();
return oauthRowToDto(
db.prepare(`SELECT * FROM oauth_connections WHERE user_id = ? AND provider = ?`).get(userId, provider),
withTokens,
);
} catch { return null; }
}
export function updateOAuthTokens(userId, provider, { accessToken, refreshToken, expiresAt }) {
try {
ensureOAuthTables();
db.prepare(`UPDATE oauth_connections SET access_token = COALESCE(?, access_token),
refresh_token = COALESCE(?, refresh_token), expires_at = COALESCE(?, expires_at), updated_at = ?
WHERE user_id = ? AND provider = ?`)
.run(accessToken ?? null, refreshToken ?? null, expiresAt ?? null, Date.now(), userId, provider);
} catch {}
}
export function deleteOAuthConnection(userId, provider) {
try {
ensureOAuthTables();
return db.prepare(`DELETE FROM oauth_connections WHERE user_id = ? AND provider = ?`).run(userId, provider);
} catch { return { changes: 0 }; }
}