From 703cbf2fdffc3bf4a29f61ac46b4bfcccc2356c8 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Wed, 30 Sep 2026 22:02:32 -0700 Subject: [PATCH] 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. --- docs/web/control-ui/sessions-and-sidebar.md | 6 ++ scripts/pr-lib/wrapper-components.txt | 7 +++ src/agents/worktrees/service-preparation.ts | 10 ++++ src/agents/worktrees/service.ts | 13 ++--- .../server-methods/projects-observed.test.ts | 55 +++++++----------- .../server-methods/projects.test-support.ts | 5 +- src/gateway/server-methods/projects.test.ts | 57 +++++++++++++++++++ src/gateway/server-methods/projects.ts | 15 ++--- src/infra/git-read-cache.ts | 12 +++- src/infra/git-read-operations.runtime.ts | 9 +++ src/infra/git-read-operations.ts | 12 ++++ src/infra/git-worker.ts | 1 + 12 files changed, 150 insertions(+), 52 deletions(-) diff --git a/docs/web/control-ui/sessions-and-sidebar.md b/docs/web/control-ui/sessions-and-sidebar.md index f2d48d73756d..87a2685665f1 100644 --- a/docs/web/control-ui/sessions-and-sidebar.md +++ b/docs/web/control-ui/sessions-and-sidebar.md @@ -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. diff --git a/scripts/pr-lib/wrapper-components.txt b/scripts/pr-lib/wrapper-components.txt index a1f56b7b30e1..ee14b8e2ec32 100644 --- a/scripts/pr-lib/wrapper-components.txt +++ b/scripts/pr-lib/wrapper-components.txt @@ -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 diff --git a/src/agents/worktrees/service-preparation.ts b/src/agents/worktrees/service-preparation.ts index 85a156cc9290..38bddde94fec 100644 --- a/src/agents/worktrees/service-preparation.ts +++ b/src/agents/worktrees/service-preparation.ts @@ -157,6 +157,16 @@ export async function resolveRepository(repoRoot: string): Promise) { const failure = await removeFailedWorktree(...args); if (failure) { diff --git a/src/agents/worktrees/service.ts b/src/agents/worktrees/service.ts index 1799f688cc28..cdbf2d774b72 100644 --- a/src/agents/worktrees/service.ts +++ b/src/agents/worktrees/service.ts @@ -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 } }); } /** diff --git a/src/gateway/server-methods/projects-observed.test.ts b/src/gateway/server-methods/projects-observed.test.ts index 8e69efc01ed2..cc1368e9e0a1 100644 --- a/src/gateway/server-methods/projects-observed.test.ts +++ b/src/gateway/server-methods/projects-observed.test.ts @@ -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[0]; - const seededSessions = vi.hoisted(() => ({ store: {} as Record, })); @@ -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[] = []; 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( - async (_checkoutPath) => { - throw new Error("checkout unavailable"); - }, - ); + const resolveRepositoryIdentity = vi.fn< + Parameters[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"); }); }); diff --git a/src/gateway/server-methods/projects.test-support.ts b/src/gateway/server-methods/projects.test-support.ts index c75cd4efee44..8714cc89078c 100644 --- a/src/gateway/server-methods/projects.test-support.ts +++ b/src/gateway/server-methods/projects.test-support.ts @@ -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, diff --git a/src/gateway/server-methods/projects.test.ts b/src/gateway/server-methods/projects.test.ts index f5e8630e26e5..0e2bedbc649a 100644 --- a/src/gateway/server-methods/projects.test.ts +++ b/src/gateway/server-methods/projects.test.ts @@ -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) { 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", diff --git a/src/gateway/server-methods/projects.ts b/src/gateway/server-methods/projects.ts index 397df62b493c..01bde0405a13 100644 --- a/src/gateway/server-methods/projects.ts +++ b/src/gateway/server-methods/projects.ts @@ -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]]), ); diff --git a/src/infra/git-read-cache.ts b/src/infra/git-read-cache.ts index ecf5878ee676..cf7869cfc37b 100644 --- a/src/infra/git-read-cache.ts +++ b/src/infra/git-read-cache.ts @@ -157,6 +157,13 @@ function createReadCache( // 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": diff --git a/src/infra/git-read-operations.runtime.ts b/src/infra/git-read-operations.runtime.ts index f30b964b75e8..0928bcdc02b2 100644 --- a/src/infra/git-read-operations.runtime.ts +++ b/src/infra/git-read-operations.runtime.ts @@ -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 { 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": diff --git a/src/infra/git-read-operations.ts b/src/infra/git-read-operations.ts index d3e53cb07acb..3cfd131bd813 100644 --- a/src/infra/git-read-operations.ts +++ b/src/infra/git-read-operations.ts @@ -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 }; diff --git a/src/infra/git-worker.ts b/src/infra/git-worker.ts index 506bda025b4d..54274b09edae 100644 --- a/src/infra/git-worker.ts +++ b/src/infra/git-worker.ts @@ -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"