diff --git a/.env.example b/.env.example index bea0b76..8f86453 100644 --- a/.env.example +++ b/.env.example @@ -13,6 +13,10 @@ COMMAND_PREFIX=! # Удалять сообщение с командой после её распознавания (нужно право ManageMessages). DELETE_COMMAND_MESSAGES=true +# Имя LiveKit-ноды из Revolt.toml, секция [hosts.livekit]. +# В стандартном self-hosted это "worldwide". +VOICE_NODE=worldwide + # ------------------------------------------------------------- Веб-панель --- PORT=3005 HOST=0.0.0.0 diff --git a/README.md b/README.md index 12518c4..4505689 100644 --- a/README.md +++ b/README.md @@ -149,7 +149,25 @@ cd /opt/stoat-mbot && docker compose up -d --build && docker compose logs -f резолвится в хост: `docker compose exec mbot node -e "fetch(process.env.STOAT_API_URL).then(r=>console.log(r.status))"`. 2. **`joining voice channel`, но нет `voice connection established`.** Значит `join_call` отдал токен, а WebSocket до LiveKit не поднялся — смотрите, доступен ли `/livekit` через внешний домен. -3. **Бот в канале, но звука нет.** Это уже медиа-трафик: LiveKit анонсирует клиентам свой адрес +3. **`Stoat считает, что бот уже в этом голосовом канале` (`AlreadyConnected`).** Зависшее + состояние в Redis после падения бота: `join_call` регистрирует участника, а выйти он не успел. + Ботам API запрещает `force_disconnect`, поэтому чистим вручную — в каталоге инстанса: + + ```bash + docker compose exec redis valkey-cli --scan --pattern 'vc:*' + ``` + + ```bash + docker compose exec redis valkey-cli DEL 'vc:' + ``` + + ```bash + docker compose exec redis valkey-cli SREM 'vc_members:' '' + ``` + + В норме состояние снимает `voice-ingress` по вебхуку от LiveKit — если ситуация повторяется + после каждого перезапуска, смотрите `docker compose logs voice-ingress`. +4. **Бот в канале, но звука нет.** Это уже медиа-трафик: LiveKit анонсирует клиентам свой адрес из `rtc.node_ip` / `use_external_ip` в `/opt/stoat/livekit.yml` и ждёт UDP на 50000-50100. Если анонсируется внешний IP, а роутер не умеет NAT loopback, пакеты от контейнера до него не дойдут. Тогда либо включите hairpin на роутере, либо запустите бота внутри compose-проекта diff --git a/src/config.ts b/src/config.ts index dc55112..1f42770 100644 --- a/src/config.ts +++ b/src/config.ts @@ -1,79 +1,81 @@ -import { readFileSync, existsSync } from "node:fs"; -import { z } from "zod"; - -// Minimal .env loader so we don't need an extra dependency. -function loadDotEnv(path = ".env"): void { - if (!existsSync(path)) return; - for (const rawLine of readFileSync(path, "utf8").split(/\r?\n/)) { - const line = rawLine.trim(); - if (!line || line.startsWith("#")) continue; - const eq = line.indexOf("="); - if (eq === -1) continue; - const key = line.slice(0, eq).trim(); - let value = line.slice(eq + 1).trim(); - if ( - (value.startsWith('"') && value.endsWith('"')) || - (value.startsWith("'") && value.endsWith("'")) - ) { - value = value.slice(1, -1); - } - if (process.env[key] === undefined) process.env[key] = value; - } -} - -loadDotEnv(); - -const schema = z.object({ - STOAT_API_URL: z.string().url(), - STOAT_BOT_TOKEN: z.string().min(1), - COMMAND_PREFIX: z.string().min(1).default("!"), - /** Remove the invoking message after a command is recognised. Needs ManageMessages. */ - DELETE_COMMAND_MESSAGES: z - .enum(["true", "false"]) - .default("true") - .transform((value) => value === "true"), - - PORT: z.coerce.number().int().positive().default(3005), - HOST: z.string().default("0.0.0.0"), - PUBLIC_URL: z.string().url(), - JWT_SECRET: z.string().min(16), - SESSION_TTL_HOURS: z.coerce.number().positive().default(168), - - YTDLP_PATH: z.string().default("yt-dlp"), - YTDLP_COOKIES: z.string().optional(), - /** Extra `--extractor-args` values, separated by ";" — e.g. youtube:player_client=default,web_safari */ - YTDLP_EXTRACTOR_ARGS: z.string().optional(), - LOCAL_MEDIA_DIR: z.string().optional(), - - DEFAULT_VOLUME: z.coerce.number().min(0).max(200).default(60), - MAX_QUEUE_SIZE: z.coerce.number().int().positive().default(500), - SEARCH_RESULT_LIMIT: z.coerce.number().int().positive().max(25).default(10), - IDLE_TIMEOUT_SECONDS: z.coerce.number().int().min(0).default(300), - DJ_ROLE_NAME: z.string().default("DJ"), - REQUIRE_DJ_ROLE: z - .enum(["true", "false"]) - .default("false") - .transform((value) => value === "true"), - - LOG_LEVEL: z.string().default("info"), - NODE_ENV: z.string().default("development"), -}); - -const parsed = schema.safeParse(process.env); - -if (!parsed.success) { - const issues = parsed.error.issues - .map((i) => ` - ${i.path.join(".")}: ${i.message}`) - .join("\n"); - console.error(`Invalid configuration, check your .env file:\n${issues}`); - process.exit(1); -} - -export const config = { - ...parsed.data, - STOAT_API_URL: parsed.data.STOAT_API_URL.replace(/\/+$/, ""), - PUBLIC_URL: parsed.data.PUBLIC_URL.replace(/\/+$/, ""), - isProduction: parsed.data.NODE_ENV === "production", -}; - -export type Config = typeof config; +import { readFileSync, existsSync } from "node:fs"; +import { z } from "zod"; + +// Minimal .env loader so we don't need an extra dependency. +function loadDotEnv(path = ".env"): void { + if (!existsSync(path)) return; + for (const rawLine of readFileSync(path, "utf8").split(/\r?\n/)) { + const line = rawLine.trim(); + if (!line || line.startsWith("#")) continue; + const eq = line.indexOf("="); + if (eq === -1) continue; + const key = line.slice(0, eq).trim(); + let value = line.slice(eq + 1).trim(); + if ( + (value.startsWith('"') && value.endsWith('"')) || + (value.startsWith("'") && value.endsWith("'")) + ) { + value = value.slice(1, -1); + } + if (process.env[key] === undefined) process.env[key] = value; + } +} + +loadDotEnv(); + +const schema = z.object({ + STOAT_API_URL: z.string().url(), + STOAT_BOT_TOKEN: z.string().min(1), + COMMAND_PREFIX: z.string().min(1).default("!"), + /** LiveKit node name from Revolt.toml ([hosts.livekit]); self-hosted default is "worldwide". */ + VOICE_NODE: z.string().min(1).default("worldwide"), + /** Remove the invoking message after a command is recognised. Needs ManageMessages. */ + DELETE_COMMAND_MESSAGES: z + .enum(["true", "false"]) + .default("true") + .transform((value) => value === "true"), + + PORT: z.coerce.number().int().positive().default(3005), + HOST: z.string().default("0.0.0.0"), + PUBLIC_URL: z.string().url(), + JWT_SECRET: z.string().min(16), + SESSION_TTL_HOURS: z.coerce.number().positive().default(168), + + YTDLP_PATH: z.string().default("yt-dlp"), + YTDLP_COOKIES: z.string().optional(), + /** Extra `--extractor-args` values, separated by ";" — e.g. youtube:player_client=default,web_safari */ + YTDLP_EXTRACTOR_ARGS: z.string().optional(), + LOCAL_MEDIA_DIR: z.string().optional(), + + DEFAULT_VOLUME: z.coerce.number().min(0).max(200).default(60), + MAX_QUEUE_SIZE: z.coerce.number().int().positive().default(500), + SEARCH_RESULT_LIMIT: z.coerce.number().int().positive().max(25).default(10), + IDLE_TIMEOUT_SECONDS: z.coerce.number().int().min(0).default(300), + DJ_ROLE_NAME: z.string().default("DJ"), + REQUIRE_DJ_ROLE: z + .enum(["true", "false"]) + .default("false") + .transform((value) => value === "true"), + + LOG_LEVEL: z.string().default("info"), + NODE_ENV: z.string().default("development"), +}); + +const parsed = schema.safeParse(process.env); + +if (!parsed.success) { + const issues = parsed.error.issues + .map((i) => ` - ${i.path.join(".")}: ${i.message}`) + .join("\n"); + console.error(`Invalid configuration, check your .env file:\n${issues}`); + process.exit(1); +} + +export const config = { + ...parsed.data, + STOAT_API_URL: parsed.data.STOAT_API_URL.replace(/\/+$/, ""), + PUBLIC_URL: parsed.data.PUBLIC_URL.replace(/\/+$/, ""), + isProduction: parsed.data.NODE_ENV === "production", +}; + +export type Config = typeof config; diff --git a/src/core/manager.ts b/src/core/manager.ts index f381967..8ec32e0 100644 --- a/src/core/manager.ts +++ b/src/core/manager.ts @@ -1,303 +1,303 @@ -import { EventEmitter } from "node:events"; -import { config } from "../config.js"; -import { logger } from "../logger.js"; -import { resolveQuery, searchTracks } from "../sources/index.js"; -import { - UserFacingError, - type LoopMode, - type PlayerSnapshot, - type Requester, - type SearchResult, - type Track, -} from "../types.js"; -import { GuildPlayer, type Notice, type PositionUpdate } from "./player.js"; -import { Revoice, type RevoiceLike } from "./revoice.js"; - -const log = logger.child({ mod: "manager" }); - -export interface VoiceChannelRef { - id: string; - name: string; -} - -export interface ServerRef { - id: string; - name: string; - iconUrl: string | null; -} - -/** - * Everything the core needs to know about the chat side of Stoat. Implemented on - * top of the bot's stoat.js client so the player itself stays testable. - */ -export interface StoatContext { - findUserVoiceChannel(serverId: string, userId: string): VoiceChannelRef | null; - getVoiceChannel(channelId: string): VoiceChannelRef | null; - listVoiceChannels(serverId: string): VoiceChannelRef[]; - getServerName(serverId: string): string | null; - listServersForUser(userId: string): Promise; - isMember(serverId: string, userId: string): Promise; - canControl(serverId: string, userId: string): Promise; - sendMessage(channelId: string, content: string): Promise; -} - -export type PlayMode = "append" | "next" | "now"; - -export interface ManagerEvents { - update: [PlayerSnapshot]; - position: [PositionUpdate]; -} - -export interface PlayOutcome extends SearchResult { - startedNow: boolean; - queuePosition: number; -} - -/** - * Owns one GuildPlayer per server and exposes the high-level operations that - * both the chat commands and the web panel call into. - */ -export class MusicManager extends EventEmitter { - private readonly players = new Map(); - private readonly revoice: RevoiceLike; - private stoat: StoatContext | null = null; - - constructor() { - super(); - this.revoice = new Revoice(config.STOAT_BOT_TOKEN, { baseURL: config.STOAT_API_URL }); - } - - attachStoat(context: StoatContext): void { - this.stoat = context; - } - - private get chat(): StoatContext { - if (!this.stoat) throw new UserFacingError("Бот ещё не подключился к Stoat"); - return this.stoat; - } - - // ---------------------------------------------------------------- players --- - - get(serverId: string): GuildPlayer | undefined { - return this.players.get(serverId); - } - - list(): GuildPlayer[] { - return [...this.players.values()]; - } - - getOrCreate(serverId: string): GuildPlayer { - const existing = this.players.get(serverId); - if (existing) return existing; - - const player = new GuildPlayer({ - serverId, - serverName: this.stoat?.getServerName(serverId) ?? null, - revoice: this.revoice, - }); - player.on("update", (snapshot) => this.emit("update", snapshot)); - player.on("position", (position) => this.emit("position", position)); - player.on("notice", (notice) => void this.deliverNotice(notice)); - this.players.set(serverId, player); - return player; - } - - private async deliverNotice(notice: Notice): Promise { - if (!notice.textChannelId || !this.stoat) return; - try { - await this.stoat.sendMessage(notice.textChannelId, notice.text); - } catch (err) { - log.warn({ err, channel: notice.textChannelId }, "failed to deliver notice"); - } - } - - async destroy(serverId: string): Promise { - const player = this.players.get(serverId); - if (!player) return; - this.players.delete(serverId); - await player.destroy(); - } - - async destroyAll(): Promise { - await Promise.allSettled([...this.players.keys()].map((id) => this.destroy(id))); - } - - // ------------------------------------------------------------ permissions --- - - async assertControl(serverId: string, userId: string): Promise { - if (!(await this.chat.canControl(serverId, userId))) { - throw new UserFacingError("Недостаточно прав для управления плеером"); - } - } - - // ---------------------------------------------------------------- actions --- - - /** Connects to the caller's voice channel (or an explicit one) and returns the player. */ - async connect( - serverId: string, - userId: string, - options: { voiceChannelId?: string | null; textChannelId?: string | null } = {}, - ): Promise { - const player = this.getOrCreate(serverId); - if (options.textChannelId) player.textChannelId = options.textChannelId; - - const target = options.voiceChannelId - ? this.chat.getVoiceChannel(options.voiceChannelId) - : (this.chat.findUserVoiceChannel(serverId, userId) ?? - (player.voiceChannelId ? this.chat.getVoiceChannel(player.voiceChannelId) : null)); - - if (!target) { - throw new UserFacingError("Зайдите в голосовой канал или укажите его явно"); - } - await player.connect(target.id, target.name); - return player; - } - - async play( - serverId: string, - requester: Requester, - query: string, - options: { mode?: PlayMode; voiceChannelId?: string | null; textChannelId?: string | null } = {}, - ): Promise { - await this.assertControl(serverId, requester.id); - const player = await this.connect(serverId, requester.id, { - voiceChannelId: options.voiceChannelId ?? null, - textChannelId: options.textChannelId ?? null, - }); - - const result = await resolveQuery(query, requester, config.MAX_QUEUE_SIZE - player.queue.length); - if (result.tracks.length === 0) throw new UserFacingError("Ничего не найдено"); - - const mode = options.mode ?? "append"; - const wasIdle = !player.current; - - if (mode === "now") { - await player.playNow(result.tracks); - return { ...result, startedNow: true, queuePosition: 0 }; - } - - player.enqueue(result.tracks, mode === "next" ? 0 : undefined); - const queuePosition = mode === "next" ? 1 : player.queue.length - result.tracks.length + 1; - await player.ensurePlaying(); - return { ...result, startedNow: wasIdle, queuePosition }; - } - - /** Queues already-resolved tracks (used by the panel's search results). */ - async enqueueTracks( - serverId: string, - requester: Requester, - tracks: Track[], - options: { mode?: PlayMode; voiceChannelId?: string | null; textChannelId?: string | null } = {}, - ): Promise { - await this.assertControl(serverId, requester.id); - const player = await this.connect(serverId, requester.id, { - voiceChannelId: options.voiceChannelId ?? null, - textChannelId: options.textChannelId ?? null, - }); - const owned = tracks.map((track) => ({ ...track, requestedBy: requester })); - const wasIdle = !player.current; - - if (options.mode === "now") { - await player.playNow(owned); - return { tracks: owned, playlist: null, startedNow: true, queuePosition: 0 }; - } - player.enqueue(owned, options.mode === "next" ? 0 : undefined); - await player.ensurePlaying(); - return { - tracks: owned, - playlist: null, - startedNow: wasIdle, - queuePosition: options.mode === "next" ? 1 : player.queue.length - owned.length + 1, - }; - } - - search(query: string, requester: Requester, limit?: number): Promise { - return searchTracks(query, requester, limit); - } - - private async require(serverId: string, userId: string): Promise { - await this.assertControl(serverId, userId); - const player = this.players.get(serverId); - if (!player) throw new UserFacingError("Плеер не запущен на этом сервере"); - return player; - } - - async pause(serverId: string, userId: string): Promise { - (await this.require(serverId, userId)).pause(); - } - - async resume(serverId: string, userId: string): Promise { - (await this.require(serverId, userId)).resume(); - } - - async togglePause(serverId: string, userId: string): Promise<"paused" | "playing"> { - const player = await this.require(serverId, userId); - if (player.snapshot().status === "paused") { - player.resume(); - return "playing"; - } - player.pause(); - return "paused"; - } - - async skip(serverId: string, userId: string, count = 1): Promise { - return (await this.require(serverId, userId)).skip(count); - } - - async stop(serverId: string, userId: string): Promise { - await (await this.require(serverId, userId)).stop(); - } - - async setVolume(serverId: string, userId: string, volume: number): Promise { - (await this.require(serverId, userId)).setVolume(volume); - } - - async setLoop(serverId: string, userId: string, mode: LoopMode): Promise { - (await this.require(serverId, userId)).setLoop(mode); - } - - async shuffle(serverId: string, userId: string): Promise { - (await this.require(serverId, userId)).shuffle(); - } - - async seek(serverId: string, userId: string, seconds: number): Promise { - await (await this.require(serverId, userId)).seek(seconds); - } - - async remove(serverId: string, userId: string, trackId: string): Promise { - return (await this.require(serverId, userId)).remove(trackId); - } - - async move(serverId: string, userId: string, trackId: string, toIndex: number): Promise { - (await this.require(serverId, userId)).move(trackId, toIndex); - } - - async clearQueue(serverId: string, userId: string): Promise { - (await this.require(serverId, userId)).clearQueue(); - } - - async leave(serverId: string, userId: string): Promise { - await (await this.require(serverId, userId)).leaveVoice(); - } - - snapshot(serverId: string): PlayerSnapshot { - const player = this.players.get(serverId); - if (player) return player.snapshot(); - return { - serverId, - serverName: this.stoat?.getServerName(serverId) ?? null, - voiceChannelId: null, - voiceChannelName: null, - textChannelId: null, - status: "idle", - current: null, - position: 0, - queue: [], - history: [], - volume: config.DEFAULT_VOLUME, - loop: "off", - shuffleUsed: false, - updatedAt: Date.now(), - }; - } -} +import { EventEmitter } from "node:events"; +import { config } from "../config.js"; +import { logger } from "../logger.js"; +import { resolveQuery, searchTracks } from "../sources/index.js"; +import { + UserFacingError, + type LoopMode, + type PlayerSnapshot, + type Requester, + type SearchResult, + type Track, +} from "../types.js"; +import { GuildPlayer, type Notice, type PositionUpdate } from "./player.js"; +import { createRevoice, type RevoiceLike } from "./revoice.js"; + +const log = logger.child({ mod: "manager" }); + +export interface VoiceChannelRef { + id: string; + name: string; +} + +export interface ServerRef { + id: string; + name: string; + iconUrl: string | null; +} + +/** + * Everything the core needs to know about the chat side of Stoat. Implemented on + * top of the bot's stoat.js client so the player itself stays testable. + */ +export interface StoatContext { + findUserVoiceChannel(serverId: string, userId: string): VoiceChannelRef | null; + getVoiceChannel(channelId: string): VoiceChannelRef | null; + listVoiceChannels(serverId: string): VoiceChannelRef[]; + getServerName(serverId: string): string | null; + listServersForUser(userId: string): Promise; + isMember(serverId: string, userId: string): Promise; + canControl(serverId: string, userId: string): Promise; + sendMessage(channelId: string, content: string): Promise; +} + +export type PlayMode = "append" | "next" | "now"; + +export interface ManagerEvents { + update: [PlayerSnapshot]; + position: [PositionUpdate]; +} + +export interface PlayOutcome extends SearchResult { + startedNow: boolean; + queuePosition: number; +} + +/** + * Owns one GuildPlayer per server and exposes the high-level operations that + * both the chat commands and the web panel call into. + */ +export class MusicManager extends EventEmitter { + private readonly players = new Map(); + private readonly revoice: RevoiceLike; + private stoat: StoatContext | null = null; + + constructor() { + super(); + this.revoice = createRevoice(config.STOAT_BOT_TOKEN, config.STOAT_API_URL, config.VOICE_NODE); + } + + attachStoat(context: StoatContext): void { + this.stoat = context; + } + + private get chat(): StoatContext { + if (!this.stoat) throw new UserFacingError("Бот ещё не подключился к Stoat"); + return this.stoat; + } + + // ---------------------------------------------------------------- players --- + + get(serverId: string): GuildPlayer | undefined { + return this.players.get(serverId); + } + + list(): GuildPlayer[] { + return [...this.players.values()]; + } + + getOrCreate(serverId: string): GuildPlayer { + const existing = this.players.get(serverId); + if (existing) return existing; + + const player = new GuildPlayer({ + serverId, + serverName: this.stoat?.getServerName(serverId) ?? null, + revoice: this.revoice, + }); + player.on("update", (snapshot) => this.emit("update", snapshot)); + player.on("position", (position) => this.emit("position", position)); + player.on("notice", (notice) => void this.deliverNotice(notice)); + this.players.set(serverId, player); + return player; + } + + private async deliverNotice(notice: Notice): Promise { + if (!notice.textChannelId || !this.stoat) return; + try { + await this.stoat.sendMessage(notice.textChannelId, notice.text); + } catch (err) { + log.warn({ err, channel: notice.textChannelId }, "failed to deliver notice"); + } + } + + async destroy(serverId: string): Promise { + const player = this.players.get(serverId); + if (!player) return; + this.players.delete(serverId); + await player.destroy(); + } + + async destroyAll(): Promise { + await Promise.allSettled([...this.players.keys()].map((id) => this.destroy(id))); + } + + // ------------------------------------------------------------ permissions --- + + async assertControl(serverId: string, userId: string): Promise { + if (!(await this.chat.canControl(serverId, userId))) { + throw new UserFacingError("Недостаточно прав для управления плеером"); + } + } + + // ---------------------------------------------------------------- actions --- + + /** Connects to the caller's voice channel (or an explicit one) and returns the player. */ + async connect( + serverId: string, + userId: string, + options: { voiceChannelId?: string | null; textChannelId?: string | null } = {}, + ): Promise { + const player = this.getOrCreate(serverId); + if (options.textChannelId) player.textChannelId = options.textChannelId; + + const target = options.voiceChannelId + ? this.chat.getVoiceChannel(options.voiceChannelId) + : (this.chat.findUserVoiceChannel(serverId, userId) ?? + (player.voiceChannelId ? this.chat.getVoiceChannel(player.voiceChannelId) : null)); + + if (!target) { + throw new UserFacingError("Зайдите в голосовой канал или укажите его явно"); + } + await player.connect(target.id, target.name); + return player; + } + + async play( + serverId: string, + requester: Requester, + query: string, + options: { mode?: PlayMode; voiceChannelId?: string | null; textChannelId?: string | null } = {}, + ): Promise { + await this.assertControl(serverId, requester.id); + const player = await this.connect(serverId, requester.id, { + voiceChannelId: options.voiceChannelId ?? null, + textChannelId: options.textChannelId ?? null, + }); + + const result = await resolveQuery(query, requester, config.MAX_QUEUE_SIZE - player.queue.length); + if (result.tracks.length === 0) throw new UserFacingError("Ничего не найдено"); + + const mode = options.mode ?? "append"; + const wasIdle = !player.current; + + if (mode === "now") { + await player.playNow(result.tracks); + return { ...result, startedNow: true, queuePosition: 0 }; + } + + player.enqueue(result.tracks, mode === "next" ? 0 : undefined); + const queuePosition = mode === "next" ? 1 : player.queue.length - result.tracks.length + 1; + await player.ensurePlaying(); + return { ...result, startedNow: wasIdle, queuePosition }; + } + + /** Queues already-resolved tracks (used by the panel's search results). */ + async enqueueTracks( + serverId: string, + requester: Requester, + tracks: Track[], + options: { mode?: PlayMode; voiceChannelId?: string | null; textChannelId?: string | null } = {}, + ): Promise { + await this.assertControl(serverId, requester.id); + const player = await this.connect(serverId, requester.id, { + voiceChannelId: options.voiceChannelId ?? null, + textChannelId: options.textChannelId ?? null, + }); + const owned = tracks.map((track) => ({ ...track, requestedBy: requester })); + const wasIdle = !player.current; + + if (options.mode === "now") { + await player.playNow(owned); + return { tracks: owned, playlist: null, startedNow: true, queuePosition: 0 }; + } + player.enqueue(owned, options.mode === "next" ? 0 : undefined); + await player.ensurePlaying(); + return { + tracks: owned, + playlist: null, + startedNow: wasIdle, + queuePosition: options.mode === "next" ? 1 : player.queue.length - owned.length + 1, + }; + } + + search(query: string, requester: Requester, limit?: number): Promise { + return searchTracks(query, requester, limit); + } + + private async require(serverId: string, userId: string): Promise { + await this.assertControl(serverId, userId); + const player = this.players.get(serverId); + if (!player) throw new UserFacingError("Плеер не запущен на этом сервере"); + return player; + } + + async pause(serverId: string, userId: string): Promise { + (await this.require(serverId, userId)).pause(); + } + + async resume(serverId: string, userId: string): Promise { + (await this.require(serverId, userId)).resume(); + } + + async togglePause(serverId: string, userId: string): Promise<"paused" | "playing"> { + const player = await this.require(serverId, userId); + if (player.snapshot().status === "paused") { + player.resume(); + return "playing"; + } + player.pause(); + return "paused"; + } + + async skip(serverId: string, userId: string, count = 1): Promise { + return (await this.require(serverId, userId)).skip(count); + } + + async stop(serverId: string, userId: string): Promise { + await (await this.require(serverId, userId)).stop(); + } + + async setVolume(serverId: string, userId: string, volume: number): Promise { + (await this.require(serverId, userId)).setVolume(volume); + } + + async setLoop(serverId: string, userId: string, mode: LoopMode): Promise { + (await this.require(serverId, userId)).setLoop(mode); + } + + async shuffle(serverId: string, userId: string): Promise { + (await this.require(serverId, userId)).shuffle(); + } + + async seek(serverId: string, userId: string, seconds: number): Promise { + await (await this.require(serverId, userId)).seek(seconds); + } + + async remove(serverId: string, userId: string, trackId: string): Promise { + return (await this.require(serverId, userId)).remove(trackId); + } + + async move(serverId: string, userId: string, trackId: string, toIndex: number): Promise { + (await this.require(serverId, userId)).move(trackId, toIndex); + } + + async clearQueue(serverId: string, userId: string): Promise { + (await this.require(serverId, userId)).clearQueue(); + } + + async leave(serverId: string, userId: string): Promise { + await (await this.require(serverId, userId)).leaveVoice(); + } + + snapshot(serverId: string): PlayerSnapshot { + const player = this.players.get(serverId); + if (player) return player.snapshot(); + return { + serverId, + serverName: this.stoat?.getServerName(serverId) ?? null, + voiceChannelId: null, + voiceChannelName: null, + textChannelId: null, + status: "idle", + current: null, + position: 0, + queue: [], + history: [], + volume: config.DEFAULT_VOLUME, + loop: "off", + shuffleUsed: false, + updatedAt: Date.now(), + }; + } +} diff --git a/src/core/player.ts b/src/core/player.ts index 81031e9..f881513 100644 --- a/src/core/player.ts +++ b/src/core/player.ts @@ -146,20 +146,29 @@ export class GuildPlayer extends EventEmitter { this.log.info({ channelId }, "joining voice channel"); const connection = await this.revoice.join(channelId); - // The room connects asynchronously inside revoice's constructor, so we wait - // for it to report readiness rather than polling a state getter. - await new Promise((resolve, reject) => { - const timer = setTimeout( - () => reject(new UserFacingError("Не удалось подключиться к голосовому каналу")), - JOIN_TIMEOUT_MS, - ); - const done = () => { - clearTimeout(timer); - resolve(); - }; - connection.once("join", done); - connection.once("roomfetched", done); - }); + try { + // The room connects asynchronously inside revoice's constructor, so we wait + // for it to report readiness rather than polling a state getter. + await new Promise((resolve, reject) => { + const timer = setTimeout( + () => reject(new UserFacingError("Не удалось подключиться к голосовому каналу")), + JOIN_TIMEOUT_MS, + ); + const done = () => { + clearTimeout(timer); + resolve(); + }; + connection.once("join", done); + connection.once("roomfetched", done); + }); + } catch (err) { + // Leave no orphan: an abandoned connection keeps the bot registered in the + // channel, and Stoat then refuses the next join with AlreadyConnected. + await connection.destroy().catch(() => {}); + connection.removeAllListeners(); + this.setStatus("idle"); + throw err; + } this.connection = connection; this.voiceReady = true; diff --git a/src/core/revoice.ts b/src/core/revoice.ts index a7dd18e..03d0c5e 100644 --- a/src/core/revoice.ts +++ b/src/core/revoice.ts @@ -1,67 +1,154 @@ -import { createRequire } from "node:module"; -import type { Readable } from "node:stream"; - -// revoice.js is CommonJS and its bundled typings lag behind the LiveKit rewrite, -// so we load it through require() and describe only the surface we rely on. -const require = createRequire(import.meta.url); - -export interface MediaPlayerLike { - readonly seconds: number; - readonly duration: number; - codecData?: { duration?: string } | null; - paused: boolean; - playing: boolean; - fProc?: { kill(signal?: string): void } | null; - originStream?: { destroy(): void } | null; - playStream(input: Readable | string, inputOptions?: string[]): Promise; - pause(): void; - resume(): void; - stop(init?: boolean): void; - destroy(): void; - setVolume(volume: number): void; - on(event: "start" | "startplay" | "buffer" | "pause" | "unpause" | "finish", listener: () => void): this; - removeAllListeners(event?: string): this; -} - -/** Voice connection state strings emitted by revoice.js (`Revoice.State`). */ -export const VOICE_STATE_OFFLINE = "off"; - -export interface VoiceConnectionLike { - channelId: string; - play(media: MediaPlayerLike): Promise; - leave(): Promise; - destroy(): Promise; - getUsers(): Array<{ id: string }>; - on(event: "join" | "leave" | "roomfetched" | "autoleave", listener: () => void): this; - on(event: "state", listener: (state: string) => void): this; - on(event: "userJoin" | "userleave" | "userLeave", listener: (user: { id: string }) => void): this; - once(event: "join" | "leave" | "roomfetched", listener: () => void): this; - removeAllListeners(event?: string): this; -} -// NB: revoice.js also exposes `connection.connected` / `isConnected()`, but both -// call `room.isConnected()` — a getter, not a method, in @livekit/rtc-node 0.13+, -// so touching them throws a TypeError. GuildPlayer tracks readiness from events. - -export interface RevoiceLike { - join(channelId: string, leaveIfEmpty?: boolean | number): Promise; - getVoiceConnection(channelId: string): VoiceConnectionLike | undefined; - connections: Map; -} - -interface RevoiceModule { - Revoice: new (token: string, apiConfig?: Record) => RevoiceLike; - MediaPlayer: new (normalisation?: boolean) => MediaPlayerLike; -} - -const revoice = require("revoice.js") as RevoiceModule; - -export const Revoice = revoice.Revoice; -export const MediaPlayer = revoice.MediaPlayer; - -/** Parses ffmpeg's `hh:mm:ss.xx` duration into seconds. */ -export function parseFfmpegDuration(value: string | undefined | null): number { - if (!value) return 0; - const parts = value.split(":").map((part) => Number.parseFloat(part)); - if (parts.some((part) => Number.isNaN(part))) return 0; - return parts.reduce((acc, part) => acc * 60 + part, 0); -} +import { createRequire } from "node:module"; +import type { Readable } from "node:stream"; +import { UserFacingError } from "../types.js"; + +// revoice.js is CommonJS and its bundled typings lag behind the LiveKit rewrite, +// so we load it through require() and describe only the surface we rely on. +const require = createRequire(import.meta.url); + +export interface MediaPlayerLike { + readonly seconds: number; + readonly duration: number; + codecData?: { duration?: string } | null; + paused: boolean; + playing: boolean; + fProc?: { kill(signal?: string): void } | null; + originStream?: { destroy(): void } | null; + playStream(input: Readable | string, inputOptions?: string[]): Promise; + pause(): void; + resume(): void; + stop(init?: boolean): void; + destroy(): void; + setVolume(volume: number): void; + on(event: "start" | "startplay" | "buffer" | "pause" | "unpause" | "finish", listener: () => void): this; + removeAllListeners(event?: string): this; +} + +/** Voice connection state strings emitted by revoice.js (`Revoice.State`). */ +export const VOICE_STATE_OFFLINE = "off"; + +export interface VoiceConnectionLike { + channelId: string; + play(media: MediaPlayerLike): Promise; + leave(): Promise; + destroy(): Promise; + getUsers(): Array<{ id: string }>; + on(event: "join" | "leave" | "roomfetched" | "autoleave", listener: () => void): this; + on(event: "state", listener: (state: string) => void): this; + on(event: "userJoin" | "userleave" | "userLeave", listener: (user: { id: string }) => void): this; + once(event: "join" | "leave" | "roomfetched", listener: () => void): this; + removeAllListeners(event?: string): this; +} +// NB: revoice.js also exposes `connection.connected` / `isConnected()`, but both +// call `room.isConnected()` — a getter, not a method, in @livekit/rtc-node 0.13+, +// so touching them throws a TypeError. GuildPlayer tracks readiness from events. + +export interface RevoiceLike { + join(channelId: string, leaveIfEmpty?: boolean | number): Promise; + getVoiceConnection(channelId: string): VoiceConnectionLike | undefined; + connections: Map; +} + +interface RevoiceModule { + Revoice: new (token: string, apiConfig?: Record) => RevoiceLike; + MediaPlayer: new (normalisation?: boolean) => MediaPlayerLike; +} + +const revoice = require("revoice.js") as RevoiceModule; + +/** How long we wait for Stoat to answer POST /join_call before giving up. */ +const JOIN_REQUEST_TIMEOUT_MS = 15_000; + +export const Revoice = revoice.Revoice; +export const MediaPlayer = revoice.MediaPlayer; + +interface RevoiceInternal extends RevoiceLike { + api: { post(path: string, body?: unknown, params?: unknown): Promise }; +} + +/** Stoat error codes we can explain better than "internal error". */ +const API_ERRORS: Record = { + AlreadyConnected: + "Stoat считает, что бот уже в этом голосовом канале. Обычно это зависшее состояние после падения — см. README, раздел про AlreadyConnected.", + NotAVoiceChannel: "Это не голосовой канал", + LiveKitUnavailable: "Голосовой сервер (LiveKit) недоступен", + UnknownNode: "LiveKit-нода не найдена в конфигурации инстанса (Revolt.toml, [hosts.livekit])", + CannotJoinCall: "В канале достигнут лимит участников", + MissingPermission: "У бота нет права Connect в этом голосовом канале", + NotFound: "Канал не найден", +}; + +function describeApiError(err: unknown): Error { + const response = (err as { response?: { status?: number; data?: { type?: string } } }).response; + const type = response?.data?.type; + if (!type) return err as Error; + const message = API_ERRORS[type] ?? `Stoat отклонил запрос: ${type}`; + return new UserFacingError(message); +} + +/** + * Wraps the revoice client so that join_call uses the configured LiveKit node + * and API failures surface Stoat's own error code instead of an axios dump. + */ +export function createRevoice(token: string, baseURL: string, node: string): RevoiceLike { + const instance = new Revoice(token, { baseURL }) as RevoiceInternal; + const post = instance.api.post.bind(instance.api); + const join = instance.join.bind(instance); + let pendingError: Error | null = null; + + instance.api.post = async (path: string, body?: unknown, params?: unknown) => { + const payload = + path.endsWith("/join_call") && typeof body === "object" && body !== null + ? { ...(body as Record), node } + : body; + try { + return await post(path, payload, params); + } catch (err) { + pendingError = describeApiError(err); + throw pendingError; + } + }; + + // revoice's join() runs an async executor inside `new Promise`, so a failing + // join_call never reaches its reject() — the promise hangs forever and the + // real error escapes as an unhandled rejection. We latch that error above and + // settle the join ourselves. + instance.join = (channelId: string, leaveIfEmpty?: boolean | number) => + new Promise((resolve, reject) => { + pendingError = null; + const startedAt = Date.now(); + let settled = false; + + const settle = (action: () => void) => { + if (settled) return; + settled = true; + clearInterval(watchdog); + action(); + }; + + const watchdog = setInterval(() => { + if (pendingError) { + const error = pendingError; + settle(() => reject(error)); + } else if (Date.now() - startedAt > JOIN_REQUEST_TIMEOUT_MS) { + settle(() => reject(new UserFacingError("Stoat не ответил на запрос подключения к каналу"))); + } + }, 50); + watchdog.unref?.(); + + join(channelId, leaveIfEmpty).then( + (connection) => settle(() => resolve(connection)), + (err: unknown) => settle(() => reject(describeApiError(err))), + ); + }); + + return instance; +} + +/** Parses ffmpeg's `hh:mm:ss.xx` duration into seconds. */ +export function parseFfmpegDuration(value: string | undefined | null): number { + if (!value) return 0; + const parts = value.split(":").map((part) => Number.parseFloat(part)); + if (parts.some((part) => Number.isNaN(part))) return 0; + return parts.reduce((acc, part) => acc * 60 + part, 0); +}