Merge pull request #3299 from deepseek-harness/worktree/fix-subagent-settings-label-color

fix(snapshot): bind fixture servers to OS-assigned ports
This commit is contained in:
Yichen Jiang 2026-08-29 13:17:27 +08:00 • committed by GitHub
commit ee309c0794
18 changed files with 654 additions and 77 deletions

View file

@ -0,0 +1,6 @@
# Bilingual-pair consistency record (docs/i18n/README.md): the git blob hash of each
# side as of the last confirmed-consistent state. Both languages carry equal authority;
# after editing either side, bring the other along and re-record with:
# pnpm run verify-translation-pairing --write .agents/notes/implemented/bug-fix/2026-08-29-windows-atomic-replace-retry.md
2026-08-29-windows-atomic-replace-retry.md: 4db5de6403be7ec39a1568a11d8877cba1ed5838
2026-08-29-windows-atomic-replace-retry.zh.md: 0138727ac0fe12af51b5a383b60859977300353c

View file

@ -0,0 +1,27 @@
# Agent Note: Retry transient Windows atomic replacements
Status: implemented
English | [中文](2026-08-29-windows-atomic-replace-retry.zh.md)
## Problem
Windows can temporarily reject a rename that replaces an existing file with `EACCES`, `EBUSY`, or `EPERM` while another system component holds the target. The cross-process writer lock orders cooperating application writers but cannot release that external handle, so treating the first error as permanent makes an otherwise valid settings or credentials update fail nondeterministically.
## Decision
`writeFileAtomic` owns replacement retry because every file-backed store needs the same guarantee. On Windows only, it retries `EACCES`, `EBUSY`, and `EPERM` up to eight times with exponential delays from 20 to 200 milliseconds. The same fully written temporary sibling remains the rename source throughout, and a caller-held writer lock remains held until `writeFileAtomic` settles.
Other error codes and other operating systems fail immediately. Exhausting the retry budget rethrows the final filesystem error after removing the temporary sibling; the existing target remains unchanged because no attempt deletes or truncates it.
## Alternatives considered
**Retry the credentials mutation.** A consumer-level retry would leave settings and future stores exposed, and replaying a read-modify-write operation can repeat work outside the atomic replacement. The shared primitive is the narrow owner of replacement-only retry.
**Delete the target before rename.** Removing the target can make readers observe an absent file and forfeits atomic replacement, so it cannot be a recovery step.
**Retry indefinitely.** A permanent permission error would then hang the writer and any lock contender. A bounded delay absorbs transient file use while preserving a predictable failure outcome.
## Consequences
A transient Windows handle can delay one replacement by at most 1.1 seconds before the final attempt fails. During that interval readers continue to see the complete old target, and success still consists of one atomic rename. Regression tests inject every retried code, permanent and non-Windows failures, and retry exhaustion; they observe rename attempts and advance fake timers rather than depending on wall-clock sleeps.

View file

@ -0,0 +1,27 @@
# Agent Note: 重试 Windows 上的瞬时原子替换失败
Status: implemented
[English](2026-08-29-windows-atomic-replace-retry.md) | 中文
## 问题
当另一个系统组件持有目标文件时,Windows 可能以 `EACCES`、`EBUSY` 或 `EPERM` 暂时拒绝替换已有文件的 rename。跨进程写锁能够排序应用内互相协作的写入方,却无法释放该外部句柄,因此把第一次错误当作永久失败会让本来有效的设置或凭据更新随机失败。
## 决策
`writeFileAtomic` 负责替换重试,因为每个文件型存储都需要相同保证。它仅在 Windows 上重试 `EACCES`、`EBUSY` 与 `EPERM`,最多八次,延迟从 20 毫秒指数增长至 200 毫秒。整个过程中,同一份已经完整写入的临时兄弟文件始终作为 rename 来源;调用方持有的写锁也会保持到 `writeFileAtomic` 结束。
其他错误码和其他操作系统会立即失败。重试预算耗尽后,函数移除临时兄弟文件并重新抛出最后一个文件系统错误;由于任何尝试都不会删除或截断现有目标,目标内容保持不变。
## 考虑过的替代方案
**重试凭据变更。** 消费方级重试仍会让设置和未来存储暴露于同一问题,而且重放一次读-修改-写操作可能重复原子替换之外的工作。共享原语是只负责替换重试的最窄所有者。
**在 rename 前删除目标。** 删除目标会让读取方观察到文件缺失,并放弃原子替换,因此不能作为恢复步骤。
**无限重试。** 永久权限错误会由此挂住写入方与所有锁竞争者。有界延迟可以吸收瞬时文件占用,同时保留可预测的失败结果。
## 后果
一个瞬时 Windows 句柄最多会让单次替换多等待 1.1 秒,随后最终尝试失败。在此期间,读取方继续看到完整的旧目标;成功仍由一次原子 rename 完成。回归测试注入每种可重试错误、永久错误、非 Windows 错误与重试耗尽,并观察 rename 尝试和推进伪时钟,而不依赖真实时间 sleep。

View file

@ -2,5 +2,5 @@
# side as of the last confirmed-consistent state. Both languages carry equal authority;
# after editing either side, bring the other along and re-record with:
# pnpm run verify-translation-pairing --write .agents/notes/implemented/testing/2026-08-24-session-log-snapshot-corpus.md
2026-08-24-session-log-snapshot-corpus.md: 328f0346554158d531dbda4b31a28277e37cc6dc
2026-08-24-session-log-snapshot-corpus.zh.md: 374ea2a1939e6e063f348e21fb74642371c345ac
2026-08-24-session-log-snapshot-corpus.md: 8b2f98e0a691e3085ff2286af3183048209f7ab8
2026-08-24-session-log-snapshot-corpus.zh.md: 19952f05c4803d50a6e3c7c987cadd9482d5e87b

View file

@ -18,6 +18,8 @@ This decision supersedes the ACP-specific placement and controller ownership in
The recorded session remains the primary input and expected output. Human-originated messages drive the selected public interface, recorded assistant chunks drive deterministic model replay, and the normalized persisted result must equal the fixture. Parent and child sessions share one typed redaction map. Committed fixtures contain relationship-preserving identity tokens and replace request system prompts and tool schemas with tokens; each distinct header class retains one explicit sidecar owner.
Scenario-owned HTTP fixtures separate the stable authority recorded in the session from their transport listener. Each fixture binds loopback port `0`, lets the operating system allocate and bind the port atomically, and maps the recorded URL or endpoint through the real provider to that listener. Any process-global transport interception matches only the recorded endpoint, is owned by the fixture fiber, and is restored before the listener closes.
Every existing ACP scenario receives a behavior-preserving destination. Ordinary one-shot behavior uses the headless profile, persistent machine control uses the SDK profile, and only ACP protocol behavior remains ACP-owned. Web scenarios driven by a recorded session join the corpus and retain their ARIA or geometry expected output as secondary evidence. Web and package tests without a recorded-session source keep owner-local expected output and stop using snapshot paths or filenames.
Workspace inputs remain scenario-local. A mutating scenario compares a complete expected final workspace that record and refresh never rewrite, so a model or tool self-report cannot satisfy the test. Existing intentional session reuse remains an explicit acyclic owner reference; the corpus adds no workspace inheritance or general fixture-merging mechanism.
@ -34,6 +36,10 @@ Workspace inputs remain scenario-local. A mutating scenario compares a complete
**Deduplicate workspaces and recorded sessions automatically.** The current workspace duplication is small and intentional locality is easier to review. Only existing semantic session reuse justifies an explicit reference.
**Bind the recorded URL's numeric port.** A stable listener port keeps transport and transcript values identical, but concurrent snapshot jobs on one host share the network namespace and race for that port.
**Probe an unused port before launching the scenario.** Releasing a probed port before the child binds it creates a time-of-check/time-of-use race. Binding port `0` inside the owning process keeps allocation and ownership atomic.
## Invariants
- Every existing recorded-session scenario has one passing replacement before its old owner is removed.
@ -43,11 +49,12 @@ Workspace inputs remain scenario-local. A mutating scenario compares a complete
- Mutating scenarios verify their final workspace externally.
- Owner-local process expectations use `*.expected.e2e.ts` and a separate built-output gate.
- Source and built adapters install replay-only packages in isolated profile fallbacks; distinct prompt-section orders keep their request headers byte-identical.
- Scenario HTTP fixtures bind OS-assigned loopback ports while preserving their recorded model-visible authorities.
- Source and built launch modes, browser replay, SDK projections, packaged Python runtime cases, documentation gates, and repository hygiene pass.
## Consequences
The corpus makes controller ownership visible: ordinary Agent behavior no longer inherits ACP protocol output, SDK and Web projections retain their interface-specific evidence, and only ACP cancellation and permission exchanges remain ACP-owned. Contributors review one normalized session diff plus the sidecars or UI expectations that add independent evidence. Adding a composition requires a manifest class pin; adding a volatile identity requires a typed relationship-preserving redaction rule rather than a broader text scrubber.
The corpus makes controller ownership visible: ordinary Agent behavior no longer inherits ACP protocol output, SDK and Web projections retain their interface-specific evidence, and only ACP cancellation and permission exchanges remain ACP-owned. Contributors review one normalized session diff plus the sidecars or UI expectations that add independent evidence. Adding a composition requires a manifest class pin; adding a volatile identity requires a typed relationship-preserving redaction rule rather than a broader text scrubber. Concurrent jobs can replay network-backed fixtures without reserving repository-wide ports, at the cost of a fixture-local mapping between the recorded authority and its transport listener.
## Risks

View file

@ -18,6 +18,8 @@ Status: implemented
录制会话仍是主要输入和预期输出。来自用户的消息驱动所选公开接口,录制的 assistant chunk 驱动确定性模型回放,规范化后的持久化结果必须等于 fixture。父会话和子会话共享同一类型化脱敏映射。提交的 fixture 使用保留关系的身份 token,并将请求 system prompt 和工具 schema 替换为 token;每个不同 header 类仍保留一个显式 sidecar 所有者。
场景拥有的 HTTP fixture 将会话中录制的稳定 authority 与传输 listener 分离。每个 fixture 在回环地址上绑定端口 `0`,由操作系统以一次原子操作分配并绑定端口,再将录制的 URL 或 endpoint 通过真实 provider 映射到该 listener。任何进程全局传输拦截只匹配录制 endpoint,由 fixture fiber 拥有,并在关闭 listener 前恢复。
每个现有 ACP 场景都获得一个保留行为的目标。普通单次行为使用 headless profile,需要持久机器控制的行为使用 SDK profile,只有 ACP 协议行为继续归 ACP 所有。由录制会话驱动的 Web 场景加入该语料,并保留其 ARIA 或几何预期输出作为辅助证据。没有录制会话来源的 Web 和包级测试保留归属方本地的预期输出,并停止使用快照路径或文件名。
Workspace 输入继续归各场景本地所有。变更文件的场景比较完整的预期最终 workspace,record 与 refresh 绝不改写该预期,因此模型或工具的自报结果无法满足测试。现有的有意会话复用继续使用显式、无环的所有者引用;语料不增加 workspace 继承或通用 fixture 合并机制。
@ -34,6 +36,10 @@ Workspace 输入继续归各场景本地所有。变更文件的场景比较完
**自动去重 workspace 和录制会话。** 当前 workspace 重复很少,有意保持本地性更易审查。只有现有的语义会话复用值得显式引用。
**直接绑定录制 URL 的数值端口。** 稳定 listener 端口使传输值与 transcript 值一致,但同一主机上的并发快照 job 共享网络命名空间,会争用该端口。
**在启动场景前探测未使用端口。** 子进程绑定前释放已探测端口会产生检查时间与使用时间竞态。在拥有该端口的进程内绑定端口 `0`,可使分配与所有权保持原子性。
## Invariants
- 每个现有录制会话场景都在移除旧所有者之前拥有一个通过的替代场景。
@ -43,11 +49,12 @@ Workspace 输入继续归各场景本地所有。变更文件的场景比较完
- 变更内容的场景从外部验证最终 workspace。
- 所属位置的进程预期使用 `*.expected.e2e.ts`,并由单独的构建产物门禁运行。
- 源码与构建适配器在隔离的 profile fallback 中安装仅回放包;不同的提示词 section 顺序值使两种模式的请求 header 保持字节一致。
- 场景 HTTP fixture 绑定由操作系统分配的回环端口,同时保留录制的模型可见 authority。
- 源码和构建启动模式、浏览器回放、SDK 投影、打包 Python 运行时场景、文档门禁和仓库卫生检查通过。
## Consequences
该语料让控制器所有权可见:普通 Agent 行为不再继承 ACP 协议输出,SDK 和 Web 投影保留各自接口专有的证据,只有 ACP 取消与权限交换仍归 ACP 所有。贡献者审查一份规范化会话差异,以及提供独立证据的 sidecar 或 UI 预期。新增组合必须提供 manifest 类别 pin;新增易变身份必须添加保留关系的带类型脱敏规则,而不是扩大文本清洗范围。
该语料让控制器所有权可见:普通 Agent 行为不再继承 ACP 协议输出,SDK 和 Web 投影保留各自接口专有的证据,只有 ACP 取消与权限交换仍归 ACP 所有。贡献者审查一份规范化会话差异,以及提供独立证据的 sidecar 或 UI 预期。新增组合必须提供 manifest 类别 pin;新增易变身份必须添加保留关系的带类型脱敏规则,而不是扩大文本清洗范围。并发 job 可以回放依赖网络的 fixture,而无需预留仓库级端口,代价是 fixture 内需要维护录制 authority 与传输 listener 的映射。
## Risks

View file

@ -2,5 +2,5 @@
# side as of the last confirmed-consistent state. Both languages carry equal authority;
# after editing either side, bring the other along and re-record with:
# pnpm run verify-translation-pairing --write packages/util/atomic-write/README.md
README.md: 22806b539faaa7668c37d863c20ffced2576bde6
README.zh.md: 1468f06aa4d46d9ea7c471bbb045314b67ae595e
README.md: 69daf671ba9d1269643533a6bb6e64462b8bee05
README.zh.md: 8a8613c673c4d12634c686cda7f2ced9957492b2

View file

@ -36,7 +36,7 @@ declare const text: string
await writeFileAtomic('/home/u/.dsh/settings.yaml', text, { mode: 0o600 })
```
Parent directories are created as needed, and readers observe either the old or the new complete content. On any failure the temporary file is removed and the failure is rethrown, so a failed replacement leaves the target untouched.
Parent directories are created as needed, and readers observe either the old or the new complete content. On Windows, transient replacement interference reported as `EACCES`, `EBUSY`, or `EPERM` is retried for a bounded interval; any remaining failure removes the temporary file and leaves the target untouched.
### Coordinating writers
@ -79,7 +79,7 @@ The package is built on one separation: the atomic commit owns the swap, and the
### Write path
`writeFileAtomic` writes a random-suffix sibling opened with exclusive create (`wx`), then renames it over the target. The exclusive open refuses to follow a symlink planted at a guessable temp path; the same-directory sibling keeps the rename on one filesystem; and the rename replaces a symlinked target itself instead of writing through to its referent.
`writeFileAtomic` writes a random-suffix sibling opened with exclusive create (`wx`), then renames it over the target. The exclusive open refuses to follow a symlink planted at a guessable temp path; the same-directory sibling keeps the rename on one filesystem; and the rename replaces a symlinked target itself instead of writing through to its referent. A Windows retry keeps the same complete sibling and uses bounded exponential backoff, so temporary use of the target by software outside the cooperative writer lock cannot turn a safe replacement into an immediate failure; the [retry decision](../../../.agents/notes/implemented/bug-fix/2026-08-29-windows-atomic-replace-retry.md) owns the rationale and rejected alternatives.
`withFileLock` creates a `<filename>.lock` sibling with `wx`. `EEXIST` identifies contention directly; `EPERM` does so only when a fresh `lstat` confirms the lock path exists, covering Windows exclusive-create behavior without hiding an unrelated permission failure. The lock records its creator's PID and is removed by the holder in a `finally`; contention backs off exponentially and fails when the per-call `waitMs` deadline (default two seconds) passes.

View file

@ -36,7 +36,7 @@ declare const text: string
await writeFileAtomic('/home/u/.dsh/settings.yaml', text, { mode: 0o600 })
```
父目录会按需创建,读取方只会观察到旧内容或完整的新内容。任何失败都会移除临时文件并重新抛出该失败,因此一次失败的替换不会改动目标文件。
父目录会按需创建,读取方只会观察到旧内容或完整的新内容。在 Windows 上,报告为 `EACCES`、`EBUSY` 或 `EPERM` 的瞬时替换干扰会在有界时间内重试;任何剩余失败都会移除临时文件,并保持目标文件不变。
### 协调写入方
@ -79,7 +79,7 @@ await withFileLock('/home/u/.dsh/settings.yaml', async () => {
### 写入路径
`writeFileAtomic` 先以独占创建(`wx`)打开一个随机后缀的同级文件并写入内容,然后 rename 到目标上。独占打开拒绝跟随预先埋在可猜测临时路径上的符号链接;同目录兄弟文件保证 rename 落在同一文件系统上;rename 替换的是符号链接目标本身,绝不写穿到其指向的文件。
`writeFileAtomic` 先以独占创建(`wx`)打开一个随机后缀的同级文件并写入内容,然后 rename 到目标上。独占打开拒绝跟随预先埋在可猜测临时路径上的符号链接;同目录兄弟文件保证 rename 落在同一文件系统上;rename 替换的是符号链接目标本身,绝不写穿到其指向的文件。Windows 重试会保留同一份完整的兄弟文件,并采用有界指数退避,因此协作式写锁之外的软件瞬时占用目标时,不会让安全替换立即失败;[重试决策](../../../.agents/notes/implemented/bug-fix/2026-08-29-windows-atomic-replace-retry.zh.md)记录了理由与被拒绝的替代方案。
`withFileLock` 以 `wx` 创建 `<filename>.lock` 同级文件。`EEXIST` 直接表示竞争;只有一次新的 `lstat` 确认锁路径存在时,`EPERM` 才表示竞争,从而兼容 Windows 的独占创建行为,又不掩盖无关的权限故障。锁记录创建者的 PID,由持有者在 `finally` 中移除;竞争按指数退避,在每次调用声明的 `waitMs` 期限(默认两秒)过后失败。

View file

@ -14,6 +14,33 @@ import { randomBytes } from 'node:crypto'
import { lstat, mkdir, rename, rm, writeFile } from 'node:fs/promises'
import { dirname } from 'node:path'
const WINDOWS_TRANSIENT_RENAME_ERRORS: ReadonlySet<string> = new Set(['EACCES', 'EBUSY', 'EPERM'])
const WINDOWS_RENAME_RETRY_INITIAL_MS = 20
const WINDOWS_RENAME_RETRY_MAX_MS = 200
const WINDOWS_RENAME_RETRY_LIMIT = 8
/** Whether Windows reported temporary interference with an atomic replacement. */
function isTransientWindowsRenameError(error: unknown): boolean {
if (process.platform !== 'win32') return false
return WINDOWS_TRANSIENT_RENAME_ERRORS.has((error as NodeJS.ErrnoException | null)?.code ?? '')
}
/** Replace the target after bounded retries for transient Windows interference. */
async function renameAtomicTemp(temp: string, filename: string): Promise<void> {
let delay = WINDOWS_RENAME_RETRY_INITIAL_MS
for (let retries = 0;; retries += 1) {
try {
await rename(temp, filename)
return
} catch (error) {
if (!isTransientWindowsRenameError(error)) throw error
if (retries >= WINDOWS_RENAME_RETRY_LIMIT) throw error
}
await new Promise(resolve => setTimeout(resolve, delay))
delay = Math.min(delay * 2, WINDOWS_RENAME_RETRY_MAX_MS)
}
}
/**
* Filesystem options for {@link writeFileAtomic}; `mode` is required so the
* permission decision stays visible at every call site.
@ -40,8 +67,10 @@ export interface WriteFileAtomicOptions {
* rename, so replacing a wider-permission file narrows it without a chmod
* race. The rename also replaces a symlinked target itself instead of writing
* through to its referent, and the same-directory sibling keeps the rename on
* one filesystem. On any failure the temp file is removed and the failure
* rethrown. Crash durability (fsync) is out of scope.
* one filesystem. Windows replacement retries transient `EACCES`, `EBUSY`,
* and `EPERM` failures for a bounded interval while the complete temp file
* remains the rename source. On any remaining failure the temp file is
* removed and the failure rethrown. Crash durability (fsync) is out of scope.
* @param filename - final path receiving the content.
* @param content - complete next file content.
* @param options - permission bits for the replacement inode.
@ -56,7 +85,7 @@ export async function writeFileAtomic(filename: string, content: string, options
const temp = `${filename}.${randomBytes(6).toString('hex')}.tmp`
try {
await writeFile(temp, content, { mode: options.mode, flag: 'wx' })
await rename(temp, filename)
await renameAtomicTemp(temp, filename)
} catch (error) {
await rm(temp, { force: true })
throw error

View file

@ -1,15 +1,28 @@
import { lstat, mkdir, mkdtemp, readFile, readdir, rm, stat, symlink, writeFile } from 'node:fs/promises'
import { lstat, mkdtemp, readFile, readdir, rm, stat, symlink, writeFile } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { dirname, join } from 'node:path'
import { afterEach, describe, expect, it, vi } from 'vitest'
import { withFileLock, writeFileAtomic } from '../src/index.ts'
const state = vi.hoisted(() => ({ failLockCreateWithEPERM: false }))
const state = vi.hoisted(() => ({
failLockCreateWithEPERM: false,
renameAttempts: 0,
renameFailures: [] as string[],
}))
vi.mock('node:fs/promises', async (importOriginal) => {
const actual = await importOriginal<typeof import('node:fs/promises')>()
return {
...actual,
rename: (async (...args: Parameters<typeof actual.rename>) => {
state.renameAttempts += 1
const code = state.renameFailures.shift()
if (code !== undefined) {
if (code === 'NO_CODE') throw new Error('injected rename failure without a code')
throw Object.assign(new Error(`${code}: injected rename failure`), { code })
}
return actual.rename(...args)
}),
writeFile: (async (path: unknown, ...rest: never[]) => {
if (state.failLockCreateWithEPERM && String(path).endsWith('.lock')) {
state.failLockCreateWithEPERM = false
@ -20,12 +33,26 @@ vi.mock('node:fs/promises', async (importOriginal) => {
}
})
afterEach(() => {
const scratchDirs: string[] = []
afterEach(async () => {
vi.useRealTimers()
vi.restoreAllMocks()
state.failLockCreateWithEPERM = false
state.renameAttempts = 0
state.renameFailures.length = 0
await Promise.all(scratchDirs.splice(0).map(dir => rm(dir, {
force: true,
maxRetries: 10,
recursive: true,
retryDelay: 20,
})))
})
async function scratch(): Promise<string> {
return mkdtemp(join(tmpdir(), 'dsh-atomic-write-'))
const dir = await mkdtemp(join(tmpdir(), 'dsh-atomic-write-'))
scratchDirs.push(dir)
return dir
}
/** Resolve once the lockfile exists, so contention is measured against a held lock. */
@ -44,9 +71,12 @@ describe('writeFileAtomic', () => {
it('creates the file and its parents with exactly the stated mode', async () => {
const dir = await scratch()
const target = join(dir, 'nested', 'deep', 'doc.yaml')
await writeFileAtomic(target, 'a: 1\n', { mode: 0o600 })
await writeFileAtomic(target, 'a: 1\n', { dirMode: 0o700, mode: 0o600 })
expect(await readFile(target, 'utf8')).toBe('a: 1\n')
if (process.platform !== 'win32') expect((await stat(target)).mode & 0o777).toBe(0o600)
if (process.platform !== 'win32') {
expect((await stat(dirname(target))).mode & 0o777).toBe(0o700)
expect((await stat(target)).mode & 0o777).toBe(0o600)
}
})
it('replaces existing content and narrows a wider-permission file to the stated mode', async () => {
@ -70,13 +100,62 @@ describe('writeFileAtomic', () => {
expect(await readFile(victim, 'utf8')).toBe('victim-content')
})
it('leaves no temp sibling and rethrows when the rename fails', async () => {
it('retries transient Windows rename interference and commits the replacement', async () => {
vi.spyOn(process, 'platform', 'get').mockReturnValue('win32')
vi.useFakeTimers()
const dir = await scratch()
const target = join(dir, 'occupied')
await mkdir(target)
await expect(writeFileAtomic(target, 'content', { mode: 0o600 })).rejects.toThrow()
const target = join(dir, 'document')
await writeFile(target, 'old')
state.renameFailures.push('EACCES', 'EBUSY', 'EPERM')
const replacement = writeFileAtomic(target, 'new', { mode: 0o600 })
await vi.waitFor(() => { expect(state.renameAttempts).toBeGreaterThan(0) })
await vi.runAllTimersAsync()
await replacement
expect(state.renameAttempts).toBe(4)
expect(await readFile(target, 'utf8')).toBe('new')
expect((await readdir(dir)).filter(entry => entry.includes('.tmp'))).toEqual([])
})
it('leaves no temp sibling after bounded Windows rename retries expire', async () => {
vi.spyOn(process, 'platform', 'get').mockReturnValue('win32')
vi.useFakeTimers()
const dir = await scratch()
const target = join(dir, 'document')
await writeFile(target, 'old')
state.renameFailures.push(...Array.from({ length: 9 }, () => 'EPERM'))
const replacement = writeFileAtomic(target, 'new', { mode: 0o600 })
await vi.waitFor(() => { expect(state.renameAttempts).toBeGreaterThan(0) })
await vi.runAllTimersAsync()
await expect(replacement).rejects.toMatchObject({ code: 'EPERM' })
expect(state.renameAttempts).toBe(9)
expect(await readFile(target, 'utf8')).toBe('old')
expect((await readdir(dir)).filter(entry => entry.includes('.tmp'))).toEqual([])
})
it('does not retry a Windows rename failure without a transient code', async () => {
vi.spyOn(process, 'platform', 'get').mockReturnValue('win32')
const dir = await scratch()
const target = join(dir, 'document')
state.renameFailures.push('NO_CODE')
await expect(writeFileAtomic(target, 'new', { mode: 0o600 })).rejects.toThrow(/without a code/)
expect(state.renameAttempts).toBe(1)
expect((await readdir(dir)).filter(entry => entry.includes('.tmp'))).toEqual([])
})
it('does not retry rename permission failures outside Windows', async () => {
vi.spyOn(process, 'platform', 'get').mockReturnValue('linux')
const dir = await scratch()
const target = join(dir, 'document')
state.renameFailures.push('EPERM')
await expect(writeFileAtomic(target, 'new', { mode: 0o600 })).rejects.toMatchObject({ code: 'EPERM' })
expect(state.renameAttempts).toBe(1)
})
})
describe('withFileLock', () => {

View file

@ -0,0 +1,203 @@
import { EventEmitter } from 'node:events'
import { Context } from '@deepseek-ai/cordis'
import { afterEach, describe, expect, it, vi } from 'vitest'
const httpMock = vi.hoisted(() => ({ createServer: vi.fn() }))
vi.mock('node:http', () => ({ createServer: httpMock.createServer }))
// Snapshot plugins are plain runtime JavaScript loaded by cordis.yml.
// @ts-expect-error The fixture intentionally has no declaration artifact.
import * as searchFixtureModule from '../snapshots/session/web-search-endpoint-guidance/web-search-error-fixture.mjs'
// @ts-expect-error The fixture intentionally has no declaration artifact.
import * as loopbackFixtureModule from '../snapshots/session/loopback-fixture-server.mjs'
const RECORDED_ENDPOINT = 'http://127.0.0.1:43118/anthropic/v1/messages'
interface FixturePlugin {
readonly name: string
readonly inject?: readonly string[]
apply(ctx: Context): Promise<void>
}
interface LoopbackFixtureOptions {
readonly label: string
readonly onCleanup: () => void
readonly onListening: (address: { port: number }) => void
readonly requestListener: () => void
}
const searchFixture = searchFixtureModule as unknown as FixturePlugin
const typedLoopbackFixtureModule = loopbackFixtureModule as unknown as {
readonly applyLoopbackServerEffect: (ctx: Context, options: LoopbackFixtureOptions) => Promise<void>
}
const { applyLoopbackServerEffect } = typedLoopbackFixtureModule
const nativeFetch = globalThis.fetch
class FixtureServer extends EventEmitter {
readonly started = Promise.withResolvers<undefined>()
listening = false
closed = false
connectionsClosed = false
unreferenced = false
private listenCallback: (() => void) | undefined
private port = 0
listen(_port: number, _host: string, callback: () => void): this {
this.listenCallback = callback
this.started.resolve(undefined)
return this
}
finishListening(port = 54321): void {
this.port = port
this.listening = true
this.listenCallback?.()
}
address(): { address: string; family: string; port: number } | null {
return this.listening ? { address: '127.0.0.1', family: 'IPv4', port: this.port } : null
}
unref(): this {
this.unreferenced = true
return this
}
close(callback: (error?: Error) => void): this {
this.listening = false
this.closed = true
callback()
return this
}
closeAllConnections(): void {
this.connectionsClosed = true
}
}
function nextServer(): FixtureServer {
const server = new FixtureServer()
httpMock.createServer.mockReturnValueOnce(server)
return server
}
function captureErrors(ctx: Context): unknown[] {
const errors: unknown[] = []
ctx.logger.error = ((error: unknown) => { errors.push(error) }) as typeof ctx.logger.error
return errors
}
async function disposeWhileStarting(fiber: { dispose(): Promise<unknown> }, server: FixtureServer): Promise<void> {
await server.started.promise
const disposal = fiber.dispose()
const settled = vi.fn()
void disposal.then(settled)
await Promise.resolve()
expect(settled).not.toHaveBeenCalled()
server.finishListening()
await disposal
}
afterEach(() => {
globalThis.fetch = nativeFetch
httpMock.createServer.mockReset()
})
describe('snapshot HTTP fixture lifecycle', () => {
it('joins search listener setup and cleanup when disposal wins the startup race', async () => {
const server = nextServer()
const ctx = new Context()
const errors = captureErrors(ctx)
const fiber = ctx.plugin(searchFixture)
await disposeWhileStarting(fiber, server)
expect(server).toMatchObject({ closed: true, connectionsClosed: true, unreferenced: true })
expect(globalThis.fetch).toBe(nativeFetch)
expect(errors).toEqual([])
})
it('runs owner cleanup and closes the listener when disposal wins the startup race', async () => {
const server = nextServer()
const ctx = new Context()
const errors = captureErrors(ctx)
const onCleanup = vi.fn()
const onListening = vi.fn()
const fiber = ctx.plugin({
name: 'loopback-fixture-lifecycle-test',
apply: testCtx => applyLoopbackServerEffect(testCtx, {
label: 'loopback-fixture-lifecycle-test',
onCleanup,
onListening,
requestListener: () => {},
}),
})
await disposeWhileStarting(fiber, server)
expect(server).toMatchObject({ closed: true, connectionsClosed: true, unreferenced: true })
expect(onListening).toHaveBeenCalledWith(expect.objectContaining({ port: 54321 }))
expect(onCleanup).toHaveBeenCalledOnce()
expect(errors).toEqual([])
})
it('maps every fetch input form and rejects another path on the recorded authority', async () => {
const server = nextServer()
const fetchMock = vi.fn(async (_input: string | URL | Request, _init?: RequestInit) => new Response('{}'))
globalThis.fetch = fetchMock
const ctx = new Context()
const fiber = ctx.plugin(searchFixture)
await server.started.promise
server.finishListening(54322)
await fiber
try {
await globalThis.fetch(RECORDED_ENDPOINT)
expect(fetchMock.mock.calls.at(-1)?.[0]).toBe('http://127.0.0.1:54322/anthropic/v1/messages')
await globalThis.fetch(new URL(RECORDED_ENDPOINT))
expect(fetchMock.mock.calls.at(-1)?.[0]).toBe('http://127.0.0.1:54322/anthropic/v1/messages')
const request = new Request(RECORDED_ENDPOINT, { method: 'POST', headers: { 'x-fixture': 'request' } })
await globalThis.fetch(request)
const mappedRequest = fetchMock.mock.calls.at(-1)?.[0]
expect(mappedRequest).toBeInstanceOf(Request)
if (!(mappedRequest instanceof Request)) throw new TypeError('mapped fetch input must be a Request')
expect(mappedRequest.url).toBe('http://127.0.0.1:54322/anthropic/v1/messages')
expect(mappedRequest.method).toBe('POST')
expect(mappedRequest.headers.get('x-fixture')).toBe('request')
const unrelated = new URL('https://example.test/')
await globalThis.fetch(unrelated)
expect(fetchMock.mock.calls.at(-1)?.[0]).toBe(unrelated)
await expect(globalThis.fetch('http://127.0.0.1:43118/unexpected'))
.rejects.toThrow('web-search-error-fixture: unexpected URL for recorded authority')
} finally {
await fiber.dispose()
}
expect(globalThis.fetch).toBe(fetchMock)
expect(server.closed).toBe(true)
})
it('preserves a later fetch wrapper while still closing the listener and reporting the ownership error', async () => {
const server = nextServer()
const fetchMock = vi.fn(async (_input: string | URL | Request, _init?: RequestInit) => new Response('{}'))
globalThis.fetch = fetchMock
const ctx = new Context()
const errors = captureErrors(ctx)
const fiber = ctx.plugin(searchFixture)
await server.started.promise
server.finishListening()
await fiber
const fixtureFetch = globalThis.fetch
const laterFetch = vi.fn((input: string | URL | Request, init?: RequestInit) => fixtureFetch(input, init))
globalThis.fetch = laterFetch
await fiber.dispose()
expect(globalThis.fetch).toBe(laterFetch)
expect(server).toMatchObject({ closed: true, connectionsClosed: true })
expect(errors.map(String).join('\n')).toContain('web-search-error-fixture: global fetch owner changed before cleanup')
})
})

View file

@ -7,6 +7,7 @@ import { homedir } from 'node:os'
import { basename, delimiter, dirname, join } from 'node:path'
import { fileURLToPath } from 'node:url'
import { describe, expect, it } from 'vitest'
import ts from 'typescript'
import {
captureExpectedWorkspaceSnapshot,
captureWorkspaceSnapshot,
@ -84,6 +85,49 @@ interface SessionLog {
readonly header: JsonObject
}
function propertyName(node: ts.PropertyName): string | undefined {
if (ts.isIdentifier(node) || ts.isStringLiteral(node) || ts.isNumericLiteral(node)) return node.text
return undefined
}
function bindsOsAssignedPort(argument: ts.Expression | undefined): boolean {
if (argument === undefined) return false
if (ts.isNumericLiteral(argument)) return Number(argument.text) === 0
if (!ts.isObjectLiteralExpression(argument)) return false
let portIsZero: boolean | undefined
for (const property of argument.properties) {
if (ts.isSpreadAssignment(property)) {
portIsZero = undefined
continue
}
if (propertyName(property.name) !== 'port') continue
portIsZero = ts.isPropertyAssignment(property)
&& ts.isNumericLiteral(property.initializer)
&& Number(property.initializer.text) === 0
}
return portIsZero === true
}
function listenerPortViolations(path: string, sourceText: string): string[] {
const source = ts.createSourceFile(path, sourceText, ts.ScriptTarget.Latest, true, ts.ScriptKind.JS)
const violations: string[] = []
const visit = (node: ts.Node): void => {
if (ts.isCallExpression(node)
&& ts.isPropertyAccessExpression(node.expression)
&& node.expression.name.text === 'listen'
&& !bindsOsAssignedPort(node.arguments[0])) {
const line = source.getLineAndCharacterOfPosition(node.getStart(source)).line + 1
const received = node.arguments[0]?.getText(source) ?? '<missing>'
violations.push(
`${path}:${line}: listener port ${received} must use listen(0, ...) or listen({ port: 0, ... })`,
)
}
ts.forEachChild(node, visit)
}
visit(source)
return violations
}
function harvested(log: SessionLog): HarvestedLog {
return {
id: String(log.header.id),
@ -507,6 +551,29 @@ describe('headless recorded-session snapshots', () => {
}
})
it('recognizes the supported OS-assigned listener forms', () => {
expect(listenerPortViolations('accepted.mjs', [
"server.listen(0, '127.0.0.1')",
"server.listen({ port: 0, host: '127.0.0.1' })",
'server.listen({ ...options, port: 0 })',
].join('\n'))).toEqual([])
expect(listenerPortViolations('fixed.mjs', 'server.listen(43118)')).toEqual([
'fixed.mjs:1: listener port 43118 must use listen(0, ...) or listen({ port: 0, ... })',
])
expect(listenerPortViolations('dynamic.mjs', 'server.listen({ port, ...options })')).toEqual([
'dynamic.mjs:1: listener port { port, ...options } must use listen(0, ...) or listen({ port: 0, ... })',
])
})
it('binds scenario HTTP fixtures only to OS-assigned ports', async () => {
const fixtureNames = (await readdir(snapshotsRoot, { recursive: true })).filter(name => name.endsWith('.mjs'))
const violations = (await Promise.all(fixtureNames.map(async (fixtureName) => listenerPortViolations(
fixtureName,
await readFile(join(snapshotsRoot, fixtureName), 'utf8'),
)))).flat()
expect(violations).toEqual([])
})
it('stores session-owned inputs with typed redaction and no ACP transcript', async () => {
for (const scenario of scenarios) {
const fixtures = await fixtureSessions(scenario)

View file

@ -0,0 +1,75 @@
/** Shared lifecycle for snapshot HTTP fixtures that bind an ephemeral loopback port. */
import { createServer } from 'node:http'
function listen(server) {
return new Promise((resolve, reject) => {
const onError = (error) => {
server.off('error', onError)
reject(error)
}
server.once('error', onError)
try {
server.listen(0, '127.0.0.1', () => {
server.off('error', onError)
resolve(undefined)
})
} catch (error) {
server.off('error', onError)
reject(error)
}
})
}
async function close(server) {
if (!server.listening) return
await new Promise((resolve, reject) => {
server.close(error => error ? reject(error) : resolve(undefined))
server.closeAllConnections()
})
}
async function cleanup(server, onCleanup, label) {
const errors = []
try {
onCleanup()
} catch (error) {
errors.push(error)
}
try {
await close(server)
} catch (error) {
errors.push(error)
}
if (errors.length === 1) throw errors[0]
if (errors.length > 1) throw new AggregateError(errors, `${label}: cleanup failed`)
}
/**
* Start a loopback server as a Cordis effect and join cleanup with its setup.
* @param ctx - Cordis context that owns the listener effect.
* @param options - Fixture callbacks and the effect label used in diagnostics.
*/
export async function applyLoopbackServerEffect(ctx, options) {
const { label, onCleanup, onListening, requestListener } = options
await ctx.effect(async () => {
const server = createServer(requestListener)
try {
await listen(server)
const address = server.address()
if (address === null || typeof address === 'string') {
throw new Error(`${label}: loopback listener has no TCP address`)
}
onListening(address)
// Snapshot fixtures must never hold the process open past protocol shutdown.
server.unref()
return () => cleanup(server, onCleanup, label)
} catch (cause) {
try {
await cleanup(server, onCleanup, label)
} catch (cleanupError) {
throw new AggregateError([cause, cleanupError], `${label}: setup and cleanup failed`)
}
throw cause
}
}, label)
}

View file

@ -1,15 +1,15 @@
/**
* Deterministic HTTP provider for the web-fetch snapshot scenario: a small
* HTML page (headings, named entities, a GFM table, nested formatting) on a
* fixed loopback port behind the real address-pinned transport. Recording and
* replay therefore exercise fetch and markdown rendering without
* external network. The port is fixed because the fetched URL is recorded.
* OS-assigned loopback port behind the real address-pinned transport. Recording
* and replay therefore exercise fetch and markdown rendering without external
* network while retaining the recorded request URL.
*/
import { createServer } from 'node:http'
import { HttpFetchProvider } from '@deepseek-ai/dsh-web-fetch-http'
import { applyLoopbackServerEffect } from '../loopback-fixture-server.mjs'
/** Fixed loopback port the scenario prompt points `web_fetch` at. */
const PORT = 43117
/** Model-visible URL retained by the recorded session. */
const RECORDED_URL = 'http://public.test:43117/menu.html'
const PAGE = `<!doctype html>
<html><head><title>Menu</title><style>.x{color:red}</style><script>ignored()</script></head>
@ -40,36 +40,52 @@ const LIMITS = {
* Register the deterministic provider and start its loopback server.
* @param ctx - Cordis context; the effect disposes the server with the fiber.
*/
export function apply(ctx) {
const server = createServer((req, res) => {
if (req.url === '/menu.html') {
res.writeHead(200, { 'content-type': 'text/html; charset=utf-8' })
res.end(PAGE)
return
}
res.writeHead(404, { 'content-type': 'text/plain; charset=utf-8' })
res.end('not found')
})
const listening = new Promise((resolve, reject) => {
server.once('error', reject)
server.listen(PORT, '127.0.0.1', () => resolve(undefined))
})
void listening.catch(() => undefined)
// The fixture must never hold the process open past protocol shutdown.
server.unref()
export async function apply(ctx) {
const readiness = Promise.withResolvers()
let transportUrl
let startupError
const resolveAddresses = async (hostname) => {
await listening
if (hostname !== 'public.test') throw new Error(`unexpected snapshot hostname: ${hostname}`)
return [{ address: '127.0.0.1', family: 4 }]
}
ctx.effect(() => async () => {
await new Promise((resolve, reject) => {
server.close(error => error ? reject(error) : resolve(undefined))
// Stop accepting first so a connection cannot arrive after the forced close.
server.closeAllConnections()
const provider = new HttpFetchProvider(LIMITS, resolveAddresses)
const unregister = ctx.web.registerFetchProvider({
id: provider.id,
available: () => provider.available(),
fetch: async (request, signal) => {
if (request.url !== RECORDED_URL) throw new Error(`unexpected snapshot URL: ${request.url}`)
await readiness.promise
if (startupError !== undefined) throw startupError
const result = await provider.fetch({ url: transportUrl.toString() }, signal)
return { ...result, url: RECORDED_URL }
},
})
try {
await applyLoopbackServerEffect(ctx, {
label: 'web-fetch-fixture-server',
requestListener: (req, res) => {
if (req.url === '/menu.html') {
res.writeHead(200, { 'content-type': 'text/html; charset=utf-8' })
res.end(PAGE)
return
}
res.writeHead(404, { 'content-type': 'text/plain; charset=utf-8' })
res.end('not found')
},
onListening: (address) => {
transportUrl = new URL(RECORDED_URL)
transportUrl.port = String(address.port)
readiness.resolve(undefined)
},
onCleanup: () => {
unregister()
},
})
}, 'web-fetch-fixture-server')
ctx.web.registerFetchProvider(new HttpFetchProvider(LIMITS, resolveAddresses))
} catch (cause) {
startupError = cause
readiness.resolve(undefined)
throw cause
}
}

View file

@ -29,4 +29,5 @@
name: '@deepseek-ai/dsh-web-search-deepseek'
config:
apiKey: snapshot-key
# The fixture maps this recorded authority to its OS-assigned listener.
baseURL: http://127.0.0.1:43118/anthropic/v1

View file

@ -7,4 +7,5 @@
name: '@deepseek-ai/dsh-web-search-deepseek'
config:
apiKey: snapshot-key
# The fixture maps this recorded authority to its OS-assigned listener.
baseURL: http://127.0.0.1:43118/anthropic/v1

View file

@ -1,32 +1,64 @@
/** Deterministic authentication failure for the search endpoint guidance snapshot. */
import { createServer } from 'node:http'
import { applyLoopbackServerEffect } from '../loopback-fixture-server.mjs'
/** Fixed loopback port recorded in the provider diagnostic. */
const PORT = 43118
/** Model-visible endpoint retained by the recorded session. */
const RECORDED_ENDPOINT = 'http://127.0.0.1:43118/anthropic/v1/messages'
const RECORDED_URL = new URL(RECORDED_ENDPOINT)
/** Cordis plugin name. */
export const name = 'web-search-error-fixture'
function requestUrl(input) {
if (typeof input === 'string') return input
if (input instanceof URL) return input.href
if (input instanceof Request) return input.url
return undefined
}
function transportInput(input, transportEndpoint) {
const url = requestUrl(input)
if (url === RECORDED_ENDPOINT) {
return input instanceof Request ? new Request(transportEndpoint, input) : transportEndpoint
}
if (url === undefined) return input
let parsed
try {
parsed = new URL(url)
} catch {
return input
}
if (parsed.host === RECORDED_URL.host) {
throw new Error(`web-search-error-fixture: unexpected URL for recorded authority: ${url}`)
}
return input
}
/** Start the local Messages endpoint and stop it with the plugin fiber. */
export async function apply(ctx) {
const server = createServer((request, response) => {
if (request.method === 'POST' && request.url === '/anthropic/v1/messages') {
response.writeHead(401, { 'content-type': 'application/json' })
response.end(JSON.stringify({ error: { message: 'invalid snapshot API key' } }))
return
}
response.writeHead(404, { 'content-type': 'text/plain; charset=utf-8' })
response.end('not found')
let restoreFetch = () => {}
await applyLoopbackServerEffect(ctx, {
label: 'web-search-error-fixture',
requestListener: (request, response) => {
if (request.method === 'POST' && request.url === '/anthropic/v1/messages') {
response.writeHead(401, { 'content-type': 'application/json' })
response.end(JSON.stringify({ error: { message: 'invalid snapshot API key' } }))
return
}
response.writeHead(404, { 'content-type': 'text/plain; charset=utf-8' })
response.end('not found')
},
onListening: (address) => {
const transportEndpoint = `http://127.0.0.1:${String(address.port)}/anthropic/v1/messages`
const originalFetch = globalThis.fetch
const fixtureFetch = async (input, init) => originalFetch(transportInput(input, transportEndpoint), init)
globalThis.fetch = fixtureFetch
restoreFetch = () => {
if (globalThis.fetch !== fixtureFetch) {
throw new Error('web-search-error-fixture: global fetch owner changed before cleanup')
}
globalThis.fetch = originalFetch
}
},
onCleanup: () => restoreFetch(),
})
await new Promise((resolve, reject) => {
server.once('error', reject)
server.listen(PORT, '127.0.0.1', () => resolve(undefined))
})
server.unref()
ctx.effect(() => async () => {
await new Promise((resolve, reject) => {
server.close(error => error ? reject(error) : resolve(undefined))
server.closeAllConnections()
})
}, 'web-search-error-fixture')
}