Add AI auto-check worker, share link message/link fields, and update journal/settings UI
This commit is contained in:
@@ -0,0 +1,218 @@
|
||||
const IDLE_MIN_MS = 2000;
|
||||
const IDLE_MAX_MS = 60000;
|
||||
const REQUEST_TIMEOUT_MS = 30000;
|
||||
const MAX_INPUT_CHARS = 2000;
|
||||
const MIN_TEXT_CHARS = 4;
|
||||
|
||||
function createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl, defaultPrompt, model }) {
|
||||
const AI_URL = aiUrl || process.env.AI_URL || 'http://text-corrector:8080';
|
||||
const MODEL = model || process.env.AI_MODEL || 'qwen2.5-1.5b-instruct';
|
||||
const DEFAULT_PROMPT = defaultPrompt || 'Ты — редактор текстов. Исправь ошибки. Верни ТОЛЬКО исправленный текст.';
|
||||
|
||||
let started = false;
|
||||
let stopped = false;
|
||||
let idleMs = IDLE_MIN_MS;
|
||||
let wake = null;
|
||||
let processing = false;
|
||||
let currentId = null;
|
||||
let startedAt = null;
|
||||
const stats = { checks: 0, corrected: 0, unchanged: 0, errors: 0, last_at: null, last_error: null };
|
||||
const attempts = new Map();
|
||||
const MAX_ATTEMPTS = 3;
|
||||
const CONFIG = {
|
||||
idle_min_ms: IDLE_MIN_MS,
|
||||
idle_max_ms: IDLE_MAX_MS,
|
||||
request_timeout_ms: REQUEST_TIMEOUT_MS,
|
||||
max_input_chars: MAX_INPUT_CHARS,
|
||||
min_text_chars: MIN_TEXT_CHARS,
|
||||
max_attempts: MAX_ATTEMPTS,
|
||||
model: MODEL,
|
||||
ai_url: AI_URL,
|
||||
};
|
||||
|
||||
function sleep(ms) {
|
||||
return new Promise((resolve) => {
|
||||
const t = setTimeout(() => { wake = null; resolve(); }, ms);
|
||||
wake = () => { clearTimeout(t); wake = null; resolve(); };
|
||||
});
|
||||
}
|
||||
|
||||
function notify() {
|
||||
if (wake) wake();
|
||||
}
|
||||
|
||||
async function isEnabled() {
|
||||
const v = await getSetting('ai_autocheck_enabled', 'true');
|
||||
return String(v) !== 'false';
|
||||
}
|
||||
|
||||
async function claimNext() {
|
||||
const client = await pool.connect();
|
||||
try {
|
||||
await client.query('BEGIN');
|
||||
const { rows } = await client.query(
|
||||
`SELECT id, description, description_original FROM entries
|
||||
WHERE ai_status = 'pending' AND deleted_at IS NULL
|
||||
ORDER BY id ASC LIMIT 1 FOR UPDATE SKIP LOCKED`
|
||||
);
|
||||
if (!rows.length) {
|
||||
await client.query('COMMIT');
|
||||
return null;
|
||||
}
|
||||
await client.query(`UPDATE entries SET ai_status = 'processing' WHERE id = $1`, [rows[0].id]);
|
||||
await client.query('COMMIT');
|
||||
return rows[0];
|
||||
} catch (e) {
|
||||
await client.query('ROLLBACK').catch(() => {});
|
||||
throw e;
|
||||
} finally {
|
||||
client.release();
|
||||
}
|
||||
}
|
||||
|
||||
async function callModel(text, systemPrompt) {
|
||||
const controller = new AbortController();
|
||||
const timer = setTimeout(() => controller.abort(), REQUEST_TIMEOUT_MS);
|
||||
try {
|
||||
const res = await fetch(`${AI_URL}/v1/chat/completions`, {
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
signal: controller.signal,
|
||||
body: JSON.stringify({
|
||||
model: MODEL,
|
||||
messages: [
|
||||
{ role: 'system', content: systemPrompt },
|
||||
{ role: 'user', content: text },
|
||||
],
|
||||
temperature: 0.1,
|
||||
max_tokens: Math.min(512, Math.max(64, text.length + 32)),
|
||||
}),
|
||||
});
|
||||
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(timer);
|
||||
}
|
||||
}
|
||||
|
||||
function sanitize(original, out) {
|
||||
if (!out) return original;
|
||||
const o = out.trim();
|
||||
if (!o) return original;
|
||||
if (o.length > original.length * 4 + 80) return original;
|
||||
return o;
|
||||
}
|
||||
|
||||
async function processOne(row) {
|
||||
const original = String(row.description_original ?? row.description ?? '').trim();
|
||||
if (original.length < MIN_TEXT_CHARS) {
|
||||
await pool.query(`UPDATE entries SET ai_status = 'skipped', ai_checked_at = now() WHERE id = $1`, [row.id]);
|
||||
return true;
|
||||
}
|
||||
const prompt = (await getSetting('ai_prompt', DEFAULT_PROMPT)) || DEFAULT_PROMPT;
|
||||
const input = original.length > MAX_INPUT_CHARS ? original.slice(0, MAX_INPUT_CHARS) : original;
|
||||
try {
|
||||
const out = await callModel(input, prompt);
|
||||
const checked = sanitize(input, out);
|
||||
const changed = checked !== input;
|
||||
await pool.query(
|
||||
`UPDATE entries SET description_ai = $1, description = $1, ai_status = 'done', ai_error = NULL, ai_checked_at = now() WHERE id = $2`,
|
||||
[changed ? checked : original, row.id]
|
||||
);
|
||||
attempts.delete(row.id);
|
||||
stats.checks++;
|
||||
if (changed) stats.corrected++; else stats.unchanged++;
|
||||
stats.last_at = new Date().toISOString();
|
||||
stats.last_error = null;
|
||||
if (logAudit) await logAudit(null, 'entry.ai.auto-check', { id: row.id, changed });
|
||||
return true;
|
||||
} catch (e) {
|
||||
const message = (e && e.message ? e.message : 'error').slice(0, 500);
|
||||
const tries = (attempts.get(row.id) || 0) + 1;
|
||||
stats.last_at = new Date().toISOString();
|
||||
stats.last_error = message;
|
||||
if (tries < MAX_ATTEMPTS) {
|
||||
attempts.set(row.id, tries);
|
||||
await pool.query(`UPDATE entries SET ai_status = 'pending', ai_error = $1 WHERE id = $2`, [message, row.id]);
|
||||
if (logAudit) await logAudit(null, 'entry.ai.retry', { id: row.id, attempt: tries, error: message });
|
||||
return false;
|
||||
}
|
||||
attempts.delete(row.id);
|
||||
stats.errors++;
|
||||
await pool.query(
|
||||
`UPDATE entries SET ai_status = 'error', ai_error = $1, ai_checked_at = now() WHERE id = $2`,
|
||||
[message, row.id]
|
||||
);
|
||||
if (logAudit) await logAudit(null, 'entry.ai.error', { id: row.id, error: message });
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
async function loop() {
|
||||
while (!stopped) {
|
||||
let row = null;
|
||||
try {
|
||||
if (!(await isEnabled())) {
|
||||
await sleep(IDLE_MAX_MS);
|
||||
continue;
|
||||
}
|
||||
row = await claimNext();
|
||||
} catch (e) {
|
||||
console.error('AI auto-check claim error:', e);
|
||||
await sleep(IDLE_MAX_MS);
|
||||
continue;
|
||||
}
|
||||
if (!row) {
|
||||
await sleep(idleMs);
|
||||
idleMs = Math.min(idleMs * 2, IDLE_MAX_MS);
|
||||
continue;
|
||||
}
|
||||
idleMs = IDLE_MIN_MS;
|
||||
processing = true;
|
||||
currentId = row.id;
|
||||
let ok = true;
|
||||
try {
|
||||
ok = await processOne(row);
|
||||
} catch (e) {
|
||||
console.error('AI auto-check process error:', e);
|
||||
} finally {
|
||||
processing = false;
|
||||
currentId = null;
|
||||
}
|
||||
if (!ok) await sleep(IDLE_MAX_MS);
|
||||
}
|
||||
}
|
||||
|
||||
async function resetStale() {
|
||||
await pool.query(`UPDATE entries SET ai_status = 'pending' WHERE ai_status = 'processing'`);
|
||||
}
|
||||
|
||||
function start() {
|
||||
if (started) return;
|
||||
started = true;
|
||||
startedAt = Date.now();
|
||||
resetStale()
|
||||
.catch((e) => console.error('AI auto-check reset error:', e))
|
||||
.finally(() => { loop().catch((e) => console.error('AI auto-check loop error:', e)); });
|
||||
}
|
||||
|
||||
function getStats() {
|
||||
return { ...stats, processing };
|
||||
}
|
||||
|
||||
function getInfo() {
|
||||
return {
|
||||
started_at: startedAt ? new Date(startedAt).toISOString() : null,
|
||||
uptime_ms: startedAt ? Date.now() - startedAt : 0,
|
||||
processing,
|
||||
current_id: currentId,
|
||||
config: CONFIG,
|
||||
...stats,
|
||||
};
|
||||
}
|
||||
|
||||
return { start, notify, getStats, getInfo };
|
||||
}
|
||||
|
||||
module.exports = { createEntryAutoChecker };
|
||||
Reference in New Issue
Block a user