Merge branch 'master' into fix/subagent-acp-onerror-containment

This commit is contained in:
Tianyi Cui
2026-07-10 23:31:17 +08:00
committed by GitHub
10 changed files with 344 additions and 63 deletions

View File

@@ -6,7 +6,7 @@ Three layers, importable separately:
- **`runScenario` (harness)** — boots the real agent bin as a subprocess via tsx (unbuilt, Loader path), drives it over ACP JSON-RPC stdio from a deterministic `input.json` script, tees raw stdout for the golden + purity check, and harvests every persisted session JSONL (parent + subagent children, primary-first) after a graceful stdin-EOF shutdown. Parameterized by `AgentUnderTest` (`binScript`, `configPath`, `tsconfigPath` — absolute paths; the subprocess cwd is a temp dir outside the repo).
- **Normalizers** — pure functions turning the two captured surfaces into stable text: `normalizeStdout` (JSON-RPC ids → first-seen sequence; UUIDs/cwd → tokens; doubles as the stdout-purity check), `normalizeSessionLog` (times zeroed, `seq` kept), and the composable `scrubRequestHeaders` (header bulk → `{{system}}`/`{{tools}}`, structure kept — [pinned-header RFC](../../../docs/rfc/implemented/testing/2026-07-06-pin-request-header-content-in-one-scenario.md)).
- **`defineAcpSnapshotSuite` (factory)** — registers the whole describe/it tree for a scenario table: per-scenario golden + re-persisted-log compares, record-mode fixture write-back, the per-header-class pin with its live uniformity guard, and the fixture guard block (no orphan scenario dirs, required files present, exactly one pin per class, pinning fixtures well-formed, non-pinning fixtures header-scrubbed). Must be called at vitest collection time.
- **`defineAcpSnapshotSuite` (factory)** — registers the whole describe/it tree for a scenario table: per-scenario golden + re-persisted-log compares, record/refresh fixture write-back, the per-header-class pin with its live uniformity guard, and the fixture guard block (no orphan scenario dirs, required files present, exactly one pin per class, pinning fixtures well-formed, non-pinning fixtures header-scrubbed). Must be called at vitest collection time.
A consuming `*.snapshot.ts` is the scenario table plus one factory call:
@@ -27,12 +27,16 @@ defineAcpSnapshotSuite({
},
snapshotsDir: join(dirname(fileURLToPath(import.meta.url)), 'snapshots'),
scenarios: SCENARIOS, // exactly one entry per header class sets pinsHeader
mode: process.env.DSH_SNAPSHOT === 'record' ? 'record' : 'replay',
mode: process.env.DSH_SNAPSHOT === 'record'
? 'record'
: process.env.DSH_SNAPSHOT === 'refresh'
? 'refresh'
: 'replay',
})
```
A scenario booting a differently-composed tree sets its own `configPath` (an overlay whose basename still ends in `cordis.yml`, so the bin's replay swap finds the sibling `*cordis.snapshot.yml`) and, when that composition changes the request header, its own `headerClass` with its own pinning scenario — the acp-agent example's Code Mode scenarios are the template.
The example also ships a `cordis.snapshot.yml` replay overlay next to its `cordis.yml` (the bin swaps them under `DSH_SNAPSHOT=replay` — [single-source replay config RFC](../../../docs/rfc/implemented/testing/2026-07-04-single-source-acp-replay-config.md)); replay fixtures are served by [`dsh-llm-replay`](../llm-replay/README.md), which this package points at via the `DSH_SNAPSHOT_*` env vars it sets on the child. Fixture roles, record/replay semantics, and scenario-table fields are documented on `Scenario` and in the [snapshot RFC](../../../docs/rfc/implemented/testing/2026-06-19-acp-snapshot-tests.md).
The example also ships a `cordis.snapshot.yml` replay overlay next to its `cordis.yml` (the bin swaps them under `DSH_SNAPSHOT=replay` — [single-source replay config RFC](../../../docs/rfc/implemented/testing/2026-07-04-single-source-acp-replay-config.md)); replay fixtures are served by [`dsh-llm-replay`](../llm-replay/README.md), which this package points at via the `DSH_SNAPSHOT_*` env vars it sets on the child. `pnpm run test:snapshot:record` calls the live LLM and rewrites the recorded scenarios' model fixtures; `pnpm run test:snapshot:refresh` stays keyless, runs the replay overlay, and rewrites stdout plus comparable session-log goldens from the committed model scripts. Fixture roles, record/replay/refresh semantics, and scenario-table fields are documented on `Scenario` and in the [snapshot RFC](../../../docs/rfc/implemented/testing/2026-06-19-acp-snapshot-tests.md).
Constraints: `suite.ts` imports vitest, so the package is importable only inside a vitest run (the harness and normalizers have no such dependency but ship from the same entry). ACP-specific by design — the harness speaks the SDK's `ClientSideConnection`. Permission round-trips are scriptable: `InputScript.permissionAnswers` is a FIFO queue of option-kind selections (`allow_once`, `reject_once`, …) the client maps to the agent-issued `optionId` at answer time; an absent or exhausted queue answers `cancelled`, and a kind the request never offered rejects the run (the agent is answered `cancelled`, so a tolerant agent cannot absorb the scenario bug).

View File

@@ -23,8 +23,11 @@
*
* `pnpm run test:snapshot:record` (DSH_SNAPSHOT=record + -u) re-records the
* `session.jsonl` fixtures against the real API and refreshes the stdout golden
* in one pass; the caller resolves that env into {@link SnapshotSuiteOptions}
* (env reading stays at the suite edge, not in this library).
* in one pass. `pnpm run test:snapshot:refresh` (DSH_SNAPSHOT=refresh) instead
* replays the committed model scripts keylessly and writes the current stdout
* + persisted-log goldens back without calling a live LLM. The caller resolves
* that env into {@link SnapshotSuiteOptions} (env reading stays at the suite
* edge, not in this library).
*
* @module @deepseek-ai/dsh-acp-snapshot/suite
*/
@@ -124,12 +127,13 @@ export interface SnapshotSuiteOptions {
/** The scenario table; exactly one entry must set `pinsHeader`. */
scenarios: Scenario[]
/**
* `replay` (keyless, the default tier) or `record` (live API; re-records the
* `recorded` scenarios' fixtures and refreshes the vitest goldens under
* `--update`). The caller derives this from `$DSH_SNAPSHOT` — env reading
* stays outside this library.
* `replay` (keyless, the default tier), `record` (live API; re-records the
* `recorded` scenarios' fixtures and refreshes the Vitest goldens under
* `--update`), or `refresh` (keyless replay that rewrites stdout goldens and
* comparable session fixtures from the replay run). The caller derives this
* from `$DSH_SNAPSHOT` — env reading stays outside this library.
*/
mode: 'replay' | 'record'
mode: 'replay' | 'record' | 'refresh'
}
/**
@@ -201,6 +205,86 @@ export function headerDeltaCount(rawLog: string): number {
.length
}
/** A literal string replacement used to carry an existing fixture's volatile value into a refreshed log. */
export interface FixtureReplacement {
/** The fresh replay-run value to replace. */
from: string
/** The existing fixture value to keep. */
to: string
}
function parseJsonlRecords(text: string): Record<string, unknown>[] {
return text.split('\n')
.filter(line => line.trim().length > 0)
.map(line => JSON.parse(line) as Record<string, unknown>)
}
/**
* Build the cross-log id/cwd replacements used by refresh write-back.
*
* @param logs The freshly harvested logs, in fixture order.
* @param fixtures The existing fixture contents, in matching order.
* @returns Literal replacements from fresh volatile values to the fixture's old values.
*/
export function refreshFixtureReplacements(logs: HarvestedLog[], fixtures: string[]): FixtureReplacement[] {
const replacements: FixtureReplacement[] = []
for (let i = 0; i < logs.length; i++) {
const fresh = parseJsonlRecords((logs[i] as HarvestedLog).content)[0]
const existing = parseJsonlRecords(fixtures[i] ?? '')[0]
for (const field of ['id', 'cwd'] as const) {
const from = fresh?.[field]
const to = existing?.[field]
if (typeof from === 'string' && typeof to === 'string' && from.length > 0 && from !== to) {
replacements.push({ from, to })
}
}
}
return replacements
}
function preserveFixtureVolatiles(record: Record<string, unknown>, existing: Record<string, unknown> | undefined): void {
if (existing === undefined || existing.type !== record.type) return
if (record.type === 'session') {
for (const field of ['id', 'createdAt', 'cwd', 'parentSession'] as const) {
if (field in record && field in existing) record[field] = existing[field]
}
return
}
if ('time' in record && 'time' in existing) record.time = existing.time
if (record.type !== 'hook/result') return
const data = record.data
const existingData = existing.data
if (
data !== null && typeof data === 'object'
&& existingData !== null && typeof existingData === 'object'
&& 'durationMs' in data && 'durationMs' in existingData
) {
(data as Record<string, unknown>).durationMs = (existingData as Record<string, unknown>).durationMs
}
}
/**
* Rewrite a fresh replay-produced log so repeated refreshes do not churn
* volatile fixture fields. Meaningful event payloads come from `fresh`; the
* existing fixture lends session ids, cwd, creation times, event times, and
* hook durations where the record shape still matches.
*
* @param fresh The newly harvested session JSONL.
* @param existing The committed fixture JSONL being refreshed.
* @param replacements Cross-log literal replacements from {@link refreshFixtureReplacements}.
* @returns The stabilized JSONL content to write back.
*/
export function stabilizeRefreshLog(fresh: string, existing: string, replacements: FixtureReplacement[]): string {
let stable = fresh
for (const { from, to } of replacements) stable = stable.split(from).join(to)
const existingRecords = parseJsonlRecords(existing)
const records = parseJsonlRecords(stable)
for (let i = 0; i < records.length; i++) {
preserveFixtureVolatiles(records[i] as Record<string, unknown>, existingRecords[i])
}
return records.map(record => JSON.stringify(record)).join('\n') + '\n'
}
/**
* Register the suite: one `describe` per scenario (the golden/log compares and
* the header-uniformity guard) plus the fixture guard block (no orphan
@@ -215,6 +299,8 @@ export function headerDeltaCount(rawLog: string): number {
export function defineAcpSnapshotSuite(options: SnapshotSuiteOptions): void {
const { agent, snapshotsDir, scenarios, mode } = options
const RECORDING = mode === 'record'
const REFRESHING = mode === 'refresh'
const childMode: 'replay' | 'record' = RECORDING ? 'record' : 'replay'
/** The class a scenario's header composition belongs to (see {@link Scenario.headerClass}). */
const classOf = (scenario: Scenario): string => scenario.headerClass ?? 'default'
@@ -238,15 +324,18 @@ export function defineAcpSnapshotSuite(options: SnapshotSuiteOptions): void {
describe(`snapshot: ${scenario.name}`, () => {
// In RECORD mode, only re-run the `recorded` (live-API) scenarios; the
// `authored` ones (sidecar-driven errors/cancel) are never re-recorded.
// REFRESH mode is replay-backed and deterministic, so it runs every
// scenario and rewrites the comparable fixtures from that replay run.
it.skipIf(RECORDING && !scenario.recorded)('matches the goldens', async () => {
const dir = join(snapshotsDir, scenario.name)
const input = JSON.parse(await readFile(join(dir, 'input.json'), 'utf8')) as InputScript
const overrideFile = join(dir, 'replay.override.json')
const workspaceDir = join(dir, 'workspace')
const childSessions = scenario.childSessions ?? 0
const comparesLog = scenario.comparesLog ?? scenario.hasModelTurn
const result = await runScenario(input, {
agent,
mode,
mode: childMode,
fixtureFile: join(dir, 'session.jsonl'),
...existsSync(overrideFile) ? { overrideFile } : {},
// In REPLAY, forward the recorded child fixtures so each subagent session
@@ -271,30 +360,47 @@ export function defineAcpSnapshotSuite(options: SnapshotSuiteOptions): void {
}
// RECORD mode (recorded model scenarios only): persist the freshly-harvested
// logs back to their fixtures — the primary to session.jsonl, each child to
// session.<n>.jsonl in harvest order. `--update` refreshes the Vitest
// goldens but NOT these fixtures, so write them here. A non-pinning
// scenario's fixtures are written header-scrubbed, so a re-record can
// never smuggle the full prompt/schema content back into every fixture.
// live logs back to their fixtures. REFRESH mode does the same from a
// keyless replay run for every comparable log, including authored
// scenarios that live record deliberately skips. The primary goes to
// session.jsonl, each child to session.<n>.jsonl in harvest order. A
// non-pinning scenario's fixtures are written header-scrubbed, so a
// re-record/refresh can never smuggle the full prompt/schema content
// back into every fixture.
const scrub = scenario.pinsHeader === true
? (log: string): string => log
: scrubRequestHeaders
if (RECORDING && scenario.recorded && scenario.hasModelTurn) {
expect(result.sessionLogs.length, 'record produced no session log to harvest').toBeGreaterThan(0)
const fixtureFiles = ['session.jsonl', ...Array.from({ length: childSessions }, (_, i) => `session.${i + 1}.jsonl`)]
const existingFixtures = REFRESHING
? await Promise.all(fixtureFiles.map(file => readFile(join(dir, file), 'utf8')))
: []
const replacements = REFRESHING ? refreshFixtureReplacements(result.sessionLogs, existingFixtures) : []
const writesSessionFixtures = (RECORDING && scenario.recorded && scenario.hasModelTurn)
|| (REFRESHING && comparesLog)
if (writesSessionFixtures) {
expect(result.sessionLogs.length, `${mode} produced no session log to harvest`).toBeGreaterThan(0)
expect(result.sessionLogs.length, `expected ${childSessions + 1} session logs (parent + children)`)
.toBe(childSessions + 1)
await writeFile(join(dir, 'session.jsonl'), scrub((result.sessionLogs[0] as HarvestedLog).content))
const primary = (result.sessionLogs[0] as HarvestedLog).content
await writeFile(join(dir, 'session.jsonl'), scrub(
REFRESHING ? stabilizeRefreshLog(primary, existingFixtures[0] as string, replacements) : primary,
))
for (let i = 1; i < result.sessionLogs.length; i++) {
await writeFile(join(dir, `session.${i}.jsonl`), scrub((result.sessionLogs[i] as HarvestedLog).content))
const child = (result.sessionLogs[i] as HarvestedLog).content
await writeFile(join(dir, `session.${i}.jsonl`), scrub(
REFRESHING ? stabilizeRefreshLog(child, existingFixtures[i] as string, replacements) : child,
))
}
}
await expect(normalizeStdout(result.rawStdout, ctx))
.toMatchFileSnapshot(join(dir, 'stdout.golden.jsonl'))
const stdout = normalizeStdout(result.rawStdout, ctx)
if (REFRESHING) {
await writeFile(join(dir, 'stdout.golden.jsonl'), stdout)
}
await expect(stdout).toMatchFileSnapshot(join(dir, 'stdout.golden.jsonl'))
// A model turn always produces a log worth comparing; a hook scenario can
// produce one without a model turn (a `rejected` turn carrying `hook/*`).
const comparesLog = scenario.comparesLog ?? scenario.hasModelTurn
if (comparesLog) {
// The harvested logs (primary-first) must match their committed fixtures
// 1:1. Each side passes through normalizeSessionLog, scrubbed against ITS
@@ -307,7 +413,6 @@ export function defineAcpSnapshotSuite(options: SnapshotSuiteOptions): void {
// reason, and config, but not its bulk content (pinned once, in the
// `pinsHeader` scenario).
expect(result.sessionLogs.length, 'this scenario must persist a session log').toBe(childSessions + 1)
const fixtureFiles = ['session.jsonl', ...Array.from({ length: childSessions }, (_, i) => `session.${i + 1}.jsonl`)]
for (let i = 0; i < fixtureFiles.length; i++) {
const harvested = scrub((result.sessionLogs[i] as HarvestedLog).content)
const fixture = scrub(await readFile(join(dir, fixtureFiles[i] as string), 'utf8'))

View File

@@ -1,11 +1,18 @@
import { cpSync, mkdtempSync } from 'node:fs'
import { cpSync, mkdtempSync, readFileSync, writeFileSync } from 'node:fs'
import { rm } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { fileURLToPath } from 'node:url'
import { afterAll, describe, expect, it } from 'vitest'
import { defineAcpSnapshotSuite, type Scenario } from '../src/index.ts'
import { childFixturePaths, fixtureContext, headerDeltaCount, normalizedHeaders } from '../src/suite.ts'
import { defineAcpSnapshotSuite, type HarvestedLog, type Scenario } from '../src/index.ts'
import {
childFixturePaths,
fixtureContext,
headerDeltaCount,
normalizedHeaders,
refreshFixtureReplacements,
stabilizeRefreshLog,
} from '../src/suite.ts'
/**
* Unit tests for the suite factory, by running it: two synthetic suites over
@@ -55,16 +62,40 @@ const RECORD_SCENARIOS: Scenario[] = [
{ name: 'rec-skip', hasModelTurn: true, recorded: false, overridden: true },
]
// Record mode mutates its snapshots dir, so run it on a throwaway copy —
// except under the documented bootstrap knob, which regenerates the committed
// fixtures/goldens in place.
// Record/refresh modes mutate their snapshots dir, so run them on throwaway
// copies — except record's documented bootstrap knob, which regenerates the
// committed record fixtures/goldens in place.
const BOOTSTRAP = process.env.ACP_SNAPSHOT_SPEC_BOOTSTRAP === '1'
const recordDir = BOOTSTRAP ? RECORD_SRC : mkdtempSync(join(tmpdir(), 'acp-snap-record-suite-'))
if (!BOOTSTRAP) cpSync(RECORD_SRC, recordDir, { recursive: true })
const refreshDir = mkdtempSync(join(tmpdir(), 'acp-snap-refresh-suite-'))
cpSync(REPLAY_DIR, refreshDir, { recursive: true })
staleRefreshFixtures(refreshDir)
afterAll(async () => {
if (!BOOTSTRAP) await rm(recordDir, { recursive: true, force: true })
await rm(refreshDir, { recursive: true, force: true })
})
function staleRefreshFixtures(dir: string): void {
writeFileSync(join(dir, 'plain-turn', 'stdout.golden.jsonl'), 'stale stdout\n')
const plainBehaviorFile = join(dir, 'plain-turn', 'behavior.json')
const plainBehavior = JSON.parse(readFileSync(plainBehaviorFile, 'utf8')) as Record<string, unknown>
plainBehavior.echoEnv = true
writeFileSync(plainBehaviorFile, `${JSON.stringify(plainBehavior, null, 2)}\n`)
writeFileSync(join(dir, 'blocked-log', 'session.jsonl'), [
'{"type":"session","id":"99999999-8888-4777-8666-555555555555","createdAt":13,"cwd":"/rec/blocked-cwd"}',
'{"type":"hook/result","seq":1,"time":13,"data":{"decision":"stale","durationMs":99}}',
'',
].join('\n'))
writeFileSync(join(dir, 'authored-error', 'session.jsonl'), [
'{"type":"session","id":"77777777-8888-4777-8666-555555555555","createdAt":13,"cwd":"/rec/error-cwd"}',
'{"type":"turn/end","seq":1,"time":9,"data":{"error":"stale"}}',
'',
].join('\n'))
}
describe('defineAcpSnapshotSuite: replay mode', () => {
defineAcpSnapshotSuite({ agent: AGENT, snapshotsDir: REPLAY_DIR, scenarios: REPLAY_SCENARIOS, mode: 'replay' })
})
@@ -75,6 +106,27 @@ describe('defineAcpSnapshotSuite: record mode', () => {
defineAcpSnapshotSuite({ agent: AGENT, snapshotsDir: recordDir, scenarios: RECORD_SCENARIOS, mode: 'record' })
})
describe('defineAcpSnapshotSuite: refresh mode', () => {
defineAcpSnapshotSuite({ agent: AGENT, snapshotsDir: refreshDir, scenarios: REPLAY_SCENARIOS, mode: 'refresh' })
})
describe('defineAcpSnapshotSuite: refresh write-back', () => {
it('rewrites stdout and comparable logs from a replay-mode child run', () => {
const stdout = readFileSync(join(refreshDir, 'plain-turn', 'stdout.golden.jsonl'), 'utf8')
expect(stdout).not.toContain('stale stdout')
expect(stdout).toContain('env:{\\"mode\\":\\"replay\\"')
expect(stdout).not.toContain('\\"mode\\":\\"refresh\\"')
const blocked = readFileSync(join(refreshDir, 'blocked-log', 'session.jsonl'), 'utf8')
expect(blocked).toContain('"decision":"block"')
expect(blocked).not.toContain('"decision":"stale"')
const authored = readFileSync(join(refreshDir, 'authored-error', 'session.jsonl'), 'utf8')
expect(authored).toContain('"error":"model exploded"')
expect(authored).not.toContain('"error":"stale"')
})
})
describe('defineAcpSnapshotSuite: registration contract', () => {
it("throws when a scenario's header class has no pinning scenario", () => {
expect(() => {
@@ -175,3 +227,56 @@ describe('headerDeltaCount', () => {
expect(headerDeltaCount(`${other}\n`)).toBe(0)
})
})
describe('refreshFixtureReplacements', () => {
it('maps fresh ids and cwd values to the existing fixture values, skipping non-replacements', () => {
const log = (content: string): HarvestedLog => ({ id: 'diagnostic', createdAt: 1, content })
const logs = [
log('{"type":"session","id":"","cwd":"/same"}\n'),
log('{"type":"session","id":"new-parent","cwd":"/new"}\n'),
log('{"type":"session","id":"new-child","cwd":"/new"}\n'),
]
const fixtures = [
'{"type":"session","id":"","cwd":"/same"}\n',
'{"type":"session","id":"old-parent","cwd":"/old"}\n',
]
expect(refreshFixtureReplacements(logs, fixtures)).toEqual([
{ from: 'new-parent', to: 'old-parent' },
{ from: '/new', to: '/old' },
])
})
})
describe('stabilizeRefreshLog', () => {
it('keeps volatile fixture fields while preserving fresh meaningful payloads', () => {
const fresh = [
'{"type":"session","id":"new-child","createdAt":200,"cwd":"/new","parentSession":"new-parent","seedLength":1}',
'{"type":"hook/result","seq":1,"time":22,"data":{"decision":"block","durationMs":37}}',
'{"type":"turn/end","seq":2,"time":33,"data":{"error":"fresh error"}}',
'{"type":"tool/result","seq":3,"time":44,"data":{"text":"new-parent in /new"}}',
'{"type":"hook/result","seq":4,"time":55,"data":{"decision":"allow","durationMs":5}}',
'',
].join('\n')
const existing = [
'{"type":"session","id":"old-child","createdAt":100,"cwd":"/old","parentSession":"old-parent","seedLength":5}',
'{"type":"hook/result","seq":1,"time":11,"data":{"decision":"stale","durationMs":99}}',
'{"type":"turn/end","seq":2,"data":{"error":"stale"}}',
'{"type":"assistant/message","seq":3,"time":12,"data":{"text":"different type"}}',
'{"type":"hook/result","seq":4,"time":13,"data":{"decision":"stale"}}',
'',
].join('\n')
expect(stabilizeRefreshLog(fresh, existing, [
{ from: 'new-parent', to: 'old-parent' },
{ from: 'new-child', to: 'old-child' },
{ from: '/new', to: '/old' },
])).toBe([
'{"type":"session","id":"old-child","createdAt":100,"cwd":"/old","parentSession":"old-parent","seedLength":1}',
'{"type":"hook/result","seq":1,"time":11,"data":{"decision":"block","durationMs":99}}',
'{"type":"turn/end","seq":2,"time":33,"data":{"error":"fresh error"}}',
'{"type":"tool/result","seq":3,"time":44,"data":{"text":"old-parent in /old"}}',
'{"type":"hook/result","seq":4,"time":13,"data":{"decision":"allow","durationMs":5}}',
'',
].join('\n'))
})
})

View File

@@ -15,6 +15,32 @@ function fakeParent(): Agent {
return { id: AgentId('workflow-parent'), options: {} } as unknown as Agent
}
// Worker-thread startup is CPU-bound (a fresh thread compiles the runtime on
// every start): on a contended CI runner it regularly blows past vitest's 5s
// default test timeout, observed repeatedly on the coverage lane.
vi.setConfig({ testTimeout: 30_000 })
/**
* `vi.waitFor` with a contention-proof default timeout: the 1s default
* flaked repeatedly on the CI coverage lane, where worker-thread cold start
* (CPU-bound — a fresh thread compiles the runtime) competes with three
* sibling vitest workers for CPU. The 10s default is for exactly those
* races — waiting for a worker to start, run its first script line, or
* deliver an async child-registration message to the host. It is NOT for a
* wait that asserts the HOST reacted PROMPTLY to something that already
* happened (a settled result, an observed worker death): those keep an
* explicit tight override below, or the generous default would silently
* accept a multi-second regression in host-side reap latency as passing
* (proven by injecting a 6s delay into one such reap and watching the
* un-overridden version of this helper still pass in ~6s).
* @param assertion - retried until it stops throwing or the timeout elapses.
* @param timeout - override for a wait that must stay deliberately tight.
* @returns resolves when the assertion passes.
*/
function waitFor(assertion: () => void, timeout = 10_000): Promise<void> {
return vi.waitFor(assertion, { timeout, interval: 50 })
}
/** The vm-context escape hatch, spelled once: real Worker tests use it to make the WORKER misbehave. */
const ESCAPE = "globalThis.constructor.constructor('return process')()"
@@ -316,7 +342,7 @@ describe('dsh-workflow-workerthread', () => {
const runEnds: WorkflowResultInfo[] = []
ctx.on('workflow/end', (_info, result) => { runEnds.push(result) })
const handle = ctx.workflows.start({ ...scripted("return await agent('long job')"), parent })
await vi.waitFor(() => { expect(provider.runs.length).toBe(1) })
await waitFor(() => { expect(provider.runs.length).toBe(1) })
handle.cancel('user stopped it')
const result = await handle.result
expect(result.stopReason).toBe('cancelled')
@@ -357,7 +383,7 @@ describe('dsh-workflow-workerthread', () => {
const controller = new AbortController()
const second = ctx.workflows.start({ ...scripted("return await agent('job')"), parent, signal: controller.signal })
await vi.waitFor(() => { expect(provider.runs.length).toBe(1) })
await waitFor(() => { expect(provider.runs.length).toBe(1) })
controller.abort()
expect((await second.result).stopReason).toBe('cancelled')
await second.dispose()
@@ -400,7 +426,7 @@ describe('dsh-workflow-workerthread', () => {
`),
parent,
})
await vi.waitFor(() => { expect(narration).toContain('started') })
await waitFor(() => { expect(narration).toContain('started') })
handle.cancel('raced the completion')
const result = await handle.result
expect(result.stopReason).toBe('cancelled')
@@ -480,7 +506,7 @@ describe('dsh-workflow-workerthread', () => {
})
const result = await handle.result
expect(result.stopReason).toBe('completed')
await vi.waitFor(() => { expect(provider.runs.length).toBe(1) })
await waitFor(() => { expect(provider.runs.length).toBe(1) })
await handle.dispose()
// Not a waitFor: by the time dispose() returns, the slow child disposal
// must already be complete (host-side registry quiescence).
@@ -525,8 +551,11 @@ describe('dsh-workflow-workerthread', () => {
const result = await handle.result
expect(result.stopReason).toBe('completed')
// BEFORE dispose(): the settlement itself must have aborted the signal —
// without it this child would stay live until dispose's terminate.
await vi.waitFor(() => { expect(aborted).toEqual(['workflow settled']) })
// without it this child would stay live until dispose's terminate. This
// is a HOST-PROMPTNESS claim, not a cold-start race — a tight explicit
// bound (unlike the file default) so a multi-second reap regression
// cannot pass by outlasting the wait.
await waitFor(() => { expect(aborted).toEqual(['workflow settled']) }, 1000)
await handle.dispose()
})
@@ -572,9 +601,9 @@ describe('dsh-workflow-workerthread', () => {
`),
parent: fakeParent(),
})
await vi.waitFor(() => { expect(starts).toBe(1) })
await waitFor(() => { expect(starts).toBe(1) })
handle.cancel('stop now')
await vi.waitFor(() => { expect(cancelled).toEqual(['stop now']) }, { timeout: 800 })
await waitFor(() => { expect(cancelled).toEqual(['stop now']) }, 800)
// The wedged worker's own completion loses to the in-flight cancel.
const result = await handle.result
expect(result.stopReason).toBe('cancelled')
@@ -602,7 +631,7 @@ describe('dsh-workflow-workerthread', () => {
`),
parent,
})
await vi.waitFor(() => { expect(provider.runs.length).toBe(1) })
await waitFor(() => { expect(provider.runs.length).toBe(1) })
const before = Date.now()
await handle.dispose()
// Bounded by the grace (plus the terminate), never by the 1.5s spin.
@@ -625,7 +654,7 @@ describe('dsh-workflow-workerthread', () => {
`),
parent,
})
await vi.waitFor(() => { expect(provider.runs.length).toBe(1) })
await waitFor(() => { expect(provider.runs.length).toBe(1) })
const handleDispose = handle.dispose()
const result = await handle.result
// The script itself settled (the wrapper's own dispose RPC found the
@@ -663,7 +692,7 @@ describe('dsh-workflow-workerthread', () => {
`),
parent,
})
await vi.waitFor(() => { expect(order.filter(entry => entry.startsWith('start:')).length).toBe(2) })
await waitFor(() => { expect(order.filter(entry => entry.startsWith('start:')).length).toBe(2) })
const fast = provider.runs.find(run => (run.request.prompt[0] as { text?: string }).text === 'fast')!
fast.settle(text('fast done'))
handle.cancel('stop now')
@@ -694,7 +723,7 @@ describe('dsh-workflow-workerthread', () => {
...scripted("await parallel([() => agent('a'), () => agent('b')])\nreturn 'unreachable'"),
parent,
})
await vi.waitFor(() => { expect(provider.runs.length).toBe(2) })
await waitFor(() => { expect(provider.runs.length).toBe(2) })
handle.cancel('user stop')
const result = await handle.result
expect(result.stopReason).toBe('cancelled')
@@ -749,7 +778,9 @@ describe('dsh-workflow-workerthread', () => {
// A worker death is a stop reason like any other: workflow/end fires
// with the error outcome — for a bus observer it is the only obituary.
expect(runEnds).toEqual([{ stopReason: 'error', error: result.error, agentsStarted: 1 }])
await vi.waitFor(() => { expect(cancelled.length).toBe(1) })
// Result already settled — this is the reap's promptness, not a
// cold-start race; tight explicit bound (see the helper's doc comment).
await waitFor(() => { expect(cancelled.length).toBe(1) }, 1000)
await handle.dispose()
}, 15_000)
@@ -770,10 +801,12 @@ describe('dsh-workflow-workerthread', () => {
expect(result.stopReason).toBe('error')
expect(result.error).toContain('worker blew up')
// The reap wound the stray child down (cancel + a CLEAN dispose).
await vi.waitFor(() => {
// Result already settled — this is the reap's promptness, not a
// cold-start race; tight explicit bound (see the helper's doc comment).
await waitFor(() => {
expect(provider.runs.length).toBe(1)
expect(provider.runs[0]!.disposed).toBe(true)
})
}, 1000)
await handle.dispose()
}, 15_000)
@@ -802,7 +835,7 @@ describe('dsh-workflow-workerthread', () => {
`),
parent,
})
await vi.waitFor(() => { expect(order.filter(entry => entry.startsWith('start:')).length).toBe(2) })
await waitFor(() => { expect(order.filter(entry => entry.startsWith('start:')).length).toBe(2) })
const fast = provider.runs.find(run => (run.request.prompt[0] as { text?: string }).text === 'fast')!
fast.settle(text('fast done'))
const result = await handle.result
@@ -837,7 +870,10 @@ describe('dsh-workflow-workerthread', () => {
const result = await handle.result
expect(result.stopReason).toBe('error')
expect(result.error).toContain('exit code 5')
await vi.waitFor(() => { expect(provider.runs[0]!.disposed).toBe(true) })
// Result already settled — this is the reap's promptness (bounded
// above the mock's fixed 300ms dispose delay, not a cold-start race);
// tight explicit bound (see the helper's doc comment).
await waitFor(() => { expect(provider.runs[0]!.disposed).toBe(true) }, 1000)
await handle.dispose()
}, 15_000)
@@ -855,7 +891,7 @@ describe('dsh-workflow-workerthread', () => {
})
const logs: string[] = []
ctx.on('workflow/log', (_info, message) => { logs.push(message) })
await vi.waitFor(() => { expect(logs).toContain('armed') })
await waitFor(() => { expect(logs).toContain('armed') })
handle.cancel('stop it')
// The grace is deliberately huge: only the worker's own death (exit 3,
// unreachable by the cancel — the script ignores hooks) settles this.