diff --git a/client/src/pages/Automations.tsx b/client/src/pages/Automations.tsx index 7d7b45e..4511214 100644 --- a/client/src/pages/Automations.tsx +++ b/client/src/pages/Automations.tsx @@ -78,6 +78,8 @@ const CODE_TEMPLATE = `// ====================================================== // ctx.tasks.setFieldValue(taskId, fieldCode, value) — установить значение поля // ctx.tasks.sendMessage(taskId, message) — сообщение в чат задачи от «Система» // (messageType 'status_change', бейдж «Статус», без push-уведомлений) +// ctx.tasks.scheduleStatusChange(taskId, targetStatusId, delayMs) — отложенный +// переход в статус через delayMs; отменяется, если задача раньше сменит статус // // Работа с пользователями: // ctx.users.list() — список пользователей организации @@ -329,6 +331,8 @@ export function AutomationsContent() { ctx.tasks.setFieldValue(taskId, fieldCode, value) ctx.tasks.sendMessage(taskId, message) — сообщение в чат задачи от «Система» (messageType 'status_change', бейдж «Статус», без push-уведомлений) + ctx.tasks.scheduleStatusChange(taskId, targetStatusId, delayMs) — отложенный + переход в статус через delayMs; отменяется, если задача раньше сменит статус Пользователи: ctx.users.list() / ctx.users.get(userId) diff --git a/migrations/0083_delayed_status_changes.sql b/migrations/0083_delayed_status_changes.sql new file mode 100644 index 0000000..4cafde8 --- /dev/null +++ b/migrations/0083_delayed_status_changes.sql @@ -0,0 +1,18 @@ +-- Отложенная смена статуса задачи (автоматизации, ctx.tasks.scheduleStatusChange). +-- Воркер каждые 5 секунд подбирает просроченные строки и переводит задачу +-- в target_status_id, только если она всё ещё в expected_status_id +-- (ручной переход в другой статус отменяет отложенное действие). +CREATE TABLE IF NOT EXISTS delayed_status_changes ( + id SERIAL PRIMARY KEY, + organization_id INTEGER NOT NULL REFERENCES organizations(id) ON DELETE CASCADE, + task_id INTEGER NOT NULL REFERENCES tasks(id) ON DELETE CASCADE, + form_id INTEGER NOT NULL, + expected_status_id INTEGER NOT NULL, + target_status_id INTEGER NOT NULL, + run_at TIMESTAMPTZ NOT NULL, + executed_at TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE INDEX IF NOT EXISTS delayed_status_changes_run_at_idx + ON delayed_status_changes (run_at) WHERE executed_at IS NULL; diff --git a/server/index.ts b/server/index.ts index a567301..bea9b43 100644 --- a/server/index.ts +++ b/server/index.ts @@ -8,6 +8,7 @@ import { notificationService } from "./services/notification.service"; import { webPushService } from "./services/web-push.service"; import { startMedScheduleWorker } from "./medschedule/worker"; import { startAutomationScheduler } from "./workers/automation-scheduler"; +import { startDelayedStatusChangeWorker } from "./workers/delayed-status-change.worker"; import { startErrorLogsRetention } from "./workers/error-logs-retention"; import { startGpsWorker } from "./gps/worker"; import { storage } from "./storage"; @@ -1142,6 +1143,10 @@ async function runStartupDataPatches() { startAutomationScheduler(); log('Automation scheduler started'); + // Start delayed status change worker (ctx.tasks.scheduleStatusChange, 5s interval) + startDelayedStatusChangeWorker(); + log('Delayed status change worker started'); + // Start error_logs retention (daily cleanup of API error log entries older than 30 days) startErrorLogsRetention(); log('Error logs retention started'); diff --git a/server/mcp.ts b/server/mcp.ts index 582794f..1e56ab7 100644 --- a/server/mcp.ts +++ b/server/mcp.ts @@ -2434,6 +2434,7 @@ Tasks: - ctx.tasks.getAssignees(taskId) — get task assignees - ctx.tasks.setFieldValue(taskId, fieldCode, value) — set a field value - ctx.tasks.sendMessage(taskId, message) — post a message to the task chat as 'Система' (messageType 'status_change', 'Статус' badge, no push notifications) +- ctx.tasks.scheduleStatusChange(taskId, targetStatusId, delayMs) — schedule a delayed status change: a worker moves the task to targetStatusId after delayMs, but only if the task is still in its current status (a manual status change cancels the delayed one). Audit, SSE, push and task.status_changed automations fire as usual. Users: - ctx.users.list() — list organization users diff --git a/server/routes/automation.routes.ts b/server/routes/automation.routes.ts index 7dfec76..06a64ae 100644 --- a/server/routes/automation.routes.ts +++ b/server/routes/automation.routes.ts @@ -32,8 +32,9 @@ const ALLOWED_STORAGE_METHODS = new Set([ 'updateUser', 'setUserProfileFieldValueFromAutomation', 'recalcUserFieldAverageFromAutomation', - // Не storage-метод: обрабатывается отдельно в обработчике storage-call (сервис сообщений) + // Не storage-методы: обрабатываются отдельно в обработчике storage-call 'sendTaskMessageFromAutomation', + 'createDelayedStatusChangeFromAutomation', ]); interface WorkerStorageCall { @@ -99,6 +100,39 @@ export function runAutomationInWorker( return; } + // Отложенная смена статуса из автоматизации: планирует переход воркером + // delayed-status-change.worker (отменяется, если задача сменит статус раньше) + if (method === 'createDelayedStatusChangeFromAutomation') { + try { + const [taskId, orgId, targetStatusId, delayMs] = args as [number, number, number, number]; + const task = await storage.getTask(Number(taskId), Number(orgId)); + if (!task) throw new Error(`Задача ${taskId} не найдена`); + const statuses = await storage.getFormStatuses(task.formId, Number(orgId)); + if (!statuses.some(s => s.id === Number(targetStatusId))) { + throw new Error(`Статус ${targetStatusId} не принадлежит форме ${task.formId}`); + } + const delay = Number(delayMs); + if (!Number.isFinite(delay) || delay < 0) throw new Error('delayMs должен быть неотрицательным числом'); + const created = await storage.createDelayedStatusChange({ + organizationId: Number(orgId), + taskId: task.id, + formId: task.formId, + expectedStatusId: task.currentStatusId, + targetStatusId: Number(targetStatusId), + runAt: new Date(Date.now() + delay), + }); + if (!settled) { + try { worker.postMessage({ type: 'storage-response', callId, result: created }); } catch { /* terminated */ } + } + } catch (err: unknown) { + const message = err instanceof Error ? err.message : String(err); + if (!settled) { + try { worker.postMessage({ type: 'storage-response', callId, error: message }); } catch { /* terminated */ } + } + } + return; + } + // Сообщение в чат задачи из автоматизации: не storage-метод, // идём через сервис сообщений (SSE, вложения, messageType 'status_change') if (method === 'sendTaskMessageFromAutomation') { diff --git a/server/storage/task-meta.storage.ts b/server/storage/task-meta.storage.ts index 8a15702..a416ddd 100644 --- a/server/storage/task-meta.storage.ts +++ b/server/storage/task-meta.storage.ts @@ -1,6 +1,7 @@ import { users, organizations, forms, tasks, safeUserColumns, type User, type SafeUser, type Organization, type Task } from "@shared/schema"; import { taskRelations, type TaskRelation } from "@shared/schema"; import { taskReminders, type TaskReminder, type InsertTaskReminder } from "@shared/schema"; +import { delayedStatusChanges, type DelayedStatusChange, type InsertDelayedStatusChange } from "@shared/schema"; import { taskAuditLog, type TaskAuditLog, type InsertTaskAuditLog } from "@shared/schema"; import { roleMembers, roles } from "@shared/schema"; import { db } from "../db"; @@ -96,6 +97,31 @@ export class TaskMetaStorage extends ContentStorage { .where(eq(taskReminders.id, id)); } + // Отложенная смена статуса (автоматизации: ctx.tasks.scheduleStatusChange) + async createDelayedStatusChange(data: InsertDelayedStatusChange): Promise { + const [row] = await db.insert(delayedStatusChanges).values(data).returning(); + return row; + } + + async getDueDelayedStatusChanges(limit = 50): Promise { + return db + .select() + .from(delayedStatusChanges) + .where(and( + sql`${delayedStatusChanges.executedAt} IS NULL`, + lte(delayedStatusChanges.runAt, new Date()), + )) + .orderBy(asc(delayedStatusChanges.runAt)) + .limit(limit); + } + + async markDelayedStatusChangeExecuted(id: number): Promise { + await db + .update(delayedStatusChanges) + .set({ executedAt: new Date() }) + .where(eq(delayedStatusChanges.id, id)); + } + // Массовые выборки пользователей — только безопасные колонки (без хэша пароля и токенов) async getUsersByRole(role: string, organizationId: number): Promise { return db diff --git a/server/workers/automation-worker-code.ts b/server/workers/automation-worker-code.ts index 2a0522a..4dba015 100644 --- a/server/workers/automation-worker-code.ts +++ b/server/workers/automation-worker-code.ts @@ -108,6 +108,11 @@ const ctx = { sendMessage: function(taskId, message) { return callStorage('sendTaskMessageFromAutomation', [taskId, organizationId, String(message)]); }, + // Отложенная смена статуса: воркер переведёт задачу в targetStatusId через delayMs, + // если она к тому моменту останется в текущем статусе + scheduleStatusChange: function(taskId, targetStatusId, delayMs) { + return callStorage('createDelayedStatusChangeFromAutomation', [taskId, organizationId, Number(targetStatusId), Number(delayMs)]); + }, }, }; diff --git a/server/workers/delayed-status-change.worker.ts b/server/workers/delayed-status-change.worker.ts new file mode 100644 index 0000000..95a657e --- /dev/null +++ b/server/workers/delayed-status-change.worker.ts @@ -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 { + 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)"); +} diff --git a/shared/schema.ts b/shared/schema.ts index 8cd1728..864fedb 100644 --- a/shared/schema.ts +++ b/shared/schema.ts @@ -2816,6 +2816,28 @@ export type TaskReminder = typeof taskReminders.$inferSelect; export type InsertTaskReminder = z.infer; export type ReminderRecipient = { type: 'user'; userId: number } | { type: 'role'; role: string } | { type: 'orgRole'; roleId: number }; +// Отложенная смена статуса задачи (автоматизации: ctx.tasks.scheduleStatusChange). +// Воркер server/workers/delayed-status-change.worker.ts раз в 5 секунд +// подбирает просроченные строки и переводит задачу в target_status_id, +// только если она всё ещё в expected_status_id (иначе — отмена). +export const delayedStatusChanges = pgTable("delayed_status_changes", { + id: serial("id").primaryKey(), + organizationId: integer("organization_id").notNull().references(() => organizations.id, { onDelete: "cascade" }), + taskId: integer("task_id").notNull().references(() => tasks.id, { onDelete: "cascade" }), + formId: integer("form_id").notNull(), + expectedStatusId: integer("expected_status_id").notNull(), + targetStatusId: integer("target_status_id").notNull(), + runAt: timestamp("run_at").notNull(), + executedAt: timestamp("executed_at"), + createdAt: timestamp("created_at").defaultNow(), +}, (table) => ({ + // Частичный индекс для воркера: ищет неисполненные, отсортированные по времени + runAtIdx: index("delayed_status_changes_run_at_idx").on(table.runAt).where(sql`${table.executedAt} IS NULL`), +})); + +export type DelayedStatusChange = typeof delayedStatusChanges.$inferSelect; +export type InsertDelayedStatusChange = typeof delayedStatusChanges.$inferInsert; + // Task audit log table (history of changes) export const taskAuditLog = pgTable("task_audit_log", { id: serial("id").primaryKey(),