import os from "node:os"; import { execFile } from "node:child_process"; /** * Jedyna brama do lokalnych procesów Pythona (segmentacja, recolor, studio). * * Każdy proces ładuje modele (torch + BiRefNet/DINO/SAM ≈ 1,5–2,5 GB RAM), * więc równoległe wywołania z workera i z tras HTTP potrafiły zapchać * 8 GB RAM do 99% i zamrozić maszynę. Zasady: * - co najwyżej VILMAL_PY_CONCURRENCY procesów naraz (domyślnie 1), reszta czeka w kolejce FIFO; * - wątki BLAS/torch ograniczone (VILMAL_PY_THREADS, domyślnie 2 z nproc); * - proces dostaje niższy priorytet systemowy, żeby UI i serwer HTTP zostały responsywne. */ const limit = Math.max(1, Number(process.env.VILMAL_PY_CONCURRENCY ?? 1)); const threads = String(Math.max(1, Number(process.env.VILMAL_PY_THREADS ?? Math.min(2, os.cpus().length)))); let active = 0; const waiting: Array<() => void> = []; async function acquire(): Promise { if (active < limit) { active++; return; } await new Promise((resolve) => waiting.push(resolve)); } function release(): void { const next = waiting.shift(); if (next) next(); else active--; } /** Stan kolejki — do statusu zadań w UI. */ export function pythonQueueState(): { active: number; waiting: number; limit: number } { return { active, waiting: waiting.length, limit }; } export interface PyResult { stdout: string; stderr: string; } export class PyError extends Error { constructor(message: string, readonly stdout: string, readonly stderr: string) { super(message); } } export async function runPython( args: string[], opts: { cwd?: string; env?: NodeJS.ProcessEnv; timeoutMs: number; maxBuffer?: number } ): Promise { await acquire(); try { const py = process.env.VILMAL_PYTHON ?? "python"; return await new Promise((resolve, reject) => { const child = execFile( py, args, { cwd: opts.cwd, env: { ...(opts.env ?? process.env), OMP_NUM_THREADS: threads, MKL_NUM_THREADS: threads, OPENBLAS_NUM_THREADS: threads, VILMAL_PY_THREADS: threads, PYTHONIOENCODING: "utf-8", }, timeout: opts.timeoutMs, maxBuffer: opts.maxBuffer ?? 16 * 1024 * 1024, windowsHide: true, }, (err, stdout, stderr) => { if (err) reject(new PyError(err.message, String(stdout ?? ""), String(stderr ?? ""))); else resolve({ stdout: String(stdout), stderr: String(stderr) }); } ); if (child.pid) { try { os.setPriority(child.pid, os.constants.priority.PRIORITY_BELOW_NORMAL); } catch { /* brak uprawnień do zmiany priorytetu — proces działa z domyślnym */ } } }); } finally { release(); } }