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