diff --git a/.changeset/fn-7205-concurrency-counter-enforcement.md b/.changeset/fn-7205-concurrency-counter-enforcement.md new file mode 100644 index 0000000000..1436396446 --- /dev/null +++ b/.changeset/fn-7205-concurrency-counter-enforcement.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Concurrency panels now prefer live engine counts so running-agent totals stay accurate. +category: fix +dev: Prefer engine-manager task stores over stale registered/default fallback stores in the dashboard live-count source; add regressions for count normalization and scoped semaphore live-limit behavior. diff --git a/packages/core/src/__tests__/live-agent-count.test.ts b/packages/core/src/__tests__/live-agent-count.test.ts index 55d310c433..b380162dda 100644 --- a/packages/core/src/__tests__/live-agent-count.test.ts +++ b/packages/core/src/__tests__/live-agent-count.test.ts @@ -1,5 +1,5 @@ import { describe, expect, it } from "vitest"; -import { countRunningAgentTasks, isRunningAgentTask } from "../live-agent-count.js"; +import { countRunningAgentTasks, deriveRunningAgentCounts, isRunningAgentTask } from "../live-agent-count.js"; import type { Task } from "../types.js"; function task(overrides: Pick & Partial>): Pick { @@ -44,4 +44,20 @@ describe("live agent count predicates", () => { task({ column: "archived" }), ])).toBe(7); }); + + it("normalizes display counts for zero, one, multi-project, unopened, and oversubscribed states", () => { + expect(deriveRunningAgentCounts({})).toEqual({ currentlyActive: 0, projectsActive: {} }); + expect(deriveRunningAgentCounts({ proj_zero: 0, proj_one: 1 })).toEqual({ + currentlyActive: 1, + projectsActive: { proj_one: 1 }, + }); + expect(deriveRunningAgentCounts({ proj_a: 2, proj_b: 4, proj_unopened: 0 })).toEqual({ + currentlyActive: 6, + projectsActive: { proj_a: 2, proj_b: 4 }, + }); + expect(deriveRunningAgentCounts({ proj_over_limit: 12, proj_negative: -3, proj_nan: Number.NaN })).toEqual({ + currentlyActive: 12, + projectsActive: { proj_over_limit: 12 }, + }); + }); }); diff --git a/packages/dashboard/src/__tests__/server.test.ts b/packages/dashboard/src/__tests__/server.test.ts index 8101b7c0a5..ff5d5b75e5 100644 --- a/packages/dashboard/src/__tests__/server.test.ts +++ b/packages/dashboard/src/__tests__/server.test.ts @@ -203,7 +203,10 @@ describe("createServer options", () => { it("registers live counts for the already-open default store by central project id", async () => { const store = createMockStore({ - listTasks: vi.fn().mockResolvedValue([{ id: "FN-1" }, { id: "FN-2" }]), + listTasks: vi.fn().mockResolvedValue([ + { id: "FN-1", column: "in-progress" }, + { id: "FN-2", column: "triage", status: "planning", paused: false }, + ]), }); const centralCore = { getDefaultProjectId: vi.fn().mockResolvedValue("proj_default"), @@ -215,13 +218,13 @@ describe("createServer options", () => { expect(source).toBeDefined(); if (!source) throw new Error("expected running-agent count source"); await expect(source(["proj_default", "proj_unopened"])).resolves.toEqual({ proj_default: 2 }); - expect(store.listTasks).toHaveBeenCalledWith({ column: "in-progress", slim: true }); + expect(store.listTasks).toHaveBeenCalledWith({ slim: true }); }); it("registers live counts for already-open engine-manager stores without starting engines", async () => { const store = createMockStore(); const engineStore = createMockStore({ - listTasks: vi.fn().mockResolvedValue([{ id: "FN-3" }]), + listTasks: vi.fn().mockResolvedValue([{ id: "FN-3", column: "in-review", status: "merging", paused: false }]), }); const getEngine = vi.fn((projectId: string) => projectId === "proj_engine" ? { getTaskStore: vi.fn(() => engineStore) } @@ -236,7 +239,40 @@ describe("createServer options", () => { await expect(source(["proj_engine", "proj_unopened"])).resolves.toEqual({ proj_engine: 1 }); expect(getEngine).toHaveBeenCalledWith("proj_engine"); expect(getEngine).toHaveBeenCalledWith("proj_unopened"); - expect(engineStore.listTasks).toHaveBeenCalledWith({ column: "in-progress", slim: true }); + expect(engineStore.listTasks).toHaveBeenCalledWith({ slim: true }); + expect(store.listTasks).not.toHaveBeenCalled(); + }); + + it("prefers the live engine-manager store over the default fallback for the default project", async () => { + const store = createMockStore({ + listTasks: vi.fn().mockResolvedValue([]), + }); + const engineStore = createMockStore({ + listTasks: vi.fn().mockResolvedValue([ + { id: "FN-4", column: "in-progress" }, + { id: "FN-5", column: "in-review", status: "reviewing", paused: false }, + ]), + }); + const getEngine = vi.fn((projectId: string) => projectId === "proj_default" + ? { getTaskStore: vi.fn(() => engineStore) } + : undefined); + const engineManager = { getEngine }; + const centralCore = { + getDefaultProjectId: vi.fn().mockResolvedValue("proj_default"), + }; + + createServer(store, { + centralCore: centralCore as unknown as CentralCore, + engineManager: engineManager as unknown as import("@fusion/engine").ProjectEngineManager, + }); + const source = getRunningAgentCountSource(); + + expect(source).toBeDefined(); + if (!source) throw new Error("expected running-agent count source"); + await expect(source(["proj_default", "proj_unopened"])).resolves.toEqual({ proj_default: 2 }); + expect(getEngine).toHaveBeenCalledWith("proj_default"); + expect(getEngine).toHaveBeenCalledWith("proj_unopened"); + expect(engineStore.listTasks).toHaveBeenCalledWith({ slim: true }); expect(store.listTasks).not.toHaveBeenCalled(); }); }); diff --git a/packages/dashboard/src/server.ts b/packages/dashboard/src/server.ts index cd9b8f242a..1588f6e11f 100644 --- a/packages/dashboard/src/server.ts +++ b/packages/dashboard/src/server.ts @@ -838,6 +838,9 @@ export function createServer(store: TaskStore, options?: ServerOptions): ReturnT FNXC:GlobalConcurrencyControls 2026-06-26-23:41: The default in-process TaskStore is already open but is intentionally not part of the secondary project-store cache. Include it by central default project id, and include any engine-manager stores already resident in memory, so live reads cover every already-open store without calling getOrCreateProjectStore(), watch(), or runtime startup paths. + + FNXC:GlobalConcurrencyControls 2026-06-28-16:48: + FN-7205 requires the footer and Command Center counters to reflect actual running top-level agents from the live runtime store. Prefer engine-manager stores over registered/default fallback stores, and only use the default store when no live engine/registered source has supplied that project, so stale bootstrap stores cannot overwrite the semaphore-facing runtime count after a slider or engine lifecycle change. */ setRunningAgentCountSource(async (projectIds) => { const requestedProjectIds = new Set(projectIds); @@ -845,9 +848,6 @@ export function createServer(store: TaskStore, options?: ServerOptions): ReturnT if (options?.engineManager) { await Promise.all(projectIds.map(async (projectId) => { - if (counts[projectId] !== undefined) { - return; - } const engine = options.engineManager?.getEngine(projectId); if (!engine) { return; @@ -857,7 +857,7 @@ export function createServer(store: TaskStore, options?: ServerOptions): ReturnT } const defaultProjectId = await options?.centralCore?.getDefaultProjectId?.(); - if (defaultProjectId && requestedProjectIds.has(defaultProjectId)) { + if (defaultProjectId && requestedProjectIds.has(defaultProjectId) && counts[defaultProjectId] === undefined) { counts[defaultProjectId] = await countRunningAgentsInStore(store); } diff --git a/packages/engine/src/__tests__/concurrency.test.ts b/packages/engine/src/__tests__/concurrency.test.ts index 590c3905ed..90c42d300a 100644 --- a/packages/engine/src/__tests__/concurrency.test.ts +++ b/packages/engine/src/__tests__/concurrency.test.ts @@ -113,6 +113,44 @@ describe("ScopedAgentSemaphore", () => { expect(shared.activeCount).toBe(0); }); + it("honors live global-limit changes across scoped project semaphores on the next acquire", async () => { + let globalLimit = 2; + const shared = new AgentSemaphore(() => globalLimit); + const projectA = new ScopedAgentSemaphore(shared); + const projectB = new ScopedAgentSemaphore(shared); + + await projectA.acquire(PRIORITY_EXECUTE); + await projectB.acquire(PRIORITY_MERGE); + expect(shared.snapshot()).toEqual({ activeCount: 2, waitingCount: 0, availableCount: 0, limit: 2 }); + + globalLimit = 1; + let acquired = false; + const waiter = projectA.acquire(PRIORITY_EXECUTE).then(() => { + acquired = true; + }); + await Promise.resolve(); + + expect(acquired).toBe(false); + expect(shared.snapshot()).toEqual({ activeCount: 2, waitingCount: 1, availableCount: 0, limit: 1 }); + + projectA.release(); + await Promise.resolve(); + expect(acquired).toBe(false); + expect(shared.snapshot()).toEqual({ activeCount: 1, waitingCount: 1, availableCount: 0, limit: 1 }); + + globalLimit = 2; + projectB.release(); + await waiter; + + expect(acquired).toBe(true); + expect(projectA.heldCount).toBe(1); + expect(projectB.heldCount).toBe(0); + expect(shared.snapshot()).toEqual({ activeCount: 1, waitingCount: 0, availableCount: 1, limit: 2 }); + + projectA.release(); + expect(shared.activeCount).toBe(0); + }); + it("reconciles only this scope's slots when another project still holds global capacity", async () => { const shared = new AgentSemaphore(3); const idleProject = new ScopedAgentSemaphore(shared);