fix(cloud): settle uploads when their workspace owner closes (#142237)

Close the response as well as the request when a transfer is cancelled, including a signal aborted before handler entry. A consumed request can already be destroyed while its uploader still waits for staging validation. Preserve live-authority checks before publishing the upload. Extracted from #134931. Both real HTTP regressions fail before repair; sibling tests, final owner tests, changed checks, and independent review pass.
This commit is contained in:
Peter Steinberger
2026-09-09 07:02:57 -07:00
committed by GitHub
parent b827d20356
commit 9956abfbe3
3 changed files with 143 additions and 4 deletions
+2
View File
@@ -419,6 +419,8 @@ Disconnected workers have no cleanup deadline. Nodes also reclaim copies when th
Completed cloud turns preserve eligible, size-bounded workspace files before the turn claim is released. Repository-only sessions accept a cumulative immutable checkpoint in the Gateway's bare artifact repository. Gateway-source sessions apply those changes to their managed worktree. Worker-turn uses its terminal worker event to create the durable pending-result fence. Remote-exec waits for workspace quiescence and enters the same reconciliation flow after the local Codex attempt. Before applying the result, the Gateway stages complete authenticated base/current manifests plus each changed resulting blob as a Git ref under `refs/openclaw/worker-results/`; deletions are represented by the manifests and need no blob. This keeps the cloud delta recoverable even if the Gateway stops during the apply without duplicating unchanged baseline content. Workspace results use Git file semantics: regular files, executable bits, symlinks, additions, changes, and deletions are retained, while empty directories and other directory modes are not. Gateway-source changes remain in the managed worktree for normal review and commit; repository-only changes remain on the node and in the accepted checkpoint.
If workspace transfer ownership closes during an upload, the Gateway disconnects the uploader promptly, including while it waits for validation after sending all bytes. The cancelled upload cannot become an accepted workspace result.
Replacement and Gateway Move restore files against the pinned base; they do not restore worker commit history, merge stages, or partial staging. After a recorded cloud publication, Gateway Move continues the local branch from that verified pushed commit while keeping later accepted file changes available for review. Review recovered conflict-marker files before continuing. When a publishable checkpoint is available, restoration marks its added files as intent-to-add, keeping added and edited contents unstaged for review. Accepted publication deletions are restored as staged index removals; any recovered file bytes remain available. Ignored recovery-only files and attachments are not enrolled for publication. If publication capture was unavailable, recovered ignored files need an explicit `git add -f` before publishing.
For each OpenClaw `worker-turn`, the Gateway binds its effective shared GitHub identity into the worker's `exec` launches, using the same [`tools.github`](/gateway/config-tools#tools-github) selection as ordinary Gateway-host exec. When that identity is available, `gh` is authenticated and HTTPS `git push` uses the `gh auth git-credential` helper. The worker checkout carries the session-owned branch name and, for GitHub repositories, an HTTPS `origin`. The agent commits and pushes directly from the worker. Reconciliation preserves file contents, not the worker's commit history, so work pushed from the worker lands on GitHub first. At every turn start, the worker fast-forwards its checkout to the session branch on `origin` when the local branch is behind, bringing in history pushed by an earlier worker; a diverged local branch is left untouched.
@@ -196,12 +196,20 @@ export function createNodeWorkspaceTransferHttpCallback(
clientAbort.signal,
AbortSignal.timeout(TRANSFER_TIMEOUT_MS),
]);
const abortRequest = () => {
const abortTransfer = () => {
// The request can be fully read and destroyed while its uploader still
// awaits staging validation on the open response.
if (!res.destroyed) {
res.destroy(signal.reason instanceof Error ? signal.reason : undefined);
}
if (!req.destroyed) {
req.destroy(signal.reason instanceof Error ? signal.reason : undefined);
}
};
signal.addEventListener("abort", abortRequest, { once: true });
signal.addEventListener("abort", abortTransfer, { once: true });
if (signal.aborted) {
abortTransfer();
}
const stillCurrent = () => !signal.aborted && service.isAuthorizationCurrent(authorization);
try {
if (route.kind === "manifest" || route.kind === "pack") {
@@ -295,7 +303,7 @@ export function createNodeWorkspaceTransferHttpCallback(
throw error;
}
} finally {
signal.removeEventListener("abort", abortRequest);
signal.removeEventListener("abort", abortTransfer);
stopWatchingDisconnect();
}
},
@@ -1,8 +1,9 @@
import fsSync from "node:fs";
import fs from "node:fs/promises";
import { createServer } from "node:http";
import { createServer, type IncomingMessage } from "node:http";
import path from "node:path";
import { afterEach, describe, expect, it, vi } from "vitest";
import { createDeferred, withTestTimeout } from "../../../test/helpers/promise.js";
import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js";
import {
createOperationalRunInstanceRef,
@@ -11,6 +12,7 @@ import {
} from "../../agents/admitted-run-context.js";
import { ensureStagedInputDirectory, stagedInputDirectory } from "../../media/staged-inputs.js";
import { runNodeWorkerWorkspaceTransfer } from "../../node-host/node-worker-transfer-client.js";
import { nodeWorkspaceTransferReconcilePath } from "../../worker/node-workspace-transfer-protocol.js";
import {
createNodeWorkspaceTransferHttpCallback,
handleNodeWorkspaceTransferHttpRequest,
@@ -20,6 +22,133 @@ import { createNodeWorkspaceTransferService } from "./node-workspace-transfer-se
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
afterEach(() => vi.restoreAllMocks());
describe("workspace upload cancellation", () => {
it.each([
{ boundary: "before handler", abort: false },
{ boundary: "before handler", abort: true },
{ boundary: "after body", abort: false },
{ boundary: "after body", abort: true },
] as const)("settles $boundary with owner abort=$abort", async ({ boundary, abort }) => {
const root = tempDirs.make("workspace-upload-cancellation-");
const localPath = path.join(root, "source");
await fs.mkdir(localPath);
const owner = new AbortController();
const service = createNodeWorkspaceTransferService({
temporaryRoot: path.join(root, "transfers"),
getOwner: () => ({
credential: { ownerEpoch: 1, sessionId: "session" },
environment: {
ownerEpoch: 1,
attachedSessionIds: ["session"],
destroyRequestedAtMs: null,
state: "attached",
},
}),
});
const { snapshot } = await service.prepareSync({
environmentId: "environment",
ownerEpoch: 1,
sessionId: "session",
generation: 1,
localPath,
isAuthorized: () => !owner.signal.aborted,
signal: owner.signal,
});
const token = service.prepareUpload("environment", snapshot.manifestRef);
const reached = createDeferred();
const release = createDeferred();
const finished = createDeferred();
if (boundary === "after body") {
const realpath = fs.realpath.bind(fs);
vi.spyOn(fs, "realpath").mockImplementation(async (...args) => {
if (typeof args[0] === "string" && path.basename(args[0]).startsWith("upload-")) {
// Staged-manifest verification happens after the reader consumes EOF.
reached.resolve();
await release.promise;
}
return await realpath(...args);
});
}
let incoming: IncomingMessage | undefined;
const callback = createNodeWorkspaceTransferHttpCallback(service);
const server = createServer((req, res) => {
incoming = req;
void handleNodeWorkspaceTransferHttpRequest({
req,
res,
clientIp: "127.0.0.1",
callback: async (request) => {
const authorized = await callback(request);
if (boundary === "before handler") {
reached.resolve();
await release.promise;
}
return authorized;
},
})
.catch((error: unknown) =>
res.destroy(error instanceof Error ? error : new Error(String(error))),
)
.finally(() => finished.resolve());
});
await new Promise<void>((resolve) => {
server.listen(0, "127.0.0.1", resolve);
});
const address = server.address();
if (!address || typeof address === "string") {
throw new Error("HTTP fixture did not bind");
}
const raw = Buffer.from(snapshot.rawManifest);
const length = Buffer.alloc(4);
length.writeUInt32BE(raw.length);
const request = fetch(
`http://127.0.0.1:${address.port}${nodeWorkspaceTransferReconcilePath("environment", snapshot.manifestRef)}`,
{
method: "POST",
headers: { authorization: `Bearer ${token}` },
body: Buffer.concat([length, raw, length, raw]),
},
).then(
async (response) => ({ status: response.status, body: await response.json() }),
() => "rejected" as const,
);
try {
await withTestTimeout(reached.promise, 2_000, "upload did not reach cancellation boundary");
if (boundary === "after body") {
expect(incoming?.readableEnded).toBe(true);
expect(incoming?.destroyed).toBe(true);
}
if (abort) {
owner.abort(new Error("Workspace transfer owner closed"));
}
if (!abort || boundary === "before handler") {
release.resolve();
}
const result = await withTestTimeout(request, 2_000, "upload response remained open");
release.resolve();
await finished.promise;
if (abort) {
expect(result).toBe("rejected");
expect(() => service.takeUpload("environment", snapshot.manifestRef)).toThrow();
} else {
expect(result).toEqual({ status: 200, body: { manifestRef: snapshot.manifestRef } });
expect(service.takeUpload("environment", snapshot.manifestRef).currentManifestRef).toBe(
snapshot.manifestRef,
);
}
} finally {
release.resolve();
server.closeAllConnections();
await request;
await finished.promise;
await new Promise<void>((resolve) => {
server.close(() => resolve());
});
await service.closeAll();
}
});
});
describe("attachment transfer revocation", () => {
it.each([
{ boundary: "blob-admitted", revoke: false },