import { spawn } from "bun"; import EventEmitter from "eventemitter3"; import { drainStream } from "./streams"; import type { Config, Download, Status } from "./types"; import { rid, sleep, appendLog } from "./utils"; // Imported appendLog import { existsSync, renameSync, readdirSync, mkdirSync } from "node:fs"; import { join, resolve } from "node:path"; // ========================================== // RETRY CONFIGURATION // ========================================== const RETRY_INTERVAL_NORMAL = 60_000; const RETRY_INTERVAL_BLOCKED = 300_000; // ========================================== export class DownloadManager extends EventEmitter { private downloads = new Map(); private activeSessions = new Set(); private chatProcesses = new Map(); private autoRestartInterval: NodeJS.Timeout | null = null; startAutoRestart(config: Config): void { if (this.autoRestartInterval) return; this.autoRestartInterval = setInterval(() => { this.checkLongRunningDownloads(config); }, 60_000); } stopAutoRestart(): void { if (this.autoRestartInterval) { clearInterval(this.autoRestartInterval); this.autoRestartInterval = null; } } private async checkLongRunningDownloads(config: Config): Promise { const now = Date.now(); const maxDuration = 120 * 60_000; for (const download of this.downloads.values()) { if (download.status === "downloading") { if (now - download.startTime.getTime() > maxDuration) { await this.killProcess(download.process); } } } } private getPythonCommand(): string[] { const cwd = process.cwd(); const winPath = join(cwd, ".venv", "Scripts", "python.exe"); if (process.platform === "win32" && existsSync(winPath)) return [winPath]; return ["python"]; } // Default delay 10000ms async start(rawUser: string, config: Config, initialDelay: number = 10000): Promise { const user = rawUser.replace(/^@/, "").trim(); for (const d of this.downloads.values()) { if (d.user === user && this.activeSessions.has(d.id)) { return d.id; } } const id = rid(); this.activeSessions.add(id); this._spawnCaptureSession(id, user, config, initialDelay); return id; } private async _spawnCaptureSession(id: string, user: string, config: Config, delay: number = 0) { if (!this.activeSessions.has(id)) return; if (delay > 0) { const msg = `[*] ⏳ Waiting ${Math.round(delay / 1000)}s before starting @${user}...`; console.log(msg); this.emit("log", msg); await sleep(delay); } if (!this.activeSessions.has(id)) return; const absOutputPath = resolve(config.outputPath); const tempRecordingDir = join(absOutputPath, ".recording"); if (!existsSync(tempRecordingDir)) { try { mkdirSync(tempRecordingDir, { recursive: true }); } catch {} } if (config.enableChat) { const chatDir = join(absOutputPath, "chats"); try { await Bun.write(join(chatDir, ".keep"), ""); } catch {} } const proc = this.spawnProcess(user, config, tempRecordingDir); const download: Download = { id, user, status: "waiting", process: proc, startTime: new Date(), outputPath: absOutputPath }; this.downloads.set(id, download); this.emitStatus(id, user, "waiting"); this.setupHandlers(proc, id, user, config, absOutputPath, tempRecordingDir); setTimeout(() => { const currentDl = this.downloads.get(id); if (currentDl && !currentDl.process?.exitCode && config.enableChat) { this.startChatCapture(user, config, absOutputPath); } }, 5000); } private startChatCapture(user: string, config: Config, absOutputPath: string): void { const chatDir = join(absOutputPath, "chats"); const chatFormat = config.chatFormat || "json"; const pythonCmd = this.getPythonCommand(); const scriptPath = join(import.meta.dir, "chat_capture.py"); const chatProc = spawn({ cmd: [...pythonCmd, scriptPath, user, chatDir, chatFormat], stdout: "pipe", stderr: "pipe", }); this.chatProcesses.set(user, chatProc.pid); if (chatProc.stdout) { const reader = chatProc.stdout.getReader(); const readStdout = async () => { try { while (true) { const { done, value } = await reader.read(); if (done) break; const text = new TextDecoder().decode(value); if (text.trim()) { this.emit("log", text.trim()); } } } catch {} }; readStdout(); } chatProc.exited.then(() => this.chatProcesses.delete(user)); } async stop(id: string): Promise { this.activeSessions.delete(id); const dl = this.downloads.get(id); if (!dl) return false; const chatPid = this.chatProcesses.get(dl.user); if (chatPid) { try { if (process.platform === "win32") spawn({ cmd: ["taskkill", "/F", "/T", "/PID", chatPid.toString()], stdout: "ignore", stderr: "ignore" }); else process.kill(chatPid); } catch {} this.chatProcesses.delete(dl.user); } await this.killProcess(dl.process); dl.status = "stopped"; this.emitStatus(id, dl.user, "stopped"); return true; } delete(id: string): boolean { this.stop(id); this.downloads.delete(id); return true; } async restart(id: string, config: Config): Promise { await this.stop(id); const dl = this.downloads.get(id); if (!dl) return null; return await this.start(dl.user, config, 0); } getAll(): Download[] { return Array.from(this.downloads.values()); } getRunning(): Download[] { return this.getAll().filter((d) => d.status === "downloading" || d.status === "waiting"); } async stopAll(): Promise { this.activeSessions.clear(); const promises: Promise[] = []; const allDownloads = Array.from(this.downloads.values()); for (const d of allDownloads) { if (d.status === "downloading" || d.status === "waiting") { promises.push(this.stop(d.id)); } } await Promise.all(promises); } async killGlobalProcesses(): Promise { this.activeSessions.clear(); const msg = "[*] 🛑 Global Kill initiated..."; console.log(msg); appendLog("system_logs.txt", msg); try { if (process.platform === "win32") { spawn({ cmd: ["taskkill", "/F", "/IM", "python.exe", "/T"], stdout: "ignore", stderr: "ignore" }); spawn({ cmd: ["taskkill", "/F", "/IM", "uv.exe", "/T"], stdout: "ignore", stderr: "ignore" }); spawn({ cmd: ["taskkill", "/F", "/IM", "ffmpeg.exe", "/T"], stdout: "ignore", stderr: "ignore" }); } else { spawn({ cmd: ["pkill", "-9", "python"], stdout: "ignore", stderr: "ignore" }); spawn({ cmd: ["pkill", "-9", "uv"], stdout: "ignore", stderr: "ignore" }); } } catch {} await sleep(1000); } private spawnProcess(user: string, config: Config, outputDir: string): ReturnType { const commandParts = config.commandPrefix.split(" "); const executable = commandParts[0]; const args = commandParts.slice(1); const finalArgs = [...args, "-output", outputDir, "-user", user]; return spawn({ cmd: [executable, ...finalArgs], stdout: "pipe", stderr: "pipe" }); } private async killProcess(proc: ReturnType | undefined): Promise { if (!proc) return; try { if (process.platform === "win32") spawn({ cmd: ["taskkill", "/F", "/T", "/PID", proc.pid.toString()], stdout: "ignore", stderr: "ignore" }); else proc.kill(); await sleep(600); } catch {} } private setupHandlers(proc: ReturnType, id: string, user: string, config: Config, finalOutputPath: string, tempDir: string): void { const dl = this.downloads.get(id); let isBlockedOrError = false; if (proc.stderr) { const reader = proc.stderr.getReader(); const readStderr = async () => { try { while (true) { const { done, value } = await reader.read(); if (done) break; const text = new TextDecoder().decode(value); if (text.trim()) { const trimmed = text.trim(); this.emit("log", trimmed); // LOGGING TO FILE (If not chat/gift) if (!trimmed.includes("💬") && !trimmed.includes("🎁") && !trimmed.includes("⚔️")) { appendLog("system_logs.txt", `[@${user}] ${trimmed}`); } if (text.includes("DEVICE_BLOCKED") || text.includes("blocked") || text.includes("Connection error") || text.includes("Max retries reached") || text.includes("disconnected") || text.includes("Socket closed")) { isBlockedOrError = true; } } if (text.includes("Opening stream") || text.includes("Recording") || text.includes("Connected") || text.includes("Started")) { if (dl && dl.status !== "downloading") { dl.status = "downloading"; this.emitStatus(id, user, "downloading"); } } } } catch {} }; readStderr(); } if (proc.stdout) void drainStream(proc.stdout).catch(() => {}); proc.exited.then((code) => { try { if (existsSync(tempDir)) { const files = readdirSync(tempDir); for (const file of files) { if (file.toLowerCase().includes(user.toLowerCase())) { const sourcePath = join(tempDir, file); const destPath = join(finalOutputPath, file); // Green message trigger const msg = `[*] 📦 Moving recording: ${file} -> Lives/`; console.log(msg); this.emit("log", msg); appendLog("system_logs.txt", msg); renameSync(sourcePath, destPath); } } } } catch (e) { console.error(`❌ Failed to move recording: ${e}`); } const chatPid = this.chatProcesses.get(user); if (chatPid) { try { if (process.platform === "win32") spawn({ cmd: ["taskkill", "/F", "/T", "/PID", chatPid.toString()], stdout: "ignore", stderr: "ignore" }); else process.kill(chatPid); } catch {} this.chatProcesses.delete(user); } if (!this.activeSessions.has(id)) { if (dl) dl.status = "stopped"; this.emitStatus(id, user, "stopped"); return; } const dlNow = this.downloads.get(id); if (dlNow) { dlNow.status = "waiting"; this.emitStatus(id, user, "waiting"); } let retryDelay = RETRY_INTERVAL_NORMAL; let logMsg = ""; if (isBlockedOrError) { retryDelay = RETRY_INTERVAL_BLOCKED; logMsg = `[*] ⚠️ Monitor: @${user} - BLOCK/ERROR/DISCONNECT Detected! Cooling down for ${retryDelay/1000}s before Retry...`; } else { logMsg = `[*] ♻️ Monitor: @${user} stream ended/offline. Infinite Retry active in ${retryDelay/1000}s...`; } console.log(logMsg); this.emit("log", logMsg); appendLog("system_logs.txt", logMsg); setTimeout(() => { if (this.activeSessions.has(id)) { this._spawnCaptureSession(id, user, config, 0); } }, retryDelay); }).catch(() => {}); } private emitStatus(id: string, user: string, status: Status): void { this.emit("statusChange", { id, user, status }); } }