239 lines
12 KiB
JavaScript
239 lines
12 KiB
JavaScript
'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 };
|