import { storage } from "../storage"; import { withSuperAdmin } from "../db"; import { eventBus } from "../routes/shared"; import { pushTaskUpdated } from "../utils/pushTaskUpdated"; import { runAutomationsByTrigger } from "../routes/automation.routes"; import { logger } from "../utils/logger"; /** * Воркер отложенной смены статуса (автоматизации: ctx.tasks.scheduleStatusChange). * * Раз в 5 секунд подбирает просроченные строки delayed_status_changes и переводит * задачу в target_status_id, ТОЛЬКО если она всё ещё находится в expected_status_id — * ручной (или иной) переход задачи в другой статус отменяет отложенное действие. * Смена статуса сопровождается аудитом, SSE (task_updated), push и триггером * автоматизаций task.status_changed (как при обычном переходе). */ const log = logger("Delayed Status Change"); const TICK_MS = 5_000; async function processOne(row: { id: number; organizationId: number; taskId: number; formId: number; expectedStatusId: number; targetStatusId: number; }): Promise { await withSuperAdmin(async () => { const task = await storage.getTask(row.taskId, row.organizationId); if (!task) { log.info(`Задача ${row.taskId} не найдена — отмена отложенного перехода ${row.id}`); await storage.markDelayedStatusChangeExecuted(row.id); return; } if (task.currentStatusId !== row.expectedStatusId) { log.info( `Задача ${row.taskId} в статусе ${task.currentStatusId} (ожидался ${row.expectedStatusId}) — отмена отложенного перехода ${row.id}` ); await storage.markDelayedStatusChangeExecuted(row.id); return; } const statuses = await storage.getFormStatuses(task.formId, row.organizationId); const target = statuses.find(s => s.id === row.targetStatusId); if (!target) { log.error(`Статус ${row.targetStatusId} не найден в форме ${task.formId} — отмена перехода ${row.id}`); await storage.markDelayedStatusChangeExecuted(row.id); return; } const from = statuses.find(s => s.id === task.currentStatusId); const updated = await storage.updateTask(row.taskId, row.organizationId, { currentStatusId: row.targetStatusId, isCompleted: target.isFinal, completedAt: target.isFinal ? (task.completedAt ?? new Date()) : null, }); await storage.addTaskAuditLog({ taskId: row.taskId, organizationId: row.organizationId, action: 'status.changed', fieldName: 'Статус', oldValue: from?.name ?? String(task.currentStatusId), newValue: target.name, changedBy: null, changedByName: 'Автоматизация', metadata: { source: 'automation', delayed: true }, }); eventBus.publishEvent({ type: 'task_updated', data: { taskId: row.taskId, formId: task.formId, task: updated }, organizationId: row.organizationId, taskId: row.taskId, }); pushTaskUpdated({ taskId: row.taskId, organizationId: row.organizationId, formId: task.formId, taskTitle: updated.title, actorId: 0, }).catch(() => {}); await storage.markDelayedStatusChangeExecuted(row.id); // Автоматизации по триггеру task.status_changed (как при обычном переходе), // fire-and-forget — ошибки скриптов не ломают выполнение runAutomationsByTrigger(row.organizationId, 'task.status_changed', { formId: task.formId, taskId: row.taskId, oldStatusId: row.expectedStatusId, newStatusId: row.targetStatusId, task: updated, }).catch(err => log.error(`task.status_changed automations error (task ${row.taskId}):`, err)); log.info(`Задача ${row.taskId}: "${from?.name}" → "${target.name}" (отложенный переход ${row.id})`); }); } async function tick(): Promise { const due = await withSuperAdmin(() => storage.getDueDelayedStatusChanges()); for (const row of due) { try { await processOne(row); } catch (err) { // Не помечаем исполненным — следующий тик повторит попытку log.error(`Ошибка отложенного перехода ${row.id} (задача ${row.taskId}):`, err); } } } export function startDelayedStatusChangeWorker(): void { setInterval(() => { tick().catch((err) => log.error("Tick error:", err)); }, TICK_MS); log.info("started (5s interval)"); }