diff --git a/db/init.sql b/db/init.sql index 4084b72..9eebcb7 100644 --- a/db/init.sql +++ b/db/init.sql @@ -54,6 +54,21 @@ ALTER TABLE entries ADD COLUMN IF NOT EXISTS ai_error TEXT; CREATE INDEX IF NOT EXISTS idx_entries_ai_pending ON entries(id) WHERE ai_status = 'pending' AND deleted_at IS NULL; +CREATE OR REPLACE FUNCTION notify_entries_changed() RETURNS trigger AS $$ +BEGIN + IF (TG_OP = 'INSERT') THEN + PERFORM pg_notify('entries_changed', json_build_object('type', 'entry_created', 'id', NEW.id)::text); + ELSIF (TG_OP = 'UPDATE' AND OLD.ai_status IS DISTINCT FROM NEW.ai_status) THEN + PERFORM pg_notify('entries_changed', json_build_object('type', 'ai_status', 'id', NEW.id, 'status', NEW.ai_status, 'error', NEW.ai_error)::text); + END IF; + RETURN NULL; +END; +$$ LANGUAGE plpgsql; + +DROP TRIGGER IF EXISTS trg_entries_notify ON entries; +CREATE TRIGGER trg_entries_notify AFTER INSERT OR UPDATE OF ai_status ON entries +FOR EACH ROW EXECUTE FUNCTION notify_entries_changed(); + CREATE TABLE IF NOT EXISTS settings ( key TEXT PRIMARY KEY, value TEXT diff --git a/db/migration.sql b/db/migration.sql index 2550774..112b0e6 100644 --- a/db/migration.sql +++ b/db/migration.sql @@ -184,3 +184,18 @@ FROM ( FROM entry_photos ORDER BY entry_id, sort_order, id ) p WHERE e.photo_path IS NULL AND p.entry_id = e.id; + +CREATE OR REPLACE FUNCTION notify_entries_changed() RETURNS trigger AS $$ +BEGIN + IF (TG_OP = 'INSERT') THEN + PERFORM pg_notify('entries_changed', json_build_object('type', 'entry_created', 'id', NEW.id)::text); + ELSIF (TG_OP = 'UPDATE' AND OLD.ai_status IS DISTINCT FROM NEW.ai_status) THEN + PERFORM pg_notify('entries_changed', json_build_object('type', 'ai_status', 'id', NEW.id, 'status', NEW.ai_status, 'error', NEW.ai_error)::text); + END IF; + RETURN NULL; +END; +$$ LANGUAGE plpgsql; + +DROP TRIGGER IF EXISTS trg_entries_notify ON entries; +CREATE TRIGGER trg_entries_notify AFTER INSERT OR UPDATE OF ai_status ON entries +FOR EACH ROW EXECUTE FUNCTION notify_entries_changed(); diff --git a/public/js/journal.js b/public/js/journal.js index ed19784..69c82f0 100644 --- a/public/js/journal.js +++ b/public/js/journal.js @@ -66,7 +66,7 @@ function aiBadge(e) { const st = e.ai_status; const label = AI_STATUS_LABELS[st] || `ИИ: ${st}`; const icon = AI_STATUS_ICONS[st] || 'sparkles'; - return ``; + return ``; } function findFile(id) { @@ -236,6 +236,27 @@ function goPage(p) { let liveSource = null; let liveReconnectTimer = null; +let aiRefreshTimer = null; +let entriesRefreshTimer = null; +function applyAiUpdate(update) { + const e = currentEntries.find(x => x.id === update.id); + if (!e) return; + e.ai_status = update.ai_status; + e.ai_error = update.ai_error || null; + if (update.ai_status === 'done') { + if (aiRefreshTimer) return; + aiRefreshTimer = setTimeout(() => { aiRefreshTimer = null; loadEntries(); }, 300); + return; + } + document.querySelectorAll(`[data-ai-badge="${update.id}"]`).forEach(badgeEl => { + badgeEl.outerHTML = aiBadge(e); + }); + renderIcons(); + const editBadge = document.getElementById('editAiStatus'); + if (editId === update.id && editBadge && editBadge.style.display !== 'none') { + renderEditAiStatus(e); + } +} function connectLiveUpdates() { if (liveSource) { liveSource.close(); liveSource = null; } const connect = () => { @@ -243,9 +264,16 @@ function connectLiveUpdates() { if (!t) return; liveSource = new EventSource(`${API}/api/events?token=${encodeURIComponent(t)}`); liveSource.addEventListener('entries_changed', () => { - showToast('Поступила новая запись — список обновлён'); - page = 1; - loadEntries(); + if (entriesRefreshTimer) return; + entriesRefreshTimer = setTimeout(() => { + entriesRefreshTimer = null; + showToast('Поступила новая запись — список обновлён'); + page = 1; + loadEntries(); + }, 300); + }); + liveSource.addEventListener('ai_status', ev => { + try { applyAiUpdate(JSON.parse(ev.data)); } catch (err) {} }); liveSource.onerror = () => { if (liveSource) { liveSource.close(); liveSource = null; } diff --git a/server.js b/server.js index 04563f9..c68c249 100644 --- a/server.js +++ b/server.js @@ -17,6 +17,39 @@ types.setTypeParser(1082, v => v); const app = express(); const pool = new Pool({ connectionString: process.env.DATABASE_URL }); +const pgClient = require('pg').Client; +const lister = new pgClient({ connectionString: process.env.DATABASE_URL }); +let listerConnected = false; +function connectLister() { + if (listerConnected) return; + listerConnected = true; + lister.connect() + .then(() => lister.query('LISTEN entries_changed')) + .catch(e => { + listerConnected = false; + console.error('LISTEN entries_changed failed:', e.message); + setTimeout(connectLister, 5000); + }); +} +connectLister(); +lister.on('error', (e) => { + console.error('LISTEN connection error:', e.message); +}); +lister.on('end', () => { + listerConnected = false; + setTimeout(connectLister, 3000); +}); +lister.on('notification', (msg) => { + let payload = null; + try { payload = JSON.parse(msg.payload || '{}'); } catch (e) { payload = null; } + const type = payload && payload.type ? payload.type : 'entry_created'; + if (type === 'ai_status') { + broadcastAiStatus(payload.id, payload.status, payload.error); + } else { + broadcastEntryChanged(); + } +}); + const cacheStore = new Map(); const SETTINGS_TTL_MS = 30 * 1000; const PUBLIC_TTL_MS = 60 * 1000; @@ -72,6 +105,13 @@ function broadcastEntryChanged() { } } +function broadcastAiStatus(entryId, status, error) { + const frame = `event: ai_status\ndata: ${JSON.stringify({ id: entryId, ai_status: status, ai_error: error || null, ts: Date.now() })}\n\n`; + for (const client of sseClients) { + try { client.write(frame); } catch (e) { sseClients.delete(client); } + } +} + app.get('/api/events', async (req, res) => { try { const token = req.headers['x-auth-token'] || req.query.token;