test(fs): close abort/concurrency/observed coverage gaps + with-key e2e

Behavioral gaps from the coverage audit (line coverage was already 100%; these
close BEHAVIOR gaps):

- fs-local: service-level writeText/editText pre-abort → FS_ABORTED (file
  unchanged); concurrent guarded-write race and mixed write-vs-edit race (one
  wins, one FS_STALE_VERSION, locks released); edit→edit version refresh at the
  provider; the replaceIfVersion post-write version matches a fresh stat. fsio:
  a mid-stream abort → FS_ABORTED (previously only pre-abort was covered).
- fs-policy: the agent-without-session owner rung ({agent:{}} → no owner →
  createIfAbsent / FS_NOT_OBSERVED); fs/write-intent first-wins (symmetric to
  the existing edit-intent test).
- tool-fs: abort-through-the-tool for read/write/edit (isError FS_ABORTED, file
  unchanged); a deterministic tool-tier concurrent-edit race via a shared read;
  the throwing-fs/observed contract (a throwing listener surfaces as isError but
  the mutation already hit disk); the replace_all edit message; parseReadArgs
  rejects fractional/NaN offset and zero/negative limit.
- dsh-fs: FsError chains a cause through ErrorOptions.

New with-key e2e (packages/fs/tool-fs/tests/fs-tools.e2e.ts, self-skips without
DEEPSEEK_API_KEY): a real model drives the real read/write/edit tools to create
→ read → edit a file, verified on disk; a second test proves a relative path
resolves against the per-session cwd (factory meta.cwd) not config.cwd. Booted
via a plain tests/harness.ts. Added dsh-agent-loop + dsh-llm-deepseek devDeps.
This commit is contained in:
Tianyi Cui
2026-07-02 20:38:06 +08:00
parent 94cbec8162
commit b3f8b4c9c6
10 changed files with 348 additions and 0 deletions

View File

@@ -196,6 +196,41 @@ describe('writeText', () => {
.rejects.toMatchObject({ code: 'FS_NOT_OBSERVED' })
expect(lockCount(fs)).toBe(0)
})
it('replaceIfVersion returns the post-write version (matches a fresh stat)', async () => {
await writeFile(join(dir, 'a.txt'), 'v1')
const target = await fs.resolve('a.txt')
const before = await versionOf(target)
// Change the byte length so the mtimeMs:size token provably differs (a
// same-size same-tick rewrite can collide — the documented version-token
// limitation; not what this test is about).
const outcome = await fs.writeText(target, 'a much longer replacement body', { kind: 'replaceIfVersion', version: before })
expect(outcome.version).not.toBe(before)
expect(outcome.version).toBe(await versionOf(target))
})
it('honors a pre-aborted signal without creating the file', async () => {
const target = await fs.resolve('aborted.txt')
await expect(fs.writeText(target, 'x', undefined, AbortSignal.abort()))
.rejects.toMatchObject({ code: 'FS_ABORTED' })
await expect(stat(join(dir, 'aborted.txt'))).rejects.toMatchObject({ code: 'ENOENT' })
expect(lockCount(fs)).toBe(0)
})
it('two concurrent guarded writes: one updates, the other is rejected as stale', async () => {
await writeFile(join(dir, 'a.txt'), 'base')
const target = await fs.resolve('a.txt')
const version = await versionOf(target)
const results = await Promise.allSettled([
fs.writeText(target, 'one', { kind: 'replaceIfVersion', version }),
fs.writeText(target, 'two', { kind: 'replaceIfVersion', version }),
])
expect(results.filter(r => r.status === 'fulfilled')).toHaveLength(1)
const rejected = results.filter(r => r.status === 'rejected')
expect(rejected).toHaveLength(1)
expect((rejected[0] as PromiseRejectedResult).reason).toMatchObject({ code: 'FS_STALE_VERSION' })
expect(lockCount(fs)).toBe(0)
})
})
describe('editText', () => {
@@ -297,6 +332,41 @@ describe('editText', () => {
expect((rejected[0] as PromiseRejectedResult).reason).toMatchObject({ code: 'FS_STALE_VERSION' })
expect(lockCount(fs)).toBe(0)
})
it('honors a pre-aborted signal without rewriting the file', async () => {
await writeFile(join(dir, 'a.txt'), 'keep')
const target = await fs.resolve('a.txt')
await expect(fs.editText(target, { oldString: 'keep', newString: 'x', replaceAll: false }, undefined, AbortSignal.abort()))
.rejects.toMatchObject({ code: 'FS_ABORTED' })
expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('keep')
expect(lockCount(fs)).toBe(0)
})
it('a successful edit refreshes the version so an immediate follow-up edit proceeds', async () => {
await writeFile(join(dir, 'a.txt'), 'one two')
const target = await fs.resolve('a.txt')
const first = await fs.editText(target, { oldString: 'one', newString: 'ONE', replaceAll: false }, { version: await versionOf(target) })
// The version the first edit returned is a valid guard for a second edit —
// no intervening re-stat needed.
const second = await fs.editText(target, { oldString: 'two', newString: 'TWO', replaceAll: false }, { version: first.version })
expect(second.replacements).toBe(1)
expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('ONE TWO')
})
it('concurrent write vs edit at the same version: one wins, the other is stale', async () => {
await writeFile(join(dir, 'a.txt'), 'base')
const target = await fs.resolve('a.txt')
const version = await versionOf(target)
const results = await Promise.allSettled([
fs.writeText(target, 'written', { kind: 'replaceIfVersion', version }),
fs.editText(target, { oldString: 'base', newString: 'edited', replaceAll: false }, { version }),
])
expect(results.filter(r => r.status === 'fulfilled')).toHaveLength(1)
const rejected = results.filter(r => r.status === 'rejected')
expect(rejected).toHaveLength(1)
expect((rejected[0] as PromiseRejectedResult).reason).toMatchObject({ code: 'FS_STALE_VERSION' })
expect(lockCount(fs)).toBe(0)
})
})
describe('symlink targetKey identity', () => {

View File

@@ -222,6 +222,22 @@ describe('streamWholeText', () => {
await writeFile(file, 'one\ntwo')
expect(await collect(streamWholeText(localTarget(file), new AbortController().signal))).toBe('one\ntwo')
})
it('translates a mid-stream abort into FS_ABORTED', async () => {
// A multi-chunk file so the stream yields more than once; abort after the
// first chunk and assert the structured code, not a raw AbortError.
const file = join(dir, 'big.txt')
await writeFile(file, 'x'.repeat(256 * 1024))
const ac = new AbortController()
const run = async (): Promise<void> => {
let seen = 0
for await (const _chunk of streamWholeText(localTarget(file), ac.signal)) {
seen += 1
if (seen === 1) ac.abort()
}
}
await expect(run()).rejects.toMatchObject({ code: 'FS_ABORTED' })
})
})
describe('writeFileAtomic — temp-file safety', () => {

View File

@@ -64,6 +64,13 @@ describe('write-intent decision', () => {
expect(await writeIntent(ctx, target('a.txt'), {})).toEqual({ kind: 'createIfAbsent' })
})
it('an actor with an agent but no session has no owner (createIfAbsent)', async () => {
// The middle optional-chain rung: agent present, session undefined ⇒ owner
// undefined ⇒ unobservable, so a write can only be a blind create.
const { ctx } = await setup()
expect(await writeIntent(ctx, target('a.txt'), { agent: {} })).toEqual({ kind: 'createIfAbsent' })
})
it('an observed target decides replaceIfVersion at the observed version', async () => {
const { ctx } = await setup()
const exec = ownerExec({})
@@ -83,6 +90,11 @@ describe('edit-intent decision', () => {
await expect(editIntent(ctx, target('a.txt'), undefined)).rejects.toMatchObject({ code: 'FS_NOT_OBSERVED' })
})
it('rejects an edit whose actor has an agent but no session (no owner)', async () => {
const { ctx } = await setup()
await expect(editIntent(ctx, target('a.txt'), { agent: {} })).rejects.toMatchObject({ code: 'FS_NOT_OBSERVED' })
})
it('returns the observed version as the CAS basis after an observation', async () => {
const { ctx } = await setup()
const exec = ownerExec({})
@@ -166,6 +178,17 @@ describe('single-slot, first-wins', () => {
await editIntent(ctx, target('a.txt'), exec)
expect(secondRan).toBe(false)
})
it('a SECOND write-intent decider registered AFTER fs-policy is not reached', async () => {
const { ctx } = await setup()
let secondRan = false
ctx.on('fs/write-intent', () => {
secondRan = true
return Promise.resolve(undefined)
})
await writeIntent(ctx, target('a.txt'), ownerExec({}))
expect(secondRan).toBe(false)
})
})
describe('disposal releases recorded state (HMR safety)', () => {

View File

@@ -108,4 +108,11 @@ describe('FsError', () => {
expect(error.name).toBe('FsError')
expect(error).toBeInstanceOf(Error)
})
it('chains an underlying cause through ErrorOptions', () => {
const root = new Error('EACCES')
const error = new FsError('cannot read', 'FS_ABORTED', { cause: root })
expect(error.cause).toBe(root)
expect(error.code).toBe('FS_ABORTED')
})
})

View File

@@ -30,10 +30,12 @@
},
"devDependencies": {
"@deepseek-ai/dsh-agent": "workspace:^",
"@deepseek-ai/dsh-agent-loop": "workspace:^",
"@deepseek-ai/dsh-fs-policy": "workspace:^",
"@deepseek-ai/dsh-fs": "workspace:^",
"@deepseek-ai/dsh-fs-local": "workspace:^",
"@deepseek-ai/dsh-llm": "workspace:^",
"@deepseek-ai/dsh-llm-deepseek": "workspace:^",
"@deepseek-ai/dsh-session": "workspace:^",
"@deepseek-ai/dsh-system-prompt": "workspace:^",
"@deepseek-ai/dsh-tools": "workspace:^",

View File

@@ -0,0 +1,84 @@
import { mkdtemp, readFile, rm } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, describe, expect, it } from 'vitest'
import type { Context } from 'cordis'
import { AgentId } from '@deepseek-ai/dsh-agent'
import { SessionId } from '@deepseek-ai/dsh-session'
import { fsHarness, waitForIdle } from './harness.ts'
/**
* With-key smoke for the filesystem tools: a REAL model drives the REAL
* read/write/edit tools (over the real local backend + policy gate), and we
* verify the WORLD — the file on disk — not the agent's self-report. This is the
* "green units, broken product" guard: mocks prove the plumbing, only a real
* model proves the tools actually work end-to-end. Key-gated (self-skips without
* DEEPSEEK_API_KEY).
*/
let ctx: Context | undefined
let workdir: string | undefined
afterEach(async () => {
await ctx?.fiber.dispose()
ctx = undefined
if (workdir !== undefined) await rm(workdir, { recursive: true, force: true })
workdir = undefined
})
const SYSTEM = 'You are a coding assistant. Use the write tool to create files, the read tool to inspect '
+ 'them, and the edit tool for literal replacements. Read a file before editing it. Keep replies terse.'
describe.skipIf(!process.env.DEEPSEEK_API_KEY)('fs tools with-key smoke', () => {
it('creates, reads, then edits a file — verified on disk', async () => {
workdir = await mkdtemp(join(tmpdir(), 'dsh-fs-e2e-'))
ctx = await fsHarness(workdir)
// agentLoop.create prepares a session with no cwd, so the provider default
// (config.cwd = workdir) is the workspace.
const agent = ctx.agentLoop.create(AgentId('fs-e2e'), { model: 'deepseek-v4-flash', systemPrompt: SYSTEM })
agent.send([{ type: 'text', text:
'Create a file named note.txt containing exactly the line: status: draft. '
+ 'Then read it back, then edit it to replace the literal word draft with final. '
+ 'Tell me when done.' }])
await waitForIdle(ctx, agent)
// Verify the WORLD: the edit landed on disk.
const content = await readFile(join(workdir, 'note.txt'), 'utf8')
expect(content).toContain('status: final')
expect(content).not.toContain('draft')
// The log records real read/write/edit tool calls (not bash).
const calls = [...agent.session.events].filter(e => e.type === 'tool/call').map(e => e.data.name)
expect(calls).toContain('write')
expect(calls).toContain('read')
expect(calls).toContain('edit')
}, 180_000)
it('resolves a relative path against the per-session cwd (factory meta.cwd)', async () => {
// config.cwd is the harness workdir, but the agent's SESSION cwd is a
// different dir; the write must land in the SESSION dir, proving the tool
// passes the per-session cwd (not the backend default).
const configDir = await mkdtemp(join(tmpdir(), 'dsh-fs-e2e-cfg-'))
workdir = configDir
const sessionDir = await mkdtemp(join(tmpdir(), 'dsh-fs-e2e-session-'))
try {
ctx = await fsHarness(configDir)
const handle = ctx.agents.create({
agentId: AgentId('fs-e2e-cwd'),
sessionId: SessionId(`fs-e2e-cwd-${Date.now()}`),
meta: { cwd: sessionDir },
agentOptions: { model: 'deepseek-v4-flash', systemPrompt: SYSTEM },
})
handle.agent.send([{ type: 'text', text:
'Use the write tool to create a file named where.txt containing exactly the line: here. Tell me when done.' }])
await waitForIdle(ctx, handle.agent)
// The file is in the SESSION dir, not the config dir.
expect(await readFile(join(sessionDir, 'where.txt'), 'utf8')).toContain('here')
await expect(readFile(join(configDir, 'where.txt'), 'utf8')).rejects.toMatchObject({ code: 'ENOENT' })
} finally {
await rm(sessionDir, { recursive: true, force: true })
}
}, 180_000)
})

View File

@@ -0,0 +1,47 @@
import { Context } from 'cordis'
import LlmService from '@deepseek-ai/dsh-llm'
import SessionStore from '@deepseek-ai/dsh-session'
import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
import ToolRegistry from '@deepseek-ai/dsh-tools'
import AgentRegistry, { type Agent } from '@deepseek-ai/dsh-agent'
import AgentLoop from '@deepseek-ai/dsh-agent-loop'
import LocalFileSystem from '@deepseek-ai/dsh-fs-local'
import * as FsPolicy from '@deepseek-ai/dsh-fs-policy'
import * as ToolFs from '@deepseek-ai/dsh-tool-fs'
import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek'
/**
* Shared harness for the fs-tools with-key e2e: a minimal real agent stack (the
* DeepSeek adapter + the real fs provider + the read-before-write/edit policy +
* the model-facing read/write/edit tools). Lives outside the *.e2e.ts pattern so
* importing it never re-registers another file's tests.
*
* `fsCwd` is the local backend's default base; a per-session cwd (set via a
* session header) overrides it, but this harness creates agents without a
* session cwd, so the provider default IS the workspace.
*/
export async function fsHarness(fsCwd: string): Promise<Context> {
const ctx = new Context()
await ctx.plugin(LlmService)
await ctx.plugin(SessionStore)
await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry)
await ctx.plugin(AgentRegistry)
await ctx.plugin(AgentLoop, { agents: [] })
await ctx.plugin(LlmDeepSeek, { models: ['deepseek-v4-flash'] })
await ctx.plugin(LocalFileSystem, { cwd: fsCwd })
await ctx.plugin(FsPolicy)
await ctx.plugin(ToolFs)
return ctx
}
export function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
return new Promise((resolve) => {
const dispose = ctx.on('agent/status', (subject, status) => {
if (subject === agent && status === 'idle') {
dispose()
resolve()
}
})
})
}

View File

@@ -338,3 +338,73 @@ describe('per-session cwd', () => {
expect(await readFile(join(sessionDir, 'code.txt'), 'utf8')).toBe('beta')
})
})
// --------------------------------------------------------------------------
// Abort-through-the-tool, tool-tier concurrency, and the fs/observed contract —
// all through ctx.tools.execute() against the REAL backend + policy.
// --------------------------------------------------------------------------
describe('signal, concurrency, and the fs/observed contract', () => {
beforeEach(async () => {
dir = await mkdtemp(join(tmpdir(), 'dsh-tool-fs-'))
ctx = new Context()
await ctx.plugin(SystemPrompt)
await ctx.plugin(ToolRegistry)
await ctx.plugin(LocalFileSystem, { cwd: dir })
await ctx.plugin(FsPolicy)
fiber = await ctx.plugin(ToolFs)
})
const session = { header: {} }
const callSig = (signal: AbortSignal, name: string, args: unknown) =>
ctx.tools.execute({ callId: CallId(`c-${++callCounter}`), name, arguments: args, agent: { session } as never, signal })
const callOwned = (name: string, args: unknown) =>
ctx.tools.execute({ callId: CallId(`c-${++callCounter}`), name, arguments: args, agent: { session } as never })
it('a pre-aborted signal makes read/write/edit return isError FS_ABORTED', async () => {
await writeFile(join(dir, 'a.txt'), 'hello')
const read = await callSig(AbortSignal.abort(), 'read', { file_path: 'a.txt' })
expect(read.isError).toBe(true)
expect(read.error).toMatchObject({ code: 'FS_ABORTED' })
const write = await callSig(AbortSignal.abort(), 'write', { file_path: 'new.txt', content: 'x' })
expect(write.isError).toBe(true)
expect(write.error).toMatchObject({ code: 'FS_ABORTED' })
await expect(readFile(join(dir, 'new.txt'), 'utf8')).rejects.toMatchObject({ code: 'ENOENT' })
// Read first (un-aborted, SAME session owner) so the edit clears the
// observation gate; then the aborted edit fails on the signal, not on
// FS_NOT_OBSERVED.
expect((await callOwned('read', { file_path: 'a.txt' })).isError).toBe(false)
const edit = await callSig(AbortSignal.abort(), 'edit', { file_path: 'a.txt', old_string: 'hello', new_string: 'bye' })
expect(edit.isError).toBe(true)
expect(edit.error).toMatchObject({ code: 'FS_ABORTED' })
expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('hello') // unchanged
})
it('two concurrent edits of the same file, same session: one wins, one FS_STALE_VERSION', async () => {
await writeFile(join(dir, 'a.txt'), 'base value here')
// One read establishes the observed version both edits guard against; then
// race two edits so both carry the SAME observed version (the barrier).
expect((await callOwned('read', { file_path: 'a.txt' })).isError).toBe(false)
const [one, two] = await Promise.all([
callOwned('edit', { file_path: 'a.txt', old_string: 'base', new_string: 'ONE', replaceAll: false }),
callOwned('edit', { file_path: 'a.txt', old_string: 'value', new_string: 'TWO', replaceAll: false }),
])
const errors = [one, two].filter(r => r.isError)
expect(errors).toHaveLength(1)
expect(errors[0]?.error).toMatchObject({ code: 'FS_STALE_VERSION' })
// The world is consistent: exactly one edit landed.
const onDisk = await readFile(join(dir, 'a.txt'), 'utf8')
expect(onDisk === 'ONE value here' || onDisk === 'base TWO here').toBe(true)
})
it('a throwing fs/observed listener surfaces as isError, but the mutation already hit disk', async () => {
// fs/observed is a plain ctx.emit AFTER the write succeeded; a throwing
// listener cannot roll the write back — it only turns the tool result into
// isError. The file must still carry the written bytes.
ctx.on('fs/observed', () => { throw new Error('recording bug') })
const result = await callOwned('write', { file_path: 'w.txt', content: 'durable' })
expect(result.isError).toBe(true)
expect(await readFile(join(dir, 'w.txt'), 'utf8')).toBe('durable')
})
})

View File

@@ -158,6 +158,20 @@ describe('read tool', () => {
expect(text(result)).toContain('offset must be a positive integer')
})
it('rejects a fractional or NaN offset, and a zero/negative limit', async () => {
const { ctx } = await setup()
for (const args of [
{ file_path: 'a.txt', offset: 1.5 },
{ file_path: 'a.txt', offset: Number.NaN },
{ file_path: 'a.txt', limit: 0 },
{ file_path: 'a.txt', limit: -3 },
]) {
const result = await call(ctx, 'read', args)
expect(result.isError, JSON.stringify(args)).toBe(true)
expect(text(result)).toMatch(/must be a positive integer/)
}
})
it('rejects a limit above the cap', async () => {
const { ctx } = await setup()
const result = await call(ctx, 'read', { file_path: 'a.txt', limit: 99999 })
@@ -291,6 +305,15 @@ describe('edit tool', () => {
expect(text(result)).toBe('The file /abs/a.txt has been updated successfully.')
})
it('formats the replace_all success message distinctly', async () => {
const { ctx, fs } = await setup()
const session = { header: {} }
fs.files.set('key:a.txt', 'a a a')
await call(ctx, 'read', { file_path: 'a.txt' }, { session })
const result = await call(ctx, 'edit', { file_path: 'a.txt', old_string: 'a', new_string: 'b', replace_all: true }, { session })
expect(text(result)).toBe('The file /abs/a.txt has been updated. All occurrences were successfully replaced.')
})
it('rejects identical old/new strings', async () => {
const { ctx } = await setup()
const result = await call(ctx, 'edit', { file_path: 'a.txt', old_string: 'x', new_string: 'x' })

6
pnpm-lock.yaml generated
View File

@@ -323,6 +323,9 @@ importers:
'@deepseek-ai/dsh-agent':
specifier: workspace:^
version: link:../../core/agent
'@deepseek-ai/dsh-agent-loop':
specifier: workspace:^
version: link:../../core/agent-loop
'@deepseek-ai/dsh-fs':
specifier: workspace:^
version: link:../fs
@@ -335,6 +338,9 @@ importers:
'@deepseek-ai/dsh-llm':
specifier: workspace:^
version: link:../../llm/llm
'@deepseek-ai/dsh-llm-deepseek':
specifier: workspace:^
version: link:../../llm/llm-deepseek
'@deepseek-ai/dsh-session':
specifier: workspace:^
version: link:../../core/session