diff --git a/.github/workflows/runner-full-stack-e2e.yml b/.github/workflows/runner-full-stack-e2e.yml index 7ea1660729..9016efa30c 100644 --- a/.github/workflows/runner-full-stack-e2e.yml +++ b/.github/workflows/runner-full-stack-e2e.yml @@ -1235,6 +1235,11 @@ jobs: with: version: 9.15.4 + # Resolve the trusted publication checkout's lockfile without lifecycle + # scripts, then install the exact result. This matches the report job's + # frozen-install preparation while keeping target-controlled code out of + # the AWS credentialed publisher. + - run: pnpm install --lockfile-only --ignore-scripts --no-frozen-lockfile - run: pnpm install --frozen-lockfile - name: Install publisher-only Chromium diff --git a/.github/workflows/runner-protocol-live-evals.yml b/.github/workflows/runner-protocol-live-evals.yml index e3eafa2da5..c3bd0e01bf 100644 --- a/.github/workflows/runner-protocol-live-evals.yml +++ b/.github/workflows/runner-protocol-live-evals.yml @@ -154,6 +154,57 @@ jobs: } >> "$GITHUB_OUTPUT" fi + target_lock: + name: Resolve target pnpm lockfile + needs: authorize + runs-on: ubuntu-latest + timeout-minutes: 10 + permissions: + contents: read + outputs: + artifact_id: ${{ steps.upload.outputs.artifact-id }} + lock_sha256: ${{ steps.lock.outputs.sha256 }} + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7 + with: + ref: ${{ needs.authorize.outputs.target_sha }} + persist-credentials: false + + - uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7 + with: + node-version: 24 + + - uses: pnpm/action-setup@0977fd99725f1db4007ccb2928dbb4e90d06cc86 # v6 + env: + NPM_CONFIG_AUDIT: "false" + NPM_CONFIG_FUND: "false" + NPM_CONFIG_UPDATE_NOTIFIER: "false" + with: + version: 9.15.4 + + - name: Resolve target lockfile without lifecycle scripts + id: lock + run: | + set -euo pipefail + pnpm install --ignore-scripts --no-frozen-lockfile --lockfile-only + test -s pnpm-lock.yaml + unexpected="$(git status --short | awk '$2 != "pnpm-lock.yaml" { print }')" + if [ -n "$unexpected" ]; then + echo "Lockfile resolution changed files other than pnpm-lock.yaml:" >&2 + echo "$unexpected" >&2 + exit 1 + fi + echo "sha256=$(sha256sum pnpm-lock.yaml | cut -d ' ' -f 1)" >> "$GITHUB_OUTPUT" + + - name: Upload resolved target lockfile + id: upload + uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7 + with: + name: runner-protocol-target-pnpm-lock-${{ github.run_id }}-${{ github.run_attempt }} + path: pnpm-lock.yaml + retention-days: 30 + if-no-files-found: error + catalog: name: Pin and fan out the direct Evalbook roster needs: authorize @@ -228,7 +279,7 @@ jobs: build_runner: name: Build portable direct-eval runner once - needs: [authorize, catalog] + needs: [authorize, target_lock, catalog] runs-on: ${{ needs.authorize.outputs.test_runner }} timeout-minutes: 30 permissions: @@ -251,6 +302,26 @@ jobs: with: version: 9.15.4 + - name: Download resolved target lockfile + uses: actions/download-artifact@3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c # v8 + with: + artifact-ids: ${{ needs.target_lock.outputs.artifact_id }} + path: ${{ runner.temp }}/runner-protocol-target-lock + + - name: Restore resolved target lockfile + env: + TARGET_SHA: ${{ needs.authorize.outputs.target_sha }} + EXPECTED_LOCK_SHA256: ${{ needs.target_lock.outputs.lock_sha256 }} + run: | + set -euo pipefail + test "$(git rev-parse HEAD)" = "$TARGET_SHA" + lock="$RUNNER_TEMP/runner-protocol-target-lock/pnpm-lock.yaml" + test -f "$lock" + test "$(find "$(dirname "$lock")" -type f | wc -l | tr -d ' ')" = 1 + test "$(sha256sum "$lock" | cut -d ' ' -f 1)" = "$EXPECTED_LOCK_SHA256" + cp "$lock" pnpm-lock.yaml + test "$(sha256sum pnpm-lock.yaml | cut -d ' ' -f 1)" = "$EXPECTED_LOCK_SHA256" + - uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7 with: node-version: 24 @@ -258,6 +329,9 @@ jobs: - run: pnpm install --frozen-lockfile --ignore-scripts + - name: Materialize the pinned OpenCode executable before packaging + run: node packages/paperclip-runner/scripts/materialize-opencode-binary.mjs + - name: Build runner CLI, daemon, and canonical attempt viewer run: | set -euo pipefail @@ -382,6 +456,40 @@ jobs: tar --extract --gzip --file runner-protocol-build.tar.gz --directory extracted test -x extracted/paperclip-runnerd + # Keep the host policy aligned with runner-full-stack-e2e.yml. This is + # deliberately before any provider credential or web-identity step. + - name: Provision Codex sandbox on the disposable trusted runner + if: matrix.provider == 'codex' || matrix.rosterId == 'protocol-live-acpx-codex-control' + run: | + node --input-type=module <<'NODE' + import { execFileSync } from "node:child_process"; + import { createHash } from "node:crypto"; + import { readFileSync, realpathSync, writeFileSync } from "node:fs"; + import { createRequire } from "node:module"; + import path from "node:path"; + if (process.platform !== "linux") process.exit(0); + let restricted = "0"; + try { restricted = readFileSync("/proc/sys/kernel/apparmor_restrict_unprivileged_userns", "utf8").trim(); } catch {} + if (restricted !== "1") process.exit(0); + const root = realpathSync(path.join(process.env.GITHUB_WORKSPACE, "runner-protocol-build/extracted/portable")); + const runnerRequire = createRequire(path.join(root, "package.json")); + const acpRequire = createRequire(runnerRequire.resolve("@agentclientprotocol/codex-acp/package.json")); + const codexRequire = createRequire(acpRequire.resolve("@openai/codex/package.json")); + const arch = process.arch === "x64" ? "x64" : process.arch === "arm64" ? "arm64" : null; + if (!arch) throw new Error("Unsupported Codex CI architecture"); + const platformPackage = codexRequire.resolve(`@openai/codex-linux-${arch}/package.json`); + const triple = arch === "x64" ? "x86_64-unknown-linux-musl" : "aarch64-unknown-linux-musl"; + const suffix = `/vendor/${triple}/bin/codex`; + const binary = realpathSync(path.join(path.dirname(platformPackage), suffix)); + if (!binary.startsWith(root + "/node_modules/.pnpm/") || !binary.endsWith(suffix) || !/^[/A-Za-z0-9_.@+\-]+$/.test(binary)) { + throw new Error("Codex executable is outside the resolved dependency tree"); + } + const name = `paperclip-e2e-codex-${createHash("sha256").update(binary).digest("hex").slice(0,16)}`; + const profilePath = path.join(process.env.RUNNER_TEMP, "paperclip-codex-userns.apparmor"); + writeFileSync(profilePath, `abi ,\ninclude \nprofile ${name} "${binary}" flags=(unconfined) {\n userns,\n}\n`, {mode:0o600, flag:"wx"}); + execFileSync("sudo", ["-n", "apparmor_parser", "-r", profilePath], {timeout:15000, stdio:"pipe"}); + NODE + - name: Prepare short-lived AgentCore web identity if: matrix.credentialName == 'AWS_AGENTCORE_OIDC' env: @@ -537,6 +645,25 @@ jobs: with: node-version: 24 + - uses: pnpm/action-setup@0977fd99725f1db4007ccb2928dbb4e90d06cc86 # v6 + env: + NPM_CONFIG_AUDIT: "false" + NPM_CONFIG_FUND: "false" + NPM_CONFIG_UPDATE_NOTIFIER: "false" + with: + version: 9.15.4 + + - name: Resolve trusted report lockfile without lifecycle scripts + run: | + set -euo pipefail + pnpm install --ignore-scripts --no-frozen-lockfile --lockfile-only + test -s pnpm-lock.yaml + + - uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7 + with: + node-version: 24 + cache: pnpm + - name: Download immutable campaign catalog uses: actions/download-artifact@3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c # v8 with: diff --git a/packages/paperclip-runner/scripts/runner-protocol-eval-workflow-security.test.mjs b/packages/paperclip-runner/scripts/runner-protocol-eval-workflow-security.test.mjs index d35fbd8bc5..2d45293ef8 100644 --- a/packages/paperclip-runner/scripts/runner-protocol-eval-workflow-security.test.mjs +++ b/packages/paperclip-runner/scripts/runner-protocol-eval-workflow-security.test.mjs @@ -51,6 +51,19 @@ test("pull request CI builds the canonical Evalbook viewer", async () => { test("resolves both repositories immutably and bounds total matrix concurrency", async () => { const workflow = await readFile(workflowPath, "utf8"); + const targetLock = workflow.slice( + workflow.indexOf(" target_lock:"), + workflow.indexOf(" catalog:"), + ); + assert.match(targetLock, /ref: \$\{\{ needs\.authorize\.outputs\.target_sha \}\}/u); + assert.match(targetLock, /pnpm install --ignore-scripts --no-frozen-lockfile --lockfile-only/u); + assert.match(targetLock, /Upload resolved target lockfile/u); + const build = workflow.slice( + workflow.indexOf(" build_runner:"), + workflow.indexOf(" eval_shard_0:"), + ); + assert.match(build, /Download resolved target lockfile/u); + assert.match(build, /Restore resolved target lockfile/u); const authorize = workflow.slice( workflow.indexOf(" authorize:"), workflow.indexOf(" catalog:"), @@ -164,3 +177,54 @@ test("publishes only the separately sanitized Evalbook through trusted OIDC code assert.doesNotMatch(publisher, /paperclipai\/paperclip-evals/u); assert.doesNotMatch(publisher, /downloaded-runner-protocol-evals/u); }); + +test("provisions the Codex userns profile before any provider credentials", async () => { + const workflow = await readFile(workflowPath, "utf8"); + const direct = workflow.slice( + workflow.indexOf(" steps: &direct_eval_steps"), + workflow.indexOf(" eval_shard_1:"), + ); + const profile = direct.indexOf("Provision Codex sandbox on the disposable trusted runner"); + const credentials = direct.indexOf("Prepare short-lived AgentCore web identity"); + assert.ok(profile >= 0 && credentials >= 0 && profile < credentials); + assert.match(direct, /apparmor_restrict_unprivileged_userns/u); + assert.match(direct, /apparmor_parser/u); + assert.match(direct, /userns,/u); + assert.match(direct, /runner-protocol-build\/extracted\/portable/u); + assert.match(direct, /matrix\.provider == 'codex'/u); + assert.match(direct, /matrix\.rosterId == 'protocol-live-acpx-codex-control'/u); + assert.doesNotMatch(direct, /matrix\.profileId/u); + const build = workflow.slice( + workflow.indexOf(" build_runner:"), + workflow.indexOf(" eval_shard_0:"), + ); + assert.match(build, /Materialize the pinned OpenCode executable before packaging/u); + assert.match(build, /materialize-opencode-binary\.mjs/u); +}); + +test("report preparation stays on the trusted lock and install mode", async () => { + const workflow = await readFile(workflowPath, "utf8"); + const report = workflow.slice( + workflow.indexOf(" report:"), + workflow.indexOf(" publish_history:"), + ); + assert.match(report, /ref: \$\{\{ github\.sha \}\}/u); + assert.doesNotMatch(report, /Download resolved target lockfile/u); + assert.doesNotMatch(report, /Restore resolved target lockfile/u); + assert.match(report, /Resolve trusted report lockfile without lifecycle scripts/u); + assert.match(report, /pnpm install --ignore-scripts --no-frozen-lockfile --lockfile-only/u); + assert.match(report, /pnpm install --frozen-lockfile --ignore-scripts\n/u); +}); + +test("trusted catalog, direct eval, and report orchestration stay on workflow revision", async () => { + const workflow = await readFile(workflowPath, "utf8"); + for (const [job, next] of [ + [" catalog:", " build_runner:"], + [" eval_shard_0:", " eval_shard_1:"], + [" report:", " publish_history:"], + ]) { + const section = workflow.slice(workflow.indexOf(job), workflow.indexOf(next)); + assert.match(section, /ref: \$\{\{ github\.sha \}\}/u, `${job} must use the trusted workflow revision`); + assert.doesNotMatch(section, /ref: \$\{\{ needs\.authorize\.outputs\.target_sha \}\}/u, `${job} must not execute target orchestration code`); + } +}); diff --git a/packages/paperclip-runner/src/backends/harness-driver-backend.test.ts b/packages/paperclip-runner/src/backends/harness-driver-backend.test.ts index b8a1ae0211..317f80e5a1 100644 --- a/packages/paperclip-runner/src/backends/harness-driver-backend.test.ts +++ b/packages/paperclip-runner/src/backends/harness-driver-backend.test.ts @@ -956,6 +956,20 @@ describe("HarnessDriverBackend", () => { await expect(iterator.next()).rejects.toThrow("provider transport lost after resolution"); }); + it("does not start provider work when Stop arrives before the first turn", async () => { + const provider = new FakeHarnessSession(); + const start = vi.spyOn(provider, "startTurn"); + const backend = new HarnessDriverBackend({ ...driver, async openSession() { return provider; } }); + const session = await backend.openSession({ + identity: { runId: "run-1", sessionId: "session-1", companyId: "company-1", issueId: "issue-1", agentId: "agent-1" }, + workingDirectory: "/workspace", + }); + await session.cancel({ reason: "user stop during startup", signal: new AbortController().signal }).cleanup; + await expect(session.startTurn({ message: { role: "user", text: "cancelled work" } })) + .rejects.toThrow("native_session_cancelled"); + expect(start).not.toHaveBeenCalled(); + }); + it("does not synthesize a fallback after explicit run cancellation", async () => { class CancelledProviderSession extends FakeHarnessSession { override async *events() { diff --git a/packages/paperclip-runner/src/backends/harness-driver-backend.ts b/packages/paperclip-runner/src/backends/harness-driver-backend.ts index 3669e58c43..135724c542 100644 --- a/packages/paperclip-runner/src/backends/harness-driver-backend.ts +++ b/packages/paperclip-runner/src/backends/harness-driver-backend.ts @@ -667,6 +667,10 @@ class HarnessNativeSession implements NativeSession { async startTurn(input: Parameters[0]) { this.#assertProtocolIntegrity(); + // Stop can arrive after session publication but before the first turn. + // The provider has no active turn to interrupt yet. Do not launch work + // whose events cancellation would suppress and leave the owner waiting. + if (this.#explicitlyCancelled) throw new Error("native_session_cancelled"); try { const started = await this.#session.startTurn(input); this.#assertProtocolIntegrity(); diff --git a/packages/plugins/sandbox-providers/daytona/src/file-sync.test.ts b/packages/plugins/sandbox-providers/daytona/src/file-sync.test.ts index 9fe2a21ed3..dd4e8f4efc 100644 --- a/packages/plugins/sandbox-providers/daytona/src/file-sync.test.ts +++ b/packages/plugins/sandbox-providers/daytona/src/file-sync.test.ts @@ -356,6 +356,88 @@ it.skipIf(!gnuTar)("extracts interleaved read-only skill directories with GNU ta } }, 30_000); +it.skipIf(!gnuTar)("uploads gzip directory archives and preserves content, executable modes, and symlinks", async () => { + const root = await fs.mkdtemp(path.join(os.tmpdir(), "paperclip-daytona-gzip-dir-")); + const source = path.join(root, "source"); + const remoteDir = path.join(root, "remote"); + const target = path.join(remoteDir, "target"); + const bin = path.join(root, "bin"); + await fs.mkdir(path.join(source, "bin"), { recursive: true }); + await fs.mkdir(remoteDir); + await fs.mkdir(bin); + await fs.symlink(gnuTar!, path.join(bin, "tar")); + await fs.writeFile(path.join(source, "bin", "tool.sh"), "#!/bin/sh\necho ok\n", { mode: 0o755 }); + await fs.symlink("bin/tool.sh", path.join(source, "tool-link")); + + try { + let uploadedArchiveBytes: Buffer | undefined; + const { sandbox } = createRealExecSandbox({ + commandEnv: { ...process.env, PATH: `${bin}:${process.env.PATH}` }, + uploadOverride: async (uploads) => { + uploadedArchiveBytes = await fs.readFile(uploads[0]!.source); + for (const upload of uploads) await fs.copyFile(upload.source, upload.destination); + return true; + }, + }); + + await performSyncIn({ + sandbox: sandbox as never, + remoteDir, + timeoutSeconds: 30, + operations: [{ + operationId: "gzip-directory", + files: [{ sourcePath: source, targetPath: target, kind: "directory" }], + }], + }); + + expect(uploadedArchiveBytes).toBeDefined(); + expect(uploadedArchiveBytes!.subarray(0, 2)).toEqual(Buffer.from([0x1f, 0x8b])); + expect(await fs.readFile(path.join(target, "bin", "tool.sh"), "utf8")).toBe("#!/bin/sh\necho ok\n"); + expect((await fs.stat(path.join(target, "bin", "tool.sh"))).mode & 0o777).toBe(0o755); + expect(await fs.readlink(path.join(target, "tool-link"))).toBe("bin/tool.sh"); + } finally { + await fs.rm(root, { recursive: true, force: true }); + } +}, 30_000); + +it.skipIf(!gnuTar)("uploads and extracts a gzip empty directory archive", async () => { + const root = await fs.mkdtemp(path.join(os.tmpdir(), "paperclip-daytona-gzip-empty-")); + const source = path.join(root, "source"); + const remoteDir = path.join(root, "remote"); + const target = path.join(remoteDir, "target"); + const bin = path.join(root, "bin"); + await fs.mkdir(source); + await fs.mkdir(remoteDir); + await fs.mkdir(bin); + await fs.symlink(gnuTar!, path.join(bin, "tar")); + + try { + let uploadedArchiveBytes: Buffer | undefined; + const { sandbox } = createRealExecSandbox({ + commandEnv: { ...process.env, PATH: `${bin}:${process.env.PATH}` }, + uploadOverride: async (uploads) => { + uploadedArchiveBytes = await fs.readFile(uploads[0]!.source); + for (const upload of uploads) await fs.copyFile(upload.source, upload.destination); + return true; + }, + }); + await performSyncIn({ + sandbox: sandbox as never, + remoteDir, + timeoutSeconds: 30, + operations: [{ + operationId: "gzip-empty-directory", + files: [{ sourcePath: source, targetPath: target, kind: "directory" }], + }], + }); + expect(uploadedArchiveBytes).toBeDefined(); + expect(uploadedArchiveBytes!.subarray(0, 2)).toEqual(Buffer.from([0x1f, 0x8b])); + expect(await fs.readdir(target)).toEqual([]); + } finally { + await fs.rm(root, { recursive: true, force: true }); + } +}, 30_000); + describe("daytona file-sync inbound zstd transport compression", () => { const cleanupDirs: string[] = []; diff --git a/packages/plugins/sandbox-providers/daytona/src/file-sync.ts b/packages/plugins/sandbox-providers/daytona/src/file-sync.ts index d76b5cd412..f235076753 100644 --- a/packages/plugins/sandbox-providers/daytona/src/file-sync.ts +++ b/packages/plugins/sandbox-providers/daytona/src/file-sync.ts @@ -151,7 +151,7 @@ async function withHostTempDir(fn: (dir: string) => Promise): Promise { } /** - * Build a host-side tarball of a directory, mirroring the runtime's own + * Build a host-side gzip-compressed tarball of a directory, mirroring the runtime's own * `createTarballFromDirectory`: archive top-level entries by name (no "." self * entry), suppress AppleDouble/xattr sidecars, honor `exclude`, and reproduce the * `followSymlinks` → `-h` mapping so the native path is observationally identical @@ -167,14 +167,17 @@ async function createHostTarball(input: { const entries = (await fs.readdir(input.localDir)).sort((left, right) => left.localeCompare(right)); if (entries.length === 0) { // An empty source is valid (blank workspace / empty asset dir). Write a valid - // empty tar (1024-byte zero EOF marker) so extraction is a clean no-op. - await fs.writeFile(input.archivePath, Buffer.alloc(1024)); + // gzip-compressed empty tar (1024-byte zero EOF marker) so extraction is a + // clean no-op and uses the same transport as non-empty directories. + await fs.writeFile(input.archivePath, await new Promise((resolve, reject) => { + zlib.gzip(Buffer.alloc(1024), (error, compressed) => error ? reject(error) : resolve(compressed)); + })); return; } await execFileAsync( "tar", [ - "-c", + "-cz", "--no-xattrs", ...(input.followSymlinks ? ["-h"] : []), "-f", @@ -845,7 +848,7 @@ async function syncInDirectoryMapping(input: { }): Promise<{ filesTransferred: number; bytesTransferred: number }> { const { sandbox, mapping, remoteDir, timeoutSeconds } = input; return withHostTempDir(async (tmp) => { - const archivePath = path.join(tmp, "sync-in.tar"); + const archivePath = path.join(tmp, "sync-in.tar.gz"); // The pack step is host-local: it builds the tarball and makes no sandbox // round trip. The `pack` span records its wall time. // `pack` span: build a tarball on the host — no sandbox round trip. @@ -862,7 +865,7 @@ async function syncInDirectoryMapping(input: { const bytesTransferred = (await fs.stat(archivePath)).size; // The tar bytes ride the native bulk channel (string source ⇒ streamed); // only the extract/cleanup control commands use exec. - const remoteTar = path.posix.join(remoteDir, scratchName(".tar")); + const remoteTar = path.posix.join(remoteDir, scratchName(".tar.gz")); // Count the serial guard round trips before the transfer, so the transfer // span records how much of the wall time is guard cost. let guardRoundTrips = 0; diff --git a/packages/plugins/sandbox-providers/daytona/src/plugin.test.ts b/packages/plugins/sandbox-providers/daytona/src/plugin.test.ts index e29ac4780c..30066a5aaa 100644 --- a/packages/plugins/sandbox-providers/daytona/src/plugin.test.ts +++ b/packages/plugins/sandbox-providers/daytona/src/plugin.test.ts @@ -4092,7 +4092,7 @@ describe("daytona native file-sync hooks", () => { expect(provision!.attributes["paperclip.sandbox.startup.provider"]).toBe("daytona"); }); - it("syncIn tars a directory mapping host-side honoring excludes and the followSymlinks flag, then extracts it in-sandbox via a single quoted tar command", async () => { + it("gzip-tars a directory mapping host-side honoring excludes and the followSymlinks flag, then extracts it in-sandbox via a single quoted tar command", async () => { const hostDir = await makeHostDir(); const sourceDir = path.join(hostDir, "assets"); await fs.mkdir(path.join(sourceDir, "keep"), { recursive: true }); @@ -4134,8 +4134,8 @@ describe("daytona native file-sync hooks", () => { expect(sandbox.fs.uploadFiles).toHaveBeenCalledTimes(1); const [uploads] = sandbox.fs.uploadFiles.mock.calls[0] as [Array<{ source: string; destination: string }>]; expect(uploads).toHaveLength(1); - expect(uploads[0].source).toMatch(/\.tar$/); - expect(path.posix.basename(uploads[0].destination)).toMatch(/^\.paperclip-upload-.*\.tar$/); + expect(uploads[0].source).toMatch(/\.tar\.gz$/); + expect(path.posix.basename(uploads[0].destination)).toMatch(/^\.paperclip-upload-.*\.tar\.gz$/); expect(uploads[0].destination.startsWith(`${REMOTE_DIR}/`)).toBe(true); // Inspect the real host tar: excluded file gone; symlink preserved AS a link. @@ -4159,7 +4159,7 @@ describe("daytona native file-sync hooks", () => { // followed by removing the scratch tar. expect(extractCommand).toContain(".paperclip-runtime/assets"); expect(extractCommand).toContain("tar -xf"); - expect(extractCommand).toMatch(/rm -f .*\.paperclip-upload-.*\.tar/); + expect(extractCommand).toMatch(/rm -f .*\.paperclip-upload-.*\.tar\.gz/); }); it("syncIn dereferences symlinks to bytes when followSymlinks is true (tar -h)", async () => { diff --git a/server/src/__tests__/tool-gateway-service.test.ts b/server/src/__tests__/tool-gateway-service.test.ts index d6691d07ca..111251fe71 100644 --- a/server/src/__tests__/tool-gateway-service.test.ts +++ b/server/src/__tests__/tool-gateway-service.test.ts @@ -381,8 +381,17 @@ describeEmbeddedPostgres("tool gateway service", () => { expect(nativeResponses).toMatchObject([{ interactionId: interaction.id, response: { status: "accepted", result: { toolAction: { status: "executed", resultSummary: expect.stringContaining("bodyLength") } } } }]); expect(wakeup).not.toHaveBeenCalled(); expect((await db.select().from(toolActionDeliveries))[0].deliveredAt).toBeNull(); + await deliveries.deliverForRun({ companyId: company.id, runId: run.id }); + expect(wakeup).not.toHaveBeenCalled(); await db.update(heartbeatRuns).set({ status: "succeeded" }).where(eq(heartbeatRuns.id, run.id)); const restarted = toolActionDeliveryService(db, { wakeup }); + await deliveries.deliverForRun({ companyId: randomUUID(), runId: run.id }); + await deliveries.deliverForRun({ companyId: company.id, runId: randomUUID() }); + expect(wakeup).not.toHaveBeenCalled(); + // The original executor's terminal cleanup must deliver a review that was + // approved while it was running, without waiting for the scheduler sweep. + await deliveries.deliverForRun({ companyId: company.id, runId: run.id }); + expect(wakeup).toHaveBeenCalledTimes(1); await Promise.all([restarted.sweepPending(), deliveries.sweepPending()]); await gateway.approveActionRequest({ companyId: company.id, actionRequestId: request.id, actor: { userId: "second-reviewer" } }); await restarted.sweepPending(); diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index d2eb2f94e9..5018ba1b5f 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -1,4 +1,5 @@ import { publicChatTaskUrl } from "./chat-task-url.js"; +import { toolActionDeliveryService } from "./tool-action-delivery.js"; import { readQueuedInteractionResponse } from "./queued-interaction-response.js"; import { AGENT_CHAT_DIRECTIVE, conversationReplay, isConversation, isConversationExecutionWake, isWaitingConversation, prepareConversationTurn, settleConversationTurn } from "./agent-conversations.js"; import { PROCESS_IDENTITY_RECORDED, recordNativeLocalProcessStop } from "./native-local-process-stop.js"; @@ -26011,6 +26012,13 @@ export function heartbeatService( if (latestRun) await resumeRemoteStopComments(latestRun).catch(err => { logger.warn({ err, runId: run.id }, "failed to resume user messages after remote Stop"); }); + if (latestRun && isHeartbeatRunTerminalStatus(latestRun.status)) { + await toolActionDeliveryService(db, { wakeup: trackWakeup }) + .deliverForRun({ companyId: run.companyId, runId: run.id }) + .catch(err => { + logger.warn({ err, runId: run.id }, "failed to deliver settled tool reviews after execution cleanup"); + }); + } await startNextQueuedRunForAgent(run.agentId); } } diff --git a/server/src/services/native-runtime/native-session-executor.test.ts b/server/src/services/native-runtime/native-session-executor.test.ts index b8e49a833a..1cc3797b21 100644 --- a/server/src/services/native-runtime/native-session-executor.test.ts +++ b/server/src/services/native-runtime/native-session-executor.test.ts @@ -4903,7 +4903,7 @@ describe("native session cancellation", () => { }); }); - it("waits for an in-flight startup handle before acknowledging Stop", async () => { + it.each([true, false])("waits for an in-flight startup handle before acknowledging Stop (runnerd=%s)", async (useRunnerd) => { const root = await mkdtemp(join(tmpdir(), "native-startup-stop-")); const previous = process.env.PAPERCLIP_RUNNER_STATE_DIR; process.env.PAPERCLIP_RUNNER_STATE_DIR = root; @@ -4927,7 +4927,7 @@ describe("native session cancellation", () => { }; }); const running = executePaperclipNativeSession({ - db: leaseDb(), execution, runnerInstanceId: "runner", useRunnerd: true, + db: leaseDb(), execution, runnerInstanceId: "runner", useRunnerd, }); const outcome = running.catch(error => error); const persistence = cancellationDb(); @@ -4935,6 +4935,9 @@ describe("native session cancellation", () => { let acknowledged = false; try { await admitted; + await expect(executePaperclipNativeSession({ + db: leaseDb(), execution, runnerInstanceId: "duplicate-runner", useRunnerd, + })).rejects.toThrow("native_session_supervisor_busy"); stopping = cancelNativeSession(execution.binding.runId, "operator Stop during startup", { db: persistence.db, scope: "run", }); @@ -4961,7 +4964,7 @@ describe("native session cancellation", () => { } }); - it("fences a late startup even when Stop reaches its acknowledgement deadline", async () => { + it.each([true, false])("fences a late startup even when Stop reaches its acknowledgement deadline (runnerd=%s)", async (useRunnerd) => { const root = await mkdtemp(join(tmpdir(), "native-late-startup-stop-")); const previous = process.env.PAPERCLIP_RUNNER_STATE_DIR; process.env.PAPERCLIP_RUNNER_STATE_DIR = root; @@ -4975,7 +4978,7 @@ describe("native session cancellation", () => { submitTurn(); throw new Error("a stopped startup must not reach prompt submission"); }); - const outcome = executePaperclipNativeSession({ db: leaseDb(), execution, runnerInstanceId: "runner", useRunnerd: true }).catch(error => error); + const outcome = executePaperclipNativeSession({ db: leaseDb(), execution, runnerInstanceId: "runner", useRunnerd }).catch(error => error); const persistence = cancellationDb(); try { await admitted; @@ -6593,6 +6596,25 @@ describe("native session bounded recovery", () => { expect(state.upsertRecoveryAction).not.toHaveBeenCalled(); }); + it.each(["pending", "acknowledged"])("preserves a %s Stop when cancellation wins before the first turn", async (dispatchState) => { + const updates: Array<{ table: unknown; values: Record }> = []; + const stop: Record = {}; + state.execute.mockReset().mockImplementationOnce(async () => { + Object.assign(stop, { nativeCancellation: { + schema: "paperclip.native-cancellation.v1", ...execution.binding, + scope: "run", reasonCode: "cancellation_run_only", dispatchState, + dispatched: true, intentAuditId: "intent", acknowledgementAuditId: "ack", + } }); + throw new Error("native_session_cancelled"); + }); + state.upsertRecoveryAction.mockClear(); + await expect(executePaperclipNativeSession({ + db: leaseDb(execution, {}, stop, updates), execution, runnerInstanceId: "stop-before-first-turn", + })).rejects.toThrow("native_cancellation_pending_recovery"); + expect(updates.some(update => update.table === heartbeatRuns && update.values.status === "failed")).toBe(false); + expect(state.upsertRecoveryAction).not.toHaveBeenCalled(); + }); + it("keeps typed integrity failure permanent even if a wrapper changes its message", () => { const failure = new NativeSessionProtocolIntegrityError( "semantic_input_digest_mismatch", diff --git a/server/src/services/native-runtime/native-session-executor.ts b/server/src/services/native-runtime/native-session-executor.ts index 904b3a6073..8a8bf41cd7 100644 --- a/server/src/services/native-runtime/native-session-executor.ts +++ b/server/src/services/native-runtime/native-session-executor.ts @@ -6980,16 +6980,27 @@ export async function executePaperclipNativeSession(input: { }, ) => Promise; }): Promise { - if (!input.useRunnerd) { - return executePaperclipNativeSessionWithinScope(input); + const runId = input.execution.binding.runId; + if (nativeSessionStartups.has(runId)) { + throw new Error("native_session_supervisor_busy"); } + // Register before the first asynchronous operation on either backend path. + // A duplicate execution must not replace the original startup handoff. + let resolveStartup!: (session: ActiveNativeSession | null) => void; + const startup: NativeSessionStartup = { + promise: new Promise(resolve => { resolveStartup = resolve; }), + resolve: session => resolveStartup(session), + }; + nativeSessionStartups.set(runId, startup); let preparedInput: typeof input = input; let cleanupStagedAttachments: () => Promise = async () => undefined; let sessionScopeId: string | null = null; let ownsSessionScope = false; let executionFailure: unknown; - let startup: NativeSessionStartup | undefined; try { + if (!input.useRunnerd) { + return await executePaperclipNativeSessionWithinScope(input); + } // The session scope is unaffected by appending server-staged attachment // descriptors. Claim it before any workspace scrub/write so a duplicate // execution cannot truncate or replace the active turn's staging inode. @@ -7002,11 +7013,6 @@ export async function executePaperclipNativeSession(input: { input.execution.binding.runId, ); ownsSessionScope = true; - let resolveStartup!: (session: ActiveNativeSession | null) => void; - const startupPromise = new Promise(resolve => { resolveStartup = resolve; }); - startup = { promise: startupPromise, resolve: resolveStartup }; - nativeSessionStartups.set(input.execution.binding.runId, startup); - // The shutdown sweep can have removed an idle owner while its remote // checkpoint is still being saved. The scope reservation also prevents // a later sweep from closing an owner this turn is about to acquire. @@ -7061,9 +7067,9 @@ export async function executePaperclipNativeSession(input: { executionFailure = error; throw error; } finally { - startup?.resolve(null); - if (startup && nativeSessionStartups.get(input.execution.binding.runId) === startup) { - nativeSessionStartups.delete(input.execution.binding.runId); + startup.resolve(null); + if (nativeSessionStartups.get(runId) === startup) { + nativeSessionStartups.delete(runId); } if ( ownsSessionScope && @@ -8303,14 +8309,28 @@ async function executePaperclipNativeSessionWithinScope( // Stop before paperclip_finish is normal. Bounded provider teardown has // finished; a missing result must not overwrite cancellation or trigger // another turn to perform completion bookkeeping. - if (protocolIntegrityFailure === null && error instanceof Error && error.message === "native_finalization_missing: session returned no semantic result") { + const stoppedBeforeFirstTurn = error instanceof Error && error.message === "native_session_cancelled"; + if (protocolIntegrityFailure === null && error instanceof Error && + (stoppedBeforeFirstTurn || error.message === "native_finalization_missing: session returned no semantic result")) { const [stoppedRun] = await input.db.select().from(heartbeatRuns).where(and( eq(heartbeatRuns.id, input.execution.binding.runId), eq(heartbeatRuns.companyId, input.execution.binding.companyId), eq(heartbeatRuns.agentId, input.execution.binding.agentId), eq(heartbeatRuns.nativeIssueId, input.execution.binding.issueId), )).limit(1); - if (stoppedRun && (hasAcknowledgedNativeStopIntent(stoppedRun) || hasAcknowledgedNativeReassignmentStopIntent(stoppedRun))) { + const stopIntent = record(stoppedRun?.resultJson?.nativeCancellation); + // The backend can reject the first turn while the Stop API is still + // recording its acknowledgement. Preserve that audited, exactly bound + // cancellation instead of racing it with a generic provider failure. + const pendingStartupStop = stoppedBeforeFirstTurn && + stopIntent.schema === "paperclip.native-cancellation.v1" && + stopIntent.companyId === input.execution.binding.companyId && + stopIntent.runId === input.execution.binding.runId && + stopIntent.issueId === input.execution.binding.issueId && + stopIntent.scope === "run" && + typeof stopIntent.intentAuditId === "string" && stopIntent.intentAuditId.length > 0 && + ["pending", "acknowledged"].includes(String(stopIntent.dispatchState)); + if (stoppedRun && (pendingStartupStop || hasAcknowledgedNativeStopIntent(stoppedRun) || hasAcknowledgedNativeReassignmentStopIntent(stoppedRun))) { const [settled] = await input.db.update(nativeRunFinalizations).set({ phase: "terminal_failure", failureCode: "native_retry_cancelled", nextAttemptAt: null, leaseOwner: null, leaseExpiresAt: null, controlDeadlineAt: null, recoveryState: null, diff --git a/server/src/services/tool-action-delivery.ts b/server/src/services/tool-action-delivery.ts index 752208f7e5..5a2d4b7f89 100644 --- a/server/src/services/tool-action-delivery.ts +++ b/server/src/services/tool-action-delivery.ts @@ -337,6 +337,29 @@ export function toolActionDeliveryService( } return { deliver, + async deliverForRun(input: { companyId: string; runId: string }) { + // A review can settle before its source run yields. Retry that run's + // receipts as soon as execution cleanup finishes; the periodic sweep + // remains the recovery path if this callback is interrupted. + const pending = await db + .select({ id: toolActionDeliveries.actionRequestId }) + .from(toolActionDeliveries) + .innerJoin(toolActionRequests, and( + eq(toolActionRequests.id, toolActionDeliveries.actionRequestId), + eq(toolActionRequests.companyId, input.companyId), + )) + .innerJoin(toolInvocations, and( + eq(toolInvocations.id, toolActionRequests.invocationId), + eq(toolInvocations.companyId, input.companyId), + eq(toolInvocations.runId, input.runId), + )) + .where(and( + eq(toolActionDeliveries.companyId, input.companyId), + isNull(toolActionDeliveries.deliveredAt), + inArray(toolActionRequests.status, terminalStatuses), + )); + for (const row of pending) await deliver(row.id); + }, async sweepPending() { let cursor: string | undefined; let scanned = 0; diff --git a/tests/runner-e2e/FIXTURES.md b/tests/runner-e2e/FIXTURES.md index f03bfda838..9b1acf8351 100644 --- a/tests/runner-e2e/FIXTURES.md +++ b/tests/runner-e2e/FIXTURES.md @@ -173,6 +173,14 @@ Screenshots are allowlisted to the exact disposable agent chat. Cleanup cancels all active runs in the isolated company, including handed-off work; usage from failed and cancelled runs must not disappear from campaign totals. + +Warm three-turn continuity grades the exact workspace file after each turn, +task completion, and sandbox/session identity. It also requires a visible +persisted final reply with each turn marker once and in order. It does not +grade exact final-reply wording; the hello +and continuation fixtures retain those exact-response checks. This separates +workspace persistence failures from model response-format variance. + `chat-hardening.ts` adds the explicit-only `agent-chat-hardening` journeys. Use the ordinary public APIs to seed source documents and blockers. Keep the answer out of the user's status/review request. Grade the exact source values, latest diff --git a/tests/runner-e2e/catalog.ts b/tests/runner-e2e/catalog.ts index 3addfd186f..11833435c6 100644 --- a/tests/runner-e2e/catalog.ts +++ b/tests/runner-e2e/catalog.ts @@ -849,16 +849,14 @@ export const daytonaWarmContinuityTask: RunnerTaskFixture = { warmTurnInstructions(3, nonce), ], buildMatchers(nonce, execution) { - const markers = ([1, 2, 3] as const).map((turn) => - warmTurnMarker(turn, nonce), - ); + // Workspace persistence is the oracle for this story. Exact response text + // formatting must not mask a valid workspace, but every warm turn still + // needs one visible marker in chronological order. Surrounding provider + // prose is allowed; the occurrence and order matchers grade only markers. + const markers = ([1, 2, 3] as const).map((turn) => warmTurnMarker(turn, nonce)); return [ - { kind: "message_exact", expected: markers[2] }, - ...markers.map( - (expected) => - ({ kind: "message_occurrences", expected, count: 1 }) as const, - ), - { kind: "message_ordered", expected: markers }, + ...markers.map((marker) => ({ kind: "message_occurrences" as const, expected: marker, count: 1 })), + { kind: "message_ordered" as const, expected: markers }, { kind: "file_exact", path: `daytona-warm-${nonce}.txt`, diff --git a/tests/runner-e2e/connection-reviews.ts b/tests/runner-e2e/connection-reviews.ts index 9808b3cd36..a92cfff203 100644 --- a/tests/runner-e2e/connection-reviews.ts +++ b/tests/runner-e2e/connection-reviews.ts @@ -47,7 +47,7 @@ export async function setupConnectionReview(input: { ); connection = { id: connected.connectionId }; } else { - await page.goto(`/${input.prefix}/apps`); + await page.goto(`/${input.prefix}/apps`, { waitUntil: "domcontentloaded" }); const connector = page .getByRole("list", { name: "Connector list" }) .getByRole("listitem") @@ -81,10 +81,14 @@ export async function setupConnectionReview(input: { }, ); expect(installed.ok()).toBe(true); - await page.goto(`/${input.prefix}/apps/${connection.id}/permissions`); - await page - .getByRole("radio", { name: "List fixture pages: Ask first" }) - .click(); + await page.goto(`/${input.prefix}/apps/${connection.id}/permissions`, { + waitUntil: "domcontentloaded", + }); + const askFirst = page.getByRole("radio", { + name: "List fixture pages: Ask first", + }); + await expect(askFirst).toBeVisible({ timeout: 30_000 }); + await askFirst.click(); return { ...provider, connectionId: connection.id, diff --git a/tests/runner-e2e/continuation-screenshot.ts b/tests/runner-e2e/continuation-screenshot.ts index 85506bcacb..d215d998ec 100644 --- a/tests/runner-e2e/continuation-screenshot.ts +++ b/tests/runner-e2e/continuation-screenshot.ts @@ -6,12 +6,26 @@ export async function captureLoadedContinuation( title: string, capture: () => Promise, timeout = 30_000, + expectedVisibleText?: string, +) { + await waitForTaskChatRendered(page, title, timeout, expectedVisibleText); + await page.evaluate(() => document.fonts.ready.then(() => undefined)); + await capture(); +} + +/** Wait for the persisted task projection before taking a browser screenshot. */ +export async function waitForTaskChatRendered( + page: Page, + title?: string, + timeout = 30_000, + expectedVisibleText?: string, ) { const thread = page.getByTestId("task-chat-thread"); await expect(thread).toBeVisible({ timeout }); await expect(thread.locator(':scope > [aria-busy="false"]')).toBeVisible({ timeout }); await expect(thread.getByTestId("task-chat-history-loading")).toHaveCount(0, { timeout }); - await expect(thread.getByRole("heading", { name: title, exact: true })).toBeVisible({ timeout }); - await page.evaluate(() => document.fonts.ready.then(() => undefined)); - await capture(); + if (title !== undefined) + await expect(thread.getByRole("heading", { name: title, exact: true })).toBeVisible({ timeout }); + if (expectedVisibleText !== undefined) + await expect(thread.getByTestId("task-chat-agent-bubble").filter({ hasText: expectedVisibleText })).toHaveCount(1, { timeout }); } diff --git a/tests/runner-e2e/everyday-flow.ts b/tests/runner-e2e/everyday-flow.ts index a0c4fc5d2c..428f1d761b 100644 --- a/tests/runner-e2e/everyday-flow.ts +++ b/tests/runner-e2e/everyday-flow.ts @@ -1,7 +1,7 @@ import { expect, type Page } from "@playwright/test"; import { spawn } from "node:child_process"; import { createHash } from "node:crypto"; -import { mkdir, readFile, stat, readdir } from "node:fs/promises"; +import { mkdir, readFile, readdir } from "node:fs/promises"; import path from "node:path"; import { isDeepStrictEqual } from "node:util"; import { pollUntil, type RunnerApi } from "./api.js"; @@ -12,6 +12,8 @@ import { import type { LiveFixtureValues } from "./live-fixtures.js"; import type { MatrixExecution } from "./types.js"; import { createTaskThroughUi, submitTaskReply } from "./user-actions.js"; +import { waitForTaskChatRendered } from "./continuation-screenshot.js"; +import { hasPersistedSource, isSavedSourceCheckpoint } from "./everyday-interruption.js"; import { setupConnectionReview } from "./connection-reviews.js"; import { pendingStoryDecision, @@ -170,13 +172,14 @@ export async function runEverydayFlow(input: Input) { parentBlockedObserved: boolean; } | undefined; async function workspaceFiles() { + const workspaceRoot = input.workspacePath; const files: Record = {}; - for (const entry of await readdir(input.workspacePath, { + for (const entry of await readdir(workspaceRoot, { withFileTypes: true, })) { if (entry.isFile() && /\.(py|md|zip)$/.test(entry.name)) files[entry.name] = createHash("sha256") - .update(await readFile(path.join(input.workspacePath, entry.name))) + .update(await readFile(path.join(workspaceRoot, entry.name))) .digest("hex"); } return files; @@ -229,8 +232,12 @@ export async function runEverydayFlow(input: Input) { } const taskUrl = (issue: StoryIssue) => `/${prefix}/issues/${issue.identifier ?? issue.id}`; + async function openTask(issue: StoryIssue) { + await page.goto(taskUrl(issue), { waitUntil: "domcontentloaded" }); + await waitForTaskChatRendered(page, String(issue.title)); + } async function openParent() { - await page.goto(taskUrl(parent!), { waitUntil: "domcontentloaded" }); + await openTask(parent!); } function observableAgentIds(state: EverydayEvidence) { return [ @@ -430,7 +437,7 @@ export async function runEverydayFlow(input: Input) { : zips[zips.length - 1]; if (!attachment) throw new Error("Selected delivery is no longer available"); const issue = ev.issues.find((i) => i.id === issueId)!; - await page.goto(taskUrl(issue), { waitUntil: "domcontentloaded" }); + await openTask(issue); const links = page.locator( `a[href*="/api/attachments/${attachment.id}/content"]`, ); @@ -469,34 +476,58 @@ export async function runEverydayFlow(input: Input) { }); await openParent(); } - async function sourceReady() { - await pollUntil({ - label: "saved source before interruption", - deadlineAt: Math.min(input.deadlineAt, Date.now() + 180_000), + async function recordSource( + label = "source-saved-before-interruption", + evidenceName = "source-before-interruption.json", + requireActive = false, + maxWaitMs = 180_000, + ) { + const filePath = path.join(input.workspacePath, "slugify.py"); + const bytes = await pollUntil({ + label: requireActive ? "active run with saved source" : "saved source after interruption", + deadlineAt: Math.min(input.deadlineAt, Date.now() + maxWaitMs), intervalMs: 500, load: async () => { await refresh(); - try { - return ( - (await stat(path.join(input.workspacePath, "slugify.py"))).size > - 0 && ev.runs.some(isActiveStoryRun) - ); - } catch { - return false; - } + const source = await readFile(filePath).catch(() => undefined); + return { active: ev.runs.some(isActiveStoryRun), source }; }, - accept: Boolean, + accept: ({ active, source }) => + requireActive + ? isSavedSourceCheckpoint(active, source) + : hasPersistedSource(source), }); - const bytes = await readFile(path.join(input.workspacePath, "slugify.py")); - await input.evidence("source-before-interruption.json", { - body: bytes.toString("utf8"), - sha256: createHash("sha256").update(bytes).digest("hex"), + if (!bytes.source) throw new Error("Saved source was not available at the controlled boundary"); + await input.evidence(evidenceName, { + body: bytes.source.toString("utf8"), + sha256: createHash("sha256").update(bytes.source).digest("hex"), }); - note("source-saved-before-interruption", { - sha256: createHash("sha256").update(bytes).digest("hex"), - bytes: bytes.length, + note(label, { + sha256: createHash("sha256").update(bytes.source).digest("hex"), + bytes: bytes.source.length, + active: bytes.active, }); } + async function sourceReady() { + await recordSource("source-saved-before-interruption", "source-before-interruption.json", true); + } + async function prepareStopBoundary() { + // Providers can finish a short first turn before the browser can click + // Stop. If that happens, submit one ordinary user follow-up through the + // composer and use that fresh run as the controlled interruption boundary. + // Save the checkpoint even if a fast provider has already finished. Do + // not shorten the normal source-creation budget to manufacture a timeout. + await recordSource("source-saved-before-interruption", "source-before-interruption.json"); + await refresh(); + if (ev.runs.some(isActiveStoryRun)) return; + await submitTaskReply(page, `${SLUGIFY_REVISION}\nContinue working until the source file is saved.`); + note("stop-boundary-continuation-submitted"); + await recordSource( + "source-saved-before-interruption", + "source-before-interruption.json", + true, + ); + } try { await mkdir(path.join(input.privateDir, "snapshots"), { recursive: true }); const revision = await runCommand("git", ["rev-parse", "HEAD"]); @@ -779,7 +810,7 @@ export async function runEverydayFlow(input: Input) { childId: child.id, activeRunIds: ev.runs.filter(isActiveStoryRun).map((r) => r.id), }); - await page.goto(taskUrl(child), { waitUntil: "domcontentloaded" }); + await openTask(child); await reply(LATE_REQUIREMENT, child); note("late-feedback-delivered-to-child", { childId: child.id }); await openParent(); @@ -804,7 +835,8 @@ export async function runEverydayFlow(input: Input) { load: refresh, accept: (s) => s.runs.some(isActiveStoryRun), }); - } else await sourceReady(); + } else if (caseId === "stop-redirect") await prepareStopBoundary(); + else await sourceReady(); const active = ev.runs.find((r) => r.status === "running"); if (!active) throw new Error( @@ -821,6 +853,7 @@ export async function runEverydayFlow(input: Input) { accept: Boolean, intervalMs: 250, }); + await recordSource("source-saved-after-interruption", "source-after-interruption.json"); stoppedWorkspace = await workspaceFiles(); note("stopped-workspace-snapshot", stoppedWorkspace); await reply( @@ -1411,13 +1444,7 @@ export async function runEverydayFlow(input: Input) { ), "No completion confirmation or unanswered interaction remains.", ); - await expect( - page - .locator( - '[data-testid="task-chat-thread"], [data-testid="thread-root"]', - ) - .first(), - ).toBeVisible(); + await waitForTaskChatRendered(page, String(parent!.title)); const latestAgentComment = ev.issues .find((i) => i.id === parent!.id) ?.comments?.filter((c: Row) => c.authorAgentId) diff --git a/tests/runner-e2e/everyday-interruption.ts b/tests/runner-e2e/everyday-interruption.ts new file mode 100644 index 0000000000..9d6611c5d4 --- /dev/null +++ b/tests/runner-e2e/everyday-interruption.ts @@ -0,0 +1,11 @@ +export function hasPersistedSource(value: Buffer | undefined): boolean { + return Boolean(value && value.length > 0); +} + +/** A Stop boundary is valid only when both facts hold at the same observation. */ +export function isSavedSourceCheckpoint( + active: boolean, + value: Buffer | undefined, +): boolean { + return active && hasPersistedSource(value); +} diff --git a/tests/runner-e2e/everyday.test.ts b/tests/runner-e2e/everyday.test.ts index b4bbdfba42..5016adf268 100644 --- a/tests/runner-e2e/everyday.test.ts +++ b/tests/runner-e2e/everyday.test.ts @@ -26,8 +26,20 @@ import { type StoryRun, type StoryIssue, } from "./everyday-observations.js"; +import { hasPersistedSource, isSavedSourceCheckpoint } from "./everyday-interruption.js"; describe("everyday workflow grader and review timing", () => { + it("requires nonempty persisted source evidence", () => { + expect(hasPersistedSource(undefined)).toBe(false); + expect(hasPersistedSource(Buffer.alloc(0))).toBe(false); + expect(hasPersistedSource(Buffer.from("source"))).toBe(true); + }); + + it("requires an active run and saved source at the same Stop checkpoint", () => { + expect(isSavedSourceCheckpoint(false, Buffer.from("source"))).toBe(false); + expect(isSavedSourceCheckpoint(true, Buffer.alloc(0))).toBe(false); + expect(isSavedSourceCheckpoint(true, Buffer.from("source"))).toBe(true); + }); it("uses base grading for review delivery and preserves late max-length cases", () => { expect(artifactGradeModeForPhase("reviewed-delivery")).toBe("base"); expect(artifactGradeModeForPhase("delegated-delivery")).toBe("max-length"); diff --git a/tests/runner-e2e/failure-classifier.ts b/tests/runner-e2e/failure-classifier.ts index f17f0fc543..b231fb9270 100644 --- a/tests/runner-e2e/failure-classifier.ts +++ b/tests/runner-e2e/failure-classifier.ts @@ -15,6 +15,10 @@ const NON_RETRYABLE_ACPX_SESSION_OPEN = /native_session_recovery_failed[\s\S]*ACPX sidecar command session\.open was rejected\s*\([^)]*\bretryable\s*=\s*false\b/i; const CANDIDATE = /(?:matcher|expected.*observed|marker|issue status|run status|runtime mode|wrong output|missing output)/i; +// File-transfer deadlines are sandbox transport failures. Do not generalize +// this to all RPC timeouts: runner protocol defects must remain visible. +const SANDBOX_TRANSFER_TIMEOUT = + /RPC call "environmentSync(?:In|Out)" timed out after \d+ms/i; export function classifyFailure(error: unknown): FailureClass { if (error instanceof ObservedStateTimeout) return error.failureClass; @@ -36,7 +40,8 @@ export function classifyFailure(error: unknown): FailureClass { NON_RETRYABLE_ACPX_SESSION_OPEN.test(message) ) return "candidate_failure"; - if (TRANSIENT.test(message)) return "transient_infrastructure"; + if (TRANSIENT.test(message) || SANDBOX_TRANSFER_TIMEOUT.test(message)) + return "transient_infrastructure"; if (CANDIDATE.test(message)) return "candidate_failure"; return "candidate_failure"; } diff --git a/tests/runner-e2e/runner.spec.ts b/tests/runner-e2e/runner.spec.ts index c93b6a8d75..73ce46984c 100644 --- a/tests/runner-e2e/runner.spec.ts +++ b/tests/runner-e2e/runner.spec.ts @@ -530,6 +530,7 @@ for (const execution of executions) { const companyRunFlow = ["continuation", "agent_chat", "everyday_workflow", "first_task"].includes(execution.task.flow); const consoleDiagnostics: Array> = []; const networkDiagnostics: Array> = []; + const pageLifecycleDiagnostics: Array> = []; let fixtures: LiveFixtureValues | undefined; let reviewProvider: Awaited> | undefined; let issue: IssueRecord | undefined; @@ -721,6 +722,22 @@ for (const execution of executions) { }); } }); + page.on("pageerror", (error) => { + pageLifecycleDiagnostics.push({ + type: "pageerror", + message: error.message, + stack: error.stack ?? null, + url: page.url(), + }); + }); + page.on("framenavigated", (frame) => { + if (frame === page.mainFrame()) + pageLifecycleDiagnostics.push({ + type: "navigation", + url: frame.url(), + at: new Date().toISOString(), + }); + }); page.on("requestfailed", (requestEvent) => { networkDiagnostics.push({ method: requestEvent.method(), @@ -2330,19 +2347,27 @@ for (const execution of executions) { const visibleAgentReplies = page .getByTestId("task-chat-thread") .getByTestId("task-chat-agent-bubble"); - const escapedMarker = marker.replace(/[.*+?^${}()|[\]\\]/g, "\\$&"); - const terminalAgentReplies = visibleAgentReplies.filter({ - hasText: new RegExp(`^\\s*${escapedMarker}\\s*$`), - }); - await expect(terminalAgentReplies).toHaveCount(1, { timeout: 30_000 }); - await expect(terminalAgentReplies.first()).toBeVisible(); - // A string-valued toHaveText assertion compares the complete rendered - // text while normalizing ordinary DOM whitespace. This keeps Markdown - // layout differences harmless without allowing prefixed, suffixed, or - // substituted provider prose to masquerade as the requested response. - await expect(terminalAgentReplies.first()).toHaveText(marker, { - useInnerText: true, - }); + if (execution.task.flow === "warm_three_turn") { + // Prove the persisted user-facing response is visible, independently + // of the byte-for-byte workspace checks and lease continuity checks. + expect(finalRunMessage.trim()).not.toBe(""); + await expect(visibleAgentReplies.filter({ hasText: finalRunMessage }).last()) + .toBeVisible({ timeout: 30_000 }); + } else { + const escapedMarker = marker.replace(/[.*+?^${}()|[\]\\]/g, "\\$&"); + const terminalAgentReplies = visibleAgentReplies.filter({ + hasText: new RegExp(`^\\s*${escapedMarker}\\s*$`), + }); + await expect(terminalAgentReplies).toHaveCount(1, { timeout: 30_000 }); + await expect(terminalAgentReplies.first()).toBeVisible(); + // A string-valued toHaveText assertion compares the complete rendered + // text while normalizing ordinary DOM whitespace. This keeps Markdown + // layout differences harmless without allowing prefixed, suffixed, or + // substituted provider prose to masquerade as the requested response. + await expect(terminalAgentReplies.first()).toHaveText(marker, { + useInnerText: true, + }); + } await expect( page.getByTestId("issue-detail-header").getByRole("button", { name: "Change status (current: Done)", @@ -2421,6 +2446,7 @@ for (const execution of executions) { { console: consoleDiagnostics, network: networkDiagnostics, + lifecycle: pageLifecycleDiagnostics, }, secrets, ); diff --git a/tests/runner-e2e/support.test.ts b/tests/runner-e2e/support.test.ts index 3be21db56e..e7c4cb95ff 100644 --- a/tests/runner-e2e/support.test.ts +++ b/tests/runner-e2e/support.test.ts @@ -705,6 +705,16 @@ describe("runner E2E run observations", () => { }); describe("runner E2E failure policy", () => { + it("classifies sandbox file-transfer RPC deadlines without hiding other RPC defects", () => { + for (const method of ["environmentSyncIn", "environmentSyncOut"]) { + expect(classifyFailure(new Error( + `Stopped waiting for everyday recover-controller settled: native execution failed native_session_interrupted: RPC call "${method}" timed out after 330000ms`, + ))).toBe("transient_infrastructure"); + } + expect(classifyFailure(new Error('RPC call "run.attach" timed out after 330000ms'))) + .toBe("candidate_failure"); + }); + it.each([ "native_session_close_unrecoverable: provider transport failed", "Provider connection closed: runner did not durably suspend before checkpoint", @@ -1331,3 +1341,58 @@ describe("persisted final response selection", () => { expect(persistedFinalRunMessage([], { id: "run-1", resultJson: { summary: "FINAL" } })).toBe(""); }); }); + +describe("warm continuity grading scope", () => { + it("checks workspace bytes, lifecycle, and ordered turn markers without exact response formatting", () => { + const execution = runnerMatrix.find((cell) => cell.task.flow === "warm_three_turn")!; + const matchers = execution.task.buildMatchers("test-nonce", execution); + expect(matchers).toContainEqual({ + kind: "message_occurrences", expected: "PAPERCLIP_E2E_WARM_T1_test-nonce", count: 1, + }); + expect(matchers).toContainEqual({ + kind: "message_occurrences", expected: "PAPERCLIP_E2E_WARM_T2_test-nonce", count: 1, + }); + expect(matchers).toContainEqual({ + kind: "message_occurrences", expected: "PAPERCLIP_E2E_WARM_T3_test-nonce", count: 1, + }); + expect(matchers).toContainEqual({ + kind: "message_ordered", + expected: [ + "PAPERCLIP_E2E_WARM_T1_test-nonce", + "PAPERCLIP_E2E_WARM_T2_test-nonce", + "PAPERCLIP_E2E_WARM_T3_test-nonce", + ], + }); + expect(matchers).toContainEqual({ + kind: "file_exact", path: "daytona-warm-test-nonce.txt", + expected: "T1-test-nonce\nT2-test-nonce\nT3-test-nonce\n", + }); + expect(matchers).toContainEqual({ kind: "issue_status", expected: "done" }); + const hello = runnerMatrix.find((cell) => cell.task.id === "hello-complete")!; + expect(hello.task.buildMatchers("test-nonce", hello).some((matcher) => matcher.kind === "message_exact")).toBe(true); + }); + + it("accepts warm-turn prose while rejecting missing, duplicate, or out-of-order markers", async () => { + const execution = runnerMatrix.find((cell) => cell.task.flow === "warm_three_turn")!; + const matchers = execution.task.buildMatchers("test-nonce", execution) + .filter((matcher) => matcher.kind.startsWith("message_")); + const passing = await evaluateMatchers(matchers, { + message: [ + "Turn one is complete: PAPERCLIP_E2E_WARM_T1_test-nonce.", + "Turn two is complete: PAPERCLIP_E2E_WARM_T2_test-nonce.", + "Turn three is complete: PAPERCLIP_E2E_WARM_T3_test-nonce.", + ].join("\n"), + }); + expect(passing.every((result) => result.passed)).toBe(true); + + const invalidMessages = [ + "PAPERCLIP_E2E_WARM_T1_test-nonce PAPERCLIP_E2E_WARM_T1_test-nonce PAPERCLIP_E2E_WARM_T2_test-nonce PAPERCLIP_E2E_WARM_T3_test-nonce", + "PAPERCLIP_E2E_WARM_T1_test-nonce PAPERCLIP_E2E_WARM_T3_test-nonce", + "PAPERCLIP_E2E_WARM_T3_test-nonce PAPERCLIP_E2E_WARM_T2_test-nonce PAPERCLIP_E2E_WARM_T1_test-nonce", + ]; + for (const message of invalidMessages) { + const results = await evaluateMatchers(matchers, { message }); + expect(results.some((result) => !result.passed)).toBe(true); + } + }); +});