121 lines
4.8 KiB
TypeScript
121 lines
4.8 KiB
TypeScript
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<void> {
|
||
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<void> {
|
||
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)");
|
||
}
|