feat: support MCP elicitation and rebuild release
This commit is contained in:
238
lib/gitea-workflow-queue.js
Normal file
238
lib/gitea-workflow-queue.js
Normal file
@@ -0,0 +1,238 @@
|
||||
'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 };
|
||||
Reference in New Issue
Block a user