feat: persistent photo job queue, worker, history with rollback
- Replace in-memory photoAiJobs Map with DB-backed photo_jobs table - Background photo worker (worker.js) with retry, backoff, stale reset - Handles both AI (Real-ESRGAN) and server-side (sharp) enhancement - Controlled via photo_worker_enabled setting - Photo job history in enhance modal with before/after thumbnails + rollback - Worker dashboard: photo jobs section with status counts, recent table, compare slider for before/after, rollback from worker UI - New endpoints: /api/photo-jobs/status|wake|enabled|requeue-failed, /api/entries/:id/photo/jobs (history), .../rollback - swapEntryPhotoFiles logs every mutation to photo_jobs table
This commit is contained in:
@@ -3,6 +3,261 @@ const IDLE_MAX_MS = 60000;
|
||||
const REQUEST_TIMEOUT_MS = parseInt(process.env.AI_REQUEST_TIMEOUT_MS || '120000', 10);
|
||||
const MAX_INPUT_CHARS = 2000;
|
||||
const MIN_TEXT_CHARS = 4;
|
||||
const path = require('path');
|
||||
const fs = require('fs');
|
||||
const crypto = require('crypto');
|
||||
const PHOTO_MAX_ATTEMPTS = 3;
|
||||
const PHOTO_AI_TIMEOUT_MS = 300000;
|
||||
|
||||
function createPhotoEnhanceWorker({ pool, getSetting, logAudit, invalidateEntries, sharp, photoAiUrl, uploadsDir, originalsDir }) {
|
||||
const AI_URL = photoAiUrl || process.env.PHOTO_AI_URL || '';
|
||||
const IDLE_MIN = 2000;
|
||||
const IDLE_MAX = 30000;
|
||||
|
||||
let started = false;
|
||||
let stopped = false;
|
||||
let idleMs = IDLE_MIN;
|
||||
let wake = null;
|
||||
let processing = false;
|
||||
let currentId = null;
|
||||
let startedAt = null;
|
||||
const stats = { jobs: 0, done: 0, errors: 0, last_at: null, last_error: null };
|
||||
const CONFIG = {
|
||||
idle_min_ms: IDLE_MIN,
|
||||
idle_max_ms: IDLE_MAX,
|
||||
ai_timeout_ms: PHOTO_AI_TIMEOUT_MS,
|
||||
max_attempts: PHOTO_MAX_ATTEMPTS,
|
||||
ai_url: AI_URL,
|
||||
};
|
||||
|
||||
function sleep(ms) {
|
||||
return new Promise((resolve) => {
|
||||
const t = setTimeout(() => { wake = null; resolve(); }, ms);
|
||||
wake = () => { clearTimeout(t); wake = null; resolve(); };
|
||||
});
|
||||
}
|
||||
|
||||
function notify() {
|
||||
if (wake) wake();
|
||||
}
|
||||
|
||||
async function isEnabled() {
|
||||
const v = await getSetting('photo_worker_enabled', 'true');
|
||||
return String(v) !== 'false';
|
||||
}
|
||||
|
||||
async function claimNext() {
|
||||
const client = await pool.connect();
|
||||
try {
|
||||
await client.query('BEGIN');
|
||||
const { rows } = await client.query(
|
||||
`SELECT j.id, j.entry_id, j.action, j.params
|
||||
FROM photo_jobs j
|
||||
WHERE j.status = 'pending'
|
||||
ORDER BY j.id ASC LIMIT 1 FOR UPDATE SKIP LOCKED`
|
||||
);
|
||||
if (!rows.length) {
|
||||
await client.query('COMMIT');
|
||||
return null;
|
||||
}
|
||||
await client.query(`UPDATE photo_jobs SET status = 'processing' WHERE id = $1`, [rows[0].id]);
|
||||
await client.query('COMMIT');
|
||||
return rows[0];
|
||||
} catch (e) {
|
||||
await client.query('ROLLBACK').catch(() => {});
|
||||
throw e;
|
||||
} finally {
|
||||
client.release();
|
||||
}
|
||||
}
|
||||
|
||||
function safeUnlinkPath(p) {
|
||||
try {
|
||||
if (!p) return;
|
||||
const abs = path.resolve(p);
|
||||
const root = path.resolve(uploadsDir);
|
||||
if (abs !== root && abs.startsWith(root + path.sep)) fs.unlinkSync(abs);
|
||||
} catch (e) { console.error('photo worker unlink:', e.message); }
|
||||
}
|
||||
|
||||
function backupOldFile(oldRelPath) {
|
||||
if (!oldRelPath) return null;
|
||||
const oldAbs = path.join(uploadsDir, String(oldRelPath).replace(/^\/+/, '').replace(/^uploads\//, ''));
|
||||
if (!fs.existsSync(oldAbs)) return null;
|
||||
const backupName = crypto.randomBytes(12).toString('hex') + (path.extname(oldAbs) || '.jpg');
|
||||
const backupPath = path.join(originalsDir, backupName);
|
||||
try {
|
||||
fs.renameSync(oldAbs, backupPath);
|
||||
return `/uploads/.originals/${backupName}`;
|
||||
} catch (e) {
|
||||
console.error('photo worker backup failed:', e.message);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
async function enhanceWithSharp(srcAbs, params) {
|
||||
if (!sharp) throw new Error('sharp недоступен на сервере');
|
||||
const p = params || {};
|
||||
const clamp = (v, min, max, def) => {
|
||||
const n = parseFloat(v);
|
||||
return Number.isFinite(n) ? Math.min(max, Math.max(min, n)) : def;
|
||||
};
|
||||
const brightness = clamp(p.brightness, 10, 300, 100);
|
||||
const contrast = clamp(p.contrast, 10, 300, 100);
|
||||
const saturate = clamp(p.saturate, 0, 300, 100);
|
||||
const sharpAmt = clamp(p.sharp, 0, 100, 0);
|
||||
const denoise = clamp(p.denoise, 0, 100, 0);
|
||||
let pipeline = sharp(srcAbs).rotate();
|
||||
if (denoise > 0) pipeline = pipeline.median(denoise > 70 ? 5 : 3);
|
||||
pipeline = pipeline.modulate({ brightness: brightness / 100, saturation: saturate / 100 });
|
||||
if (contrast !== 100) {
|
||||
const a = contrast / 100;
|
||||
pipeline = pipeline.linear(a, 128 * (1 - a));
|
||||
}
|
||||
if (sharpAmt > 0) {
|
||||
pipeline = pipeline.sharpen({ sigma: 0.5 + (sharpAmt / 100) * 1.5, m1: 0, m2: 1 + sharpAmt / 50 });
|
||||
}
|
||||
const newName = crypto.randomBytes(12).toString('hex') + '.jpg';
|
||||
await pipeline.jpeg({ quality: 92, mozjpeg: true }).toFile(path.join(uploadsDir, newName));
|
||||
return `/uploads/${newName}`;
|
||||
}
|
||||
|
||||
async function runAiEnhance(srcAbs) {
|
||||
if (!AI_URL) throw new Error('PHOTO_AI_URL не настроен');
|
||||
const buf = fs.readFileSync(srcAbs);
|
||||
const fd = new FormData();
|
||||
fd.append('image', new Blob([buf], { type: 'image/jpeg' }), 'photo.jpg');
|
||||
fd.append('scale', '2');
|
||||
const resp = await fetch(AI_URL.replace(/\/+$/, '') + '/enhance', {
|
||||
method: 'POST',
|
||||
body: fd,
|
||||
signal: AbortSignal.timeout(PHOTO_AI_TIMEOUT_MS),
|
||||
});
|
||||
if (!resp.ok) throw new Error('AI service error: ' + resp.status);
|
||||
const out = Buffer.from(await resp.arrayBuffer());
|
||||
const newName = crypto.randomBytes(12).toString('hex') + '.jpg';
|
||||
fs.writeFileSync(path.join(uploadsDir, newName), out);
|
||||
async function applyResult(job, photoPath, newPath) {
|
||||
const beforePath = backupOldFile(photoPath);
|
||||
await pool.query('UPDATE entries SET photo_path = $1 WHERE id = $2', [newPath, job.entry_id]);
|
||||
await pool.query('UPDATE entry_photos SET photo_path = $1 WHERE entry_id = $2 AND photo_path = $3', [newPath, job.entry_id, photoPath]);
|
||||
const oldThumb = path.join(uploadsDir, '.thumbs', path.basename(photoPath).replace(/\.[^.]+$/, '') + '.webp');
|
||||
safeUnlinkPath(oldThumb);
|
||||
await pool.query(
|
||||
`UPDATE photo_jobs SET status = 'done', before_path = $1, after_path = $2, error = NULL, finished_at = now() WHERE id = $3`,
|
||||
[beforePath, newPath, job.id]
|
||||
);
|
||||
if (logAudit) await logAudit(null, 'photo.job.done', { entry_id: job.entry_id, job_id: job.id, action: job.action, before_path: beforePath, after_path: newPath });
|
||||
if (invalidateEntries) invalidateEntries();
|
||||
}
|
||||
async function processOne(job) {
|
||||
stats.jobs++;
|
||||
const { rows } = await pool.query('SELECT photo_path FROM entries WHERE id = $1', [job.entry_id]);
|
||||
if (!rows.length || !rows[0].photo_path) {
|
||||
await pool.query(
|
||||
`UPDATE photo_jobs SET status = 'error', error = 'У записи нет фото', finished_at = now() WHERE id = $1`,
|
||||
[job.id]
|
||||
);
|
||||
stats.errors++;
|
||||
stats.last_at = new Date().toISOString();
|
||||
stats.last_error = 'У записи нет фото';
|
||||
return true;
|
||||
}
|
||||
const photoPath = rows[0].photo_path;
|
||||
const srcAbs = path.join(uploadsDir, String(photoPath).replace(/^\/+/, '').replace(/^uploads\//, ''));
|
||||
try {
|
||||
const newPath = job.action === 'ai' ? await runAiEnhance(srcAbs) : await enhanceWithSharp(srcAbs, job.params);
|
||||
await applyResult(job, photoPath, 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]
|
||||
);
|
||||
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 };
|
||||
}
|
||||
|
||||
stats.errors++;
|
||||
if (logAudit) await logAudit(null, 'photo.job.error', { entry_id: job.entry_id, job_id: job.id, error: message });
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
return `/uploads/${newName}`;
|
||||
}
|
||||
|
||||
|
||||
function createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl, defaultPrompt, model }) {
|
||||
const AI_URL = aiUrl || process.env.AI_URL || 'http://text-corrector:8080';
|
||||
@@ -258,4 +513,6 @@ function createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl, defaultProm
|
||||
return { start, notify, getStats, getInfo };
|
||||
}
|
||||
|
||||
module.exports = { createEntryAutoChecker };
|
||||
module.exports = { createEntryAutoChecker, createPhotoEnhanceWorker };
|
||||
|
||||
(function photoEnhanceWorkerImpl() {})();
|
||||
|
||||
Reference in New Issue
Block a user