- Implement feedback management system in admin interface - Add feedback processing and AI-based features - Update user interface with new feedback components - Modify server-side API endpoints for feedback handling - Enhance worker.js for background feedback processing - Update deployment configuration for feedback components Co-authored-by: openhands <openhands@all-hands.dev>
1353 lines
48 KiB
JavaScript
1353 lines
48 KiB
JavaScript
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 };
|
||
}
|
||
|
||
const LESSON_TOPIC_POSITION_LABELS = {
|
||
first: 'первое занятие модуля',
|
||
middle: 'промежуточное занятие модуля',
|
||
last: 'последнее занятие модуля',
|
||
};
|
||
|
||
function lessonTopicPosition(topic) {
|
||
const m = /(\d+)\s*\/\s*(\d+)\s*$/.exec(String(topic || '').trim());
|
||
if (!m) return null;
|
||
const n = parseInt(m[1], 10);
|
||
const total = parseInt(m[2], 10);
|
||
if (!Number.isInteger(n) || !Number.isInteger(total) || total < 1 || n < 1 || n > total) return null;
|
||
const kind = n === total ? 'last' : n === 1 ? 'first' : 'middle';
|
||
return { kind, n, total, label: `${LESSON_TOPIC_POSITION_LABELS[kind]} (${n} из ${total})` };
|
||
}
|
||
|
||
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.topic, 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 topic = String(row.topic || '').trim();
|
||
const position = lessonTopicPosition(topic);
|
||
const context = [
|
||
`Группа: ${row.group_name || '—'}`,
|
||
`Дата занятия: ${date}`,
|
||
time ? `Время занятия: ${time}` : '',
|
||
topic ? `Тема занятия: ${topic}` : '',
|
||
position ? `Позиция темы: ${position.label}` : '',
|
||
'',
|
||
'Текст отчёта:',
|
||
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 };
|
||
}
|
||
|
||
function createFeedbackWriter(opts) {
|
||
const pool = opts.pool;
|
||
const getSetting = opts.getSetting;
|
||
const logAudit = opts.logAudit;
|
||
const aiUrl = (opts.aiUrl || process.env.AI_URL || 'http://text-corrector:8080');
|
||
const defaultPrompt = opts.defaultPrompt || '';
|
||
const bus = opts.bus;
|
||
const onDone = opts.onDone || (() => {});
|
||
|
||
const AI_MODEL = process.env.AI_MODEL || 'qwen2.5-1.5b-instruct-q4_k_m.gguf';
|
||
const TIMEOUT_MS = opts.timeoutMs || parseInt(process.env.FEEDBACK_AI_TIMEOUT_MS || '120000', 10) || 120000;
|
||
const MAX_ATTEMPTS = 3;
|
||
const IDLE_MIN_MS = 500;
|
||
const IDLE_MAX_MS = 30000;
|
||
const MIN_CHARS = 1;
|
||
const MAX_INPUT_CHARS = 6000;
|
||
const TEXT_MAX = 5000;
|
||
|
||
let running = false;
|
||
let timer = null;
|
||
let stopped = false;
|
||
let idleMs = IDLE_MIN_MS;
|
||
let inFlight = false;
|
||
|
||
function notify() {
|
||
if (!running || inFlight) return;
|
||
clearTimeout(timer);
|
||
timer = setTimeout(tick, 50);
|
||
}
|
||
|
||
if (bus && typeof bus.subscribe === 'function') {
|
||
try { bus.subscribe(notify); } catch (e) {}
|
||
}
|
||
|
||
async function getSettingValue(key, def) {
|
||
try {
|
||
const v = await getSetting(key, def);
|
||
return v;
|
||
} catch (e) {
|
||
return def;
|
||
}
|
||
}
|
||
|
||
async function resolveActiveProfile(profileId) {
|
||
try {
|
||
const active = await getSettingValue('ai_active_profile', 'native');
|
||
const profilesRaw = await getSettingValue('ai_profiles', '[]');
|
||
let profiles = [];
|
||
try { profiles = JSON.parse(profilesRaw); } catch (e) {}
|
||
if (profileId) {
|
||
const p = profiles.find(x => String(x.id) === String(profileId));
|
||
if (p) return { type: 'profile', profile: p };
|
||
}
|
||
if (String(active) !== 'native') {
|
||
const p = profiles.find(x => String(x.id) === String(active));
|
||
if (p) return { type: 'profile', profile: p };
|
||
}
|
||
return { type: 'native', model: AI_MODEL, baseUrl: aiUrl };
|
||
} catch (e) {
|
||
return { type: 'native', model: AI_MODEL, baseUrl: aiUrl };
|
||
}
|
||
}
|
||
|
||
function normalizeOpenAiBase(baseUrl) {
|
||
const b = (baseUrl || '').replace(/\/$/, '');
|
||
return b.includes('/v1') ? b : `${b}/v1`;
|
||
}
|
||
|
||
async function callModel(systemPrompt, userText, profileId) {
|
||
const text = String(userText || '');
|
||
if (text.length > MAX_INPUT_CHARS) {
|
||
return text.slice(0, TEXT_MAX);
|
||
}
|
||
const resolved = await resolveActiveProfile(profileId);
|
||
const controller = new AbortController();
|
||
const timerAbort = setTimeout(() => controller.abort(), TIMEOUT_MS);
|
||
try {
|
||
let url, headers, body;
|
||
if (resolved.type === 'profile') {
|
||
const p = resolved.profile;
|
||
const base = normalizeOpenAiBase(p.base_url || aiUrl);
|
||
url = `${base}/chat/completions`;
|
||
headers = { 'Content-Type': 'application/json' };
|
||
if (p.api_key) headers.Authorization = `Bearer ${p.api_key}`;
|
||
body = {
|
||
model: p.model || AI_MODEL,
|
||
messages: [
|
||
{ role: 'system', content: systemPrompt },
|
||
{ role: 'user', content: text }
|
||
],
|
||
temperature: 0.1,
|
||
max_tokens: p.max_tokens || Math.min(2048, Math.max(256, text.length * 2))
|
||
};
|
||
} else {
|
||
const base = normalizeOpenAiBase(resolved.baseUrl || aiUrl);
|
||
url = `${base}/chat/completions`;
|
||
headers = { 'Content-Type': 'application/json' };
|
||
body = {
|
||
model: resolved.model || AI_MODEL,
|
||
messages: [
|
||
{ role: 'system', content: systemPrompt },
|
||
{ role: 'user', content: text }
|
||
],
|
||
temperature: 0.1,
|
||
max_tokens: Math.min(2048, Math.max(256, text.length * 2))
|
||
};
|
||
}
|
||
const res = await fetch(url, { method: 'POST', headers, body: JSON.stringify(body), signal: controller.signal });
|
||
if (!res.ok) throw new Error(`AI service error: ${res.status}`);
|
||
const data = await res.json();
|
||
return data.choices?.[0]?.message?.content?.trim() || '';
|
||
} finally {
|
||
clearTimeout(timerAbort);
|
||
}
|
||
}
|
||
|
||
function normalizeOutput(s) {
|
||
let out = String(s || '').trim();
|
||
if (out.startsWith('```') && out.endsWith('```')) {
|
||
out = out.slice(3, -3).trim();
|
||
}
|
||
out = out.replace(/^(Итоговый текст:|Отзыв:)\s*/i, '');
|
||
out = out.replace(/^"|"$/g, '');
|
||
out = out.replace(/\n{3,}/g, '\n\n');
|
||
return out;
|
||
}
|
||
|
||
async function isEnabled() {
|
||
try {
|
||
const v = await getSettingValue('feedback_ai_enabled', 'true');
|
||
return String(v) !== 'false';
|
||
} catch (e) {
|
||
return true;
|
||
}
|
||
}
|
||
|
||
async function claimNext() {
|
||
const client = await pool.connect();
|
||
try {
|
||
await client.query('BEGIN');
|
||
const q = `SELECT f.*, g.name AS group_name, g.branch_id AS group_branch_id
|
||
FROM feedbacks f
|
||
LEFT JOIN groups g ON g.id = f.group_id
|
||
WHERE f.ai_status = 'pending' AND f.deleted_at IS NULL
|
||
ORDER BY f.id ASC
|
||
FOR UPDATE OF f SKIP LOCKED
|
||
LIMIT 1`;
|
||
const r = await client.query(q);
|
||
if (r.rows.length === 0) {
|
||
await client.query('ROLLBACK');
|
||
return null;
|
||
}
|
||
const row = r.rows[0];
|
||
await client.query(`UPDATE feedbacks SET ai_status='processing', updated_at=now() WHERE id=$1`, [row.id]);
|
||
await client.query('COMMIT');
|
||
row.branch_id = row.branch_id ?? row.group_branch_id ?? null;
|
||
return row;
|
||
} catch (e) {
|
||
try { await client.query('ROLLBACK'); } catch (e2) {}
|
||
throw e;
|
||
} finally {
|
||
client.release();
|
||
}
|
||
}
|
||
|
||
async function processOne(row) {
|
||
if (row.text && String(row.text).trim().length >= MIN_CHARS) {
|
||
// allow overwrite only if empty? keep as done? but generate on request
|
||
}
|
||
let context = '';
|
||
context += `Группа: ${row.group_name || ''}\n`;
|
||
context += `Дата занятия: ${row.feedback_date || ''}\n`;
|
||
context += `Резидент: ${row.resident || ''}\n`;
|
||
context += `Тема занятия: ${row.topic || ''}\n`;
|
||
context += `Прошлый отзыв: ${row.past_review || 'нет'}\n`;
|
||
let shorts = [];
|
||
try {
|
||
const ids = Array.isArray(row.short_message_ids) ? row.short_message_ids : [];
|
||
if (ids.length) {
|
||
const sr = await pool.query(`SELECT resident, lesson_date, topic, message FROM short_messages WHERE id = ANY($1::int[]) AND deleted_at IS NULL ORDER BY id`, [ids]);
|
||
shorts = sr.rows;
|
||
}
|
||
} catch (e) {}
|
||
if (shorts.length) {
|
||
context += 'Короткие сообщения (выбраны):\n';
|
||
for (const s of shorts) {
|
||
context += `- ${s.lesson_date || ''} – ${s.topic || '-'} – ${s.message || ''}\n`;
|
||
}
|
||
}
|
||
const systemPrompt = defaultPrompt + (context ? '\n\nКОНТЕКСТ:\n' + context : '');
|
||
let out = await callModel(systemPrompt, row.text && row.text.trim() ? row.text : (context || 'Сформируй отзыв по контексту'), row.ai_profile || null);
|
||
let norm = normalizeOutput(out);
|
||
if (!norm) return { status: 'error', error: 'Пустой ответ от модели' };
|
||
if (norm.length > TEXT_MAX) norm = norm.slice(0, TEXT_MAX);
|
||
const oldText = (row.text || '').trim();
|
||
if (oldText && norm.trim() === oldText) {
|
||
return { status: 'skipped', text: oldText };
|
||
}
|
||
return { status: 'done', text: norm, aiText: norm };
|
||
}
|
||
|
||
async function finish(row, res) {
|
||
const client = await pool.connect();
|
||
try {
|
||
await client.query('BEGIN');
|
||
const sets = [];
|
||
const params = [];
|
||
let i = 1;
|
||
sets.push(`ai_status=$${i++}`); params.push(res.status);
|
||
sets.push(`ai_checked_at=now()`);
|
||
if (res.error) { sets.push(`ai_error=$${i++}`); params.push(res.error); } else { sets.push(`ai_error=NULL`); }
|
||
if (res.status === 'done' && res.text) { sets.push(`text=$${i++}`); params.push(res.text); }
|
||
sets.push(`updated_at=now()`);
|
||
params.push(row.id);
|
||
await client.query(`UPDATE feedbacks SET ${sets.join(', ')} WHERE id=$${i}`, params);
|
||
await client.query('COMMIT');
|
||
} catch (e) {
|
||
try { await client.query('ROLLBACK'); } catch (e2) {}
|
||
throw e;
|
||
} finally {
|
||
client.release();
|
||
}
|
||
try { onDone({ row, status: res.status, text: res.text || row.text, aiText: res.aiText, error: res.error }); } catch (e) {}
|
||
}
|
||
|
||
async function tick() {
|
||
if (stopped || !running || inFlight) return;
|
||
const enabled = await isEnabled();
|
||
if (!enabled) {
|
||
idleMs = Math.min(idleMs * 2, IDLE_MAX_MS);
|
||
timer = setTimeout(tick, idleMs);
|
||
return;
|
||
}
|
||
inFlight = true;
|
||
try {
|
||
let attempts = 0;
|
||
while (attempts < MAX_ATTEMPTS) {
|
||
attempts++;
|
||
const row = await claimNext();
|
||
if (!row) break;
|
||
let res;
|
||
try {
|
||
res = await processOne(row);
|
||
} catch (e) {
|
||
res = { status: 'error', error: e.message };
|
||
}
|
||
if (res.status === 'error' && attempts < MAX_ATTEMPTS) {
|
||
const client = await pool.connect();
|
||
try {
|
||
await client.query('BEGIN');
|
||
await client.query(`UPDATE feedbacks SET ai_status='pending', updated_at=now(), ai_error=$2 WHERE id=$1`, [row.id, res.error || 'error']);
|
||
await client.query('COMMIT');
|
||
} catch (e2) {
|
||
try { await client.query('ROLLBACK'); } catch (e3) {}
|
||
} finally {
|
||
client.release();
|
||
}
|
||
idleMs = IDLE_MIN_MS;
|
||
continue;
|
||
}
|
||
await finish(row, res);
|
||
idleMs = IDLE_MIN_MS;
|
||
}
|
||
} catch (e) {
|
||
idleMs = Math.min(idleMs * 2, IDLE_MAX_MS);
|
||
} finally {
|
||
inFlight = false;
|
||
timer = setTimeout(tick, idleMs);
|
||
}
|
||
}
|
||
|
||
async function resetStale() {
|
||
try {
|
||
await pool.query(`UPDATE feedbacks SET ai_status='pending', ai_error=NULL, ai_checked_at=NULL WHERE ai_status='processing'`);
|
||
} catch (e) {}
|
||
}
|
||
|
||
function start() {
|
||
if (running) return;
|
||
running = true;
|
||
stopped = false;
|
||
idleMs = IDLE_MIN_MS;
|
||
resetStale().finally(() => { tick(); });
|
||
}
|
||
|
||
function stop() {
|
||
running = false;
|
||
stopped = true;
|
||
if (timer) { clearTimeout(timer); timer = null; }
|
||
}
|
||
|
||
function getInfo() {
|
||
return { running, inFlight, idleMs };
|
||
}
|
||
|
||
return { start, stop, notify, getInfo };
|
||
}
|
||
|
||
module.exports = { createEntryAutoChecker, createPhotoEnhanceWorker, createLessonReportChecker, createFeedbackWriter, lessonTopicPosition };
|