import type { StreamState, WatchTarget } from '@stream-control/shared'; import { DEFAULT_WATCH_SETTINGS, fetchStripchatStatus, preemptionVerdict, } from '@stream-control/shared'; import { config } from './config.ts'; import { sessionsRepo, targetsRepo } from './db.ts'; import { hub } from './hub.ts'; import { notifyTargetLive } from './pushover.ts'; import { startTargetRecording } from './recorder.ts'; /** * Tentatives d'amorçage automatique par diffusion. Les causes d'échec restantes * une fois les garde-fous passés (OBS injoignable, par exemple) sont surtout * persistantes : réessayer indéfiniment toutes les dix secondes noierait le * journal sans rien changer. */ const MAX_AUTO_ATTEMPTS = 3; /** * Sonde périodiquement les profils surveillés et signale les passages en direct. * * Côté serveur et non côté agent : un profil se surveille indépendamment de * toute machine d'enregistrement, et une seule requête suffit quel que soit le * nombre d'agents. */ class Watchlist { private timer: NodeJS.Timeout | null = null; private running = false; /** Tentatives déjà faites pour la diffusion en cours, par profil. */ private readonly autoAttempts = new Map(); start(): void { if (this.timer) return; this.timer = setInterval(() => void this.tick(), config.watchlistIntervalMs); this.timer.unref?.(); void this.tick(); } stop(): void { if (this.timer) clearInterval(this.timer); this.timer = null; } /** Sonde immédiate d'un seul profil, déclenchée depuis le dashboard. */ async checkOne(id: string): Promise { const target = targetsRepo.get(id); if (!target) return null; await this.probe(target); return targetsRepo.get(id); } private async tick(): Promise { if (this.running) return; // un cycle lent ne doit pas s'empiler sur le suivant this.running = true; try { // Du plus prioritaire au moins prioritaire : c'est l'ordre de sondage qui // décide qui tente sa chance en premier sur une VM libre. Sonder au // hasard laisserait un profil mineur s'en emparer, pour se faire couper // dans la foulée par son aîné — le résultat serait le même, au prix d'un // fichier de quelques secondes. const targets = targetsRepo .list() .sort((a, b) => b.priority - a.priority || a.username.localeCompare(b.username)); // Séquentiel et espacé : une poignée de profils ne justifie pas de // marteler l'API en parallèle. for (const target of targets) { await this.probe(target); await delay(250); } } finally { this.running = false; } } private async probe(target: WatchTarget): Promise { let next: StreamState; let raw: string | null = null; let error: string | null = null; let avatarUrl = target.avatarUrl; let statusChangedAt = target.statusChangedAt; try { const result = await fetchStripchatStatus( target.username, DEFAULT_WATCH_SETTINGS.privateStatuses, ); next = result.state; raw = result.raw; // On conserve la dernière valeur connue si la plateforme ne la renvoie pas. avatarUrl = result.avatarUrl ?? target.avatarUrl; statusChangedAt = result.statusChangedAt ?? target.statusChangedAt; } catch (err) { // Une sonde en échec ne change pas l'état connu : on garde le dernier // verdict fiable plutôt que d'annoncer un faux passage hors-ligne. error = err instanceof Error ? err.message : String(err); next = target.state; } const changed = next !== target.state; // Diffusion en cours : sans condition sur `changed`, car l'ouverture est // idempotente. C'est ce qui rattrape un profil déjà en direct au démarrage // du serveur — aucune transition ne se produira pour lui, et sa diffusion // resterait sinon absente de la frise jusqu'à la suivante. if (next === 'public') { sessionsRepo.open(target.id, statusChangedAt ?? Date.now()); } // Fin de diffusion : on fige le stream qui vient de se terminer. Son début // est le statusChangedAt d'avant la transition, que la plateforme vient de // remplacer par celui du nouveau statut. let lastLiveStartedAt = target.lastLiveStartedAt; let lastLiveEndedAt = target.lastLiveEndedAt; if (changed && target.state === 'public') { lastLiveStartedAt = target.statusChangedAt ?? target.stateSince; lastLiveEndedAt = Date.now(); sessionsRepo.close(target.id, lastLiveEndedAt); } targetsRepo.updateState(target.id, { state: next, rawStatus: raw ?? target.rawStatus, stateSince: changed ? Date.now() : target.stateSince, lastError: error, avatarUrl, statusChangedAt, lastLiveStartedAt, lastLiveEndedAt, }); const updated = targetsRepo.get(target.id); if (!updated) return; hub.publishTarget(updated); this.maybeAutoRecord(updated); if (changed) { const name = updated.label ?? updated.username; hub.log(null, 'info', `Veille : ${name} est passé « ${labelOf(next)} »`); // Seul le passage effectif au flux public déclenche une notification. if (next === 'public' && target.state !== 'public' && updated.notify) { hub.publishTargetLive(updated); void notifyTargetLive(updated.label ?? updated.username, updated.url); } } } /** * Déclenche l'enregistrement d'un profil en automatisme, une fois qu'il est en * direct depuis assez longtemps. * * Le délai d'amorçage n'est pas une précaution de style : un modèle qui sort * d'un show privé repasse « public » quelques secondes avant de se remettre en * place, et un flux qui redémarre alterne parfois plusieurs fois. Enregistrer * sur la première lecture produirait des fichiers de dix secondes. */ private maybeAutoRecord(target: WatchTarget): void { if (target.state !== 'public') { // Diffusion terminée : la suivante aura droit à ses propres tentatives. this.autoAttempts.delete(target.id); return; } if (!target.autoRecord || !target.agentId) return; if ((this.autoAttempts.get(target.id) ?? 0) >= MAX_AUTO_ATTEMPTS) return; // `statusChangedAt` vient de la plateforme et vaut mieux que notre première // observation : il survit à un redémarrage du serveur. const liveSince = target.statusChangedAt ?? target.stateSince; if (Date.now() - liveSince < config.autoRecordDelayMs) return; if (!hub.isOnline(target.agentId)) return; // VM occupée : par défaut la capture en place l'emporte, sauf si ce profil // a reçu le droit d'interrompre et le rang pour le faire. Le verdict est // recalculé par startTargetRecording — le refaire ici évite seulement de // consommer une des trois tentatives à chaque cycle pour une interruption // qui restera refusée tant que rien ne change. if (hub.statusOf(target.agentId).recording) { const verdict = preemptionVerdict(target, hub.recordingTarget(target.agentId)); if (!verdict.allowed) return; } const agentId = target.agentId; const name = target.label ?? target.username; // Marqué avant l'appel : la séquence de capture dure plusieurs secondes, et // le cycle suivant ne doit pas en lancer une seconde en parallèle. const attempt = (this.autoAttempts.get(target.id) ?? 0) + 1; this.autoAttempts.set(target.id, attempt); void startTargetRecording(target) .then(() => { this.autoAttempts.set(target.id, MAX_AUTO_ATTEMPTS); hub.log( agentId, 'info', `Automatisme : enregistrement de ${name} démarré`, Date.now(), 'capture.started', ); }) .catch((err: unknown) => { const message = err instanceof Error ? err.message : String(err); const giveUp = attempt >= MAX_AUTO_ATTEMPTS; hub.log( agentId, 'warn', `Automatisme : démarrage de ${name} en échec (${attempt}/${MAX_AUTO_ATTEMPTS}) — ${message}` + (giveUp ? '. Abandon jusqu\'à la prochaine diffusion.' : ''), Date.now(), 'command.failed', ); }); } } function labelOf(state: StreamState): string { switch (state) { case 'public': return 'en direct'; case 'private': return 'en show privé'; case 'offline': return 'hors-ligne'; default: return 'inconnu'; } } function delay(ms: number): Promise { return new Promise((resolve) => setTimeout(resolve, ms)); } export const watchlist = new Watchlist();