Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
74 changes: 73 additions & 1 deletion v2/e2e/dispatch.e2e.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,15 @@
import { execFile } from "node:child_process";
import { cp, mkdir, mkdtemp, readFile, realpath, stat, writeFile } from "node:fs/promises";
import {
chmod,
cp,
mkdir,
mkdtemp,
readFile,
readdir,
realpath,
stat,
writeFile,
} from "node:fs/promises";
import { tmpdir } from "node:os";
import { dirname, join } from "node:path";
import { promisify } from "node:util";
Expand Down Expand Up @@ -349,6 +359,68 @@ describe("crew start", () => {
expect(dispatch["fixture:LOW-1"].reason).toBe("slots-full");
});

it("reserves one shared slot across concurrent starts for distinct tasks", async () => {
const fixture = await createDispatchFixture({
maximumInProgress: 1,
tasks: [task({ id: "A-1", repositories: [] }), task({ id: "B-1", repositories: [] })],
});
const barrierDirectory = join(fixture.root, "which-barrier");
const releasePath = join(barrierDirectory, "release");
const firstReadyPath = join(barrierDirectory, "first-ready");
const secondReadyPath = join(barrierDirectory, "second-ready");
await mkdir(barrierDirectory);
const whichPath = join(barrierDirectory, "which");
await writeFile(
whichPath,
[
"#!/bin/sh",
'if [ "$1" = "codex" ]; then',
' touch "$FAKE_WHICH_READY"',
' while [ ! -e "$FAKE_WHICH_RELEASE" ]; do sleep 0.01; done',
"fi",
'exec /usr/bin/which "$@"',
"",
].join("\n"),
);
await chmod(whichPath, 0o755);
const environment = {
...fixture.environment,
FAKE_WHICH_RELEASE: releasePath,
PATH: `${barrierDirectory}:${fixture.environment["PATH"]}`,
};

const firstStart = runCrew({
arguments: ["start", "A-1"],
environment: { ...environment, FAKE_WHICH_READY: firstReadyPath },
});
const secondStart = runCrew({
arguments: ["start", "B-1"],
environment: { ...environment, FAKE_WHICH_READY: secondReadyPath },
});
await Promise.all([waitForPath(firstReadyPath), waitForPath(secondReadyPath)]);
await writeFile(releasePath, "release\n");

const results = await Promise.all([firstStart, secondStart]);
const records = (await readdir(fixture.runsDirectory)).filter((entry) =>
entry.endsWith(".json"),
);
expect(results.filter((result) => result.stdout.includes("Started fixture:"))).toHaveLength(1);
expect(records).toHaveLength(1);

const verdicts = Object.values(
JSON.parse(await readFile(join(dirname(fixture.runsDirectory), "dispatch.json"), "utf8")),
);
const updates = (await readFile(fixture.updatesPath, "utf8")).trim().split("\n");
const calls = (await readFile(fixture.cmuxCallsPath, "utf8"))
.trim()
.split("\n")
.map((line) => JSON.parse(line));

expect(verdicts).toContainEqual(expect.objectContaining({ reason: "slots-full" }));
expect(updates).toHaveLength(1);
expect(calls.filter((call) => call.arguments[0] === "new-workspace")).toHaveLength(1);
});

it("skips a blocked task unless an explicit start uses force", async () => {
const fixture = await createDispatchFixture({
tasks: [task({ blocked: true, repositories: [] })],
Expand Down
26 changes: 14 additions & 12 deletions v2/src/core/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -325,7 +325,6 @@ async function start(input: {
agentProfile: run.record.agentProfile,
canonicalTaskId: run.record.canonicalTaskId,
}));
let activeCount = active.length;
input.onProgress?.({
active,
force: input.force,
Expand Down Expand Up @@ -395,21 +394,24 @@ async function start(input: {
}
continue;
}
if (!input.force && activeCount >= runtime.config.orchestrator.maximumInProgress) {
const workspaceDirectory = runtime.workspaces.workspaceDirectory({ slug });
const reservation = await runtime.runs.reserveDispatch({
agentProfile: profileName,
canonicalTaskId,
force: input.force,
maximumInProgress: runtime.config.orchestrator.maximumInProgress,
repositories: task.repositories,
workspaceDirectory,
});
if (reservation.type === "full") {
await skipTask({
canonicalTaskId,
detail: "concurrency limit reached",
reason: "slots-full",
});
continue;
}
const workspaceDirectory = runtime.workspaces.workspaceDirectory({ slug });
const provisioning = await runtime.runs.beginDispatch({
agentProfile: profileName,
canonicalTaskId,
repositories: task.repositories,
workspaceDirectory,
});
const provisioning = reservation.run;
const claim = await runtime.registry.update({
canonicalTaskId,
event: { runId: provisioning.record.runId, type: "claimed" },
Expand All @@ -430,9 +432,10 @@ async function start(input: {
input.onProgress?.({
agentProfile: profileName,
canonicalTaskId,
forced: input.force && activeCount >= runtime.config.orchestrator.maximumInProgress,
forced:
input.force && reservation.activeCount >= runtime.config.orchestrator.maximumInProgress,
maximum: runtime.config.orchestrator.maximumInProgress,
slot: activeCount + 1,
slot: reservation.activeCount + 1,
type: "dispatching",
});
try {
Expand Down Expand Up @@ -465,7 +468,6 @@ async function start(input: {
runtime,
});
started.push(canonicalTaskId);
activeCount += 1;
} catch (error) {
const reason =
error instanceof PrepareWorktreeError ? prepareFailureReason(error) : errorMessage(error);
Expand Down
29 changes: 28 additions & 1 deletion v2/src/run/index.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { describe, expect, it } from "vitest";
import { describe, expect, it, vi } from "vitest";
import { RunModule, seedWorkspaceTrust, withFileLock } from "./index.js";

describe("RunModule lifecycle", () => {
Expand Down Expand Up @@ -370,6 +370,33 @@ describe("RunModule lifecycle", () => {
});

describe("RunModule concurrency", () => {
it("reserves dispatch without colliding with a matching task lock name", async () => {
const stateRoot = await mkdtemp(join(tmpdir(), "groundcrew-v2-run-admission-lock-"));
const runs = new RunModule({
environment: process.env,
presenterName: "cmux",
stateRoot,
});
const currentTime = Date.now();
const mockDateNow = vi.spyOn(Date, "now").mockReturnValue(currentTime + 31_000);

try {
const reservation = await runs.reserveDispatch({
agentProfile: "codex",
canonicalTaskId: "dispatch:admission",
force: false,
maximumInProgress: 1,
repositories: [],
workspaceDirectory: join(stateRoot, "workspace"),
});

expect(reservation).toMatchObject({ type: "reserved" });
expect(mockDateNow).not.toHaveBeenCalled();
} finally {
mockDateNow.mockRestore();
}
});

it("creates one durable run when initial writers race", async () => {
const stateRoot = await mkdtemp(join(tmpdir(), "groundcrew-v2-run-race-"));
const runs = new RunModule({
Expand Down
49 changes: 49 additions & 0 deletions v2/src/run/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,14 @@ export interface PresentedWorkspaceSnapshot {
};
}

export type DispatchReservation =
| { readonly activeCount: number; readonly type: "full" }
| {
readonly activeCount: number;
readonly run: ProvisioningRun;
readonly type: "reserved";
};

export class RunModule {
readonly #store: RunStore;
readonly #environment: NodeJS.ProcessEnv;
Expand Down Expand Up @@ -159,6 +167,24 @@ export class RunModule {
return new ProvisioningRunHandle({ module: this, record });
}

public async reserveDispatch(input: {
readonly canonicalTaskId: string;
readonly agentProfile: string;
readonly workspaceDirectory: string;
readonly repositories: readonly string[];
readonly force: boolean;
readonly maximumInProgress: number;
}): Promise<DispatchReservation> {
const reservation = await this.#store.reserveDispatch(input);
return reservation.record === undefined
? { activeCount: reservation.activeCount, type: "full" }
: {
activeCount: reservation.activeCount,
run: new ProvisioningRunHandle({ module: this, record: reservation.record }),
type: "reserved",
};
}

public async findBySlug(input: { readonly slug: string }): Promise<RunHandle | undefined> {
const record = await this.#store.getBySlug(input);
return record === undefined ? undefined : createRunHandle({ module: this, record });
Expand Down Expand Up @@ -598,6 +624,29 @@ class RunStore {
});
}

public async reserveDispatch(input: {
readonly canonicalTaskId: string;
readonly agentProfile: string;
readonly workspaceDirectory: string;
readonly repositories: readonly string[];
readonly force: boolean;
readonly maximumInProgress: number;
}): Promise<{ readonly activeCount: number; readonly record?: RunRecord | undefined }> {
return await withFileLock({
operation: async () => {
const activeCount = (await this.list()).filter(
(record) => record.state === "provisioning" || record.state === "running",
).length;
if (!input.force && activeCount >= input.maximumInProgress) {
return { activeCount };
}
const record = await this.create(input);
return { activeCount, record };
},
path: join(this.#runsDirectory, ".locks", "dispatch-admission.lock"),
});
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

public async getBySlug(input: { readonly slug: string }): Promise<RunRecord | undefined> {
try {
return JSON.parse(await readFile(this.path(input), "utf8")) as RunRecord;
Expand Down