diff --git a/src/gateway/worker-environments/workspace-accepted-sync.test.ts b/src/gateway/worker-environments/workspace-accepted-sync.test.ts index 919536c876da..b3de97e8045f 100644 --- a/src/gateway/worker-environments/workspace-accepted-sync.test.ts +++ b/src/gateway/worker-environments/workspace-accepted-sync.test.ts @@ -17,9 +17,10 @@ import { isAcceptedWorkspacePublicationIndeterminateError, } from "./workspace-accepted-publication.js"; import { - createAcceptedWorkspacePublisherFactory, + createAcceptedWorkspacePublisherFactory as createAcceptedWorkspacePublisherFactoryRaw, recoverAcceptedWorkspacePublication, } from "./workspace-accepted-sync.js"; +import { createWorkspaceReconcileMetrics } from "./workspace-hash-memo.js"; import { serializeWorkerWorkspaceManifest, type WorkerWorkspaceManifest, @@ -69,6 +70,40 @@ function settlement(outcome: "begun" | "rolled-back" | "applied" | "committed"): return result({ stdout: `${JSON.stringify({ version: 1, outcome })}\n` }); } +function createAcceptedWorkspacePublisherFactory( + params: Omit< + Parameters[0], + "hashMemo" | "metrics" + >, +) { + const runWorkspaceCommand = params.runWorkspaceCommand; + return createAcceptedWorkspacePublisherFactoryRaw({ + ...params, + hashMemo: new Map(), + metrics: createWorkspaceReconcileMetrics(), + runWorkspaceCommand: async (command) => { + const response = await runWorkspaceCommand(command); + const returnedRef = response.stdout.trim(); + if (command.argv.at(-1) !== "memo-v1" || !/^sha256:[a-f0-9]{64}$/u.test(returnedRef)) { + return response; + } + return result({ + stdout: `${JSON.stringify({ + version: 1, + manifestRef: returnedRef, + memo: [], + metrics: { + contentHashCount: 0, + contentHashDurationMs: 0, + memoHitCount: 0, + totalDurationMs: 0, + }, + })}\n`, + }); + }, + }); +} + describe("accepted workspace publication", () => { it.skipIf(process.platform === "win32")( "waits for the staging receiver group before promoting its inodes live", diff --git a/src/gateway/worker-environments/workspace-accepted-sync.ts b/src/gateway/worker-environments/workspace-accepted-sync.ts index a4abc636bf5d..2f197baad8e3 100644 --- a/src/gateway/worker-environments/workspace-accepted-sync.ts +++ b/src/gateway/worker-environments/workspace-accepted-sync.ts @@ -10,13 +10,14 @@ import { parseAcceptedWorkspaceSettlement, type AcceptedWorkspaceSettlementOutcome, } from "./workspace-accepted-publication.js"; +import type { WorkspaceHashMemo, WorkspaceReconcileMetrics } from "./workspace-hash-memo.js"; import { serializeWorkerWorkspaceManifest, type WorkerWorkspaceManifest, } from "./workspace-manifest.js"; import { changedPaths, manifestNodes } from "./workspace-reconcile.js"; import { - parseManifestRef, + captureRemoteWorkspaceManifest, WORKER_WORKSPACE_RSYNC_DESTINATION, workerAcceptedWorkspaceRsyncReceiverPath, workerWorkspaceCommandSucceeded, @@ -67,6 +68,8 @@ function createAcceptedWorkspacePublisher(params: { localPath: string; remoteWorkspaceDir: string; remoteManifest: WorkerWorkspaceManifest; + hashMemo: WorkspaceHashMemo; + metrics: WorkspaceReconcileMetrics; }) { return async (accepted: { manifestRef: string; @@ -96,21 +99,14 @@ function createAcceptedWorkspacePublisher(params: { } const verifyAcceptedWorkspace = async () => { - const verified = await params.runWorkspaceCommand({ - transportRetry: "idempotent", - argv: [ - "node", - "-e", - REMOTE_WORKSPACE_MANIFEST_JS, - params.remoteWorkspaceDir, - accepted.manifest.baseCommit ?? "", - ...(accepted.manifest.baseCommit ? ["eligible", acceptedDigest] : []), - ], + const verifiedRef = await captureRemoteWorkspaceManifest({ + runWorkspaceCommand: params.runWorkspaceCommand, + remoteWorkspaceDir: params.remoteWorkspaceDir, + baseCommit: accepted.manifest.baseCommit, + priorManifestDigests: accepted.manifest.baseCommit ? [acceptedDigest] : [], + hashMemo: params.hashMemo, + metrics: params.metrics, }); - if (!workerWorkspaceCommandSucceeded(verified)) { - throw workspaceSyncError(verified); - } - const verifiedRef = parseManifestRef(verified.stdout.trim()); if (verifiedRef !== accepted.manifestRef) { throw new Error( `Worker workspace does not match its accepted manifest: expected ${accepted.manifestRef}, got ${verifiedRef}`, diff --git a/src/gateway/worker-environments/workspace-finalize.test.ts b/src/gateway/worker-environments/workspace-finalize.test.ts index 6f912077da5f..02e708f1c094 100644 --- a/src/gateway/worker-environments/workspace-finalize.test.ts +++ b/src/gateway/worker-environments/workspace-finalize.test.ts @@ -1,29 +1,33 @@ import { describe, expect, it, vi } from "vitest"; -import { verifyReconciledWorkspaceFinal } from "./workspace-finalize.js"; +import { + registerWorkspaceReconcileReporter, + verifyReconciledWorkspaceFinal, +} from "./workspace-finalize.js"; describe("final worker workspace fences", () => { it("rechecks remote and local stability after the final quiescence renewal", async () => { const log: string[] = []; - await verifyReconciledWorkspaceFinal( - { - manifestRef: "sha256:" + "a".repeat(64), - changed: true, - verifyStable: async () => { - log.push("remote"); - }, - verifyLocalStable: async () => { - log.push("local"); - }, + const reconciliation = { + manifestRef: "sha256:" + "a".repeat(64), + changed: true, + verifyStable: async () => { + log.push("remote"); }, - { - assertActive: async () => { - log.push("quiescence"); - }, - resume: async () => {}, + verifyLocalStable: async () => { + log.push("local"); }, - ); + }; + const outcomes: string[] = []; + registerWorkspaceReconcileReporter(reconciliation, (outcome) => outcomes.push(outcome)); + await verifyReconciledWorkspaceFinal(reconciliation, { + assertActive: async () => { + log.push("quiescence"); + }, + resume: async () => {}, + }); expect(log).toEqual(["remote", "local", "quiescence", "remote", "local"]); + expect(outcomes).toEqual(["succeeded"]); }); it("rejects a remote write observed after the final quiescence renewal", async () => { diff --git a/src/gateway/worker-environments/workspace-finalize.ts b/src/gateway/worker-environments/workspace-finalize.ts index ef025b812bae..5b2b35cfcf76 100644 --- a/src/gateway/worker-environments/workspace-finalize.ts +++ b/src/gateway/worker-environments/workspace-finalize.ts @@ -31,40 +31,74 @@ const runRetryableFinalFenceStep = async (operation: () => Promise): Promi const runResultPreservingFinalFenceStep = async (operation: () => Promise): Promise => await runFinalFenceStep(operation, "preserve-result"); +type WorkspaceReconcileOutcome = "failed" | "succeeded"; + +const workspaceReconcileReporters = new WeakMap< + WorkerWorkspaceReconcileResult, + (outcome: WorkspaceReconcileOutcome) => void +>(); + +export function registerWorkspaceReconcileReporter( + reconciliation: WorkerWorkspaceReconcileResult, + reporter: (outcome: WorkspaceReconcileOutcome) => void, +): void { + workspaceReconcileReporters.set(reconciliation, reporter); +} + +function reportWorkspaceReconcile( + reconciliation: WorkerWorkspaceReconcileResult, + outcome: WorkspaceReconcileOutcome, +): void { + const reporter = workspaceReconcileReporters.get(reconciliation); + workspaceReconcileReporters.delete(reconciliation); + reporter?.(outcome); +} + /** Rechecks both owners after renewing the remote quiescence lease. */ export async function verifyReconciledWorkspaceFinal( reconciliation: WorkerWorkspaceReconcileResult, quiescence: WorkerWorkspaceQuiescence, ): Promise { - if (reconciliation.applyPreparedStagedResult && reconciliation.publishStagedResult) { - try { - // Fence the prepared remote capture before quiescence renewal can enroll late writers. - await runRetryableFinalFenceStep(async () => await reconciliation.verifyStable()); - // Renew quiescence and freeze any writers that appeared after the prepared capture. - await runRetryableFinalFenceStep(async () => await quiescence.assertActive()); - // Keep this fence: a late writer can mutate before renewal enrolls and SIGSTOPs it. - await runRetryableFinalFenceStep(async () => await reconciliation.verifyStable()); - await reconciliation.applyPreparedStagedResult(); - await reconciliation.verifyLocalStable(); - // Renew after apply so lease expiry cannot race the final publish gate. - await runResultPreservingFinalFenceStep(async () => await quiescence.assertActive()); - // Recheck the remote owner after apply before publishing the prepared result. - await runResultPreservingFinalFenceStep(async () => await reconciliation.verifyStable()); - await runResultPreservingFinalFenceStep(async () => await reconciliation.verifyLocalStable()); - await reconciliation.publishStagedResult(); - return reconciliation.getAppliedWorkspaceResult?.(); - } catch (error) { - await reconciliation.discardPreparedStagedResult?.().catch(() => undefined); - throw error; + let succeeded = false; + try { + if (reconciliation.applyPreparedStagedResult && reconciliation.publishStagedResult) { + try { + // Fence the prepared remote capture before quiescence renewal can enroll late writers. + await runRetryableFinalFenceStep(async () => await reconciliation.verifyStable()); + // Renew quiescence and freeze any writers that appeared after the prepared capture. + await runRetryableFinalFenceStep(async () => await quiescence.assertActive()); + // Keep this fence: a late writer can mutate before renewal enrolls and SIGSTOPs it. + await runRetryableFinalFenceStep(async () => await reconciliation.verifyStable()); + await reconciliation.applyPreparedStagedResult(); + await reconciliation.verifyLocalStable(); + // Renew after apply so lease expiry cannot race the final publish gate. + await runResultPreservingFinalFenceStep(async () => await quiescence.assertActive()); + // Recheck the remote owner after apply before publishing the prepared result. + await runResultPreservingFinalFenceStep(async () => await reconciliation.verifyStable()); + await runResultPreservingFinalFenceStep( + async () => await reconciliation.verifyLocalStable(), + ); + await reconciliation.publishStagedResult(); + const applied = reconciliation.getAppliedWorkspaceResult?.(); + succeeded = true; + return applied; + } catch (error) { + await reconciliation.discardPreparedStagedResult?.().catch(() => undefined); + throw error; + } } + const runFenceStep = reconciliation.changed + ? runResultPreservingFinalFenceStep + : runRetryableFinalFenceStep; + await runFenceStep(async () => await reconciliation.verifyStable()); + await runFenceStep(async () => await reconciliation.verifyLocalStable()); + await runFenceStep(async () => await quiescence.assertActive()); + await runFenceStep(async () => await reconciliation.verifyStable()); + await runFenceStep(async () => await reconciliation.verifyLocalStable()); + const applied = reconciliation.getAppliedWorkspaceResult?.(); + succeeded = true; + return applied; + } finally { + reportWorkspaceReconcile(reconciliation, succeeded ? "succeeded" : "failed"); } - const runFenceStep = reconciliation.changed - ? runResultPreservingFinalFenceStep - : runRetryableFinalFenceStep; - await runFenceStep(async () => await reconciliation.verifyStable()); - await runFenceStep(async () => await reconciliation.verifyLocalStable()); - await runFenceStep(async () => await quiescence.assertActive()); - await runFenceStep(async () => await reconciliation.verifyStable()); - await runFenceStep(async () => await reconciliation.verifyLocalStable()); - return reconciliation.getAppliedWorkspaceResult?.(); } diff --git a/src/gateway/worker-environments/workspace-hash-memo.test.ts b/src/gateway/worker-environments/workspace-hash-memo.test.ts new file mode 100644 index 000000000000..aafd47c000cf --- /dev/null +++ b/src/gateway/worker-environments/workspace-hash-memo.test.ts @@ -0,0 +1,236 @@ +import { createHash } from "node:crypto"; +import fs from "node:fs/promises"; +import path from "node:path"; +import { afterEach, describe, expect, it, vi } from "vitest"; +import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js"; +import { runCommandWithTimeout } from "../../process/exec.js"; +import { + createWorkspaceReconcileMetrics, + MAX_WORKSPACE_HASH_MEMO_BYTES, + recordRemoteWorkspaceHashMetrics, + serializeRemoteWorkspaceHashMemo, + withWorkspaceHashMemo, +} from "./workspace-hash-memo.js"; +import { MAX_RECONCILIATION_ENTRIES, type WorkerWorkspaceManifest } from "./workspace-manifest.js"; +import { preflightWorkspaceApply, readActualWorkspaceManifest } from "./workspace-reconcile.js"; +import { REMOTE_WORKSPACE_MANIFEST_JS } from "./workspace-sync-scripts.js"; + +const tempDirs = useAutoCleanupTempDirTracker(afterEach); +afterEach(() => vi.restoreAllMocks()); + +function hashMetrics() { + return { + contentHashCount: 0, + contentHashDurationMs: 0, + memoHitCount: 0, + }; +} + +describe("workspace hash memo", () => { + it("reuses content hashes only within one reconcile stat identity", async () => { + const root = await fs.realpath(tempDirs.make("openclaw-workspace-hash-memo-")); + const target = path.join(root, "same-size.txt"); + await fs.writeFile(target, "alpha"); + const memo = new Map(); + const metrics = hashMetrics(); + let replacedManifestRef = ""; + await withWorkspaceHashMemo( + memo, + async () => { + const first = await readActualWorkspaceManifest({ root, baseCommit: null }); + const unchanged = await withWorkspaceHashMemo( + memo, + async () => await readActualWorkspaceManifest({ root, baseCommit: null }), + ); + expect(unchanged.manifestRef).toBe(first.manifestRef); + expect(metrics).toMatchObject({ contentHashCount: 1, memoHitCount: 1 }); + + await fs.writeFile(target, "bravo"); + await fs.utimes(target, new Date(), new Date(Date.now() + 1_000)); + const changed = await readActualWorkspaceManifest({ root, baseCommit: null }); + expect(changed.manifestRef).not.toBe(first.manifestRef); + expect(metrics.contentHashCount).toBe(2); + + const replacement = path.join(root, "replacement.txt"); + await fs.writeFile(replacement, "cider"); + await fs.rename(replacement, target); + const replaced = await readActualWorkspaceManifest({ root, baseCommit: null }); + expect(replaced.manifestRef).not.toBe(changed.manifestRef); + expect(metrics.contentHashCount).toBe(3); + replacedManifestRef = replaced.manifestRef; + }, + metrics, + ); + + const nextReconcileMetrics = hashMetrics(); + const nextReconcile = await withWorkspaceHashMemo( + new Map(), + async () => await readActualWorkspaceManifest({ root, baseCommit: null }), + nextReconcileMetrics, + ); + expect(nextReconcile.manifestRef).toBe(replacedManifestRef); + expect(nextReconcileMetrics).toMatchObject({ contentHashCount: 1, memoHitCount: 0 }); + }); + + it("reuses local workspace nodes within one preflight but not across fences", async () => { + const root = await fs.realpath(tempDirs.make("openclaw-workspace-preflight-memo-")); + await fs.writeFile(path.join(root, "parent"), "base"); + const baseContent = Buffer.from("base"); + const currentContent = Buffer.from("worker"); + const base: WorkerWorkspaceManifest = { + version: 1, + baseCommit: null, + entries: [ + { + path: "parent", + type: "file", + mode: 0o644, + size: baseContent.length, + sha256: createHash("sha256").update(baseContent).digest("hex"), + }, + ], + directories: [], + }; + const current: WorkerWorkspaceManifest = { + version: 1, + baseCommit: null, + entries: [ + { + path: "parent/child.txt", + type: "file", + mode: 0o644, + size: currentContent.length, + sha256: createHash("sha256").update(currentContent).digest("hex"), + }, + { + path: "parent/sibling.txt", + type: "file", + mode: 0o644, + size: currentContent.length, + sha256: createHash("sha256").update(currentContent).digest("hex"), + }, + ], + directories: ["parent"], + }; + const metrics = hashMetrics(); + const open = vi.spyOn(fs, "open"); + const parentPath = path.join(root, "parent"); + const parentSnapshots = () => open.mock.calls.filter(([file]) => file === parentPath).length; + + const first = await withWorkspaceHashMemo( + new Map(), + async () => await preflightWorkspaceApply({ root, base, current }), + metrics, + ); + expect([...first.applyPaths].toSorted()).toEqual([ + "parent", + "parent/child.txt", + "parent/sibling.txt", + ]); + expect(metrics.contentHashCount).toBe(1); + expect(parentSnapshots()).toBe(1); + + await withWorkspaceHashMemo( + new Map(), + async () => await preflightWorkspaceApply({ root, base, current }), + metrics, + ); + expect(metrics.contentHashCount).toBe(2); + expect(parentSnapshots()).toBe(2); + }); + + it("aggregates remote metrics and bounds a maximum-entry memo envelope", () => { + const aggregate = createWorkspaceReconcileMetrics(); + recordRemoteWorkspaceHashMetrics(aggregate, { + contentHashCount: 7, + contentHashDurationMs: 11, + memoHitCount: 13, + totalDurationMs: 17, + }); + recordRemoteWorkspaceHashMetrics(aggregate, { + contentHashCount: 19, + contentHashDurationMs: 23, + memoHitCount: 29, + totalDurationMs: 31, + }); + expect(aggregate).toMatchObject({ + remoteContentHashCount: 26, + remoteMemoHitCount: 42, + remoteHashDurationMs: 34, + remoteManifestDurationMs: 48, + }); + + const uint64 = "18446744073709551615"; + const memo = new Map(); + for (let index = 0; index < MAX_RECONCILIATION_ENTRIES; index += 1) { + const inode = String(index).padStart(20, "0"); + memo.set( + `worker:${uint64}:${inode}:${uint64}:${uint64}:${uint64}`, + index.toString(16).padStart(64, "0"), + ); + } + const serializedMemo = serializeRemoteWorkspaceHashMemo(memo); + const envelopeBytes = Buffer.byteLength( + `${JSON.stringify({ + version: 1, + manifestRef: `sha256:${"f".repeat(64)}`, + memo: JSON.parse(serializedMemo), + metrics: { + contentHashCount: MAX_RECONCILIATION_ENTRIES, + contentHashDurationMs: Number.MAX_SAFE_INTEGER, + memoHitCount: MAX_RECONCILIATION_ENTRIES, + totalDurationMs: Number.MAX_SAFE_INTEGER, + }, + })}\n`, + ); + expect(envelopeBytes).toBeLessThan(MAX_WORKSPACE_HASH_MEMO_BYTES); + expect(MAX_WORKSPACE_HASH_MEMO_BYTES - envelopeBytes).toBeGreaterThan(3 * 1024 * 1024); + }); + + it("reuses hashes only for matching stat identities in one remote reconcile", async () => { + const root = tempDirs.make("openclaw-remote-manifest-memo-"); + const home = path.join(root, "home"); + let workspace = path.join(root, "workspace"); + await Promise.all([fs.mkdir(home), fs.mkdir(workspace)]); + workspace = await fs.realpath(workspace); + const target = path.join(workspace, "same-size.txt"); + await fs.writeFile(target, "alpha"); + const env = { ...process.env, HOME: home }; + type MemoResponse = { + manifestRef: string; + memo: [string, string][]; + metrics: { contentHashCount: number; memoHitCount: number }; + }; + const capture = async (memo: [string, string][]): Promise => { + const result = await runCommandWithTimeout( + [process.execPath, "-e", REMOTE_WORKSPACE_MANIFEST_JS, workspace, "", "memo-v1"], + { timeoutMs: 10_000, baseEnv: env, input: JSON.stringify(memo) }, + ); + expect(result).toMatchObject({ code: 0, stderr: "" }); + return JSON.parse(result.stdout) as MemoResponse; + }; + + const first = await capture([]); + expect(first.metrics).toMatchObject({ contentHashCount: 1, memoHitCount: 0 }); + const unchanged = await capture(first.memo); + expect(unchanged.manifestRef).toBe(first.manifestRef); + expect(unchanged.metrics).toMatchObject({ contentHashCount: 0, memoHitCount: 1 }); + + await fs.writeFile(target, "bravo"); + await fs.utimes(target, new Date(), new Date(Date.now() + 1_000)); + const changed = await capture(unchanged.memo); + expect(changed.manifestRef).not.toBe(first.manifestRef); + expect(changed.metrics).toMatchObject({ contentHashCount: 1, memoHitCount: 0 }); + + const replacement = path.join(workspace, "replacement.txt"); + await fs.writeFile(replacement, "cider"); + await fs.rename(replacement, target); + const replaced = await capture(changed.memo); + expect(replaced.manifestRef).not.toBe(changed.manifestRef); + expect(replaced.metrics).toMatchObject({ contentHashCount: 1, memoHitCount: 0 }); + + const nextReconcile = await capture([]); + expect(nextReconcile.manifestRef).toBe(replaced.manifestRef); + expect(nextReconcile.metrics).toMatchObject({ contentHashCount: 1, memoHitCount: 0 }); + }); +}); diff --git a/src/gateway/worker-environments/workspace-hash-memo.ts b/src/gateway/worker-environments/workspace-hash-memo.ts new file mode 100644 index 000000000000..0d699b274328 --- /dev/null +++ b/src/gateway/worker-environments/workspace-hash-memo.ts @@ -0,0 +1,111 @@ +import { AsyncLocalStorage } from "node:async_hooks"; + +type WorkspaceHashMetrics = { + contentHashCount: number; + contentHashDurationMs: number; + memoHitCount: number; +}; + +export type WorkspaceHashMemo = Map; + +export type WorkspaceReconcileMetrics = { + gateway: WorkspaceHashMetrics; + remoteManifestCalls: number; + remoteContentHashCount: number; + remoteMemoHitCount: number; + remoteHashDurationMs: number; + remoteManifestDurationMs: number; + remoteManifestWallDurationMs: number; + localReconciliationDurationMs: number; +}; + +type RemoteWorkspaceHashMetrics = WorkspaceHashMetrics & { totalDurationMs: number }; + +export const MAX_WORKSPACE_HASH_MEMO_BYTES = 8 * 1024 * 1024; + +type WorkspaceHashContext = { + memo: WorkspaceHashMemo; + metrics?: WorkspaceHashMetrics; +}; + +const workspaceHashContext = new AsyncLocalStorage(); + +export function createWorkspaceReconcileMetrics(): WorkspaceReconcileMetrics { + return { + gateway: { + contentHashCount: 0, + contentHashDurationMs: 0, + memoHitCount: 0, + }, + remoteManifestCalls: 0, + remoteContentHashCount: 0, + remoteMemoHitCount: 0, + remoteHashDurationMs: 0, + remoteManifestDurationMs: 0, + remoteManifestWallDurationMs: 0, + localReconciliationDurationMs: 0, + }; +} + +export function activeWorkspaceHashContext(): WorkspaceHashContext | undefined { + return workspaceHashContext.getStore(); +} + +export async function withWorkspaceHashMemo( + memo: WorkspaceHashMemo, + operation: () => Promise, + metrics?: WorkspaceHashMetrics, +): Promise { + const active = workspaceHashContext.getStore(); + const inheritedMetrics = metrics ?? active?.metrics; + if (active?.memo === memo && active.metrics === inheritedMetrics) { + return await operation(); + } + return await workspaceHashContext.run({ memo, metrics: inheritedMetrics }, operation); +} + +export async function withWorkspaceHashContext(operation: () => Promise): Promise { + const active = workspaceHashContext.getStore(); + return await withWorkspaceHashMemo(active?.memo ?? new Map(), operation, active?.metrics); +} + +export function serializeRemoteWorkspaceHashMemo(memo: WorkspaceHashMemo): string { + const serialized = JSON.stringify( + [...memo] + .filter(([identity]) => identity.startsWith("worker:")) + .toSorted(([left], [right]) => left.localeCompare(right)), + ); + if (Buffer.byteLength(serialized) > MAX_WORKSPACE_HASH_MEMO_BYTES) { + throw new Error("Workspace hash memo exceeds its byte limit"); + } + return serialized; +} + +export function recordRemoteWorkspaceHashMetrics( + aggregate: WorkspaceReconcileMetrics, + metrics: RemoteWorkspaceHashMetrics, +): void { + aggregate.remoteContentHashCount += metrics.contentHashCount; + aggregate.remoteMemoHitCount += metrics.memoHitCount; + aggregate.remoteHashDurationMs += metrics.contentHashDurationMs; + aggregate.remoteManifestDurationMs += metrics.totalDurationMs; +} + +export async function measureLocalWorkspaceReconciliation( + metrics: WorkspaceReconcileMetrics, + operation: () => Promise, +): Promise { + const startedAt = performance.now(); + try { + return await operation(); + } finally { + metrics.localReconciliationDurationMs += performance.now() - startedAt; + } +} + +export function workspaceStatIdentity( + owner: "gateway" | "worker", + stats: { dev: bigint; ino: bigint; size: bigint; mtimeNs: bigint; ctimeNs: bigint }, +): string { + return `${owner}:${stats.dev}:${stats.ino}:${stats.size}:${stats.mtimeNs}:${stats.ctimeNs}`; +} diff --git a/src/gateway/worker-environments/workspace-quiescence.ts b/src/gateway/worker-environments/workspace-quiescence.ts new file mode 100644 index 000000000000..e20762854df0 --- /dev/null +++ b/src/gateway/worker-environments/workspace-quiescence.ts @@ -0,0 +1,117 @@ +import path from "node:path"; +import type { SpawnResult } from "../../process/exec.js"; +import type { WorkerWorkspaceCommand, WorkerWorkspaceQuiescence } from "./tunnel-contract.js"; +import { + REMOTE_WORKSPACE_QUIESCE_JS, + REMOTE_WORKSPACE_RENEW_QUIESCENCE_JS, + REMOTE_WORKSPACE_RESUME_JS, +} from "./workspace-quiescence-scripts.js"; +import { + waitForQuiescenceRenewal, + workerWorkspaceCommandSucceeded, + workspaceSyncError, +} from "./workspace-sync-helpers.js"; + +const WORKSPACE_QUIESCENCE_TIMEOUT_MS = 12 * 60_000; +const WORKSPACE_QUIESCENCE_RENEW_INTERVAL_MS = 4 * 60_000; + +export function createWorkerWorkspaceQuiescence(params: { + ownerSignal: AbortSignal; + sharedHost: boolean; + runWorkspaceCommand: (command: WorkerWorkspaceCommand) => Promise; +}): (remoteWorkspaceDir: string) => Promise { + return async (remoteWorkspaceDir) => { + if (!path.posix.isAbsolute(remoteWorkspaceDir)) { + throw new Error("Worker workspace quiescence path must be absolute"); + } + const hostMode = params.sharedHost ? "shared-host" : "dedicated"; + const run = async (argv: string[]) => { + const result = await params.runWorkspaceCommand({ transportRetry: "never", argv }); + if (!workerWorkspaceCommandSucceeded(result)) { + throw workspaceSyncError(result); + } + return result; + }; + const result = await run([ + "node", + "-e", + REMOTE_WORKSPACE_QUIESCE_JS, + remoteWorkspaceDir, + String(WORKSPACE_QUIESCENCE_TIMEOUT_MS), + hostMode, + ]); + const acknowledgement = /^quiesced ([a-f0-9]{32})$/u.exec(result.stdout.trim()); + if (!acknowledgement) { + throw new Error("Worker workspace quiescence returned an invalid acknowledgement"); + } + const nonce = acknowledgement[1]!; + let resumed = false; + let renewalFailure: unknown; + const renewalAbort = new AbortController(); + const abortRenewal = () => renewalAbort.abort(params.ownerSignal.reason); + params.ownerSignal.addEventListener("abort", abortRenewal, { once: true }); + let renewalQueue = Promise.resolve(); + const renew = (validationMode: "heartbeat" | "final") => { + const operation = renewalQueue.then(async () => { + const renewedResult = await run([ + "node", + "-e", + REMOTE_WORKSPACE_RENEW_QUIESCENCE_JS, + remoteWorkspaceDir, + nonce, + String(WORKSPACE_QUIESCENCE_TIMEOUT_MS), + validationMode, + hostMode, + ]); + if (renewedResult.stdout.trim() !== `renewed ${nonce}`) { + throw new Error( + "Worker workspace quiescence renewal returned an invalid acknowledgement", + ); + } + }); + renewalQueue = operation.catch(() => undefined); + return operation; + }; + const renewalLoop = (async () => { + while (!renewalAbort.signal.aborted) { + if ( + !(await waitForQuiescenceRenewal( + renewalAbort.signal, + WORKSPACE_QUIESCENCE_RENEW_INTERVAL_MS, + )) + ) { + return; + } + try { + await renew("heartbeat"); + } catch (error) { + renewalFailure = error; + return; + } + } + })(); + return { + assertActive: async () => { + if (resumed) { + throw new Error("Worker workspace quiescence was already released"); + } + if (renewalFailure) { + throw new Error("Worker workspace quiescence renewal failed", { + cause: renewalFailure, + }); + } + await renew("final"); + }, + resume: async () => { + if (resumed) { + return; + } + params.ownerSignal.removeEventListener("abort", abortRenewal); + renewalAbort.abort(); + await renewalLoop; + await run(["node", "-e", REMOTE_WORKSPACE_RESUME_JS, remoteWorkspaceDir, nonce]); + resumed = true; + }, + }; + }; +} diff --git a/src/gateway/worker-environments/workspace-reconcile-apply.ts b/src/gateway/worker-environments/workspace-reconcile-apply.ts index 5cbfaf13a912..d93a74b69f86 100644 --- a/src/gateway/worker-environments/workspace-reconcile-apply.ts +++ b/src/gateway/worker-environments/workspace-reconcile-apply.ts @@ -1,6 +1,11 @@ import { createHash, randomBytes } from "node:crypto"; import fs from "node:fs/promises"; import { isAcceptedWorkspacePublicationIndeterminateError } from "./workspace-accepted-publication.js"; +import { + activeWorkspaceHashContext, + withWorkspaceHashContext, + withWorkspaceHashMemo, +} from "./workspace-hash-memo.js"; import { MAX_RECONCILIATION_ENTRIES, type WorkerWorkspaceManifest, @@ -46,6 +51,15 @@ export async function applyStagedWorkerWorkspace(params: { conflictPaths: string[]; }) => Promise; }): Promise { + return await withWorkspaceHashContext( + async () => await applyStagedWorkerWorkspaceWithMemo(params), + ); +} + +async function applyStagedWorkerWorkspaceWithMemo( + params: Parameters[0], +): Promise { + const { memo: hashMemo, metrics } = activeWorkspaceHashContext()!; const root = await fs.realpath(params.root); const preserveDirectories = new Set(reconciliationDirectories(params.current.directories)); // Git workspaces must keep the eligibility boundary established at dispatch. @@ -53,6 +67,26 @@ export async function applyStagedWorkerWorkspace(params: { const includePaths = params.current.baseCommit ? new Set([...manifestNodes(params.base).keys(), ...manifestNodes(params.current).keys()]) : undefined; + const createApplyResult = ( + actual: Awaited>, + conflictPaths: string[], + ): WorkerWorkspaceApplyResult => ({ + ...actual, + conflictPaths, + verifyLocalStable: async () => + await withWorkspaceHashMemo( + hashMemo, + async () => + await assertActualWorkspaceManifest({ + root, + expectedRef: actual.manifestRef, + baseCommit: actual.manifest.baseCommit, + preserveDirectories, + includePaths, + }), + metrics, + ), + }); const preflight = await preflightWorkspaceApply({ root, base: params.base, @@ -76,18 +110,7 @@ export async function applyStagedWorkerWorkspace(params: { const conflictPaths = retainedConflictPaths(preflight, preflight.applyPaths); await params.publishAcceptedManifest?.({ ...actual, conflictPaths }); params.journal.commit(actual.manifestRef); - return { - ...actual, - conflictPaths, - verifyLocalStable: async () => - await assertActualWorkspaceManifest({ - root, - expectedRef: actual.manifestRef, - baseCommit: actual.manifest.baseCommit, - preserveDirectories, - includePaths, - }), - }; + return createApplyResult(actual, conflictPaths); } const baseByPath = new Map( reconciliationEntries(params.base.entries).map((entry) => [entry.path, entry]), @@ -197,18 +220,7 @@ export async function applyStagedWorkerWorkspace(params: { const conflictPaths = retainedConflictPaths(finalPreflight, preflight.applyPaths); await params.publishAcceptedManifest?.({ ...actual, conflictPaths }); params.journal.commit(actual.manifestRef); - return { - ...actual, - conflictPaths, - verifyLocalStable: async () => - await assertActualWorkspaceManifest({ - root, - expectedRef: actual.manifestRef, - baseCommit: actual.manifest.baseCommit, - preserveDirectories, - includePaths, - }), - }; + return createApplyResult(actual, conflictPaths); } catch (error) { // Transport or settlement timeouts are observation evidence, never authority // for an inverse operation; recovery owns restoring both sides. diff --git a/src/gateway/worker-environments/workspace-reconcile-core.ts b/src/gateway/worker-environments/workspace-reconcile-core.ts index 4b4288aaa106..89753dbe74aa 100644 --- a/src/gateway/worker-environments/workspace-reconcile-core.ts +++ b/src/gateway/worker-environments/workspace-reconcile-core.ts @@ -2,6 +2,7 @@ import { createHash } from "node:crypto"; import fs from "node:fs/promises"; import path from "node:path"; import { FsSafeError, root as openFsSafeRoot } from "../../infra/fs-safe.js"; +import { activeWorkspaceHashContext, withWorkspaceHashMemo } from "./workspace-hash-memo.js"; import { MAX_RECONCILIATION_ENTRIES, MAX_RECONCILIATION_TOTAL_BYTES, @@ -324,6 +325,7 @@ export async function inspectAcceptedWorkerWorkspace(params: { current: WorkerWorkspaceManifest; }): Promise { const root = await fs.realpath(params.root); + const { memo: hashMemo, metrics } = activeWorkspaceHashContext() ?? {}; const preserveDirectories = new Set(reconciliationDirectories(params.current.directories)); const actual = await readActualWorkspaceManifest({ root, @@ -341,16 +343,20 @@ export async function inspectAcceptedWorkerWorkspace(params: { const conflictPaths = params.allowAdvancedLocalState ? retainedConflictPaths(preflight) : preflight.conflictPaths; + const verifyLocalStable = async () => + await assertActualWorkspaceManifest({ + root, + expectedRef: actual.manifestRef, + baseCommit: actual.manifest.baseCommit, + preserveDirectories, + }); return { ...actual, conflictPaths, verifyLocalStable: async () => - await assertActualWorkspaceManifest({ - root, - expectedRef: actual.manifestRef, - baseCommit: actual.manifest.baseCommit, - preserveDirectories, - }), + hashMemo + ? await withWorkspaceHashMemo(hashMemo, verifyLocalStable, metrics) + : await verifyLocalStable(), }; } @@ -526,6 +532,18 @@ export async function preflightWorkspaceApply(params: { const applyPaths = new Set(); const conflicts = new Set(); const blockingConflicts = new Set(); + // Node snapshots may be shared only inside this pass. Separate preflight + // calls are concurrency fences and must stat paths again. + const localNodes = new Map>(); + const localNode = (entryPath: string): Promise => { + const existing = localNodes.get(entryPath); + if (existing) { + return existing; + } + const node = localWorkspaceNode(params.root, entryPath); + localNodes.set(entryPath, node); + return node; + }; for (const entryPath of paths) { if (hasPathAncestor(blockingConflicts, entryPath)) { continue; @@ -551,7 +569,7 @@ export async function preflightWorkspaceApply(params: { const baseAncestor = baseNodes.get(ancestor); const currentAncestor = currentNodes.get(ancestor); if (!baseAncestor && !currentAncestor) { - const localAncestor = await localWorkspaceNode(params.root, ancestor); + const localAncestor = await localNode(ancestor); if (localAncestor && localAncestor.type !== "directory") { conflicts.add(ancestor); blockingConflicts.add(ancestor); @@ -560,7 +578,7 @@ export async function preflightWorkspaceApply(params: { } continue; } - const localAncestor = await localWorkspaceNode(params.root, ancestor); + const localAncestor = await localNode(ancestor); const localStructurallyMatchesBase = localAncestor?.type === "directory" && baseAncestor?.type === "directory" ? true @@ -588,7 +606,7 @@ export async function preflightWorkspaceApply(params: { baseAncestor && baseAncestor.type !== "directory" && !sameEntry(baseAncestor, currentNodes.get(ancestor)) && - sameEntry(await localWorkspaceNode(params.root, ancestor), baseAncestor) + sameEntry(await localNode(ancestor), baseAncestor) ) { replacedBaseAncestor = true; break; @@ -597,7 +615,7 @@ export async function preflightWorkspaceApply(params: { if (replacedBaseAncestor) { local = undefined; } else { - local = await localWorkspaceNode(params.root, entryPath); + local = await localNode(entryPath); if ( local?.type === "directory" && (!baseNodes.has(entryPath) || !currentNodes.has(entryPath)) && diff --git a/src/gateway/worker-environments/workspace-reconcile-fs.ts b/src/gateway/worker-environments/workspace-reconcile-fs.ts index 79bb161db91f..a2eb0b52d758 100644 --- a/src/gateway/worker-environments/workspace-reconcile-fs.ts +++ b/src/gateway/worker-environments/workspace-reconcile-fs.ts @@ -4,6 +4,7 @@ import fs from "node:fs/promises"; import path from "node:path"; import { isPathInside, resolveOpenedFileRealPathForHandle } from "../../infra/fs-safe.js"; import { runCommandBuffered } from "../../process/exec.js"; +import { activeWorkspaceHashContext, workspaceStatIdentity } from "./workspace-hash-memo.js"; import { gitFileMode, MAX_RECONCILIATION_FILE_BYTES, @@ -26,44 +27,54 @@ async function readOpenedWorkspaceFile(params: { expectedPath: string; root?: string; }): Promise { - const before = await params.handle.stat(); + const { memo: hashMemo, metrics } = activeWorkspaceHashContext() ?? {}; + const before = await params.handle.stat({ bigint: true }); const realPath = await resolveOpenedFileRealPathForHandle(params.handle, params.expectedPath); if (!before.isFile() || (params.root && !isPathInside(params.root, realPath))) { throw new Error("Gateway workspace file changed while it was being read"); } - if (before.size > MAX_RECONCILIATION_FILE_BYTES) { + if (before.size > BigInt(MAX_RECONCILIATION_FILE_BYTES)) { return { type: "unsupported" }; } - const hash = createHash("sha256"); - const buffer = Buffer.allocUnsafe(64 * 1024); - let size = 0; - for (;;) { - const { bytesRead } = await params.handle.read(buffer, 0, buffer.length, size); - if (bytesRead === 0) { - break; + const identity = workspaceStatIdentity("gateway", before); + let sha256 = hashMemo?.get(identity); + let size = Number(before.size); + if (sha256) { + if (metrics) { + metrics.memoHitCount += 1; } - size += bytesRead; - if (size > MAX_RECONCILIATION_FILE_BYTES) { - return { type: "unsupported" }; + } else { + const hashStartedAt = performance.now(); + const hash = createHash("sha256"); + const buffer = Buffer.allocUnsafe(64 * 1024); + size = 0; + for (;;) { + const { bytesRead } = await params.handle.read(buffer, 0, buffer.length, size); + if (bytesRead === 0) { + break; + } + size += bytesRead; + if (size > MAX_RECONCILIATION_FILE_BYTES) { + return { type: "unsupported" }; + } + hash.update(buffer.subarray(0, bytesRead)); + } + sha256 = hash.digest("hex"); + if (metrics) { + metrics.contentHashCount += 1; + metrics.contentHashDurationMs += performance.now() - hashStartedAt; } - hash.update(buffer.subarray(0, bytesRead)); } - const after = await params.handle.stat(); - if ( - after.size !== size || - after.size !== before.size || - after.mtimeMs !== before.mtimeMs || - after.ctimeMs !== before.ctimeMs || - after.ino !== before.ino || - after.dev !== before.dev - ) { + const after = await params.handle.stat({ bigint: true }); + if (after.size !== BigInt(size) || workspaceStatIdentity("gateway", after) !== identity) { throw new Error("Gateway workspace file changed while it was being read"); } + hashMemo?.set(identity, sha256); return { type: "file", - mode: gitFileMode(after.mode & 0o777), + mode: gitFileMode(Number(after.mode & 0o777n)), size, - sha256: hash.digest("hex"), + sha256, }; } @@ -77,7 +88,11 @@ export async function readWorkspaceFileSnapshot( constants.O_RDONLY | constants.O_NOFOLLOW | constants.O_NONBLOCK, ); try { - return await readOpenedWorkspaceFile({ handle, expectedPath: absolute, root }); + return await readOpenedWorkspaceFile({ + handle, + expectedPath: absolute, + root, + }); } finally { await handle.close(); } diff --git a/src/gateway/worker-environments/workspace-result-staging.ts b/src/gateway/worker-environments/workspace-result-staging.ts index 4e87402957d3..92f1a8bebdbd 100644 --- a/src/gateway/worker-environments/workspace-result-staging.ts +++ b/src/gateway/worker-environments/workspace-result-staging.ts @@ -4,6 +4,11 @@ import os from "node:os"; import path from "node:path"; import { runCommandBuffered, runCommandWithTimeout } from "../../process/exec.js"; import type { WorkerWorkspaceReconcileRequest } from "./tunnel-contract.js"; +import { + activeWorkspaceHashContext, + withWorkspaceHashContext, + withWorkspaceHashMemo, +} from "./workspace-hash-memo.js"; import { MAX_RECONCILIATION_ENTRIES, MAX_RECONCILIATION_FILE_BYTES, @@ -440,6 +445,14 @@ export async function applyStagedWorkerWorkspaceResult(params: { conflictPaths: string[]; }) => Promise; }): Promise { + return await withWorkspaceHashContext( + async () => await applyStagedWorkerWorkspaceResultWithMemo(params), + ); +} + +async function applyStagedWorkerWorkspaceResultWithMemo( + params: Parameters[0], +): Promise { const root = await fs.realpath(params.root); const staged = await loadStagedWorkerWorkspace(root, params.stagedResultRef); if (params.alreadyAccepted || staged.baseManifestRef !== params.expectedBaseManifestRef) { @@ -523,6 +536,9 @@ async function prepareRequestedWorkerWorkspaceResult(params: { throw new Error("Cloud workspace durable result staging was not requested"); } const candidateRef = preparedWorkerWorkspaceResultRef(stagedResult.ref); + const active = activeWorkspaceHashContext(); + const hashMemo = active?.memo ?? new Map(); + const metrics = active?.metrics; let appliedWorkspaceResult: WorkerWorkspaceApplyResult | undefined; await stageWorkerWorkspaceResult({ root: params.request.localPath, @@ -536,13 +552,18 @@ async function prepareRequestedWorkerWorkspaceResult(params: { return { applyPreparedStagedResult: async () => { const root = await ensureWorkerWorkspaceResultRepository(params.request.localPath); - appliedWorkspaceResult = await applyStagedWorkerWorkspaceResult({ - root, - stagedResultRef: candidateRef, - expectedBaseManifestRef: params.request.baseManifestRef, - journal: params.request.journal, - publishAcceptedManifest: params.publishAcceptedManifest, - }); + appliedWorkspaceResult = await withWorkspaceHashMemo( + hashMemo, + async () => + await applyStagedWorkerWorkspaceResult({ + root, + stagedResultRef: candidateRef, + expectedBaseManifestRef: params.request.baseManifestRef, + journal: params.request.journal, + publishAcceptedManifest: params.publishAcceptedManifest, + }), + metrics, + ); }, getAppliedWorkspaceResult: () => appliedWorkspaceResult, verifyLocalStable: async () => { diff --git a/src/gateway/worker-environments/workspace-sync-helpers.ts b/src/gateway/worker-environments/workspace-sync-helpers.ts index f2138f58287e..363d2c70d554 100644 --- a/src/gateway/worker-environments/workspace-sync-helpers.ts +++ b/src/gateway/worker-environments/workspace-sync-helpers.ts @@ -3,6 +3,7 @@ import fs from "node:fs/promises"; import path from "node:path"; import { setTimeout as delay } from "node:timers/promises"; import { isRecord } from "@openclaw/normalization-core/record-coerce"; +import { z } from "zod"; import { redactSensitiveText } from "../../logging/redact.js"; import type { CommandOptions, SpawnResult } from "../../process/exec.js"; import { @@ -12,9 +13,39 @@ import { workerSshRemoteCommand, } from "./ssh.js"; import type { WorkerWorkspaceCommand, WorkerWorkspaceSyncRequest } from "./tunnel-contract.js"; +import { + recordRemoteWorkspaceHashMetrics, + serializeRemoteWorkspaceHashMemo, + type WorkspaceHashMemo, + type WorkspaceReconcileMetrics, +} from "./workspace-hash-memo.js"; +import { MAX_RECONCILIATION_ENTRIES } from "./workspace-manifest.js"; import { REMOTE_WORKSPACE_MANIFEST_JS } from "./workspace-sync-scripts.js"; const MANIFEST_REF_PATTERN = /^sha256:[a-f0-9]{64}$/u; +const WORKER_HASH_IDENTITY_PATTERN = /^worker:\d+:\d+:\d+:\d+:\d+$/u; +const SHA256_PATTERN = /^[a-f0-9]{64}$/u; +const remoteWorkspaceManifestEnvelopeSchema = z + .object({ + version: z.literal(1), + manifestRef: z.string().regex(MANIFEST_REF_PATTERN), + memo: z + .array( + z.tuple([z.string().regex(WORKER_HASH_IDENTITY_PATTERN), z.string().regex(SHA256_PATTERN)]), + ) + .max(MAX_RECONCILIATION_ENTRIES), + metrics: z + .object({ + contentHashCount: z.number().finite().nonnegative(), + contentHashDurationMs: z.number().finite().nonnegative(), + memoHitCount: z.number().finite().nonnegative(), + totalDurationMs: z.number().finite().nonnegative(), + }) + .strict(), + }) + .strict(); +const INBOUND_QUOTA_INITIAL_POLL_MS = 25; +const INBOUND_QUOTA_MAX_POLL_MS = 250; export const WORKER_WORKSPACE_RSYNC_DESTINATION = "openclaw-rsync-destination"; export type WorkerWorkspaceActionsOptions = { @@ -200,33 +231,55 @@ export async function resolveRemoteWorkspaceManifest( ); } -export async function verifyRemoteWorkspaceManifest(params: { +export async function captureRemoteWorkspaceManifest(params: { runWorkspaceCommand: (command: WorkerWorkspaceCommand) => Promise; remoteWorkspaceDir: string; baseCommit: string | null; - baseDigest: string; - expectedRef: string; -}): Promise { - const expectedDigest = params.expectedRef.slice("sha256:".length); - const verified = await params.runWorkspaceCommand({ - transportRetry: "idempotent", - argv: [ - "node", - "-e", - REMOTE_WORKSPACE_MANIFEST_JS, - params.remoteWorkspaceDir, - params.baseCommit ?? "", - // Seed both manifests so a deleted path recreated under a new ignore rule - // still invalidates the fence. - ...(params.baseCommit ? ["eligible", expectedDigest, params.baseDigest] : []), - ], - }); - if (!workerWorkspaceCommandSucceeded(verified)) { - throw workspaceSyncError(verified); + priorManifestDigests: readonly string[]; + hashMemo: WorkspaceHashMemo; + metrics: WorkspaceReconcileMetrics; +}): Promise { + params.metrics.remoteManifestCalls += 1; + const startedAt = performance.now(); + const captured = await params + .runWorkspaceCommand({ + transportRetry: "idempotent", + argv: [ + "node", + "-e", + REMOTE_WORKSPACE_MANIFEST_JS, + params.remoteWorkspaceDir, + params.baseCommit ?? "", + ...(params.baseCommit ? ["eligible"] : []), + ...params.priorManifestDigests, + "memo-v1", + ], + input: serializeRemoteWorkspaceHashMemo(params.hashMemo), + }) + .finally(() => { + params.metrics.remoteManifestWallDurationMs += performance.now() - startedAt; + }); + if (!workerWorkspaceCommandSucceeded(captured)) { + throw workspaceSyncError(captured); } - if (parseManifestRef(verified.stdout.trim()) !== params.expectedRef) { - throw new Error("Cloud workspace changed during final reconciliation"); + let response; + try { + response = remoteWorkspaceManifestEnvelopeSchema.parse(JSON.parse(captured.stdout)); + } catch (error) { + throw new Error("Worker workspace manifest returned an invalid memo response", { + cause: error, + }); } + for (const identity of params.hashMemo.keys()) { + if (identity.startsWith("worker:")) { + params.hashMemo.delete(identity); + } + } + for (const [identity, sha256] of response.memo) { + params.hashMemo.set(identity, sha256); + } + recordRemoteWorkspaceHashMetrics(params.metrics, response.metrics); + return response.manifestRef; } export async function probeWorkspaceGitMode(params: { @@ -400,7 +453,10 @@ export async function runBoundedInboundRsync(params: { () => true, ); let quotaError: Error | undefined; - while (!(await Promise.race([transferSettled, delay(25).then(() => false)]))) { + let pollIntervalMs = INBOUND_QUOTA_INITIAL_POLL_MS; + // Rsync reports logical updates, not partial files or retry residue. Back off + // the canonical tree scan, then always recheck once more before acceptance. + while (!(await Promise.race([transferSettled, delay(pollIntervalMs).then(() => false)]))) { const usage = await inboundDirectoryUsage(params.destinationRoot, { bytes: params.totalByteLimit, entries: params.entryLimit, @@ -412,6 +468,7 @@ export async function runBoundedInboundRsync(params: { quotaAbort.abort(quotaError); break; } + pollIntervalMs = Math.min(pollIntervalMs * 2, INBOUND_QUOTA_MAX_POLL_MS); } let result: SpawnResult; try { diff --git a/src/gateway/worker-environments/workspace-sync-local.ts b/src/gateway/worker-environments/workspace-sync-local.ts index 5f18eeab4793..05abe56bbc1e 100644 --- a/src/gateway/worker-environments/workspace-sync-local.ts +++ b/src/gateway/worker-environments/workspace-sync-local.ts @@ -317,3 +317,27 @@ export async function createGitTransferList(params: { }); return outputPath; } + +export async function filterExistingGitTransferList(params: { + gitRoot: string; + preparedListPath: string; + outputPath: string; +}): Promise { + const output = await fs.open(params.outputPath, "wx", 0o600); + try { + for await (const file of readNulFile(params.preparedListPath)) { + const stats = await fs.lstat(path.join(params.gitRoot, file)).catch((error: unknown) => { + if (hasNodeErrorCode(error, "ENOENT")) { + return undefined; + } + throw error; + }); + if (stats?.isFile() || stats?.isSymbolicLink()) { + await output.write(`${file}\0`); + } + } + } finally { + await output.close(); + } + return params.outputPath; +} diff --git a/src/gateway/worker-environments/workspace-sync-scripts.ts b/src/gateway/worker-environments/workspace-sync-scripts.ts index 4c35d698372a..b70e8aa328b7 100644 --- a/src/gateway/worker-environments/workspace-sync-scripts.ts +++ b/src/gateway/worker-environments/workspace-sync-scripts.ts @@ -4,6 +4,8 @@ import { } from "./workspace-manifest-remote-script.js"; export { REMOTE_WORKSPACE_ACCEPTED_TRANSACTION_JS } from "./workspace-accepted-remote-script.js"; export { REMOTE_GIT_WORKSPACE_RETRY_RESET_JS } from "./workspace-mutation-remote-script.js"; +import { MAX_WORKSPACE_HASH_MEMO_BYTES, workspaceStatIdentity } from "./workspace-hash-memo.js"; +import { MAX_RECONCILIATION_ENTRIES } from "./workspace-manifest.js"; import { DERIVED_WORKSPACE_DIRECTORY_NAMES, DERIVED_WORKSPACE_FILE_NAMES, @@ -68,16 +70,46 @@ const DERIVED_WORKSPACE_DIRECTORY_NAMES = ${JSON.stringify(DERIVED_WORKSPACE_DIR const DERIVED_WORKSPACE_FILE_NAMES = ${JSON.stringify(DERIVED_WORKSPACE_FILE_NAMES)}; const DERIVED_WORKSPACE_FILE_SUFFIXES = ${JSON.stringify(DERIVED_WORKSPACE_FILE_SUFFIXES)}; const isDerivedWorkspacePath = ${isDerivedWorkspacePath.toString()}; +const workspaceStatIdentity = ${workspaceStatIdentity.toString()}; +const MAX_RECONCILIATION_ENTRIES = ${MAX_RECONCILIATION_ENTRIES}; +const MAX_HASH_MEMO_BYTES = ${MAX_WORKSPACE_HASH_MEMO_BYTES}; const root = fs.realpathSync(process.argv[1]); const requestedBaseCommit = process.argv[2] || null; const eligibleOnly = process.argv[3] === "eligible"; const requestedManifestDigest = process.argv[3] === "resolve" ? process.argv[4] : null; const publishedManifestDigest = process.argv[3] === "publish" ? process.argv[4] : null; -const priorManifestDigests = [...new Set(process.argv.slice(4).filter(Boolean))]; +const memoMode = process.argv.at(-1) === "memo-v1"; +const priorManifestDigests = [ + ...new Set(process.argv.slice(4).filter((value) => value && value !== "memo-v1")), +]; const entriesByPath = new Map(); +const usedHashMemo = new Map(); +const metrics = { contentHashCount: 0, contentHashDurationMs: 0, memoHitCount: 0 }; +const startedAt = performance.now(); function fail(message) { throw new Error(message); } +function readHashMemo() { + if (!memoMode) return new Map(); + const raw = fs.readFileSync(0, "utf8"); + if (Buffer.byteLength(raw) > MAX_HASH_MEMO_BYTES) { + fail("workspace hash memo exceeds its byte limit"); + } + let entries; + try { + entries = JSON.parse(raw); + } catch { + fail("invalid workspace hash memo"); + } + if ( + !Array.isArray(entries) || + entries.length > MAX_RECONCILIATION_ENTRIES + ) { + fail("invalid workspace hash memo"); + } + return new Map(entries); +} +const hashMemo = readHashMemo(); ${REMOTE_WORKSPACE_MANIFEST_CANONICAL_JS} function addEntry(relative) { if ( @@ -213,12 +245,40 @@ async function hashFiles() { if (entry.type !== "file") { continue; } - const hash = crypto.createHash("sha256"); - const stream = fs.createReadStream(path.join(root, entry.path)); - for await (const chunk of stream) { - hash.update(chunk); + const absolute = path.join(root, entry.path); + const handle = await fs.promises.open( + absolute, + fs.constants.O_RDONLY | fs.constants.O_NOFOLLOW | fs.constants.O_NONBLOCK, + ); + try { + const before = await handle.stat({ bigint: true }); + if (!before.isFile()) fail("worker workspace file changed while it was being read"); + const identity = workspaceStatIdentity("worker", before); + let sha256 = hashMemo.get(identity); + if (sha256) { + metrics.memoHitCount += 1; + } else { + const hashStartedAt = performance.now(); + const hash = crypto.createHash("sha256"); + const stream = handle.createReadStream({ autoClose: false }); + for await (const chunk of stream) { + hash.update(chunk); + } + sha256 = hash.digest("hex"); + metrics.contentHashCount += 1; + metrics.contentHashDurationMs += performance.now() - hashStartedAt; + } + const after = await handle.stat({ bigint: true }); + if (workspaceStatIdentity("worker", after) !== identity) { + fail("worker workspace file changed while it was being read"); + } + entry.mode = Number(after.mode & 0o777n); + entry.size = Number(after.size); + entry.sha256 = sha256; + usedHashMemo.set(identity, sha256); + } finally { + await handle.close(); } - entry.sha256 = hash.digest("hex"); } return entries; } @@ -267,7 +327,20 @@ async function main() { const baseCommit = requestedBaseCommit; const manifest = serializeManifest(baseCommit, entries); const digest = publishManifest(manifestRoot, manifest); - process.stdout.write("sha256:" + digest + "\n"); + const manifestRef = "sha256:" + digest; + const measured = { ...metrics, totalDurationMs: performance.now() - startedAt }; + if (memoMode) { + process.stdout.write(JSON.stringify({ + version: 1, + manifestRef, + memo: [...usedHashMemo].sort((left, right) => + left[0] < right[0] ? -1 : left[0] > right[0] ? 1 : 0, + ), + metrics: measured, + }) + "\n"); + } else { + process.stdout.write(manifestRef + "\n"); + } } main().catch((error) => { process.stderr.write(String(error && error.stack ? error.stack : error) + "\n"); diff --git a/src/gateway/worker-environments/workspace-sync-tunnel.test.ts b/src/gateway/worker-environments/workspace-sync-tunnel.test.ts index fcda18c3515d..5af9be88af78 100644 --- a/src/gateway/worker-environments/workspace-sync-tunnel.test.ts +++ b/src/gateway/worker-environments/workspace-sync-tunnel.test.ts @@ -517,6 +517,11 @@ describe("worker tunnel manager", () => { entry.argv[0] === "rsync" && entry.argv.some((arg) => arg.startsWith("--files-from=")), ); expect(transfers.map((entry) => rsyncArgvPort(entry.argv))).toEqual([2222, 22]); + const fileLists = transfers.map((entry) => + entry.argv.find((arg) => arg.startsWith("--files-from="))!.slice(13), + ); + expect(new Set(fileLists.map((file) => path.dirname(file))).size).toBe(1); + expect(fileLists.map((file) => path.basename(file))).toEqual(["attempt-0", "attempt-1"]); for (const transfer of transfers) { expect(transfer.argv).toContain("--delete-delay"); expect(transfer.argv).not.toContain("--delete-excluded"); diff --git a/src/gateway/worker-environments/workspace-sync.test.ts b/src/gateway/worker-environments/workspace-sync.test.ts index 7eeb4d0abdc4..66114a1b2388 100644 --- a/src/gateway/worker-environments/workspace-sync.test.ts +++ b/src/gateway/worker-environments/workspace-sync.test.ts @@ -5,6 +5,7 @@ import { afterEach, describe, expect, it, vi } from "vitest"; import type { CommandOptions, SpawnResult } from "../../process/exec.js"; import type { PreparedWorkerSsh } from "./ssh.js"; import { rsyncArgvPort, sshArgvPort } from "./worker-ssh-argv.test-support.js"; +import { runBoundedInboundRsync } from "./workspace-sync-helpers.js"; import { createWorkerWorkspaceRsyncTransport } from "./workspace-sync-transport.js"; import { createWorkerWorkspaceActions } from "./workspace-sync.js"; @@ -189,3 +190,65 @@ describe("worker workspace rsync transport retry", () => { } }); }); + +describe("bounded inbound workspace transfer", () => { + it("aborts an in-flight transfer when the destination crosses quota", async () => { + const destinationRoot = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-rsync-quota-")); + let transferSignal: AbortSignal | undefined; + try { + const runTask = vi.fn(async (_argv: string[], options: CommandOptions) => { + transferSignal = options.signal; + await fs.writeFile(path.join(destinationRoot, "oversized"), "over quota"); + return await new Promise((_resolve, reject) => { + const abort = () => { + const reason = options.signal?.reason; + reject(reason instanceof Error ? reason : new Error("aborted")); + }; + options.signal?.addEventListener("abort", abort, { once: true }); + if (options.signal?.aborted) { + abort(); + } + }); + }); + + await expect( + runBoundedInboundRsync({ + argv: ["rsync"], + destinationRoot, + entryLimit: 10, + totalByteLimit: 1, + ownerSignal: new AbortController().signal, + runTask, + timeoutMs: 10_000, + }), + ).rejects.toThrow("inbound transfer exceeds"); + expect(transferSignal?.aborted).toBe(true); + } finally { + await fs.rm(destinationRoot, { recursive: true, force: true }); + } + }); + + it("rejects a completed over-quota transfer in the authoritative final scan", async () => { + const destinationRoot = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-rsync-final-quota-")); + try { + const runTask = vi.fn(async () => { + await fs.writeFile(path.join(destinationRoot, "oversized"), "over quota"); + return result(); + }); + + await expect( + runBoundedInboundRsync({ + argv: ["rsync"], + destinationRoot, + entryLimit: 10, + totalByteLimit: 1, + ownerSignal: new AbortController().signal, + runTask, + timeoutMs: 10_000, + }), + ).rejects.toThrow("inbound transfer exceeds"); + } finally { + await fs.rm(destinationRoot, { recursive: true, force: true }); + } + }); +}); diff --git a/src/gateway/worker-environments/workspace-sync.ts b/src/gateway/worker-environments/workspace-sync.ts index 9f094b6dd279..b4ac9a7c7e83 100644 --- a/src/gateway/worker-environments/workspace-sync.ts +++ b/src/gateway/worker-environments/workspace-sync.ts @@ -2,6 +2,7 @@ import { randomBytes } from "node:crypto"; import fs from "node:fs/promises"; import os from "node:os"; import path from "node:path"; +import { createSubsystemLogger } from "../../logging/subsystem.js"; import type { CommandOptions, SpawnResult } from "../../process/exec.js"; import { type PreparedWorkerSsh, runWorkerSshCandidates, workerSshCommandOptions } from "./ssh.js"; import { @@ -17,12 +18,16 @@ import { createAcceptedWorkspacePublisherFactory, recoverAcceptedWorkspacePublication, } from "./workspace-accepted-sync.js"; -import { DERIVED_WORKSPACE_RSYNC_EXCLUDES } from "./workspace-path-exclusions.js"; +import { registerWorkspaceReconcileReporter } from "./workspace-finalize.js"; import { - REMOTE_WORKSPACE_QUIESCE_JS, - REMOTE_WORKSPACE_RENEW_QUIESCENCE_JS, - REMOTE_WORKSPACE_RESUME_JS, -} from "./workspace-quiescence-scripts.js"; + createWorkspaceReconcileMetrics, + MAX_WORKSPACE_HASH_MEMO_BYTES, + measureLocalWorkspaceReconciliation, + withWorkspaceHashMemo, + type WorkspaceReconcileMetrics, +} from "./workspace-hash-memo.js"; +import { DERIVED_WORKSPACE_RSYNC_EXCLUDES } from "./workspace-path-exclusions.js"; +import { createWorkerWorkspaceQuiescence } from "./workspace-quiescence.js"; import { applyStagedWorkerWorkspace, assertWorkspaceMatchesManifest, @@ -39,6 +44,7 @@ import { workerWorkspaceTransferPaths, } from "./workspace-result-staging.js"; import { + captureRemoteWorkspaceManifest, createWorkerWorkspaceRsyncReceiverPathFactory, parseManifestRef, parseRemoteWorkspaceSetup, @@ -47,8 +53,6 @@ import { resolveRemoteWorkspaceManifest, stableWorkerPathComponent, validateWorkspaceSyncRequest, - verifyRemoteWorkspaceManifest, - waitForQuiescenceRenewal, WORKER_WORKSPACE_RSYNC_DESTINATION, workerWorkspaceCommandSucceeded as success, workerWorkspaceRsyncRemoteCommand, @@ -57,7 +61,11 @@ import { workspaceSyncError, type WorkerWorkspaceActionsOptions, } from "./workspace-sync-helpers.js"; -import { createGitTransferList, runLocalCommandToFile } from "./workspace-sync-local.js"; +import { + createGitTransferList, + filterExistingGitTransferList, + runLocalCommandToFile, +} from "./workspace-sync-local.js"; import { REMOTE_GIT_WORKSPACE_RETRY_RESET_JS, REMOTE_GIT_WORKSPACE_SETUP_SCRIPT, @@ -68,14 +76,13 @@ import { createWorkerWorkspaceRsyncTransport } from "./workspace-sync-transport. const REMOTE_SETUP_TIMEOUT_MS = 20_000; const WORKSPACE_TIMEOUT_MS = 10 * 60_000; -const WORKSPACE_QUIESCENCE_TIMEOUT_MS = 12 * 60_000; -const WORKSPACE_QUIESCENCE_RENEW_INTERVAL_MS = 4 * 60_000; // Relative to the canonical worker $HOME owned by REMOTE_WORKSPACE_SETUP_SCRIPT; // rsync targets must use the returned absolute directory, never this relative path. const REMOTE_WORKSPACE_ROOT = ".openclaw-worker/workspaces"; const REMOTE_GIT_PACK_NAME = ".openclaw-base.pack"; const GIT_COMMIT_PATTERN = /^[a-f0-9]{40}(?:[a-f0-9]{24})?$/u; const INBOUND_RSYNC_BW_LIMIT_KIB = 65_536; +const workspaceSyncLog = createSubsystemLogger("gateway/worker-workspace"); /** Binds workspace commands and synchronization to one connected tunnel owner. */ export function createWorkerWorkspaceActions( @@ -117,12 +124,22 @@ export function createWorkerWorkspaceActions( const signal = command.signal ? AbortSignal.any([options.ownerSignal, command.signal]) : options.ownerSignal; + const commandOptions = (remainingTimeoutMs: number): CommandOptions => { + const base = workerSshCommandOptions({ + input: command.input, + timeoutMs: remainingTimeoutMs, + signal, + }); + return command.argv.at(-1) === "memo-v1" + ? { ...base, maxOutputBytes: MAX_WORKSPACE_HASH_MEMO_BYTES } + : base; + }; // Exit 255 does not prove whether the remote command was accepted, so stateful // commands must stay pinned to one transport attempt. if (command.transportRetry === "never") { return await runTask( workerWorkspaceSshArgv(prepared, command.argv), - workerSshCommandOptions({ input: command.input, timeoutMs, signal }), + commandOptions(timeoutMs), ); } return await runWorkerSshCandidates( @@ -131,119 +148,16 @@ export function createWorkerWorkspaceActions( async (port, remainingTimeoutMs) => await runTask( workerWorkspaceSshArgv(prepared, command.argv, port), - workerSshCommandOptions({ - input: command.input, - timeoutMs: remainingTimeoutMs, - signal, - }), + commandOptions(remainingTimeoutMs), ), ); }; - const quiesceWorkspace = async (remoteWorkspaceDir: string) => { - if (!path.posix.isAbsolute(remoteWorkspaceDir)) { - throw new Error("Worker workspace quiescence path must be absolute"); - } - const result = await runWorkspaceCommand({ - transportRetry: "never", - argv: [ - "node", - "-e", - REMOTE_WORKSPACE_QUIESCE_JS, - remoteWorkspaceDir, - String(WORKSPACE_QUIESCENCE_TIMEOUT_MS), - options.sharedHost === true ? "shared-host" : "dedicated", - ], - }); - if (!success(result)) { - throw workspaceSyncError(result); - } - const acknowledgement = /^quiesced ([a-f0-9]{32})$/u.exec(result.stdout.trim()); - if (!acknowledgement) { - throw new Error("Worker workspace quiescence returned an invalid acknowledgement"); - } - const nonce = acknowledgement[1]!; - let resumed = false; - let renewalFailure: unknown; - const renewalAbort = new AbortController(); - const abortRenewal = () => renewalAbort.abort(options.ownerSignal.reason); - options.ownerSignal.addEventListener("abort", abortRenewal, { once: true }); - let renewalQueue = Promise.resolve(); - const renew = (validationMode: "heartbeat" | "final") => { - const operation = renewalQueue.then(async () => { - const renewedResult = await runWorkspaceCommand({ - transportRetry: "never", - argv: [ - "node", - "-e", - REMOTE_WORKSPACE_RENEW_QUIESCENCE_JS, - remoteWorkspaceDir, - nonce, - String(WORKSPACE_QUIESCENCE_TIMEOUT_MS), - validationMode, - options.sharedHost === true ? "shared-host" : "dedicated", - ], - }); - if (!success(renewedResult)) { - throw workspaceSyncError(renewedResult); - } - if (renewedResult.stdout.trim() !== `renewed ${nonce}`) { - throw new Error( - "Worker workspace quiescence renewal returned an invalid acknowledgement", - ); - } - }); - renewalQueue = operation.catch(() => undefined); - return operation; - }; - const renewalLoop = (async () => { - while (!renewalAbort.signal.aborted) { - if ( - !(await waitForQuiescenceRenewal( - renewalAbort.signal, - WORKSPACE_QUIESCENCE_RENEW_INTERVAL_MS, - )) - ) { - return; - } - try { - await renew("heartbeat"); - } catch (error) { - renewalFailure = error; - return; - } - } - })(); - return { - assertActive: async () => { - if (resumed) { - throw new Error("Worker workspace quiescence was already released"); - } - if (renewalFailure) { - throw new Error("Worker workspace quiescence renewal failed", { - cause: renewalFailure, - }); - } - await renew("final"); - }, - resume: async () => { - if (resumed) { - return; - } - options.ownerSignal.removeEventListener("abort", abortRenewal); - renewalAbort.abort(); - await renewalLoop; - const resumedResult = await runWorkspaceCommand({ - transportRetry: "never", - argv: ["node", "-e", REMOTE_WORKSPACE_RESUME_JS, remoteWorkspaceDir, nonce], - }); - if (!success(resumedResult)) { - throw workspaceSyncError(resumedResult); - } - resumed = true; - }, - }; - }; + const quiesceWorkspace = createWorkerWorkspaceQuiescence({ + ownerSignal: options.ownerSignal, + sharedHost: options.sharedHost === true, + runWorkspaceCommand, + }); const syncWorkspaceImpl = async ( request: WorkerWorkspaceSyncRequest, @@ -288,7 +202,7 @@ export function createWorkerWorkspaceActions( remoteRelative, }; const mutationReceiverPath = createWorkerWorkspaceRsyncReceiverPathFactory(receiverContext); - let prepareGitTransferList: (() => Promise) | undefined; + let gitTransferListPath: string | undefined; if (mode === "git") { const [canonicalRequestPath, canonicalGitRoot] = await Promise.all([ fs.realpath(request.localPath), @@ -301,14 +215,12 @@ export function createWorkerWorkspaceActions( throw new Error("Worker workspace git base is not a commit id"); } - let transferAttempt = 0; - prepareGitTransferList = async () => - await createGitTransferList({ - gitRoot, - temporaryDirectory: path.join(temporaryDirectory, `transfer-${transferAttempt++}`), - signal: options.ownerSignal, - timeoutMs: WORKSPACE_TIMEOUT_MS, - }); + gitTransferListPath = await createGitTransferList({ + gitRoot, + temporaryDirectory: path.join(temporaryDirectory, "transfer"), + signal: options.ownerSignal, + timeoutMs: WORKSPACE_TIMEOUT_MS, + }); const objectListPath = path.join(temporaryDirectory, "base-objects"); const packPath = path.join(temporaryDirectory, "base.pack"); @@ -396,7 +308,9 @@ export function createWorkerWorkspaceActions( `${prepared.scpTarget}:${WORKER_WORKSPACE_RSYNC_DESTINATION}`, ]; let retryingGitTransfer = false; - const transfer = prepareGitTransferList + let transferAttempt = 0; + const preparedGitTransferListPath = gitTransferListPath; + const transfer = preparedGitTransferListPath ? await runWorkerSshCandidates( prepared, WORKSPACE_TIMEOUT_MS, @@ -435,7 +349,14 @@ export function createWorkerWorkspaceActions( ); } } - const fileListPath = await prepareGitTransferList(); + const fileListPath = await filterExistingGitTransferList({ + gitRoot, + preparedListPath: preparedGitTransferListPath, + outputPath: path.join( + path.dirname(preparedGitTransferListPath), + `attempt-${transferAttempt++}`, + ), + }); const result = await runTask( transferArgv(workerWorkspaceRsyncRemoteCommand(prepared, port), fileListPath), commandOptions(), @@ -473,8 +394,9 @@ export function createWorkerWorkspaceActions( } }; - const reconcileWorkspaceImpl = async ( + const reconcileWorkspaceRun = async ( request: WorkerWorkspaceReconcileRequest, + metrics: WorkspaceReconcileMetrics, ): Promise => { if (!path.isAbsolute(request.localPath) || !path.posix.isAbsolute(request.remoteWorkspaceDir)) { throw new Error("Worker workspace reconcile paths must be absolute"); @@ -484,6 +406,11 @@ export function createWorkerWorkspaceActions( await recoverWorkerWorkspaceReconciliation({ root: request.localPath, journal: pending }); request.journal.abort(); } + const hashMemo = new Map(); + const runLocalReconciliation = (operation: () => Promise): Promise => + measureLocalWorkspaceReconciliation(metrics, () => + withWorkspaceHashMemo(hashMemo, operation, metrics.gateway), + ); const baseDigest = await resolveRemoteWorkspaceManifest( runWorkspaceCommand, request.remoteWorkspaceDir, @@ -504,6 +431,8 @@ export function createWorkerWorkspaceActions( receiverEntryPath, localPath: request.localPath, remoteWorkspaceDir: request.remoteWorkspaceDir, + hashMemo, + metrics, }); try { await fs.mkdir(stagingRoot, { mode: 0o700 }); @@ -539,177 +468,161 @@ export function createWorkerWorkspaceActions( runWorkspaceCommand, remoteWorkspaceDir: request.remoteWorkspaceDir, }); - const verifyStable = async (expectedRef: string): Promise => - await verifyRemoteWorkspaceManifest({ + const verifyStable = async (expectedRef: string): Promise => { + const expectedDigest = expectedRef.slice("sha256:".length); + const observedRef = await captureRemoteWorkspaceManifest({ runWorkspaceCommand, remoteWorkspaceDir: request.remoteWorkspaceDir, baseCommit: base.baseCommit, - baseDigest, - expectedRef, + // Seed both manifests so a recreated path under a new ignore rule + // still invalidates the late-writer fence. + priorManifestDigests: base.baseCommit ? [expectedDigest, baseDigest] : [], + hashMemo, + metrics, }); - const currentResult = await runWorkspaceCommand({ - transportRetry: "idempotent", - argv: [ - "node", - "-e", - REMOTE_WORKSPACE_MANIFEST_JS, - request.remoteWorkspaceDir, - base.baseCommit ?? "", - ...(base.baseCommit ? ["eligible"] : []), - ...(base.baseCommit ? [baseDigest] : []), - ], - }); - if (!success(currentResult)) { - throw workspaceSyncError(currentResult); - } - const currentRef = parseManifestRef(currentResult.stdout.trim()); - if (currentRef === request.baseManifestRef) { - const { expectedRemoteRef, publishAcceptedManifest } = acceptedWorkspacePublisher( - base, - currentRef, - ); - await verifyStable(currentRef); - const stagedResult = request.stagedResult - ? await workerWorkspaceResultStaging.prepareRequestedWorkerWorkspaceResult({ - request, - stagingRoot, - currentManifestRef: currentRef, - baseManifestRaw: baseRaw, - currentManifestRaw: baseRaw, - publishAcceptedManifest, - }) - : undefined; - let appliedWorkspaceResult: WorkerWorkspaceApplyResult | undefined; - if (!stagedResult) { - appliedWorkspaceResult = await applyStagedWorkerWorkspace({ - root: request.localPath, - stagingRoot, - baseManifestRef: request.baseManifestRef, - currentManifestRef: currentRef, - base, - current: base, - journal: request.journal, - publishAcceptedManifest, - }); + if (observedRef !== expectedRef) { + throw new Error("Cloud workspace changed during final reconciliation"); } - return { - get manifestRef() { - return expectedRemoteRef(); - }, - changed: false, - verifyStable: async () => await verifyStable(expectedRemoteRef()), - verifyLocalStable: async () => - await (appliedWorkspaceResult?.verifyLocalStable() ?? - assertWorkspaceResultStable({ root: request.localPath, base, current: base })), - getAppliedWorkspaceResult: () => appliedWorkspaceResult, - ...stagedResult, - }; - } - const currentDigest = currentRef.slice("sha256:".length); - const currentManifestPath = path.join(manifestRoot, `${currentDigest}.json`); - const currentManifestTransfer = await runBoundedInboundRsync({ - prepared, - argv: (rsyncSsh) => [ - "rsync", - "--archive", - "--no-recursive", - "--checksum", - `--max-size=${MAX_RECONCILIATION_FILE_BYTES}`, - `--bwlimit=${INBOUND_RSYNC_BW_LIMIT_KIB}`, - "-e", - rsyncSsh, - "--", - `${prepared.scpTarget}:.openclaw-worker/manifests/${currentDigest}.json`, - currentManifestPath, - ], - destinationRoot: manifestRoot, - entryLimit: 1, - totalByteLimit: MAX_RECONCILIATION_FILE_BYTES, + }; + const currentRef = await captureRemoteWorkspaceManifest({ + runWorkspaceCommand, + remoteWorkspaceDir: request.remoteWorkspaceDir, + baseCommit: base.baseCommit, + priorManifestDigests: base.baseCommit ? [baseDigest] : [], + hashMemo, + metrics, }); - if (!success(currentManifestTransfer)) { - throw workspaceSyncError(currentManifestTransfer); - } - const currentRaw = await readTransferredManifest(currentManifestPath); - const current = parseWorkerWorkspaceManifest(currentRaw, currentRef); - const { expectedRemoteRef, publishAcceptedManifest } = acceptedWorkspacePublisher( - current, - currentRef, - ); - const transferPaths = workerWorkspaceTransferPaths(current, base); - const transferPathSet = new Set(transferPaths); - if (transferPaths.length > 0) { - await fs.writeFile(transferListPath, Buffer.from(`${transferPaths.join("\0")}\0`), { - mode: 0o600, - }); - const resultTransfer = await runBoundedInboundRsync({ + const changed = currentRef !== request.baseManifestRef; + let current = base; + let currentRaw = baseRaw; + if (changed) { + const currentDigest = currentRef.slice("sha256:".length); + const currentManifestPath = path.join(manifestRoot, `${currentDigest}.json`); + const currentManifestTransfer = await runBoundedInboundRsync({ prepared, argv: (rsyncSsh) => [ "rsync", "--archive", + "--no-recursive", "--checksum", `--max-size=${MAX_RECONCILIATION_FILE_BYTES}`, `--bwlimit=${INBOUND_RSYNC_BW_LIMIT_KIB}`, - "--from0", - `--files-from=${transferListPath}`, "-e", rsyncSsh, "--", - `${prepared.scpTarget}:${request.remoteWorkspaceDir}/`, - `${stagingRoot}/`, + `${prepared.scpTarget}:.openclaw-worker/manifests/${currentDigest}.json`, + currentManifestPath, ], - destinationRoot: stagingRoot, - entryLimit: MAX_RECONCILIATION_ENTRIES * 2, - totalByteLimit: MAX_RECONCILIATION_TOTAL_BYTES, + destinationRoot: manifestRoot, + entryLimit: 1, + totalByteLimit: MAX_RECONCILIATION_FILE_BYTES, }); - if (!success(resultTransfer)) { - throw workspaceSyncError(resultTransfer); + if (!success(currentManifestTransfer)) { + throw workspaceSyncError(currentManifestTransfer); } + currentRaw = await readTransferredManifest(currentManifestPath); + current = parseWorkerWorkspaceManifest(currentRaw, currentRef); + } + const { expectedRemoteRef, publishAcceptedManifest } = acceptedWorkspacePublisher( + current, + currentRef, + ); + if (changed) { + const transferPaths = workerWorkspaceTransferPaths(current, base); + const transferPathSet = new Set(transferPaths); + if (transferPaths.length > 0) { + await fs.writeFile(transferListPath, Buffer.from(`${transferPaths.join("\0")}\0`), { + mode: 0o600, + }); + const resultTransfer = await runBoundedInboundRsync({ + prepared, + argv: (rsyncSsh) => [ + "rsync", + "--archive", + "--checksum", + `--max-size=${MAX_RECONCILIATION_FILE_BYTES}`, + `--bwlimit=${INBOUND_RSYNC_BW_LIMIT_KIB}`, + "--from0", + `--files-from=${transferListPath}`, + "-e", + rsyncSsh, + "--", + `${prepared.scpTarget}:${request.remoteWorkspaceDir}/`, + `${stagingRoot}/`, + ], + destinationRoot: stagingRoot, + entryLimit: MAX_RECONCILIATION_ENTRIES * 2, + totalByteLimit: MAX_RECONCILIATION_TOTAL_BYTES, + }); + if (!success(resultTransfer)) { + throw workspaceSyncError(resultTransfer); + } + } + await assertWorkspaceMatchesManifest({ + root: stagingRoot, + manifest: current, + entries: current.entries.filter((entry) => transferPathSet.has(entry.path)), + }); } - await assertWorkspaceMatchesManifest({ - root: stagingRoot, - manifest: current, - entries: current.entries.filter((entry) => transferPathSet.has(entry.path)), - }); // Catch additions, deletions, and writes that raced the inbound transfer. // Stop performs this check once more after local acceptance, directly // before destroying the remote owner. await verifyStable(currentRef); - const stagedResult = request.stagedResult - ? await workerWorkspaceResultStaging.prepareRequestedWorkerWorkspaceResult({ - request, - stagingRoot, - currentManifestRef: currentRef, - baseManifestRaw: baseRaw, - currentManifestRaw: currentRaw, - publishAcceptedManifest, - }) + const preparedStagedResult = request.stagedResult + ? await runLocalReconciliation( + async () => + await workerWorkspaceResultStaging.prepareRequestedWorkerWorkspaceResult({ + request, + stagingRoot, + currentManifestRef: currentRef, + baseManifestRaw: baseRaw, + currentManifestRaw: currentRaw, + publishAcceptedManifest, + }), + ) + : undefined; + const stagedResult = preparedStagedResult + ? { + ...preparedStagedResult, + applyPreparedStagedResult: async () => + await runLocalReconciliation( + async () => await preparedStagedResult.applyPreparedStagedResult(), + ), + verifyLocalStable: async () => + await runLocalReconciliation( + async () => await preparedStagedResult.verifyLocalStable(), + ), + } : undefined; let appliedWorkspaceResult: WorkerWorkspaceApplyResult | undefined; if (!stagedResult) { - appliedWorkspaceResult = await applyStagedWorkerWorkspace({ - root: request.localPath, - stagingRoot, - baseManifestRef: request.baseManifestRef, - currentManifestRef: currentRef, - base, - current, - journal: request.journal, - publishAcceptedManifest, - }); + appliedWorkspaceResult = await runLocalReconciliation( + async () => + await applyStagedWorkerWorkspace({ + root: request.localPath, + stagingRoot, + baseManifestRef: request.baseManifestRef, + currentManifestRef: currentRef, + base, + current, + journal: request.journal, + publishAcceptedManifest, + }), + ); } return { get manifestRef() { return expectedRemoteRef(); }, - changed: true, + changed, verifyStable: async () => await verifyStable(expectedRemoteRef()), verifyLocalStable: async () => - appliedWorkspaceResult - ? await appliedWorkspaceResult.verifyLocalStable() - : await assertWorkspaceResultStable({ root: request.localPath, base, current }), - ...(appliedWorkspaceResult - ? { getAppliedWorkspaceResult: () => appliedWorkspaceResult } - : {}), + await runLocalReconciliation( + async () => + await (appliedWorkspaceResult?.verifyLocalStable() ?? + assertWorkspaceResultStable({ root: request.localPath, base, current })), + ), + getAppliedWorkspaceResult: () => appliedWorkspaceResult, ...stagedResult, }; } finally { @@ -717,6 +630,28 @@ export function createWorkerWorkspaceActions( } }; + const reconcileWorkspaceImpl = async ( + request: WorkerWorkspaceReconcileRequest, + ): Promise => { + const metrics = createWorkspaceReconcileMetrics(); + const startedAt = performance.now(); + const report = (outcome: "failed" | "succeeded") => { + workspaceSyncLog.debug("worker workspace reconcile completed", { + outcome, + durationMs: performance.now() - startedAt, + ...metrics, + }); + }; + try { + const reconciliation = await reconcileWorkspaceRun(request, metrics); + registerWorkspaceReconcileReporter(reconciliation, report); + return reconciliation; + } catch (error) { + report("failed"); + throw error; + } + }; + return { quiesceWorkspace, reconcileWorkspace: (request) => track(reconcileWorkspaceImpl(request)),