import fs from 'node:fs'; import path from 'node:path'; import { DatabaseSync } from 'node:sqlite'; import type { LogEntry, LogLevel, ObsSettings, Platform, RecordingSettings, StreamState, BrowserSettings, WatchSettings, WatchTarget, } from '@stream-control/shared'; import { DEFAULT_BROWSER_SETTINGS, DEFAULT_OBS_SETTINGS, DEFAULT_RECORDING_SETTINGS, DEFAULT_WATCH_SETTINGS, normalizeBrowserSettings, normalizeRecordingSettings, normalizeWatchSettings, safeJsonParse, stripchatProfileUrl, } from '@stream-control/shared'; import { config } from './config.ts'; fs.mkdirSync(path.dirname(config.dbPath), { recursive: true }); export const db = new DatabaseSync(config.dbPath); db.exec('PRAGMA journal_mode = WAL'); db.exec('PRAGMA foreign_keys = ON'); db.exec(` CREATE TABLE IF NOT EXISTS agents ( id TEXT PRIMARY KEY, name TEXT NOT NULL, hostname TEXT, platform TEXT NOT NULL DEFAULT 'unknown', agent_version TEXT, token_hash TEXT NOT NULL, obs_host TEXT NOT NULL DEFAULT '127.0.0.1', obs_port INTEGER NOT NULL DEFAULT 4455, obs_password TEXT NOT NULL DEFAULT '', auto_connect INTEGER NOT NULL DEFAULT 1, notes TEXT, created_at INTEGER NOT NULL, last_seen_at INTEGER ); CREATE TABLE IF NOT EXISTS logs ( id INTEGER PRIMARY KEY AUTOINCREMENT, agent_id TEXT, level TEXT NOT NULL, message TEXT NOT NULL, ts INTEGER NOT NULL ); CREATE INDEX IF NOT EXISTS idx_logs_ts ON logs (ts DESC); -- Profils surveillés : le serveur sonde leur statut et signale les passages -- en direct. Indépendant des agents : on peut veiller sans rien enregistrer. CREATE TABLE IF NOT EXISTS watch_targets ( id TEXT PRIMARY KEY, provider TEXT NOT NULL DEFAULT 'stripchat', username TEXT NOT NULL, label TEXT, agent_id TEXT, notify INTEGER NOT NULL DEFAULT 1, state TEXT NOT NULL DEFAULT 'unknown', raw_status TEXT, state_since INTEGER NOT NULL, last_checked_at INTEGER, last_error TEXT, created_at INTEGER NOT NULL ); CREATE UNIQUE INDEX IF NOT EXISTS idx_targets_identity ON watch_targets (provider, username); `); /** Migrations additives : sûres à rejouer sur une base déjà peuplée. */ function addColumnIfMissing(table: string, column: string, definition: string): void { const columns = db.prepare(`PRAGMA table_info(${table})`).all() as unknown as Array<{ name: string; }>; if (columns.some((entry) => entry.name === column)) return; db.exec(`ALTER TABLE ${table} ADD COLUMN ${column} ${definition}`); } // Surveillance du stream source : stockée en JSON, le schéma évolue plus vite // que la table (nouveaux fournisseurs, nouveaux statuts). addColumnIfMissing('agents', 'watch_json', 'TEXT'); addColumnIfMissing('agents', 'browser_json', 'TEXT'); // Preset d'enregistrement : réglable VM par VM, absent = aucune intervention // de l'agent sur les réglages d'OBS. addColumnIfMissing('agents', 'recording_json', 'TEXT'); // Enrichissement des profils surveillés : photo et historique de diffusion. addColumnIfMissing('watch_targets', 'avatar_url', 'TEXT'); addColumnIfMissing('watch_targets', 'status_changed_at', 'INTEGER'); addColumnIfMissing('watch_targets', 'last_live_started_at', 'INTEGER'); addColumnIfMissing('watch_targets', 'last_live_ended_at', 'INTEGER'); export interface AgentRow { id: string; name: string; hostname: string | null; platform: string; agent_version: string | null; token_hash: string; obs_host: string; obs_port: number; obs_password: string; auto_connect: number; watch_json: string | null; browser_json: string | null; recording_json: string | null; notes: string | null; created_at: number; last_seen_at: number | null; } export interface AgentRecord { id: string; name: string; hostname: string | null; platform: Platform; agentVersion: string | null; tokenHash: string; obs: ObsSettings; autoConnectObs: boolean; watch: WatchSettings; browser: BrowserSettings; recording: RecordingSettings; notes: string | null; createdAt: number; lastSeenAt: number | null; } function toRecord(row: AgentRow): AgentRecord { return { id: row.id, name: row.name, hostname: row.hostname, platform: (row.platform as Platform) ?? 'unknown', agentVersion: row.agent_version, tokenHash: row.token_hash, obs: { host: row.obs_host || DEFAULT_OBS_SETTINGS.host, port: Number(row.obs_port) || DEFAULT_OBS_SETTINGS.port, password: row.obs_password ?? '', }, autoConnectObs: Number(row.auto_connect) === 1, watch: normalizeWatchSettings( row.watch_json ? safeJsonParse(row.watch_json) : DEFAULT_WATCH_SETTINGS, ), browser: normalizeBrowserSettings( row.browser_json ? safeJsonParse(row.browser_json) : DEFAULT_BROWSER_SETTINGS, ), recording: normalizeRecordingSettings( row.recording_json ? safeJsonParse(row.recording_json) : DEFAULT_RECORDING_SETTINGS, ), notes: row.notes, createdAt: Number(row.created_at), lastSeenAt: row.last_seen_at === null ? null : Number(row.last_seen_at), }; } const stmts = { listAgents: db.prepare('SELECT * FROM agents ORDER BY name COLLATE NOCASE'), getAgent: db.prepare('SELECT * FROM agents WHERE id = ?'), getAgentByTokenHash: db.prepare('SELECT * FROM agents WHERE token_hash = ?'), insertAgent: db.prepare(` INSERT INTO agents (id, name, hostname, platform, agent_version, token_hash, obs_host, obs_port, obs_password, auto_connect, watch_json, browser_json, recording_json, notes, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) `), updateIdentity: db.prepare(` UPDATE agents SET hostname = ?, platform = ?, agent_version = ?, last_seen_at = ? WHERE id = ? `), updateSettings: db.prepare(` UPDATE agents SET name = ?, obs_host = ?, obs_port = ?, obs_password = ?, auto_connect = ?, watch_json = ?, browser_json = ?, recording_json = ?, notes = ? WHERE id = ? `), touchAgent: db.prepare('UPDATE agents SET last_seen_at = ? WHERE id = ?'), rotateToken: db.prepare('UPDATE agents SET token_hash = ? WHERE id = ?'), deleteAgent: db.prepare('DELETE FROM agents WHERE id = ?'), insertLog: db.prepare('INSERT INTO logs (agent_id, level, message, ts) VALUES (?, ?, ?, ?)'), recentLogs: db.prepare(` SELECT l.id, l.agent_id, l.level, l.message, l.ts, a.name AS agent_name FROM logs l LEFT JOIN agents a ON a.id = l.agent_id ORDER BY l.id DESC LIMIT ? `), pruneLogs: db.prepare(` DELETE FROM logs WHERE id NOT IN (SELECT id FROM logs ORDER BY id DESC LIMIT ?) `), }; export const agentsRepo = { list(): AgentRecord[] { return (stmts.listAgents.all() as unknown as AgentRow[]).map(toRecord); }, get(id: string): AgentRecord | null { const row = stmts.getAgent.get(id) as unknown as AgentRow | undefined; return row ? toRecord(row) : null; }, findByTokenHash(tokenHash: string): AgentRecord | null { const row = stmts.getAgentByTokenHash.get(tokenHash) as unknown as AgentRow | undefined; return row ? toRecord(row) : null; }, create(input: { id: string; name: string; tokenHash: string; hostname?: string | null; platform?: Platform; agentVersion?: string | null; obs?: Partial; autoConnectObs?: boolean; watch?: WatchSettings; browser?: BrowserSettings; recording?: RecordingSettings; notes?: string | null; }): AgentRecord { const obs = { ...DEFAULT_OBS_SETTINGS, ...input.obs }; stmts.insertAgent.run( input.id, input.name, input.hostname ?? null, input.platform ?? 'unknown', input.agentVersion ?? null, input.tokenHash, obs.host, obs.port, obs.password, input.autoConnectObs === false ? 0 : 1, JSON.stringify(normalizeWatchSettings(input.watch ?? DEFAULT_WATCH_SETTINGS)), JSON.stringify(normalizeBrowserSettings(input.browser ?? DEFAULT_BROWSER_SETTINGS)), JSON.stringify(normalizeRecordingSettings(input.recording ?? DEFAULT_RECORDING_SETTINGS)), input.notes ?? null, Date.now(), ); const created = agentsRepo.get(input.id); if (!created) throw new Error(`Échec de création de l'agent ${input.id}`); return created; }, updateIdentity( id: string, identity: { hostname: string | null; platform: Platform; agentVersion: string | null }, ): void { stmts.updateIdentity.run( identity.hostname, identity.platform, identity.agentVersion, Date.now(), id, ); }, updateSettings( id: string, settings: { name: string; obs: ObsSettings; autoConnectObs: boolean; watch: WatchSettings; browser: BrowserSettings; recording: RecordingSettings; notes: string | null; }, ): void { stmts.updateSettings.run( settings.name, settings.obs.host, settings.obs.port, settings.obs.password, settings.autoConnectObs ? 1 : 0, JSON.stringify(normalizeWatchSettings(settings.watch)), JSON.stringify(normalizeBrowserSettings(settings.browser)), JSON.stringify(normalizeRecordingSettings(settings.recording)), settings.notes, id, ); }, touch(id: string): void { stmts.touchAgent.run(Date.now(), id); }, rotateToken(id: string, tokenHash: string): void { stmts.rotateToken.run(tokenHash, id); }, remove(id: string): void { stmts.deleteAgent.run(id); }, }; // --- Profils surveillés ------------------------------------------------------ interface TargetRow { id: string; provider: string; username: string; label: string | null; agent_id: string | null; notify: number; state: string; raw_status: string | null; state_since: number; last_checked_at: number | null; last_error: string | null; created_at: number; avatar_url: string | null; status_changed_at: number | null; last_live_started_at: number | null; last_live_ended_at: number | null; } const num = (value: number | null): number | null => (value === null ? null : Number(value)); function toTarget(row: TargetRow): WatchTarget { return { id: row.id, provider: (row.provider as WatchTarget['provider']) ?? 'stripchat', username: row.username, label: row.label, url: stripchatProfileUrl(row.username), agentId: row.agent_id, notify: Number(row.notify) === 1, state: (row.state as StreamState) ?? 'unknown', rawStatus: row.raw_status, stateSince: Number(row.state_since), lastCheckedAt: num(row.last_checked_at), lastError: row.last_error, createdAt: Number(row.created_at), avatarUrl: row.avatar_url, statusChangedAt: num(row.status_changed_at), lastLiveStartedAt: num(row.last_live_started_at), lastLiveEndedAt: num(row.last_live_ended_at), }; } const targetStmts = { list: db.prepare('SELECT * FROM watch_targets ORDER BY username COLLATE NOCASE'), get: db.prepare('SELECT * FROM watch_targets WHERE id = ?'), getByName: db.prepare('SELECT * FROM watch_targets WHERE provider = ? AND username = ?'), insert: db.prepare(` INSERT INTO watch_targets (id, provider, username, label, agent_id, notify, state, state_since, created_at) VALUES (?, ?, ?, ?, ?, ?, 'unknown', ?, ?) `), updateSettings: db.prepare( 'UPDATE watch_targets SET label = ?, agent_id = ?, notify = ? WHERE id = ?', ), updateState: db.prepare(` UPDATE watch_targets SET state = ?, raw_status = ?, state_since = ?, last_checked_at = ?, last_error = ?, avatar_url = ?, status_changed_at = ?, last_live_started_at = ?, last_live_ended_at = ? WHERE id = ? `), remove: db.prepare('DELETE FROM watch_targets WHERE id = ?'), }; export const targetsRepo = { list(): WatchTarget[] { return (targetStmts.list.all() as unknown as TargetRow[]).map(toTarget); }, get(id: string): WatchTarget | null { const row = targetStmts.get.get(id) as unknown as TargetRow | undefined; return row ? toTarget(row) : null; }, findByUsername(provider: string, username: string): WatchTarget | null { const row = targetStmts.getByName.get(provider, username) as unknown as TargetRow | undefined; return row ? toTarget(row) : null; }, create(input: { id: string; username: string; provider?: string; label?: string | null; agentId?: string | null; notify?: boolean; }): WatchTarget { const now = Date.now(); targetStmts.insert.run( input.id, input.provider ?? 'stripchat', input.username, input.label ?? null, input.agentId ?? null, input.notify === false ? 0 : 1, now, now, ); const created = targetsRepo.get(input.id); if (!created) throw new Error(`Échec de création du profil surveillé ${input.username}`); return created; }, updateSettings( id: string, settings: { label: string | null; agentId: string | null; notify: boolean }, ): void { targetStmts.updateSettings.run( settings.label, settings.agentId, settings.notify ? 1 : 0, id, ); }, updateState( id: string, state: { state: StreamState; rawStatus: string | null; stateSince: number; lastError: string | null; avatarUrl: string | null; statusChangedAt: number | null; lastLiveStartedAt: number | null; lastLiveEndedAt: number | null; }, ): void { targetStmts.updateState.run( state.state, state.rawStatus, state.stateSince, Date.now(), state.lastError, state.avatarUrl, state.statusChangedAt, state.lastLiveStartedAt, state.lastLiveEndedAt, id, ); }, remove(id: string): void { targetStmts.remove.run(id); }, }; interface LogRow { id: number; agent_id: string | null; agent_name: string | null; level: string; message: string; ts: number; } export const logsRepo = { append(agentId: string | null, level: LogLevel, message: string, ts = Date.now()): LogEntry { const info = stmts.insertLog.run(agentId, level, message, ts); if (Math.random() < 0.02) stmts.pruneLogs.run(config.logRetention); const agent = agentId ? agentsRepo.get(agentId) : null; return { id: Number(info.lastInsertRowid), agentId, agentName: agent?.name ?? null, level, message, ts, }; }, recent(limit = 200): LogEntry[] { const rows = stmts.recentLogs.all(limit) as unknown as LogRow[]; return rows .map((row) => ({ id: Number(row.id), agentId: row.agent_id, agentName: row.agent_name, level: row.level as LogLevel, message: row.message, ts: Number(row.ts), })) .reverse(); }, };