'use strict'; /** * Gitea Workflow 调度器:同仓库串行、跨仓库并行。 * * 调度器只编排领域任务,不负责 clone、Codex 或 Gitea REST;这些通过 * `runner(task, context)` 注入,runner 返回 `{ state, ...patch }` 即可更新任务。 */ const { EventEmitter } = require('events'); const domain = require('./gitea-workflow-domain'); function sleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)); } function taskOrder(task) { const commentCreatedAt = task?.metadata?.event?.comment?.createdAt || task?.metadata?.event?.comment?.created_at || task?.metadata?.event?.receivedAt; const timestamp = Date.parse(commentCreatedAt || task?.createdAt || '') || Number.MAX_SAFE_INTEGER; return `${String(timestamp).padStart(16, '0')}:${String(task?.deliveryId || task?.taskId || '')}`; } class GiteaWorkflowQueue extends EventEmitter { constructor(options = {}) { super(); if (!options.store) throw new TypeError('GiteaWorkflowQueue 需要注入 store'); this.store = options.store; this.clock = typeof options.now === 'function' ? options.now : Date.now; this.defaultRunner = typeof options.runner === 'function' ? options.runner : async () => ({ state: domain.TASK_STATES.SUCCEEDED }); this.pending = []; this.active = new Map(); this.repoActive = new Map(); this.runners = new Map(); this.paused = Boolean(options.paused || this.store.getControl?.().paused); this.draining = false; this.started = false; if (options.autoRecover !== false) this.recover(); } recover() { const tasks = this.store.recover({ maxRestartRetries: 1, now: this.clock }); this.pending = tasks.filter((task) => task.state === domain.TASK_STATES.QUEUED || task.state === domain.TASK_STATES.RETRY_WAIT) .sort((a, b) => taskOrder(a).localeCompare(taskOrder(b))) .map((task) => task.taskId); this.started = true; this.drain(); return tasks; } setRunner(taskId, runner) { if (typeof runner === 'function') this.runners.set(String(taskId), runner); return this; } enqueue(taskOrId, runner) { const taskId = typeof taskOrId === 'string' ? taskOrId : taskOrId?.taskId; if (!taskId) return Promise.reject(new TypeError('enqueue 缺少 taskId')); let task = this.store.getTask(taskId); if (!task && typeof taskOrId === 'object') task = this.store.createTask(taskOrId); if (!task) return Promise.reject(Object.assign(new Error('找不到任务'), { code: 'task_not_found' })); if (domain.TERMINAL_STATES.has(task.state)) return Promise.resolve({ task: this.store.getTask(taskId), skipped: true }); if (task.state !== domain.TASK_STATES.QUEUED) { task = this.store.transitionTask(taskId, domain.TASK_STATES.QUEUED); } if (!this.pending.includes(taskId) && !this.active.has(taskId)) this.pending.push(taskId); this.pending.sort((leftId, rightId) => taskOrder(this.store.getTask(leftId)).localeCompare(taskOrder(this.store.getTask(rightId)))); if (typeof runner === 'function') this.runners.set(taskId, runner); this.emit('queued', this.store.getTask(taskId)); this.drain(); return this.waitForTerminal(taskId); } waitForTerminal(taskId, timeoutMs = 0) { const existing = this.store.getTask(taskId); if (existing && (domain.TERMINAL_STATES.has(existing.state) || existing.state === domain.TASK_STATES.WAITING_USER)) { return Promise.resolve({ task: existing }); } return new Promise((resolve, reject) => { let timer = null; const check = (task) => { if (!task || task.taskId !== taskId) return; if (domain.TERMINAL_STATES.has(task.state) || task.state === domain.TASK_STATES.WAITING_USER) { cleanup(); resolve({ task }); } }; const cleanup = () => { this.off('updated', check); if (timer) clearTimeout(timer); }; this.on('updated', check); if (timeoutMs > 0) timer = setTimeout(() => { cleanup(); reject(new Error('等待任务完成超时')); }, timeoutMs); check(this.store.getTask(taskId)); }); } pause(input = {}) { this.paused = true; this.store.setPaused(true, input); this.emit('paused', this.store.getControl()); return this.store.getControl(); } resume(input = {}) { this.paused = false; this.store.setPaused(false, input); this.emit('resumed', this.store.getControl()); this.drain(); return this.store.getControl(); } cancelQueued(taskId, input = {}) { const task = this.store.getTask(taskId); if (!task) return { ok: false, code: 'task_not_found' }; if (task.state !== domain.TASK_STATES.QUEUED && task.state !== domain.TASK_STATES.RETRY_WAIT && task.state !== domain.TASK_STATES.BLOCKED_WORKSPACE && task.state !== domain.TASK_STATES.WAITING_USER) { return { ok: false, code: 'task_not_queued', task }; } this.pending = this.pending.filter((id) => id !== String(taskId)); const updated = this.store.transitionTask(taskId, domain.TASK_STATES.CANCELLED, { errorCode: input.reason || 'cancelled_by_operator' }); this.store.appendAudit({ taskId, sessionKey: task.sessionKey, repoKey: task.repoKey, actor: input.actor || 'admin', action: 'task.cancelled', fromState: task.state, toState: updated.state, metadata: { reason: input.reason || null } }); this.emit('updated', updated); return { ok: true, task: updated }; } abortRunning(taskId, input = {}) { const entry = this.active.get(String(taskId)); const task = this.store.getTask(taskId); if (!task) return { ok: false, code: 'task_not_found' }; if (!entry || !domain.RUNNING_STATES.has(task.state)) { if (task.state === domain.TASK_STATES.ABORTED || task.state === domain.TASK_STATES.CANCELLED) return { ok: true, task }; return { ok: false, code: 'task_not_running', task }; } const updated = this.store.transitionTask(taskId, domain.TASK_STATES.ABORTING, { errorCode: input.reason || 'abort_requested' }); entry.abortRequested = true; try { entry.controller.abort(new Error(input.reason || '管理员请求中止')); } catch {} this.store.appendAudit({ taskId, sessionKey: task.sessionKey, repoKey: task.repoKey, actor: input.actor || 'admin', action: 'task.abort_requested', fromState: task.state, toState: updated.state, turnId: task.turnId, metadata: { reason: input.reason || null } }); this.emit('updated', updated); return { ok: true, task: updated }; } getStatus() { return { paused: this.paused, running: this.active.size, queued: this.pending.length, activeTaskIds: [...this.active.keys()], }; } async drain() { if (this.draining || this.paused) return; this.draining = true; try { while (!this.paused) { const index = this.pending.findIndex((id) => { const task = this.store.getTask(id); return task && (task.state === domain.TASK_STATES.QUEUED || task.state === domain.TASK_STATES.RETRY_WAIT) && !this.repoActive.has(task.repoKey || task.taskId) && (!task.nextRetryAt || new Date(task.nextRetryAt).getTime() <= this.clock()); }); if (index < 0) break; const taskId = this.pending.splice(index, 1)[0]; this.run(taskId).catch(() => undefined); } } finally { this.draining = false; } } async run(taskId) { const current = this.store.getTask(taskId); if (!current) return; const repoLockKey = current.repoKey || current.taskId; this.repoActive.set(repoLockKey, taskId); const controller = new AbortController(); this.active.set(taskId, { controller, abortRequested: false, startedAt: domain.iso(this.clock) }); let task = this.store.transitionTask(taskId, domain.TASK_STATES.PREPARING); task = this.store.upsertTask({ ...task, attempt: Number(task.attempt || 0) + 1 }); this.emit('updated', task); this.store.appendAudit({ taskId, sessionKey: task.sessionKey, repoKey: task.repoKey, action: 'task.preparing', fromState: current.state, toState: task.state, deliveryId: task.deliveryId }); try { task = this.store.transitionTask(taskId, domain.TASK_STATES.RUNNING); this.emit('updated', task); this.store.appendAudit({ taskId, sessionKey: task.sessionKey, repoKey: task.repoKey, action: 'task.running', fromState: domain.TASK_STATES.PREPARING, toState: task.state, turnId: task.turnId }); const runner = this.runners.get(taskId) || this.defaultRunner; const result = await runner(this.store.getTask(taskId), { signal: controller.signal, queue: this, status: this.getStatus() }); const patch = result && typeof result === 'object' ? result : {}; let desired = patch.state || (controller.signal.aborted ? domain.TASK_STATES.ABORTED : domain.TASK_STATES.SUCCEEDED); if (patch.waitingUser) desired = domain.TASK_STATES.WAITING_USER; if (patch.replyConfirmed) { desired = patch.usedRestFallback ? domain.TASK_STATES.SUCCEEDED_WITH_REST_FALLBACK : domain.TASK_STATES.SUCCEEDED; } if (desired === domain.TASK_STATES.ABORTED && task.state !== domain.TASK_STATES.ABORTING) { task = this.store.transitionTask(taskId, domain.TASK_STATES.ABORTING); } if (desired === domain.TASK_STATES.WAITING_USER) { task = this.store.transitionTask(taskId, desired, patch); } else if (domain.canTransition(task.state, desired)) { task = this.store.transitionTask(taskId, desired, patch); } else if (task.state === domain.TASK_STATES.RUNNING && (desired === domain.TASK_STATES.SUCCEEDED || desired === domain.TASK_STATES.SUCCEEDED_WITH_REST_FALLBACK)) { task = this.store.transitionTask(taskId, domain.TASK_STATES.VERIFYING_REPLY, patch); task = this.store.transitionTask(taskId, desired, patch); } else if (task.state === domain.TASK_STATES.RUNNING && desired === domain.TASK_STATES.FAILED_REPLY) { task = this.store.transitionTask(taskId, domain.TASK_STATES.VERIFYING_REPLY, patch); task = this.store.transitionTask(taskId, desired, patch); } else { task = this.store.transitionTask(taskId, controller.signal.aborted ? domain.TASK_STATES.ABORTED : domain.TASK_STATES.FAILED, { errorCode: 'invalid_runner_state' }); } if (patch.turn) this.store.upsertTurn(patch.turn); this.emit('updated', task); this.store.appendAudit({ taskId, sessionKey: task.sessionKey, repoKey: task.repoKey, action: `task.${task.state}`, fromState: domain.TASK_STATES.RUNNING, toState: task.state, turnId: task.turnId, errorCode: task.errorCode }); if (task.state === domain.TASK_STATES.RETRY_WAIT) { this.pending.push(task.taskId); } } catch (error) { const latest = this.store.getTask(taskId); const target = latest.state === domain.TASK_STATES.ABORTING ? domain.TASK_STATES.ABORTED : (error.code === 'blocked_workspace' ? domain.TASK_STATES.BLOCKED_WORKSPACE : domain.TASK_STATES.FAILED); task = this.store.transitionTask(taskId, target, { errorCode: error.code || 'runner_failed', errorMessage: error.message || String(error) }); this.emit('updated', task); this.store.appendAudit({ taskId, sessionKey: task.sessionKey, repoKey: task.repoKey, action: `task.${target}`, fromState: latest.state, toState: target, errorCode: task.errorCode }); } finally { this.active.delete(taskId); this.repoActive.delete(repoLockKey); this.runners.delete(taskId); this.drain(); } } } function createGiteaWorkflowQueue(options) { return new GiteaWorkflowQueue(options); } module.exports = { GiteaWorkflowQueue, createGiteaWorkflowQueue, sleep };