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 crypto = require('crypto'); const { textDiff, FIELD_LABELS } = require('./diff'); const PHOTO_MAX_ATTEMPTS = 3; const PHOTO_AI_TIMEOUT_MS = Math.max(1000, parseInt(process.env.PHOTO_AI_TIMEOUT_MS || '300000', 10) || 300000); const PHOTO_AI_FACE_TIMEOUT_MS = Math.max(PHOTO_AI_TIMEOUT_MS, parseInt(process.env.PHOTO_AI_FACE_TIMEOUT_MS || '600000', 10) || 600000); const PHOTO_AI_DEFAULT_MODEL = 'x2plus'; const PHOTO_AI_DEFAULT_FACE_MODEL = String(process.env.PHOTO_AI_FACE_MODEL || 'gfpgan').trim() || 'gfpgan'; const PHOTO_AI_DEFAULT_STRENGTH = 0.7; const PHOTO_AI_ACTIONS = new Set(['ai', 'ai_face', 'ai_upscale']); const PHOTO_SOFT_MAX_RETRIES = Math.max(1, parseInt(process.env.PHOTO_AI_SOFT_MAX_RETRIES || '60', 10) || 60); const PHOTO_SOFT_BACKOFF_MS = Math.max(1000, parseInt(process.env.PHOTO_AI_SOFT_BACKOFF_MS || '10000', 10) || 10000); const PHOTO_SOFT_BACKOFF_MAX_MS = Math.max(PHOTO_SOFT_BACKOFF_MS, parseInt(process.env.PHOTO_AI_SOFT_BACKOFF_MAX_MS || '300000', 10) || 300000); function isTimeoutFailure(e) { const name = e && e.name; return name === 'TimeoutError' || name === 'AbortError'; } function isSoftStatus(status) { return status === 408 || status === 425 || status === 429 || status >= 500; } function failureReason(e) { const parts = []; let cur = e; for (let i = 0; cur && i < 5; i++) { const code = cur.code || (cur.errors && cur.errors.code); if (code && !parts.includes(code)) parts.push(code); cur = cur.cause; } return parts.join(', '); } function softFailure(message, e) { const reason = e ? failureReason(e) || (e.message || '') : ''; const err = new Error(reason ? `${message} (${reason})` : message); err.soft = true; return err; } function createPhotoEnhanceWorker({ pool, getSetting, logAudit, invalidateEntries, sharp, photoAiUrl, storage, bus, notifyEvent, faceTimeoutMs, defaultModel, defaultFaceModel }) { const AI_URL = photoAiUrl || process.env.PHOTO_AI_URL || ''; const IDLE_MIN = 2000; const IDLE_MAX = 30000; const AI_TIMEOUT_MS = PHOTO_AI_TIMEOUT_MS; const AI_FACE_TIMEOUT_MS = Math.max(AI_TIMEOUT_MS, parseInt(faceTimeoutMs, 10) || PHOTO_AI_FACE_TIMEOUT_MS); const DEFAULT_MODEL = String(defaultModel || PHOTO_AI_DEFAULT_MODEL).trim() || PHOTO_AI_DEFAULT_MODEL; const DEFAULT_FACE_MODEL = String(defaultFaceModel || PHOTO_AI_DEFAULT_FACE_MODEL).trim() || PHOTO_AI_DEFAULT_FACE_MODEL; async function sendNotification(entryId, payload) { if (!notifyEvent) return; try { await notifyEvent(entryId, payload); } catch (e) { console.error('Photo worker notification failed:', e.message); } } 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, soft_retries: 0, last_at: null, last_error: null }; const softTries = new Map(); const CONFIG = { idle_min_ms: IDLE_MIN, idle_max_ms: IDLE_MAX, ai_timeout_ms: AI_TIMEOUT_MS, face_timeout_ms: AI_FACE_TIMEOUT_MS, max_attempts: PHOTO_MAX_ATTEMPTS, soft_max_retries: PHOTO_SOFT_MAX_RETRIES, soft_backoff_ms: PHOTO_SOFT_BACKOFF_MS, soft_backoff_max_ms: PHOTO_SOFT_BACKOFF_MAX_MS, ai_url: AI_URL, default_model: DEFAULT_MODEL, face_model: DEFAULT_FACE_MODEL, }; function sleep(ms) { return new Promise((resolve) => { const t = setTimeout(() => { wake = null; resolve(); }, ms); wake = () => { clearTimeout(t); wake = null; resolve(); }; }); } function wakeLocal() { if (wake) wake(); } function notify() { wakeLocal(); if (bus) bus.publish(); } if (bus) bus.subscribe(wakeLocal); 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(); } } async function enhanceWithSharp(srcKey, 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); const srcPath = await storage.localize(srcKey); if (!srcPath) throw new Error('Файл фото не найден'); let pipeline = sharp(srcPath).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'; const out = await pipeline.jpeg({ quality: 92, mozjpeg: true }).toBuffer(); await storage.put(newName, out); return `/uploads/${newName}`; } function aiParams(job) { const p = job && job.params && typeof job.params === 'object' ? job.params : {}; const face = String(p.face || (job && job.action === 'ai_face' ? 'face' : 'off')).trim().toLowerCase() || 'off'; const strength = Number(p.strength); return { model: String(p.model || DEFAULT_MODEL).trim().toLowerCase() || DEFAULT_MODEL, face, face_model: String(p.face_model || DEFAULT_FACE_MODEL).trim().toLowerCase() || DEFAULT_FACE_MODEL, strength: Number.isFinite(strength) ? strength : PHOTO_AI_DEFAULT_STRENGTH, timeout_ms: face === 'off' ? AI_TIMEOUT_MS : AI_FACE_TIMEOUT_MS, }; } function aiMeta(source) { const o = source && typeof source === 'object' ? source : {}; const num = (v) => (Number.isFinite(Number(v)) && v !== null && v !== '' ? Number(v) : null); return { model: o.model ? String(o.model) : null, face: o.face ? String(o.face) : null, face_model: o.face_model ? String(o.face_model) : null, faces_found: num(o.faces_found), device: o.device ? String(o.device) : null, elapsed_ms: num(o.elapsed_ms), warnings: Array.isArray(o.warnings) ? o.warnings.map((w) => String(w).slice(0, 300)).slice(0, 5) : [], }; } async function readErrorBody(resp) { try { const text = (await resp.text()).trim(); if (!text) return null; let data = null; try { data = JSON.parse(text); } catch (e) {} if (data && typeof data === 'object' && !Array.isArray(data)) { return String(data.error || data.detail || '').slice(0, 400) || null; } return text.slice(0, 400); } catch (e) { return null; } } async function runAiEnhance(srcKey, job) { if (!AI_URL) throw new Error('PHOTO_AI_URL не настроен'); const opt = aiParams(job); const buf = await storage.getBuffer(srcKey); if (!buf) throw new Error('Файл фото не найден'); const fd = new FormData(); fd.append('image', new Blob([buf], { type: 'image/jpeg' }), 'photo.jpg'); fd.append('scale', '2'); fd.append('model', opt.model); fd.append('face', opt.face); fd.append('face_model', opt.face_model); fd.append('strength', String(opt.strength)); let resp; try { resp = await fetch(AI_URL.replace(/\/+$/, '') + '/enhance', { method: 'POST', headers: { Accept: 'application/json' }, body: fd, signal: AbortSignal.timeout(opt.timeout_ms), }); } catch (e) { if (isTimeoutFailure(e)) { throw new Error(`ИИ-сервис не ответил за ${Math.round(opt.timeout_ms / 1000)} с`); } throw softFailure('ИИ-сервис недоступен', e); } if (!resp.ok) { const retryAfter = resp.headers.get('retry-after'); const detail = `ИИ-сервис ответил ${resp.status}${retryAfter ? ` (Retry-After: ${retryAfter})` : ''}`; const reason = await readErrorBody(resp); const message = reason ? `${detail}: ${reason}` : detail; if (isSoftStatus(resp.status)) throw softFailure(message); throw new Error(message); } const contentType = resp.headers.get('content-type') || ''; if (!contentType.includes('application/json')) { const out = Buffer.from(await resp.arrayBuffer()); return { out, meta: aiMeta({ model: resp.headers.get('x-photo-ai-model') || opt.model, face: resp.headers.get('x-photo-ai-face') || opt.face, face_model: opt.face === 'off' ? null : opt.face_model, device: resp.headers.get('x-photo-ai-device'), elapsed_ms: resp.headers.get('x-photo-ai-elapsed-ms'), }), }; } let data; try { data = await resp.json(); } catch (e) { throw new Error('ИИ-сервис вернул нечитаемый JSON'); } if (!data || typeof data !== 'object' || Array.isArray(data)) throw new Error('ИИ-сервис вернул неожиданный JSON'); if (data.ok === false) throw new Error(String(data.error || 'ИИ-сервис вернул ошибку').slice(0, 500)); if (typeof data.image_base64 !== 'string' || !data.image_base64) throw new Error('ИИ-сервис не вернул изображение'); let out; try { out = Buffer.from(data.image_base64, 'base64'); } catch (e) { throw new Error('ИИ-сервис вернул изображение в нечитаемом base64'); } if (!out.length) throw new Error('ИИ-сервис вернул пустое изображение'); return { out, meta: aiMeta({ model: data.model || opt.model, face: data.face || opt.face, face_model: data.face_model || (opt.face === 'off' ? null : opt.face_model), faces_found: data.faces_found, device: data.device, elapsed_ms: data.elapsed_ms, warnings: data.warnings, }), }; } async function enhanceWithAi(srcKey, job) { const { out, meta } = await runAiEnhance(srcKey, job); const newName = crypto.randomBytes(12).toString('hex') + '.jpg'; await storage.put(newName, out); return { path: `/uploads/${newName}`, meta }; } async function applyResult(job, newPath, meta) { const m = meta || {}; 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, model: m.model, face: m.face, face_model: m.face_model, device: m.device, faces_found: m.faces_found, elapsed_ms: m.elapsed_ms, warnings: m.warnings && m.warnings.length ? m.warnings : null, }); } const title = job.action === 'ai_face' ? 'Фото обработано нейросетью с восстановлением лиц: {student}' : job.action === 'ai_upscale' ? 'Фото обработано нейросетью (апскейл): {student}' : job.action === 'ai' ? 'Фото обработано нейросетью: {student}' : 'Фото улучшено на сервере: {student}'; const warn = m.warnings && m.warnings.length ? ' · ' + m.warnings.join('; ').slice(0, 300) : ''; await sendNotification(job.entry_id, { type: 'photo.job.done', title, body: `Группа {group} · запись #${job.entry_id} · результат ждёт применения${warn}`, target: { job_id: job.id, action: job.action, after_path: newPath, model: m.model, face: m.face, face_model: m.face_model, device: m.device, faces_found: m.faces_found, elapsed_ms: m.elapsed_ms, }, }); } async function softRetry(job, message) { const n = (softTries.get(job.id) || 0) + 1; if (n <= PHOTO_SOFT_MAX_RETRIES) { softTries.set(job.id, n); stats.soft_retries++; await pool.query(`UPDATE photo_jobs SET status = 'pending', error = $1 WHERE id = $2`, [message, job.id]); if (logAudit && (n === 1 || n % 10 === 0)) { await logAudit(null, 'photo.job.soft_retry', { entry_id: job.entry_id, job_id: job.id, soft_attempt: n, soft_limit: PHOTO_SOFT_MAX_RETRIES, error: message, }); } const delay = Math.min(PHOTO_SOFT_BACKOFF_MS * Math.pow(2, n - 1), PHOTO_SOFT_BACKOFF_MAX_MS); return { ok: false, delay }; } softTries.delete(job.id); const text = `${message} — ИИ-сервис недоступен, мягкие повторы исчерпаны (${PHOTO_SOFT_MAX_RETRIES}), задание остановлено`; await pool.query( `UPDATE photo_jobs SET status = 'error', error = $1, finished_at = now() WHERE id = $2`, [text, job.id] ); stats.errors++; stats.last_error = text; if (logAudit) { await logAudit(null, 'photo.job.error', { entry_id: job.entry_id, job_id: job.id, error: text, soft_attempts: PHOTO_SOFT_MAX_RETRIES }); } await sendNotification(job.entry_id, { type: 'photo.job.error', title: 'Ошибка обработки фото: {student}', body: `Группа {group} · запись #${job.entry_id} · ${text}`, target: { job_id: job.id, action: job.action, error: text }, }); return { ok: true, delay: 0 }; } 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 = 'У записи нет фото'; await sendNotification(job.entry_id, { type: 'photo.job.error', title: 'Ошибка обработки фото: {student}', body: `Группа {group} · запись #${job.entry_id} · у записи нет фото для обработки`, target: { job_id: job.id, action: job.action }, }); return { ok: true, delay: 0 }; } const photoPath = rows[0].photo_path; try { const result = PHOTO_AI_ACTIONS.has(job.action) ? await enhanceWithAi(photoPath, job) : { path: await enhanceWithSharp(photoPath, job.params), meta: null }; await applyResult(job, result.path, result.meta); softTries.delete(job.id); stats.done++; stats.last_at = new Date().toISOString(); stats.last_error = null; return { ok: true, delay: 0 }; } catch (e) { const message = (e && e.message ? e.message : 'error').slice(0, 500); stats.last_at = new Date().toISOString(); stats.last_error = message; if (e && e.soft) return await softRetry(job, 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 { ok: false, delay: 0 }; } softTries.delete(job.id); 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 }); await sendNotification(job.entry_id, { type: 'photo.job.error', title: 'Ошибка обработки фото: {student}', body: `Группа {group} · запись #${job.entry_id} · ${message}`, target: { job_id: job.id, action: job.action, error: message }, }); return { ok: true, delay: 0 }; } } 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 res = { ok: true, delay: 0 }; try { res = await processOne(job); } catch (e) { console.error('Photo worker process error:', e); } finally { processing = false; currentId = null; } if (!res || !res.ok) await sleep((res && res.delay) || 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, bus }) { 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 wakeLocal() { if (wake) wake(); } function notify() { wakeLocal(); if (bus) bus.publish(); } if (bus) bus.subscribe(wakeLocal); 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) { const d = textDiff(String(row.description ?? ''), changed ? checked : original); await logAudit(null, 'entry.ai.auto-check', { id: row.id, source: 'ai', changed, fields: changed ? ['description'] : [], changes: changed ? [{ field: 'description', label: FIELD_LABELS.description, stats: d.stats, diff: d.segments, truncated: d.truncated }] : [] }); } 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 }; } function createLessonReportChecker({ pool, getSetting, logAudit, aiUrl, model, defaultPrompt, bus, onDone }) { 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 || ''; const TIMEOUT_MS = Math.max(5000, parseInt(process.env.LESSON_AI_TIMEOUT_MS || '120000', 10) || 120000); const MIN_CHARS = 20; const MAX_INPUT_CHARS = 4000; const TEXT_MAX = 5000; const MAX_ATTEMPTS = 3; 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, formatted: 0, unchanged: 0, errors: 0, last_at: null, last_error: null }; const attempts = new Map(); 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(); if (bus) bus.publish(); } if (bus) bus.subscribe(() => notify()); async function isEnabled() { return String(await getSetting('lesson_ai_enabled', 'true')) !== 'false'; } async function claimNext() { const client = await pool.connect(); try { await client.query('BEGIN'); const { rows } = await client.query( `SELECT lr.id, lr.text, lr.text_original, lr.lesson_date, lr.lesson_time, lr.author_id, g.name AS group_name, g.branch_id FROM lesson_reports lr JOIN groups g ON g.id = lr.group_id WHERE lr.ai_status = 'pending' ORDER BY lr.id ASC LIMIT 1 FOR UPDATE OF lr SKIP LOCKED` ); if (!rows.length) { await client.query('COMMIT'); return null; } await client.query(`UPDATE lesson_reports 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) { const b = String(base || '').trim().replace(/\/+$/, ''); if (!/^https?:\/\//i.test(b)) return null; if (!/\/v1$/i.test(b)) return b + '/v1'; return b; } async function callModel(userText, 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, userText.length * 2 + 512)); } else { url = `${AI_URL.replace(/\/+$/, '')}/v1/chat/completions`; model = MODEL; maxTokens = Math.min(1536, Math.max(512, userText.length + 512)); } const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), 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: userText }, ], 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]; return (choice && choice.message && choice.message.content || '').trim(); } finally { clearTimeout(timer); } } function normalizeOutput(out) { let s = String(out || '').trim(); if (!s) return ''; s = s.replace(/^```[a-zA-Z]*\s*/, '').replace(/```$/, '').trim(); s = s.replace(/^(?:итоговый|результативный|готовый|итог)\s*текст\s*(?:отчёта|отчета|report)?\s*[:—-]\s*/i, '').trim(); s = s.replace(/^["'«”]+/, '').replace(/["'«»"]+$/, '').trim(); s = s.replace(/\r\n/g, '\n').replace(/\n{3,}/g, '\n\n'); return s; } async function finish(row, status, text, error, aiText) { await pool.query( `UPDATE lesson_reports SET text = $1, text_ai = $2, ai_status = $3, ai_error = $4, ai_checked_at = now(), updated_at = now() WHERE id = $5`, [text, aiText, status, error || null, row.id] ); attempts.delete(row.id); stats.last_at = new Date().toISOString(); stats.last_error = error || null; if (onDone) await onDone({ row, status, text, aiText, error: error || null }); } async function processOne(row) { const original = String(row.text_original || row.text || '').trim(); if (original.length < MIN_CHARS) { await finish(row, 'skipped', original, null, null); stats.unchanged++; return true; } const prompt = (await getSetting('lesson_ai_prompt', DEFAULT_PROMPT)) || DEFAULT_PROMPT; const date = String(row.lesson_date).slice(0, 10); const time = row.lesson_time ? String(row.lesson_time).slice(0, 5) : ''; const input = original.length > MAX_INPUT_CHARS ? original.slice(0, MAX_INPUT_CHARS) : original; const context = [ `Группа: ${row.group_name || '—'}`, `Дата занятия: ${date}`, time ? `Время занятия: ${time}` : '', '', 'Текст отчёта:', input, ].filter(l => l !== '').join('\n'); try { const out = normalizeOutput(await callModel(context, prompt)); if (!out || out.length > TEXT_MAX) { const reason = out ? 'Модель вернула слишком длинный текст' : 'Модель вернула пустой ответ'; await finish(row, 'error', original, reason, null); stats.errors++; return true; } const changed = out !== original; await finish(row, changed ? 'done' : 'skipped', changed ? out : original, null, changed ? out : null); stats.checks++; if (changed) stats.formatted++; else stats.unchanged++; return true; } catch (e) { const message = (e && e.message ? e.message : 'error').slice(0, 500); const tries = (attempts.get(row.id) || 0) + 1; if (tries < MAX_ATTEMPTS) { attempts.set(row.id, tries); await pool.query(`UPDATE lesson_reports SET ai_status = 'pending', ai_error = $1 WHERE id = $2`, [message, row.id]); if (logAudit) await logAudit(null, 'lesson_report.ai.retry', { id: row.id, attempt: tries, error: message }); return false; } await finish(row, 'error', original, message, null); stats.errors++; if (logAudit) await logAudit(null, 'lesson_report.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('Lesson AI claim error:', e.message); 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('Lesson AI process error:', e.message); } finally { processing = false; currentId = null; } if (!ok) await sleep(IDLE_MAX_MS); } } function resetStale() { return pool.query(`UPDATE lesson_reports SET ai_status = 'pending' WHERE ai_status = 'processing'`); } function start() { if (started) return; started = true; startedAt = Date.now(); resetStale() .catch((e) => console.error('Lesson AI reset error:', e.message)) .finally(() => { loop().catch((e) => console.error('Lesson AI loop error:', e.message)); }); } function stop() { stopped = true; notify(); } function getInfo() { return { started_at: startedAt ? new Date(startedAt).toISOString() : null, uptime_ms: startedAt ? Date.now() - startedAt : 0, processing, current_id: currentId, config: { idle_min_ms: IDLE_MIN_MS, idle_max_ms: IDLE_MAX_MS, timeout_ms: TIMEOUT_MS, model: MODEL, ai_url: AI_URL }, ...stats, }; } return { start, stop, notify, getInfo }; } module.exports = { createEntryAutoChecker, createPhotoEnhanceWorker, createLessonReportChecker };