Добавлена система уведомлений о системных и фоновых событиях (новые записи журнала, обработка фото, авто-проверка текста, блокировки IP, бэкапы). - backend (server.js, worker.js): - каталог NOTIFY_TYPES с метаданными и уровнями - таблицы notifications и notification_reads в db/init.sql и db/migration.sql - SSE-стрим GET /api/notifications/stream через Redis pub/sub с in-memory fallback - REST API: список, счётчик непрочитанных, отметка о прочтении, удаление, очистка - настройки уведомлений в settings (notify_enabled, notify_retention_days, notify_<тип>) - автоматическая очистка старых уведомлений по расписанию - frontend: - колокольчик со счётчиком непрочитанных в шапке (admin.js) - страница списка уведомлений public/notifications.html и public/js/notifications.js - секция настроек уведомлений в public/settings.html и public/js/settings.js - стили для уведомлений в public/admin.css - тесты и документация: - добавлены проверки в api.smoketest.js - обновлены README.md и AGENTS.md
544 lines
19 KiB
JavaScript
544 lines
19 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 path = require('path');
|
||
const fs = require('fs');
|
||
const crypto = require('crypto');
|
||
const { textDiff, FIELD_LABELS } = require('./diff');
|
||
const PHOTO_MAX_ATTEMPTS = 3;
|
||
const PHOTO_AI_TIMEOUT_MS = 300000;
|
||
|
||
function createPhotoEnhanceWorker({ pool, getSetting, logAudit, invalidateEntries, sharp, photoAiUrl, uploadsDir, storage, bus, notifyEvent }) {
|
||
const AI_URL = photoAiUrl || process.env.PHOTO_AI_URL || '';
|
||
const IDLE_MIN = 2000;
|
||
const IDLE_MAX = 30000;
|
||
|
||
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, 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 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 outPath = path.join(uploadsDir, newName);
|
||
await pipeline.jpeg({ quality: 92, mozjpeg: true }).toFile(outPath);
|
||
await storage.persist(newName, outPath);
|
||
return `/uploads/${newName}`;
|
||
}
|
||
|
||
async function runAiEnhance(srcKey) {
|
||
if (!AI_URL) throw new Error('PHOTO_AI_URL не настроен');
|
||
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');
|
||
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';
|
||
const outPath = path.join(uploadsDir, newName);
|
||
fs.writeFileSync(outPath, out);
|
||
await storage.persist(newName, outPath);
|
||
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 });
|
||
await sendNotification(job.entry_id, {
|
||
type: 'photo.job.done',
|
||
title: job.action === 'ai' ? 'Фото обработано нейросетью: {student}' : 'Фото улучшено на сервере: {student}',
|
||
body: `Группа {group} · запись #${job.entry_id} · результат ждёт применения`,
|
||
target: { 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 = 'У записи нет фото';
|
||
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 true;
|
||
}
|
||
const photoPath = rows[0].photo_path;
|
||
try {
|
||
const newPath = job.action === 'ai' ? await runAiEnhance(photoPath) : await enhanceWithSharp(photoPath, 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 });
|
||
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 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, 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 };
|
||
}
|
||
|
||
module.exports = { createEntryAutoChecker, createPhotoEnhanceWorker };
|