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` * (`|`) : 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 }; } }