Files
stream-control/packages/agent/src/watcher.ts
2026-08-11 00:52:33 +02:00

285 lines
9.1 KiB
TypeScript

import { EventEmitter } from 'node:events';
import type { LogLevel, StreamState, WatchSettings, WatchState } from '@stream-control/shared';
import { emptyWatchState, mapStreamStatus } from '@stream-control/shared';
import { sendHotkey } from './hotkey.ts';
const PROBE_TIMEOUT_MS = 8000;
const BROWSER_UA =
'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0 Safari/537.36';
export interface ProbeResult {
/** Statut brut renvoyé par la plateforme. */
raw: string;
state: StreamState;
}
/**
* Sonde Stripchat.
*
* Le champ autoritatif est `user.user.status` sur
* `/api/front/v2/models/username/{pseudo}/cam`. Valeurs relevées en production :
* `public`, `private`, `p2p`, `groupShow`, `idle`.
* (`cam.privateMode` existe mais reste vide y compris pendant un privé : ne pas
* s'en servir.)
*/
export async function probeStripchat(
username: string,
privateStatuses: string[],
): Promise<ProbeResult> {
const url = `https://fr.stripchat.com/api/front/v2/models/username/${encodeURIComponent(
username,
)}/cam`;
const response = await fetch(url, {
headers: { 'user-agent': BROWSER_UA, accept: 'application/json' },
signal: AbortSignal.timeout(PROBE_TIMEOUT_MS),
});
if (response.status === 404) return { raw: 'notFound', state: 'offline' };
if (!response.ok) throw new Error(`API Stripchat : HTTP ${response.status}`);
const payload = (await response.json()) as {
user?: { user?: { status?: string; isLive?: boolean } };
};
const raw = payload?.user?.user?.status;
if (typeof raw !== 'string') {
throw new Error('Réponse Stripchat inattendue : user.user.status absent');
}
return { raw, state: mapStreamStatus(raw, privateStatuses) };
}
/** Ce que le surveillant doit pouvoir demander à OBS. */
export interface WatchActions {
recordState(): Promise<{ active: boolean; paused: boolean }>;
pauseRecording(): Promise<void>;
resumeRecording(): Promise<void>;
}
/**
* Surveille le statut du stream source et met l'enregistrement en pause pendant
* les shows privés, puis le reprend et rappelle le plein écran à la reprise du
* flux public.
*
* Deux garde-fous délibérés :
* - la mise en pause exige N lectures « privé » consécutives, la reprise agit
* immédiatement (une fausse pause perd du contenu, une fausse reprise ne
* coûte que quelques secondes d'écran d'attente) ;
* - une erreur de sonde ne déclenche jamais d'action : on conserve l'état connu.
*/
export class StreamWatcher extends EventEmitter {
private settings: WatchSettings;
private state: WatchState;
private timer: NodeJS.Timeout | null = null;
private fullscreenTimer: NodeJS.Timeout | null = null;
private ticking = false;
/** Le flux public a été interrompu : il faudra rappeler le plein écran. */
private fullscreenPending = false;
constructor(
settings: WatchSettings,
private readonly actions: WatchActions,
) {
super();
this.settings = settings;
this.state = emptyWatchState(settings);
}
get snapshot(): WatchState {
return { ...this.state };
}
applySettings(settings: WatchSettings): void {
const restart =
settings.enabled !== this.settings.enabled ||
settings.username !== this.settings.username ||
settings.provider !== this.settings.provider ||
settings.pollIntervalMs !== this.settings.pollIntervalMs;
const identityChanged =
settings.username !== this.settings.username || settings.provider !== this.settings.provider;
this.settings = settings;
this.state.enabled = settings.enabled;
this.state.provider = settings.provider;
this.state.username = settings.username;
if (identityChanged) {
this.state.state = 'unknown';
this.state.rawStatus = undefined;
this.state.since = Date.now();
this.state.pendingConfirmations = 0;
this.state.autoPaused = false;
this.fullscreenPending = false;
}
if (restart) this.start();
}
start(): void {
this.stop();
if (!this.settings.enabled || !this.settings.username) return;
this.emit(
'log',
'info',
`Surveillance de « ${this.settings.username} » (${this.settings.provider}) toutes les ${
this.settings.pollIntervalMs / 1000
} s`,
);
void this.tick();
this.timer = setInterval(() => void this.tick(), this.settings.pollIntervalMs);
this.timer.unref?.();
}
stop(): void {
if (this.timer) clearInterval(this.timer);
this.timer = null;
if (this.fullscreenTimer) clearTimeout(this.fullscreenTimer);
this.fullscreenTimer = null;
}
/** Sonde immédiate, utilisée par l'action `watch.check`. */
async checkNow(): Promise<WatchState> {
await this.tick();
return this.snapshot;
}
private async tick(): Promise<void> {
if (this.ticking) return; // une sonde lente ne doit pas s'empiler
if (!this.settings.enabled || !this.settings.username) return;
this.ticking = true;
try {
const result = await probeStripchat(this.settings.username, this.settings.privateStatuses);
this.state.lastCheckedAt = Date.now();
this.state.lastError = undefined;
this.state.rawStatus = result.raw;
await this.transition(result.state);
} catch (err) {
const message = err instanceof Error ? err.message : String(err);
this.state.lastCheckedAt = Date.now();
// Sonde en échec : on ne met surtout pas l'enregistrement en pause.
if (this.state.lastError !== message) {
this.emit('log', 'warn', `Sonde « ${this.settings.username} » en échec : ${message}`);
}
this.state.lastError = message;
} finally {
this.ticking = false;
}
}
private async transition(next: StreamState): Promise<void> {
const previous = this.state.state;
if (next !== previous) {
this.state.state = next;
this.state.since = Date.now();
this.emit(
'log',
'info',
`Stream « ${this.settings.username} » : ${label(previous)}${label(next)} (${this.state.rawStatus})`,
);
}
this.state.pendingConfirmations = next === 'private' ? this.state.pendingConfirmations + 1 : 0;
if (next === 'private') {
// Le lecteur quitte le plein écran dès que l'overlay de show privé apparaît.
this.fullscreenPending = true;
await this.handlePrivate();
return;
}
if (next === 'public') await this.handlePublic();
// 'offline' / 'unknown' : on ne touche à rien, l'opérateur reste maître.
}
private async handlePrivate(): Promise<void> {
if (!this.settings.pauseOnPrivate || this.state.autoPaused) return;
if (this.state.pendingConfirmations < this.settings.confirmations) return;
try {
const record = await this.actions.recordState();
if (!record.active || record.paused) return;
await this.actions.pauseRecording();
this.state.autoPaused = true;
this.emit('log', 'info', 'Show privé détecté : enregistrement mis en pause');
} catch (err) {
this.emit(
'log',
'error',
`Mise en pause automatique impossible : ${err instanceof Error ? err.message : String(err)}`,
);
}
}
private async handlePublic(): Promise<void> {
if (this.state.autoPaused && this.settings.resumeOnPublic) {
try {
const record = await this.actions.recordState();
if (record.active && record.paused) {
await this.actions.resumeRecording();
this.emit('log', 'info', 'Flux public rétabli : enregistrement repris');
}
} catch (err) {
this.emit(
'log',
'error',
`Reprise automatique impossible : ${err instanceof Error ? err.message : String(err)}`,
);
} finally {
this.state.autoPaused = false;
}
}
if (this.fullscreenPending && this.settings.fullscreen.enabled) {
this.fullscreenPending = false;
this.scheduleFullscreen();
} else {
this.fullscreenPending = false;
}
}
/** Laisse au lecteur le temps de recharger le flux public avant d'envoyer la touche. */
private scheduleFullscreen(): void {
if (this.fullscreenTimer) clearTimeout(this.fullscreenTimer);
this.fullscreenTimer = setTimeout(() => {
this.fullscreenTimer = null;
void this.restoreFullscreen();
}, this.settings.fullscreen.delayMs);
this.fullscreenTimer.unref?.();
}
async restoreFullscreen(): Promise<{ method: string; window?: string }> {
const { key, windowMatch } = this.settings.fullscreen;
try {
const result = await sendHotkey({ key, windowMatch });
this.emit(
'log',
'info',
`Plein écran rappelé : touche « ${key} » envoyée à « ${result.window ?? windowMatch} »`,
);
return result;
} catch (err) {
const message = err instanceof Error ? err.message : String(err);
this.emit('log', 'warn', `Rappel du plein écran impossible : ${message}`);
throw err;
}
}
}
function label(state: StreamState): string {
switch (state) {
case 'public':
return 'public';
case 'private':
return 'privé';
case 'offline':
return 'hors-ligne';
default:
return 'inconnu';
}
}
export type { LogLevel };