feat(automations): ctx.tasks.scheduleStatusChange — отложенная смена статуса (воркер 5с, отмена при ручном переходе)
This commit is contained in:
120
server/workers/delayed-status-change.worker.ts
Normal file
120
server/workers/delayed-status-change.worker.ts
Normal file
@@ -0,0 +1,120 @@
|
||||
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)");
|
||||
}
|
||||
Reference in New Issue
Block a user