perf(gateway): coalesce observed project discovery (#162016)

Concurrent observed-project lists each ran their own git discovery per
project, so a reconnect burst multiplied git launches by the client count
(879 git spawns in one boot minute on Team with 11 reconnecting clients).

Share one bounded git discovery pass across concurrent callers and let
later requests refresh the facts instead of recomputing them. The
projects.list reply contract is unchanged.

Testbox fixture (80 repositories, 10 concurrent lists): git launches
1,280 -> 128, wall time 3,196 -> 609 ms, handler-thread CPU 2,136 -> 385 ms.
This commit is contained in:
Peter Steinberger 2026-09-30 22:02:32 -07:00 • committed by GitHub
parent 58432fd365
commit 703cbf2fdf
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
12 changed files with 150 additions and 52 deletions

View file

@ -557,6 +557,12 @@ checkout directory's name.
Registering the same resolved repository root again returns its existing project ID and
display name. Passing a different `name` does not rename an existing project.
`projects.list` returns recorded projects without probing Git. Operators with
`operator.write` can request `{"includeObserved":true}` to discover additional
checkouts from visible sessions and managed worktrees. Concurrent discovery of
the same checkout set shares one bounded Git pass; subsequent requests read
fresh repository metadata, including external remote changes and moved checkouts.
**Projects from GitHub.** Search the same picker or paste a GitHub HTTPS or `git@github.com` repository URL. For a remote destination, creation records that source and the runner fetches it during dispatch; no Gateway project clone is required. For Gateway execution, the picker clones into the Gateway-managed projects area. Recent repository sources retain their URL without inventing a local path. Public repository search and cloning work anonymously. Private remote checkout uses the effective shared `tools.github` identity; the discovery credential below only grants picker access. For affiliated and private repositories, prefer the explicit `gateway.controlUi.github.token` SecretRef so this service access has a clear runtime owner. When it is omitted, the Gateway still uses its shipped `GH_TOKEN` then `GITHUB_TOKEN` fallback from the shared process environment. When it is explicit, its exact environment or store name is excluded from agent execution without clearing unrelated native GitHub CLI variables. Search requires `operator.read`, cloning requires `operator.write`, and deleting a Gateway-managed cloned checkout requires `operator.admin`. Clone deletion refuses while a live session or managed worktree still references the checkout. SecretRef ownership is not an OS-user security boundary; use a sandbox, dedicated host, or dedicated OS user when same-account processes are not trusted.
Use the Effort menu to choose Fast Mode before creating a session. New Session persists that choice before the first local or remote turn starts.

View file

@ -261,8 +261,15 @@ src/agents/worktrees/filesystem-refs.native.ts
src/agents/worktrees/git-path-inventory.ts
src/agents/worktrees/git-worktree-operations.runtime.ts
src/agents/worktrees/git.ts
src/agents/worktrees/owner.ts
src/agents/worktrees/provisioned-file-inspection.ts
src/agents/worktrees/registry-read.kernel.ts
src/agents/worktrees/registry-read.ts
src/agents/worktrees/registry.ts
src/agents/worktrees/repository-paths.ts
src/agents/worktrees/run-lease-owner.ts
src/agents/worktrees/run-lease-store.kernel.ts
src/agents/worktrees/service-preparation.ts
src/agents/worktrees/snapshot-exact-state.ts
src/agents/worktrees/snapshot-index-objects.ts
src/agents/worktrees/snapshot-inventory.ts

View file

@ -157,6 +157,16 @@ export async function resolveRepository(repoRoot: string): Promise<ResolvedRepos
return await resolveRepositoryFromRealPath(requested, repoRoot);
}
export async function resolveRepositoryIdentity(repoRoot: string) {
const resolved = await resolveRepository(repoRoot);
return {
checkoutRoot: resolved.sourceRoot,
repoRoot: resolved.repoRoot,
originUrl: resolved.originUrl,
fingerprint: resolved.fingerprint,
};
}
export async function cleanupFailedCreate(...args: Parameters<typeof removeFailedWorktree>) {
const failure = await removeFailedWorktree(...args);
if (failure) {

View file

@ -90,6 +90,7 @@ import {
resetFailedWorktreeAdd,
resolveRepository,
resolveRepositoryFromRealPath,
resolveRepositoryIdentity,
runSetupScript,
validateName,
withWorktreeSource,
@ -744,13 +745,11 @@ export class ManagedWorktreeService {
originUrl: string;
fingerprint: string;
}> {
const resolved = await resolveRepository(repoRoot);
return {
checkoutRoot: resolved.sourceRoot,
repoRoot: resolved.repoRoot,
originUrl: resolved.originUrl,
fingerprint: resolved.fingerprint,
};
return await resolveRepositoryIdentity(repoRoot);
}
async resolveRepositoryIdentities(roots: string[]) {
return await runGitReadOperation({ type: "repository.identities", input: { roots } });
}
/**

View file

@ -8,8 +8,6 @@ import type { SessionEntry } from "../../config/sessions.js";
import { createProjectsHandlers } from "./projects.js";
import type { GatewayClient, GatewayRequestContext, RespondFn } from "./types.js";
type ProjectWorktreeService = Parameters<typeof createProjectsHandlers>[0];
const seededSessions = vi.hoisted(() => ({
store: {} as Record<string, SessionEntry>,
}));
@ -73,7 +71,13 @@ async function listObservedProjects(params: {
};
client?: GatewayClient;
}) {
const handlers = createProjectsHandlers(params.service as never);
const handlers = createProjectsHandlers({
listRegistryRecords: params.service.listRegistryRecords,
resolveRepositoryIdentities: (roots: string[]) =>
Promise.all(
roots.map((root) => params.service.resolveRepositoryIdentity(root).catch(() => undefined)),
),
} as never);
const responses: Parameters<RespondFn>[] = [];
const cfg = { agents: { list: [{ id: "main", default: true }] } };
await handlers["projects.list"]?.({
@ -99,35 +103,23 @@ beforeEach(() => {
});
describe("projects.list observed projects", () => {
it("deduplicates probes and overlaps bounded work without changing result order", async () => {
it("deduplicates admitted paths and preserves result order when a checkout is unavailable", async () => {
seededSessions.store = Object.fromEntries(
Array.from({ length: 5_000 }, (_, index) => [
`agent:main:session-${index}`,
{ sessionId: `session-${index}`, updatedAt: index, execCwd: `/repos/${index % 8}` },
]),
);
let active = 0;
let peak = 0;
const resolveRepositoryIdentity = vi.fn(async (checkoutPath: string) => {
active += 1;
peak = Math.max(peak, active);
try {
// Newer candidates finish later; output must retain admission order.
for (let step = 0; step < Number(checkoutPath.at(-1)); step += 1) {
await Promise.resolve();
}
if (checkoutPath === "/repos/3") {
throw new Error("checkout unavailable");
}
return {
checkoutRoot: checkoutPath.replace("/repos/", "/physical/"),
repoRoot: "/physical/main",
originUrl: "https://example.test/project.git",
fingerprint: "project",
};
} finally {
active -= 1;
if (checkoutPath === "/repos/3") {
throw new Error("checkout unavailable");
}
return {
checkoutRoot: checkoutPath.replace("/repos/", "/physical/"),
repoRoot: "/physical/main",
originUrl: "https://example.test/project.git",
fingerprint: "project",
};
});
await expect(
@ -149,9 +141,6 @@ describe("projects.list observed projects", () => {
expect(new Set(resolveRepositoryIdentity.mock.calls.map(([checkout]) => checkout)).size).toBe(
8,
);
expect(peak).toBeGreaterThan(1);
expect(peak).toBeLessThanOrEqual(4);
expect(active).toBe(0);
});
it.each([["operator.write"], ["operator.admin"]])(
@ -349,11 +338,11 @@ describe("projects.list observed projects", () => {
{ sessionId: `session-${index}`, updatedAt: index, execCwd: `/repos/${index}` },
]),
);
const resolveRepositoryIdentity = vi.fn<ProjectWorktreeService["resolveRepositoryIdentity"]>(
async (_checkoutPath) => {
throw new Error("checkout unavailable");
},
);
const resolveRepositoryIdentity = vi.fn<
Parameters<typeof listObservedProjects>[0]["service"]["resolveRepositoryIdentity"]
>(async (_checkoutPath) => {
throw new Error("checkout unavailable");
});
await expect(
listObservedProjects({
@ -362,7 +351,7 @@ describe("projects.list observed projects", () => {
).resolves.toEqual([]);
expect(resolveRepositoryIdentity).toHaveBeenCalledTimes(PROJECTS_LIST_MAX_IDENTITY_PROBES);
expect(resolveRepositoryIdentity.mock.calls.length).toBeLessThanOrEqual(rawCandidateLimit);
expect(resolveRepositoryIdentity.mock.calls[0]?.[0]).toBe(`/repos/${rawCandidateLimit + 4}`);
expect(resolveRepositoryIdentity).toHaveBeenCalledWith(`/repos/${rawCandidateLimit + 4}`);
expect(resolveRepositoryIdentity).not.toHaveBeenCalledWith("/repos/0");
});
});

View file

@ -22,8 +22,9 @@ export const resolveRepositoryIdentity = vi.fn(async (checkoutPath: string) => (
}));
export const projectsHandlers = createProjectsHandlers({
listRegistryRecords,
resolveRepositoryIdentity,
} as never);
resolveRepositoryIdentities: (roots: string[]) =>
Promise.all(roots.map(resolveRepositoryIdentity)),
});
export async function initializeRepository(
root: string,

View file

@ -12,6 +12,7 @@ import {
import * as transcriptWorker from "../../config/sessions/session-transcript-worker-runtime.js";
import type { OpenClawConfig } from "../../config/types.openclaw.js";
import { sha256HexPrefixCore } from "../../infra/crypto-digest.js";
import * as spawnDiagnostics from "../../process/spawn-diagnostics.js";
import { registerProjectRegistry } from "../../projects/project-registry.js";
import { registerClonedProjectRegistry } from "../../projects/project-registry.test-support.js";
import { SecretSurfaceUnavailableError } from "../../secrets/runtime-degraded-state.js";
@ -45,6 +46,62 @@ function withProjectState(run: (state: OpenClawTestState) => Promise<void>) {
return withOpenClawTestState({ layout: "state-only", prefix: "projects-rpc-" }, run);
}
test("projects.list coalesces concurrent observed Git discovery and refreshes later reads", async () => {
await withProjectState(async (state) => {
const repo = await initializeRepository(state.root);
replaceSessionEntrySync(
{ agentId: "main", sessionKey: "agent:main:observed" },
{ sessionId: "observed", updatedAt: 1, execCwd: repo },
);
const cfg = { agents: { entries: { main: { workspace: state.workspaceDir } } } };
const list = () =>
invokeProjectMethod(
"projects.list",
{ includeObserved: true },
cfg,
["operator.write"],
undefined,
registeredProjectsHandlers,
);
using spawns = vi.spyOn(spawnDiagnostics, "recordChildProcessSpawn");
const single = await list();
expect(single).toMatchObject({
ok: true,
payload: {
observedProjects: [
{ name: "registered", originUrl: "https://github.com/openclaw/openclaw.git" },
],
},
});
const gitSpawns = () =>
spawns.mock.calls.filter(([command]) =>
/^git(?:\.exe|\.cmd)?$/i.test(path.win32.basename(command)),
).length;
const onePass = gitSpawns();
expect(onePass).toBeGreaterThan(0);
spawns.mockClear();
const concurrent = await Promise.all(Array.from({ length: 10 }, list));
for (const result of concurrent) {
expect(result).toEqual(single);
}
expect(gitSpawns()).toBeLessThanOrEqual(onePass);
await execFileAsync("git", [
"-C",
repo,
"remote",
"set-url",
"origin",
"https://example.test/changed.git",
]);
expect(await list()).toMatchObject({
ok: true,
payload: { observedProjects: [{ originUrl: "https://example.test/changed.git" }] },
});
await fs.rm(repo, { recursive: true });
expect(await list()).toMatchObject({ ok: true, payload: { observedProjects: [] } });
});
});
test.each([
{
failure: "rate limit",

View file

@ -34,7 +34,6 @@ import {
} from "../../projects/project-registry.js";
import { isTrustedSecretSurfaceUnavailableError } from "../../secrets/runtime-degraded-state.js";
import { readCurrentUserProfileAliases } from "../../state/user-profile-list.js";
import { runTasksWithConcurrency } from "../../utils/run-with-concurrency.js";
import { readGatewayAccessRevision } from "../gateway-access-revision.js";
import {
CONTROL_UI_GITHUB_CREDENTIAL_UNAVAILABLE_MESSAGE,
@ -53,7 +52,7 @@ import { assertValidParams } from "./validation.js";
type ProjectWorktreeService = Pick<
ManagedWorktreeService,
"listRegistryRecords" | "resolveRepositoryIdentity"
"listRegistryRecords" | "resolveRepositoryIdentities"
>;
type ProjectCandidate = {
@ -255,18 +254,16 @@ async function listObservedProjects(
});
}
// Admit the same newest-first distinct paths before overlapping Git work. Keep facts
// request-local: session/registry revisions cannot detect external Git metadata edits.
// Admit newest-first paths before canonicalizing the shared discovery pass's key.
diagnostics?.mark("identityProbes");
const probePaths = [
...new Set(
rawCandidates.map((raw) => (raw.kind === "worktree" ? raw.repoRoot : raw.checkoutPath)),
),
].slice(0, PROJECTS_LIST_MAX_IDENTITY_PROBES);
const { results } = await runTasksWithConcurrency({
tasks: probePaths.map((checkoutPath) => () => service.resolveRepositoryIdentity(checkoutPath)),
limit: 4,
});
]
.slice(0, PROJECTS_LIST_MAX_IDENTITY_PROBES)
.toSorted();
const results = probePaths.length ? await service.resolveRepositoryIdentities(probePaths) : [];
const identities = new Map(
probePaths.map((checkoutPath, index) => [checkoutPath, results[index]]),
);

View file

@ -157,6 +157,13 @@ function createReadCache<Input, Output>(
// a five-minute fallback observes working-tree edits made outside OpenClaw.
function createReadCaches() {
return {
identities: createReadCache(
(input: GitReadOperations["repository.identities"]["input"], signal) =>
runGitWorkerOperation({ type: "repository.identities", input }, { signal }),
// Identity includes Git config and worktree relocation inputs without a complete revision.
// Share only pending passes so later discovery always sees external changes.
0,
),
context: createReadCache(
(input: GitReadOperations["checkout.context"]["input"], signal) =>
runGitWorkerOperation({ type: "checkout.context", input }, { signal }),
@ -246,8 +253,11 @@ export function runGitReadOperation(operation: GitReadOperation, options?: GitRe
if (state.closing) {
return Promise.reject(new Error("Git reads are unavailable while the Gateway is restarting"));
}
const { context, branchFacts, diff, branches, baseline } = (state.caches ??= createReadCaches());
const { context, branchFacts, diff, branches, baseline, identities } = (state.caches ??=
createReadCaches());
switch (operation.type) {
case "repository.identities":
return identities.read(operation.input, options);
case "checkout.revision":
return runGitWorkerOperation(operation, options);
case "checkout.context":

View file

@ -1,4 +1,5 @@
import { readRepositoryBranches } from "../agents/worktrees/branches.runtime.js";
import { resolveRepositoryIdentity } from "../agents/worktrees/service-preparation.js";
import {
readCheckoutGitContext,
readCheckoutGitRevision,
@ -8,12 +9,20 @@ import {
collectCheckoutDiff,
collectCheckoutDiffBaseline,
} from "../sessions/session-diff.runtime.js";
import { runTasksWithConcurrency } from "../utils/run-with-concurrency.js";
import type { GitReadOperation, GitReadOperationResult } from "./git-read-operations.js";
export async function executeGitReadOperation(
operation: GitReadOperation,
): Promise<GitReadOperationResult> {
switch (operation.type) {
case "repository.identities":
return (
await runTasksWithConcurrency({
tasks: operation.input.roots.map((root) => () => resolveRepositoryIdentity(root)),
limit: 4,
})
).results;
case "checkout.revision":
return readCheckoutGitRevision(operation.input);
case "checkout.context":

View file

@ -28,6 +28,18 @@ export type GitCheckoutDiffInput = { cwd: string; baseCommit?: string } & (
);
export type GitReadOperations = {
"repository.identities": {
input: { roots: string[] };
output: Array<
| {
checkoutRoot: string;
repoRoot: string;
originUrl: string;
fingerprint: string;
}
| undefined
>;
};
"checkout.revision": { input: { root: string; includeIndex: boolean }; output: string | null };
"checkout.context": { input: { root: string }; output: GitCheckoutContext | null };
"checkout.diff": { input: GitCheckoutDiffInput; output: Omit<SessionsDiffResult, "sessionKey"> };

View file

@ -73,6 +73,7 @@ function poolFor(state: GitWorkerRuntime, command: GitWorkerCommand): GitPool {
: command.type.startsWith("workspace.")
? "workspace"
: command.type === "repository.branches" ||
command.type === "repository.identities" ||
command.type === "checkout.context" ||
command.type === "checkout.revision"
? "reads"