Files
cc-web/lib/gitea-workflow-queue.js

239 lines
12 KiB
JavaScript
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

'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 };