Files
NewTube/server/db.mjs
T
bruno df53445efe feat(oauth): import favoris/abos Google (YouTube) et Twitch vers abonnements/likes
- Backend: server/oauth.mjs (auth URL, echange code, refresh, fetch
  subscriptions/likes Google, follows Twitch), table oauth_connections
  + migration, routes /api/oauth/:provider/{url,callback,preview,import},
  /api/oauth/{status,connections}
- Front: /library/import (connect, preview, import), service
  oauth-import, lien depuis Abonnements
- Config: GOOGLE_CLIENT_ID/SECRET/REDIRECT_URI, TWITCH_REDIRECT_URI,
  OAUTH_APP_BASE_URL documentes (.env.example, docker-compose),
  passthrough docker-compose.local.yml, roadmap README cochee
2026-09-26 18:47:02 -04:00

1463 lines
64 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
const db = new Database(dbFile);
db.pragma('foreign_keys = ON');
// 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);
}
// 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 {}
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 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 --------------------
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 };
}
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 = `%${q.trim()}%`;
conds.push(`(query LIKE ? OR COALESCE(filters_json,'') LIKE ?)`);
params.push(like, like);
}
if (provider && String(provider).trim()) {
const likeP = `%${String(provider).trim()}%`;
conds.push(`COALESCE(filters_json,'') LIKE ?`);
params.push(likeP);
}
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;
// 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, provider, videoId, title || null, thumbnail || null, watched_at, progressSeconds, durationSeconds, (typeof lastPositionSeconds === 'number' ? lastPositionSeconds : null), now
);
// Return the row id
const row = db.prepare(`SELECT * FROM watch_history WHERE user_id = ? AND provider = ? AND video_id = ?`).get(userId, provider, 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)); }
if (provider && String(provider).trim()) {
conds.push('provider = ?');
params.push(String(provider).trim().toLowerCase());
}
if (typeof q === 'string' && q.trim().length > 0) {
const like = `%${q.trim()}%`;
conds.push(`(COALESCE(title,'') LIKE ? OR provider LIKE ? OR video_id LIKE ?)`);
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)); }
if (provider && String(provider).trim()) {
conds.push('provider = ?');
params.push(String(provider).trim().toLowerCase());
}
if (lang && String(lang).trim()) {
conds.push('lang = ?');
params.push(String(lang).trim().slice(0, 12).toLowerCase());
}
if (typeof q === 'string' && q.trim().length > 0) {
const like = `%${q.trim()}%`;
conds.push(`(COALESCE(title,'') LIKE ? OR video_id LIKE ? OR lang LIKE ? OR COALESCE(lines_json,'') LIKE ?)`);
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()) {
db.prepare(`DELETE FROM transcript_history WHERE user_id = ? AND provider = ?`).run(userId, String(provider).trim().toLowerCase());
return;
}
db.prepare(`DELETE FROM transcript_history WHERE user_id = ?`).run(userId);
}
// -------------------- Likes (via tags/video_tags) --------------------
function ensureTag(userId, name) {
const tag = db.prepare(`SELECT id FROM tags WHERE user_id = ? AND name = ?`).get(userId, name);
if (tag && tag.id) return tag.id;
const id = cryptoRandomId();
db.prepare(`INSERT INTO tags (id, user_id, name, created_at) VALUES (?, ?, ?, ?)`)
.run(id, userId, name, nowIso());
return id;
}
export function likeVideo({ userId, provider, videoId, title, thumbnail }) {
const tagId = ensureTag(userId, 'like');
// Upsert-like behavior; ignore if exists
db.prepare(`INSERT OR IGNORE INTO video_tags (user_id, provider, video_id, tag_id, created_at) VALUES (?, ?, ?, ?, ?)`)
.run(userId, provider, videoId, tagId, nowIso());
// Also update the watch_history table with the title and thumbnail
if (title || thumbnail) {
upsertWatchHistory({ userId, provider, videoId, title, thumbnail });
}
return { provider, video_id: videoId };
}
export function unlikeVideo({ userId, provider, videoId }) {
const tag = db.prepare(`SELECT id FROM tags WHERE user_id = ? AND name = ?`).get(userId, 'like');
if (!tag) return { removed: false };
const info = db.prepare(`DELETE FROM video_tags WHERE user_id = ? AND provider = ? AND video_id = ? AND tag_id = ?`)
.run(userId, provider, videoId, tag.id);
return { removed: (info.changes || 0) > 0 };
}
export function isVideoLiked({ userId, provider, videoId }) {
const tag = db.prepare(`SELECT id FROM tags WHERE user_id = ? AND name = ?`).get(userId, 'like');
if (!tag) return false;
const row = db.prepare(`SELECT 1 FROM video_tags WHERE user_id = ? AND provider = ? AND video_id = ? AND tag_id = ?`)
.get(userId, provider, videoId, tag.id);
return Boolean(row);
}
export function listLikedVideos({ userId, limit = 100, q }) {
try {
console.log(`[listLikedVideos] Récupération des vidéos aimées pour l'utilisateur ${userId}`);
// Vérifier que la table tags existe
const tableExists = db.prepare(`
SELECT name FROM sqlite_master
WHERE type='table' AND name='tags'
`).get();
if (!tableExists) {
console.error('[listLikedVideos] La table tags n\'existe pas');
return [];
}
// Récupérer le tag 'like' de l'utilisateur
const tag = db.prepare(`
SELECT id FROM tags
WHERE user_id = ? AND name = ?
`).get(userId, 'like');
if (!tag) {
console.log(`[listLikedVideos] Aucun tag 'like' trouvé pour l'utilisateur ${userId}`);
return [];
}
console.log(`[listLikedVideos] Tag ID: ${tag.id}`);
// Vérifier que la table video_tags existe
const videoTagsExists = db.prepare(`
SELECT name FROM sqlite_master
WHERE type='table' AND name='video_tags'
`).get();
if (!videoTagsExists) {
console.error('[listLikedVideos] La table video_tags n\'existe pas');
return [];
}
// Récupérer les vidéos aimées avec les métadonnées de l'historique
// Note: La colonne thumbnail n'existe pas dans la table watch_history
const hasQ = typeof q === 'string' && q.trim().length > 0;
const like = `%${(q || '').trim()}%`;
const base = `
SELECT
vt.provider,
vt.video_id,
vt.created_at,
COALESCE(wh.title, '') AS title,
COALESCE(wh.thumbnail, '') AS thumbnail,
wh.last_watched_at AS last_watched_at
FROM video_tags vt
LEFT JOIN watch_history wh
ON wh.user_id = vt.user_id
AND wh.provider = vt.provider
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 (COALESCE(wh.title,'') LIKE ? OR vt.provider LIKE ? OR vt.video_id LIKE ?)
${orderLimit}`
: `${base} ${orderLimit}`;
console.log('[listLikedVideos] Exécution de la requête:', query.replace(/\s+/g, ' ').trim());
const rows = hasQ
? db.prepare(query).all(userId, tag.id, like, like, like, limit)
: db.prepare(query).all(userId, tag.id, limit);
console.log(`[listLikedVideos] ${rows.length} vidéos trouvées`);
return rows;
} catch (error) {
console.error('[listLikedVideos] Erreur:', error.message);
console.error(error.stack);
throw error; // Renvoyer l'erreur pour qu'elle soit gérée par le routeur
}
}
// -------------------- Playlists --------------------
export function recordPlaylistMetric({ userId, playlistId, action, meta }) {
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 };
}
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';
const info = db.prepare(`DELETE FROM playlists WHERE id = ?`).run(id);
recordPlaylistMetric({ userId, playlistId: id, action: 'delete' });
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);
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,
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,
};
}
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,
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, url, subs_count, verified, last_refreshed_at)
VALUES (@provider, @external_id, @title, @handle, @avatar_url, @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,
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,
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();
export function getCachedYoutubeSearch(qHash) {
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; }
}
export function setCachedYoutubeSearch(qHash, q, items, source, ttlMs) {
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() {
try { db.prepare(`DELETE FROM youtube_search_cache WHERE expires_at <= ?`).run(Date.now()); } catch {}
}
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; }
}
// -------------------- 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);`);
} 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;
}
return dto;
}
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 }; }
}