feat(journal): live AI status updates in entry cards via Postgres NOTIFY trigger + SSE ai_status events
This commit is contained in:
+15
@@ -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 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 (
|
CREATE TABLE IF NOT EXISTS settings (
|
||||||
key TEXT PRIMARY KEY,
|
key TEXT PRIMARY KEY,
|
||||||
value TEXT
|
value TEXT
|
||||||
|
|||||||
@@ -184,3 +184,18 @@ FROM (
|
|||||||
FROM entry_photos ORDER BY entry_id, sort_order, id
|
FROM entry_photos ORDER BY entry_id, sort_order, id
|
||||||
) p
|
) p
|
||||||
WHERE e.photo_path IS NULL AND p.entry_id = e.id;
|
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();
|
||||||
|
|||||||
+32
-4
@@ -66,7 +66,7 @@ function aiBadge(e) {
|
|||||||
const st = e.ai_status;
|
const st = e.ai_status;
|
||||||
const label = AI_STATUS_LABELS[st] || `ИИ: ${st}`;
|
const label = AI_STATUS_LABELS[st] || `ИИ: ${st}`;
|
||||||
const icon = AI_STATUS_ICONS[st] || 'sparkles';
|
const icon = AI_STATUS_ICONS[st] || 'sparkles';
|
||||||
return `<span class="ai-badge ai-badge-ic ai-badge-${esc(st)}" title="${esc(label)}" role="img" aria-label="${esc(label)}"><i data-lucide="${icon}"></i></span>`;
|
return `<span class="ai-badge ai-badge-ic ai-badge-${esc(st)}" data-ai-badge="${e.id}" title="${esc(label)}" role="img" aria-label="${esc(label)}"><i data-lucide="${icon}"></i></span>`;
|
||||||
}
|
}
|
||||||
|
|
||||||
function findFile(id) {
|
function findFile(id) {
|
||||||
@@ -236,6 +236,27 @@ function goPage(p) {
|
|||||||
|
|
||||||
let liveSource = null;
|
let liveSource = null;
|
||||||
let liveReconnectTimer = 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() {
|
function connectLiveUpdates() {
|
||||||
if (liveSource) { liveSource.close(); liveSource = null; }
|
if (liveSource) { liveSource.close(); liveSource = null; }
|
||||||
const connect = () => {
|
const connect = () => {
|
||||||
@@ -243,9 +264,16 @@ function connectLiveUpdates() {
|
|||||||
if (!t) return;
|
if (!t) return;
|
||||||
liveSource = new EventSource(`${API}/api/events?token=${encodeURIComponent(t)}`);
|
liveSource = new EventSource(`${API}/api/events?token=${encodeURIComponent(t)}`);
|
||||||
liveSource.addEventListener('entries_changed', () => {
|
liveSource.addEventListener('entries_changed', () => {
|
||||||
showToast('Поступила новая запись — список обновлён');
|
if (entriesRefreshTimer) return;
|
||||||
page = 1;
|
entriesRefreshTimer = setTimeout(() => {
|
||||||
loadEntries();
|
entriesRefreshTimer = null;
|
||||||
|
showToast('Поступила новая запись — список обновлён');
|
||||||
|
page = 1;
|
||||||
|
loadEntries();
|
||||||
|
}, 300);
|
||||||
|
});
|
||||||
|
liveSource.addEventListener('ai_status', ev => {
|
||||||
|
try { applyAiUpdate(JSON.parse(ev.data)); } catch (err) {}
|
||||||
});
|
});
|
||||||
liveSource.onerror = () => {
|
liveSource.onerror = () => {
|
||||||
if (liveSource) { liveSource.close(); liveSource = null; }
|
if (liveSource) { liveSource.close(); liveSource = null; }
|
||||||
|
|||||||
@@ -17,6 +17,39 @@ types.setTypeParser(1082, v => v);
|
|||||||
const app = express();
|
const app = express();
|
||||||
const pool = new Pool({ connectionString: process.env.DATABASE_URL });
|
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 cacheStore = new Map();
|
||||||
const SETTINGS_TTL_MS = 30 * 1000;
|
const SETTINGS_TTL_MS = 30 * 1000;
|
||||||
const PUBLIC_TTL_MS = 60 * 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) => {
|
app.get('/api/events', async (req, res) => {
|
||||||
try {
|
try {
|
||||||
const token = req.headers['x-auth-token'] || req.query.token;
|
const token = req.headers['x-auth-token'] || req.query.token;
|
||||||
|
|||||||
Reference in New Issue
Block a user