Files
stream-control/packages/server/src/db.ts
jeanotx32 1dd3959f29
Some checks failed
release / build (push) Successful in 29s
release / verify-windows (push) Failing after 1m0s
Feat : Log recorder
2026-08-11 22:34:45 +02:00

531 lines
16 KiB
TypeScript

import fs from 'node:fs';
import path from 'node:path';
import { DatabaseSync } from 'node:sqlite';
import type {
AgentEvent,
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,
isAgentEvent,
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');
// Nature de l'évènement journalisé : classe l'entrée dans l'historique d'une VM.
// Nul sur les entrées écrites avant cette colonne, et sur celles des agents non
// mis à jour — l'interface s'en accommode.
addColumnIfMissing('logs', 'event', '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<WatchSettings>(row.watch_json) : DEFAULT_WATCH_SETTINGS,
),
browser: normalizeBrowserSettings(
row.browser_json ? safeJsonParse<BrowserSettings>(row.browser_json) : DEFAULT_BROWSER_SETTINGS,
),
recording: normalizeRecordingSettings(
row.recording_json
? safeJsonParse<RecordingSettings>(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, event) VALUES (?, ?, ?, ?, ?)',
),
recentLogs: db.prepare(`
SELECT l.id, l.agent_id, l.level, l.message, l.ts, l.event, a.name AS agent_name
FROM logs l LEFT JOIN agents a ON a.id = l.agent_id
ORDER BY l.id DESC LIMIT ?
`),
agentLogs: db.prepare(`
SELECT l.id, l.agent_id, l.level, l.message, l.ts, l.event, a.name AS agent_name
FROM logs l LEFT JOIN agents a ON a.id = l.agent_id
WHERE 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<ObsSettings>;
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;
event: string | null;
}
function toLogEntry(row: LogRow): LogEntry {
return {
id: Number(row.id),
agentId: row.agent_id,
agentName: row.agent_name,
level: row.level as LogLevel,
message: row.message,
ts: Number(row.ts),
event: isAgentEvent(row.event) ? row.event : null,
};
}
export const logsRepo = {
append(
agentId: string | null,
level: LogLevel,
message: string,
ts = Date.now(),
event: AgentEvent | null = null,
): LogEntry {
const info = stmts.insertLog.run(agentId, level, message, ts, event);
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,
event,
};
},
recent(limit = 200): LogEntry[] {
const rows = stmts.recentLogs.all(limit) as unknown as LogRow[];
return rows.map(toLogEntry).reverse();
},
/** Historique d'une VM, du plus ancien au plus récent. */
forAgent(agentId: string, limit = 300): LogEntry[] {
const rows = stmts.agentLogs.all(agentId, limit) as unknown as LogRow[];
return rows.map(toLogEntry).reverse();
},
};