feat: 编排器跨项目全局并发上限闸(tsk_erQHNYdn6i9g)
新增用户级设置项 globalConcurrency(settings 表,独立于新建项目默认值 来源 concurrency),缺省 0 = 不限,向后兼容无需 DB 迁移。 - store: UserSettings 加 globalConcurrency 字段 + 默认 0;新增 globalInflightCount()(跨所有项目 started run 总数,与 inflightTaskIds 同口径) - orchestrator.claimTick: 叠加全局闸,与 per-project 闸串联(都过才领); 轮初查一次全局在途、轮内手动 ++,达上限即 return 跨所有项目停止领取 - api: PUT /api/settings 支持 globalConcurrency(空/null 归一为 0) - web: 侧栏「全局设置」入口 + 模态,编辑/回显/校验(≥0 整数) - test: 新增全局闸达上限跨项目阻止 / 未达上限正常领取 / 边界 / 轮初短路 用例 Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -137,6 +137,8 @@ export function buildServer(opts: ApiOptions): FastifyInstance {
|
||||
if (b.budgetUsd !== undefined) patch.budgetUsd = b.budgetUsd === null || b.budgetUsd === '' ? null : Number(b.budgetUsd);
|
||||
if (b.budgetPeriod !== undefined) patch.budgetPeriod = b.budgetPeriod as 'day' | 'month';
|
||||
if (b.model !== undefined) patch.model = b.model === null || b.model === '' ? null : String(b.model);
|
||||
if (b.globalConcurrency !== undefined)
|
||||
patch.globalConcurrency = b.globalConcurrency === null || b.globalConcurrency === '' ? 0 : Number(b.globalConcurrency);
|
||||
return store.putSettings(patch);
|
||||
});
|
||||
|
||||
|
||||
@@ -209,9 +209,16 @@ export function createOrchestrator(store: Store, log: OrchestratorLogger, deps:
|
||||
/**
|
||||
* 领取轮:对每个 active 且 autonomy≠manual 的项目,并发闸 = inflightTaskIds(有 started run 的任务,
|
||||
* executor + planner 通用)。在 active<concurrency 时按 claimable 领新任务并 spawn worker。
|
||||
*
|
||||
* 全局并发闸(settings.globalConcurrency,0/缺省=不限)与 per-project 闸【串联】:两闸都过才领。
|
||||
* 全局在途总数(globalActive)轮初查一次、轮内手动 ++(claimOne 同步起新 started run),命中即 return——
|
||||
* 达上限后任何项目都不应再领,故 return(跨所有项目停止)而非 break(只停当前项目)。
|
||||
*/
|
||||
function claimTick(): void {
|
||||
try {
|
||||
const cap = store.getSettings().globalConcurrency; // 0/缺省 = 不限
|
||||
let globalActive = store.globalInflightCount(); // 轮初的全局在途总数
|
||||
if (cap > 0 && globalActive >= cap) return; // 全局闸已满 → 整轮跨项目都不领
|
||||
for (const p of store.listProjects()) {
|
||||
if (p.status !== 'active' || p.autonomy === 'manual') continue;
|
||||
const inflight = store.inflightTaskIds(p.id);
|
||||
@@ -219,9 +226,11 @@ export function createOrchestrator(store: Store, log: OrchestratorLogger, deps:
|
||||
if (active >= p.concurrency) continue;
|
||||
|
||||
for (const { task, score } of claimable(p, inflight)) {
|
||||
if (active >= p.concurrency) break;
|
||||
if (active >= p.concurrency) break; // per-project 闸
|
||||
if (cap > 0 && globalActive >= cap) return; // 全局闸:达上限即跨所有项目停止领取
|
||||
claimOne(p, task, score);
|
||||
active++;
|
||||
globalActive++; // 本轮内手动累加(claimOne 已起新 started run)
|
||||
}
|
||||
}
|
||||
} catch (e) {
|
||||
|
||||
@@ -30,11 +30,14 @@ export interface UserSettings {
|
||||
budgetUsd: number | null;
|
||||
budgetPeriod: 'day' | 'month';
|
||||
model: string | null;
|
||||
/** 跨所有项目的在途 run 总数上限(全局并发闸);0 = 不限。与 concurrency(新建项目默认值)相互独立。 */
|
||||
globalConcurrency: number;
|
||||
}
|
||||
const SETTINGS_DEFAULT: UserSettings = {
|
||||
autonomy: 'manual', concurrency: 1, maxRetries: 2, timeoutMs: 1_800_000,
|
||||
autoApprovePlan: false, autoApproveExec: false,
|
||||
budgetUsd: null, budgetPeriod: 'month', model: null,
|
||||
globalConcurrency: 0,
|
||||
};
|
||||
const id = (prefix: string): string => `${prefix}_${nanoid(12)}`;
|
||||
|
||||
@@ -918,6 +921,11 @@ export class Store {
|
||||
return new Set(rows.map((r) => r.tid));
|
||||
}
|
||||
|
||||
/** 跨所有项目的在途 run 总数(全局并发闸用;与 inflightTaskIds 同口径=started run,覆盖 executor + planner)。 */
|
||||
globalInflightCount(): number {
|
||||
return (this.db.prepare(`SELECT COUNT(*) AS n FROM runs WHERE status = 'started'`).get() as { n: number }).n;
|
||||
}
|
||||
|
||||
/** 所有 started run + 其任务(reap / ingest / reconcile 通用,覆盖 executor + planner + 残留复审 run)。 */
|
||||
liveRunsWithTask(): Array<{ task: Task; run: Run }> {
|
||||
const runs = this.db.prepare(`SELECT * FROM runs WHERE status = 'started' ORDER BY started_at`).all() as RunRow[];
|
||||
|
||||
Reference in New Issue
Block a user