From 3790ca2f136d26e7a3a0b6dffeec969811a8e1e9 Mon Sep 17 00:00:00 2001 From: Dotta <34892728+cryppadotta@users.noreply.github.com> Date: Mon, 21 Sep 2026 12:50:50 -0500 Subject: [PATCH] fix(runner): repair approval and Stop races and eval infrastructure (#13750) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work. > - Runner tasks must continue after approval and stop when the user presses Stop. > - Live evals found races at approval delivery and provider startup. > - Browser readiness and CI setup errors also hid the actual task results. > - This pull request fixes those races and the related test infrastructure. > - Regression tests and saved live reports show which cases now pass. ## Linked Issues or Issue Description Companion eval definitions PR: https://github.com/paperclipai/paperclip-evals/pull/25 (AgentCore paused and provider/environment infrastructure). Related: #13741 now supplies the late-startup Stop fence and warm-attachment recovery; this PR retains that fence and extends startup tracking and regression coverage to both native backend paths. #13539 introduced queued approvals during active runs. #13738 fixes child assignment, task replies, and warm process continuity and is already in the base. #13291 concerns automatic continuation of interrupted legacy sandbox runs; this PR fixes native startup cancellation and does not change that recovery policy. **What happened?** An accepted service approval could wait after its source run stopped. Stop could return success before the provider handle existed. Work could then start after Stop, or a cancelled run could be recorded as failed. Some E2E tests also failed on unloaded browser content or irrelevant reply wording. Runner CI could fail before model work because of dependency or sandbox setup. **Expected behavior** Deliver each settled approval once after its source run stops. Do not start work after an acknowledged Stop. Preserve the audited cancellation. Test the intended product behavior with a ready browser and verified runtime dependencies. **Steps to reproduce** 1. Approve a service request while its source run is active. Let the run finish. Check that its result starts one continuation. 2. Delay provider startup. Press Stop before its handle is available. Check cancellation, then submit `/new`. 3. Run the browser, warm-workspace, and Stop-and-redirect cases from the linked report. **Paperclip version or commit** The branch includes master at `9d19f98b5`. The report records the original source for each focused attempt. **Deployment mode** Isolated local development instances and disposable Daytona sandboxes. ## What Changed - Deliver settled tool-action results for the exact company and source run during final cleanup. Keep the existing idempotent receipt and periodic recovery sweep. - Wait for startup to hand off its provider handle before acknowledging Stop. Reject first-turn admission after cancellation. Preserve a matching audited pending or acknowledged cancellation. - Wait for mounted task history and connector controls in browser tests. Record failure evidence. Grade workspace contents and process continuity separately from exact reply wording. Require each warm-turn marker once and in order, allowing surrounding prose. - Stop-and-redirect now checks that the source file exists and work is active before Stop. - Resolve target dependency locks in an uncredentialed CI job. Verify the lock artifact hash. Keep orchestration and publication on the trusted workflow revision. - Materialize the pinned OpenCode executable and configure the exact Codex executable's user-namespace profile before provider credentials are available. - Compress Daytona directory uploads with gzip. Preserve files, executable modes, symlinks, empty directories, and confinement checks. - Classify file-transfer RPC deadlines as infrastructure. Keep unrelated runner RPC failures visible. ## Verification - [Focused live report with screenshots and original attempts](https://pages.paperclip.ing/runner-reliability-20260921/): 14 of 15 selected Product E2E cases pass across the recorded revisions. Claude and Codex Stop → `/new`, Claude service approval, delegation, both hiring/reuse cases, and native Daytona warm continuity pass. - Two credentialed Runner smoke cases pass. These are not full protocol coverage. - E2E harness after the master merge: 429 tests pass. E2E and server TypeScript checks pass. - Daytona plugin: 239 tests pass, 6 skipped. Plugin TypeScript build passes. The compression test fails against the old code and passes with the change. - Runner backend/runtime regression group: 161 tests pass. Cancellation/startup selection: 26 tests pass. Approval delivery: 34 real-database tests pass. - Workflow security: 7 tests pass. Both edited workflows pass actionlint. Runner TypeScript and Rust builds pass. - After merging master, all 389 native executor tests pass, including both native backend paths and late startup after the Stop deadline. - Post-merge `pnpm -r typecheck` and `pnpm build` pass. The monolithic local `pnpm test:run` was interrupted to integrate master and is inconclusive. The [hosted CI test partitions](https://github.com/paperclipai/paperclip/actions/runs/35620461738) pass on `50a3e43822bcba1e0d07b1b45b0be91cbf9312da`. An unchanged sandbox callback schema test initially received HTTP 503. It passed five isolated local runs, its full local test file, and one failed-job CI retry. No assertion was weakened. ## Risks - Stop can wait for the bounded startup handoff. If it cannot settle, the existing pending-recovery state remains instead of a false acknowledgement. - Immediate approval delivery must remain idempotent across cleanup and recovery sweeps. Tests cover duplicate delivery and company/run boundaries. - The workflow changes still need hosted Linux verification. They retain the trusted workflow and credential boundaries. - Gzip reduces the observed provider upload from about 1.8 GB to 663 MB. It does not yet fix the remaining Claude Daytona transfer timeout. That recovery test never reached Claude, so recovery remains unverified. Use a matching image with the verified provider package preinstalled for the next recovery test; retain cold-upload coverage separately. - The report preserves diagnostic runs with missing source metadata and marks them as such. It does not claim a new full-suite pass. - This PR adds no new prompt policy or historical status reconciliation. ## Model Used OpenAI GPT-6 through Codex performed the primary implementation and review. The exact primary backend model ID is not exposed in this session. OpenAI `gpt-5.6-luna` assisted with bounded infrastructure work and verification. The agents used repository tools, code execution, and browser tests. The exact backend revision and context-window size are not exposed in this session. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` / `Refs #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip Co-authored-by: OpenAI GPT-6 --- .github/workflows/runner-full-stack-e2e.yml | 5 + .../workflows/runner-protocol-live-evals.yml | 129 +++++++++++++++++- ...r-protocol-eval-workflow-security.test.mjs | 64 +++++++++ .../backends/harness-driver-backend.test.ts | 14 ++ .../src/backends/harness-driver-backend.ts | 4 + .../daytona/src/file-sync.test.ts | 82 +++++++++++ .../daytona/src/file-sync.ts | 15 +- .../daytona/src/plugin.test.ts | 8 +- .../__tests__/tool-gateway-service.test.ts | 9 ++ server/src/services/heartbeat.ts | 8 ++ .../native-session-executor.test.ts | 30 +++- .../native-runtime/native-session-executor.ts | 46 +++++-- server/src/services/tool-action-delivery.ts | 23 ++++ tests/runner-e2e/FIXTURES.md | 8 ++ tests/runner-e2e/catalog.ts | 16 +-- tests/runner-e2e/connection-reviews.ts | 14 +- tests/runner-e2e/continuation-screenshot.ts | 20 ++- tests/runner-e2e/everyday-flow.ts | 95 ++++++++----- tests/runner-e2e/everyday-interruption.ts | 11 ++ tests/runner-e2e/everyday.test.ts | 12 ++ tests/runner-e2e/failure-classifier.ts | 7 +- tests/runner-e2e/runner.spec.ts | 52 +++++-- tests/runner-e2e/support.test.ts | 65 +++++++++ 23 files changed, 644 insertions(+), 93 deletions(-) create mode 100644 tests/runner-e2e/everyday-interruption.ts 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); + } + }); +});