const IDLE_MIN_MS = 2000; const IDLE_MAX_MS = 60000; const REQUEST_TIMEOUT_MS = parseInt(process.env.AI_REQUEST_TIMEOUT_MS || '120000', 10); const MAX_INPUT_CHARS = 2000; const MIN_TEXT_CHARS = 4; const path = require('path'); const fs = require('fs'); const crypto = require('crypto'); const PHOTO_MAX_ATTEMPTS = 3; const PHOTO_AI_TIMEOUT_MS = 300000; function createPhotoEnhanceWorker({ pool, getSetting, logAudit, invalidateEntries, sharp, photoAiUrl, uploadsDir, originalsDir }) { const AI_URL = photoAiUrl || process.env.PHOTO_AI_URL || ''; const IDLE_MIN = 2000; const IDLE_MAX = 30000; let started = false; let stopped = false; let idleMs = IDLE_MIN; let wake = null; let processing = false; let currentId = null; let startedAt = null; const stats = { jobs: 0, done: 0, errors: 0, last_at: null, last_error: null }; const CONFIG = { idle_min_ms: IDLE_MIN, idle_max_ms: IDLE_MAX, ai_timeout_ms: PHOTO_AI_TIMEOUT_MS, max_attempts: PHOTO_MAX_ATTEMPTS, ai_url: AI_URL, }; function sleep(ms) { return new Promise((resolve) => { const t = setTimeout(() => { wake = null; resolve(); }, ms); wake = () => { clearTimeout(t); wake = null; resolve(); }; }); } function notify() { if (wake) wake(); } async function isEnabled() { const v = await getSetting('photo_worker_enabled', 'true'); return String(v) !== 'false'; } async function claimNext() { const client = await pool.connect(); try { await client.query('BEGIN'); const { rows } = await client.query( `SELECT j.id, j.entry_id, j.action, j.params FROM photo_jobs j WHERE j.status = 'pending' ORDER BY j.id ASC LIMIT 1 FOR UPDATE SKIP LOCKED` ); if (!rows.length) { await client.query('COMMIT'); return null; } await client.query(`UPDATE photo_jobs SET status = 'processing' WHERE id = $1`, [rows[0].id]); await client.query('COMMIT'); return rows[0]; } catch (e) { await client.query('ROLLBACK').catch(() => {}); throw e; } finally { client.release(); } } function safeUnlinkPath(p) { try { if (!p) return; const abs = path.resolve(p); const root = path.resolve(uploadsDir); if (abs !== root && abs.startsWith(root + path.sep)) fs.unlinkSync(abs); } catch (e) { console.error('photo worker unlink:', e.message); } } function backupOldFile(oldRelPath) { if (!oldRelPath) return null; const oldAbs = path.join(uploadsDir, String(oldRelPath).replace(/^\/+/, '').replace(/^uploads\//, '')); if (!fs.existsSync(oldAbs)) return null; const backupName = crypto.randomBytes(12).toString('hex') + (path.extname(oldAbs) || '.jpg'); const backupPath = path.join(originalsDir, backupName); try { fs.renameSync(oldAbs, backupPath); return `/uploads/.originals/${backupName}`; } catch (e) { console.error('photo worker backup failed:', e.message); return null; } } async function enhanceWithSharp(srcAbs, params) { if (!sharp) throw new Error('sharp недоступен на сервере'); const p = params || {}; const clamp = (v, min, max, def) => { const n = parseFloat(v); return Number.isFinite(n) ? Math.min(max, Math.max(min, n)) : def; }; const brightness = clamp(p.brightness, 10, 300, 100); const contrast = clamp(p.contrast, 10, 300, 100); const saturate = clamp(p.saturate, 0, 300, 100); const sharpAmt = clamp(p.sharp, 0, 100, 0); const denoise = clamp(p.denoise, 0, 100, 0); let pipeline = sharp(srcAbs).rotate(); if (denoise > 0) pipeline = pipeline.median(denoise > 70 ? 5 : 3); pipeline = pipeline.modulate({ brightness: brightness / 100, saturation: saturate / 100 }); if (contrast !== 100) { const a = contrast / 100; pipeline = pipeline.linear(a, 128 * (1 - a)); } if (sharpAmt > 0) { pipeline = pipeline.sharpen({ sigma: 0.5 + (sharpAmt / 100) * 1.5, m1: 0, m2: 1 + sharpAmt / 50 }); } const newName = crypto.randomBytes(12).toString('hex') + '.jpg'; await pipeline.jpeg({ quality: 92, mozjpeg: true }).toFile(path.join(uploadsDir, newName)); return `/uploads/${newName}`; } async function runAiEnhance(srcAbs) { if (!AI_URL) throw new Error('PHOTO_AI_URL не настроен'); const buf = fs.readFileSync(srcAbs); const fd = new FormData(); fd.append('image', new Blob([buf], { type: 'image/jpeg' }), 'photo.jpg'); fd.append('scale', '2'); const resp = await fetch(AI_URL.replace(/\/+$/, '') + '/enhance', { method: 'POST', body: fd, signal: AbortSignal.timeout(PHOTO_AI_TIMEOUT_MS), }); if (!resp.ok) throw new Error('AI service error: ' + resp.status); const out = Buffer.from(await resp.arrayBuffer()); const newName = crypto.randomBytes(12).toString('hex') + '.jpg'; fs.writeFileSync(path.join(uploadsDir, newName), out); return `/uploads/${newName}`; } async function applyResult(job, newPath) { await pool.query( `UPDATE photo_jobs SET status = 'done', after_path = $1, error = NULL, finished_at = now() WHERE id = $2`, [newPath, job.id] ); if (logAudit) await logAudit(null, 'photo.job.preview', { entry_id: job.entry_id, job_id: job.id, action: job.action, after_path: newPath }); } async function processOne(job) { stats.jobs++; const { rows } = await pool.query('SELECT photo_path FROM entries WHERE id = $1', [job.entry_id]); if (!rows.length || !rows[0].photo_path) { await pool.query( `UPDATE photo_jobs SET status = 'error', error = 'У записи нет фото', finished_at = now() WHERE id = $1`, [job.id] ); stats.errors++; stats.last_at = new Date().toISOString(); stats.last_error = 'У записи нет фото'; return true; } const photoPath = rows[0].photo_path; const srcAbs = path.join(uploadsDir, String(photoPath).replace(/^\/+/, '').replace(/^uploads\//, '')); try { const newPath = job.action === 'ai' ? await runAiEnhance(srcAbs) : await enhanceWithSharp(srcAbs, job.params); await applyResult(job, newPath); stats.done++; stats.last_at = new Date().toISOString(); stats.last_error = null; return true; } catch (e) { const message = (e && e.message ? e.message : 'error').slice(0, 500); stats.last_at = new Date().toISOString(); stats.last_error = message; const { rows: cur } = await pool.query('SELECT attempts FROM photo_jobs WHERE id = $1', [job.id]); const tries = ((cur[0] && cur[0].attempts) || 0) + 1; if (tries < PHOTO_MAX_ATTEMPTS) { await pool.query(`UPDATE photo_jobs SET status = 'pending', attempts = $1, error = $2 WHERE id = $3`, [tries, message, job.id]); if (logAudit) await logAudit(null, 'photo.job.retry', { entry_id: job.entry_id, job_id: job.id, attempt: tries, error: message }); return false; } await pool.query( `UPDATE photo_jobs SET status = 'error', attempts = $1, error = $2, finished_at = now() WHERE id = $3`, [tries, message, job.id] ); stats.errors++; if (logAudit) await logAudit(null, 'photo.job.error', { entry_id: job.entry_id, job_id: job.id, error: message }); return true; } } async function loop() { while (!stopped) { let job = null; try { if (!(await isEnabled())) { await sleep(IDLE_MAX); continue; } job = await claimNext(); } catch (e) { console.error('Photo worker claim error:', e); await sleep(IDLE_MAX); continue; } if (!job) { await sleep(idleMs); idleMs = Math.min(idleMs * 2, IDLE_MAX); continue; } idleMs = IDLE_MIN; processing = true; currentId = job.id; let ok = true; try { ok = await processOne(job); } catch (e) { console.error('Photo worker process error:', e); } finally { processing = false; currentId = null; } if (!ok) await sleep(IDLE_MAX); } } async function resetStale() { await pool.query(`UPDATE photo_jobs SET status = 'pending' WHERE status = 'processing'`); } function start() { if (started) return; started = true; startedAt = Date.now(); resetStale() .catch((e) => console.error('Photo worker reset error:', e)) .finally(() => { loop().catch((e) => console.error('Photo worker loop error:', e)); }); } function getInfo() { return { started_at: startedAt ? new Date(startedAt).toISOString() : null, uptime_ms: startedAt ? Date.now() - startedAt : 0, processing, current_id: currentId, config: CONFIG, ...stats, }; } return { start, notify, getInfo }; } function createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl, defaultPrompt, model }) { const AI_URL = aiUrl || process.env.AI_URL || 'http://text-corrector:8080'; const MODEL = model || process.env.AI_MODEL || 'qwen2.5-1.5b-instruct-q4_k_m.gguf'; const DEFAULT_PROMPT = defaultPrompt || process.env.AI_PROMPT || 'Ты — редактор текстов. Исправь ТОЛЬКО грамматические, орфографические и пунктуационные ошибки в тексте. Приведи к правильному регистру буквы. НЕ меняй слова, структуру предложений, стиль или смысл текста. Верни ТОЛЬКО исправленный текст без пояснений.'; let started = false; let stopped = false; let idleMs = IDLE_MIN_MS; let wake = null; let processing = false; let currentId = null; let startedAt = null; const stats = { checks: 0, corrected: 0, unchanged: 0, errors: 0, last_at: null, last_error: null }; const attempts = new Map(); const MAX_ATTEMPTS = 3; const CONFIG = { idle_min_ms: IDLE_MIN_MS, idle_max_ms: IDLE_MAX_MS, request_timeout_ms: REQUEST_TIMEOUT_MS, max_input_chars: MAX_INPUT_CHARS, min_text_chars: MIN_TEXT_CHARS, max_attempts: MAX_ATTEMPTS, model: MODEL, ai_url: AI_URL, }; function sleep(ms) { return new Promise((resolve) => { const t = setTimeout(() => { wake = null; resolve(); }, ms); wake = () => { clearTimeout(t); wake = null; resolve(); }; }); } function notify() { if (wake) wake(); } async function isEnabled() { const v = await getSetting('ai_autocheck_enabled', 'true'); return String(v) !== 'false'; } async function claimNext() { const client = await pool.connect(); try { await client.query('BEGIN'); const { rows } = await client.query( `SELECT id, description, description_original FROM entries WHERE ai_status = 'pending' AND deleted_at IS NULL ORDER BY id ASC LIMIT 1 FOR UPDATE SKIP LOCKED` ); if (!rows.length) { await client.query('COMMIT'); return null; } await client.query(`UPDATE entries SET ai_status = 'processing' WHERE id = $1`, [rows[0].id]); await client.query('COMMIT'); return rows[0]; } catch (e) { await client.query('ROLLBACK').catch(() => {}); throw e; } finally { client.release(); } } async function resolveActiveProfile() { try { const active = await getSetting('ai_active_profile', 'native'); if (active && active !== 'native') { const raw = await getSetting('ai_profiles', '[]'); const list = JSON.parse(raw || '[]'); if (Array.isArray(list)) { return list.find(p => p && p.id === active && p.base_url && p.model) || null; } } } catch (e) {} return null; } function normalizeOpenAiBase(base) { let b = String(base || '').trim().replace(/\/+$/, ''); if (!/^https?:\/\//i.test(b)) return null; if (!/\/v1$/i.test(b)) b += '/v1'; return b; } async function callModel(text, systemPrompt) { const profile = await resolveActiveProfile(); const headers = { 'Content-Type': 'application/json' }; let url; let model; let maxTokens; if (profile) { const base = normalizeOpenAiBase(profile.base_url); if (!base) throw new Error('Некорректный base_url профиля ИИ'); url = `${base}/chat/completions`; model = profile.model; if (profile.api_key) headers.Authorization = `Bearer ${profile.api_key}`; maxTokens = parseInt(profile.max_tokens, 10) || Math.min(4096, Math.max(1024, text.length * 2 + 512)); } else { url = `${AI_URL.replace(/\/+$/, '')}/v1/chat/completions`; model = MODEL; maxTokens = Math.min(256, Math.max(128, text.length + 64)); } const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), REQUEST_TIMEOUT_MS); try { const res = await fetch(url, { method: 'POST', headers, signal: controller.signal, body: JSON.stringify({ model, messages: [ { role: 'system', content: systemPrompt }, { role: 'user', content: text }, ], temperature: 0.1, max_tokens: maxTokens, }), }); if (!res.ok) throw new Error(`AI service error: ${res.status}`); const data = await res.json(); const choice = data.choices && data.choices[0]; const content = (choice && choice.message && choice.message.content || '').trim(); if (!content) { throw new Error(`Модель вернула пустой ответ (finish_reason=${choice && choice.finish_reason || 'unknown'})`); } return content; } finally { clearTimeout(timer); } } function sanitize(original, out) { if (!out) return original; const o = out.trim(); if (!o) return original; if (o.length > original.length * 4 + 80) return original; return o; } async function processOne(row) { const original = String(row.description_original ?? row.description ?? '').trim(); if (original.length < MIN_TEXT_CHARS) { await pool.query(`UPDATE entries SET ai_status = 'skipped', ai_checked_at = now() WHERE id = $1`, [row.id]); return true; } const prompt = (await getSetting('ai_prompt', DEFAULT_PROMPT)) || DEFAULT_PROMPT; const input = original.length > MAX_INPUT_CHARS ? original.slice(0, MAX_INPUT_CHARS) : original; try { const out = await callModel(input, prompt); const checked = sanitize(input, out); const changed = checked !== input; await pool.query( `UPDATE entries SET description_ai = $1, description = $1, ai_status = 'done', ai_error = NULL, ai_checked_at = now() WHERE id = $2`, [changed ? checked : original, row.id] ); attempts.delete(row.id); stats.checks++; if (changed) stats.corrected++; else stats.unchanged++; stats.last_at = new Date().toISOString(); stats.last_error = null; if (logAudit) await logAudit(null, 'entry.ai.auto-check', { id: row.id, changed }); return true; } catch (e) { const message = (e && e.message ? e.message : 'error').slice(0, 500); const tries = (attempts.get(row.id) || 0) + 1; stats.last_at = new Date().toISOString(); stats.last_error = message; if (tries < MAX_ATTEMPTS) { attempts.set(row.id, tries); await pool.query(`UPDATE entries SET ai_status = 'pending', ai_error = $1 WHERE id = $2`, [message, row.id]); if (logAudit) await logAudit(null, 'entry.ai.retry', { id: row.id, attempt: tries, error: message }); return false; } attempts.delete(row.id); stats.errors++; await pool.query( `UPDATE entries SET ai_status = 'error', ai_error = $1, ai_checked_at = now() WHERE id = $2`, [message, row.id] ); if (logAudit) await logAudit(null, 'entry.ai.error', { id: row.id, error: message }); return true; } } async function loop() { while (!stopped) { let row = null; try { if (!(await isEnabled())) { await sleep(IDLE_MAX_MS); continue; } row = await claimNext(); } catch (e) { console.error('AI auto-check claim error:', e); await sleep(IDLE_MAX_MS); continue; } if (!row) { await sleep(idleMs); idleMs = Math.min(idleMs * 2, IDLE_MAX_MS); continue; } idleMs = IDLE_MIN_MS; processing = true; currentId = row.id; let ok = true; try { ok = await processOne(row); } catch (e) { console.error('AI auto-check process error:', e); } finally { processing = false; currentId = null; } if (!ok) await sleep(IDLE_MAX_MS); } } async function resetStale() { await pool.query(`UPDATE entries SET ai_status = 'pending' WHERE ai_status = 'processing'`); } function start() { if (started) return; started = true; startedAt = Date.now(); resetStale() .catch((e) => console.error('AI auto-check reset error:', e)) .finally(() => { loop().catch((e) => console.error('AI auto-check loop error:', e)); }); } function getStats() { return { ...stats, processing }; } function getInfo() { return { started_at: startedAt ? new Date(startedAt).toISOString() : null, uptime_ms: startedAt ? Date.now() - startedAt : 0, processing, current_id: currentId, config: CONFIG, ...stats, }; } return { start, notify, getStats, getInfo }; } module.exports = { createEntryAutoChecker, createPhotoEnhanceWorker };