integration.spec.ts 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515
  1. /**
  2. * End-to-end tool-registry tests against the real local backend. The policy deployment verifies
  3. * observed-state and guarded mutation; the bare deployment proves unconditional tools have no
  4. * policy-service dependency. Assertions read files back byte-for-byte rather than trusting tool
  5. * messages.
  6. */
  7. import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
  8. import { mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'
  9. import { tmpdir } from 'node:os'
  10. import { join } from 'node:path'
  11. import { Context } from '@deepseek-ai/cordis'
  12. import { ToolCallId } from '@deepseek-ai/dsh-llm'
  13. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  14. import ToolRuntime, { TOOL_ABORTED_BEFORE_DISPATCH } from '@deepseek-ai/dsh-tools'
  15. import { LocalFileSystem } from '@deepseek-ai/dsh-fs-local'
  16. import * as FsPolicy from '@deepseek-ai/dsh-fs-observation-policy'
  17. import * as ToolFs from '@deepseek-ai/dsh-tool-fs'
  18. const testToolSignal = new AbortController().signal
  19. let dir: string
  20. let ctx: Context
  21. let fiber: Awaited<ReturnType<Context['plugin']>>
  22. // No header cwd: sessionCwd returns undefined and the provider's configured test dir applies.
  23. const session = { header: {} }
  24. let callCounter = 0
  25. function call(name: string, args: unknown) {
  26. return ctx.tools.execute({
  27. signal: testToolSignal,
  28. callId: ToolCallId(`call-${++callCounter}`),
  29. name,
  30. arguments: args,
  31. agent: { session } as never,
  32. })
  33. }
  34. function text(result: { content: { type: string; text?: string }[] }): string {
  35. return result.content.filter(b => b.type === 'text').map(b => b.text).join('')
  36. }
  37. function notObservedDiagnostic(path: string): string {
  38. return `Error: cannot modify "${path}": file has not been read — read the file, then retry`
  39. }
  40. afterEach(async () => {
  41. await fiber.dispose()
  42. await rm(dir, { recursive: true, force: true })
  43. })
  44. // --------------------------------------------------------------------------
  45. // DEFAULT deployment: the policy gate plugin is loaded.
  46. // --------------------------------------------------------------------------
  47. describe('default deployment (with dsh-fs-observation-policy)', () => {
  48. beforeEach(async () => {
  49. dir = await mkdtemp(join(tmpdir(), 'dsh-tool-fs-'))
  50. ctx = new Context()
  51. await ctx.plugin(SystemPrompt)
  52. await ctx.plugin(ToolRuntime)
  53. await ctx.plugin(LocalFileSystem, { cwd: dir })
  54. await ctx.plugin(FsPolicy)
  55. fiber = await ctx.plugin(ToolFs)
  56. })
  57. describe('write → disk', () => {
  58. it('creates a file with exactly the requested bytes', async () => {
  59. const result = await call('write', { file_path: 'new.txt', content: 'line one\nline two\n' })
  60. expect(result.isError).toBe(false)
  61. expect(await readFile(join(dir, 'new.txt'), 'utf8')).toBe('line one\nline two\n')
  62. })
  63. it('rejects overwriting an existing file without reading it first', async () => {
  64. await writeFile(join(dir, 'a.txt'), 'original')
  65. const result = await call('write', { file_path: 'a.txt', content: 'clobber' })
  66. expect(result.isError).toBe(true)
  67. expect(result.error).toMatchObject({ info: { code: 'FS_NOT_OBSERVED' } })
  68. expect(text(result)).toBe(notObservedDiagnostic(join(dir, 'a.txt')))
  69. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('original')
  70. })
  71. it('allows overwriting after a read', async () => {
  72. await writeFile(join(dir, 'a.txt'), 'original')
  73. expect((await call('read', { file_path: 'a.txt' })).isError).toBe(false)
  74. const result = await call('write', { file_path: 'a.txt', content: 'replaced' })
  75. expect(result.isError).toBe(false)
  76. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('replaced')
  77. })
  78. it('rejects a full overwrite when the file changed since the read (stale)', async () => {
  79. await writeFile(join(dir, 'a.txt'), 'original')
  80. await call('read', { file_path: 'a.txt' })
  81. await writeFile(join(dir, 'a.txt'), 'changed-externally') // out-of-band change
  82. const result = await call('write', { file_path: 'a.txt', content: 'replaced' })
  83. expect(result.isError).toBe(true)
  84. expect(result.error).toMatchObject({ info: { code: 'FS_STALE_VERSION' } })
  85. // The model-facing text names the remedy, not just the condition.
  86. expect(text(result)).toContain('file changed since it was read')
  87. expect(text(result)).toContain('re-read the file, then retry')
  88. })
  89. it('the stale remedy is actionable: re-reading the changed file unblocks the retried write', async () => {
  90. await writeFile(join(dir, 'a.txt'), 'original')
  91. await call('read', { file_path: 'a.txt' })
  92. await writeFile(join(dir, 'a.txt'), 'changed-externally') // out-of-band change
  93. const stale = await call('write', { file_path: 'a.txt', content: 'replaced' })
  94. expect(stale.isError).toBe(true)
  95. expect(stale.error).toMatchObject({ info: { code: 'FS_STALE_VERSION' } })
  96. // Follow the remedy: re-read (refreshes the observed version), then retry.
  97. expect((await call('read', { file_path: 'a.txt' })).isError).toBe(false)
  98. const retried = await call('write', { file_path: 'a.txt', content: 'replaced' })
  99. expect(retried.isError).toBe(false)
  100. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('replaced')
  101. })
  102. })
  103. describe('read', () => {
  104. it('returns line-numbered content', async () => {
  105. await writeFile(join(dir, 'a.txt'), 'alpha\nbeta')
  106. const result = await call('read', { file_path: 'a.txt' })
  107. expect(text(result)).toContain('1: alpha')
  108. expect(text(result)).toContain('2: beta')
  109. expect(text(result)).toContain('(End of file - total 2 lines)')
  110. })
  111. it('reports a binary file as an error', async () => {
  112. await writeFile(join(dir, 'bin'), Buffer.from([0x00, 0x01, 0x02]))
  113. const result = await call('read', { file_path: 'bin' })
  114. expect(result.isError).toBe(true)
  115. expect(result.error).toMatchObject({ info: { code: 'FS_NOT_TEXT' } })
  116. })
  117. it('paginates a multi-line file with offset/limit', async () => {
  118. await writeFile(join(dir, 'a.txt'), 'one\ntwo\nthree\nfour')
  119. const result = await call('read', { file_path: 'a.txt', offset: 2, limit: 2 })
  120. expect(text(result)).toContain('2: two')
  121. expect(text(result)).toContain('3: three')
  122. expect(text(result)).toContain('(Showing lines 2-3 of 4. Use offset=4 to continue.)')
  123. })
  124. })
  125. describe('edit → disk', () => {
  126. it('applies a unique literal replacement after a read', async () => {
  127. await writeFile(join(dir, 'a.txt'), 'hello world')
  128. await call('read', { file_path: 'a.txt' })
  129. const result = await call('edit', { file_path: 'a.txt', old_string: 'world', new_string: 'there' })
  130. expect(result.isError).toBe(false)
  131. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('hello there')
  132. })
  133. it('rejects an edit before any read, leaving the file untouched', async () => {
  134. await writeFile(join(dir, 'a.txt'), 'hello world')
  135. const result = await call('edit', { file_path: 'a.txt', old_string: 'world', new_string: 'there' })
  136. expect(result.isError).toBe(true)
  137. expect(result.error).toMatchObject({ info: { code: 'FS_NOT_OBSERVED' } })
  138. expect(text(result)).toBe(notObservedDiagnostic(join(dir, 'a.txt')))
  139. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('hello world')
  140. })
  141. it('lets a WINDOWED read authorize an edit when the file is unchanged (freshness, not full-view)', async () => {
  142. // A file with more lines than the read window; read only the first line.
  143. const lines = Array.from({ length: 20 }, (_, i) => `line ${i + 1}`)
  144. await writeFile(join(dir, 'a.txt'), lines.join('\n'))
  145. const read = await call('read', { file_path: 'a.txt', offset: 1, limit: 1 })
  146. expect(read.isError).toBe(false)
  147. expect(text(read)).toContain('(Showing lines 1-1 of 20')
  148. // Editing a line OUTSIDE the window is authorized because the file is unchanged.
  149. const result = await call('edit', { file_path: 'a.txt', old_string: 'line 12', new_string: 'LINE 12' })
  150. expect(result.isError).toBe(false)
  151. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe(lines.map(l => l === 'line 12' ? 'LINE 12' : l).join('\n'))
  152. })
  153. it('rejects an edit when the file changed since the windowed read (stale before matching)', async () => {
  154. await writeFile(join(dir, 'a.txt'), 'hello world')
  155. await call('read', { file_path: 'a.txt', offset: 1, limit: 1 })
  156. await writeFile(join(dir, 'a.txt'), 'goodbye') // out-of-band change removes 'world'
  157. const result = await call('edit', { file_path: 'a.txt', old_string: 'world', new_string: 'there' })
  158. expect(result.isError).toBe(true)
  159. expect(result.error).toMatchObject({ info: { code: 'FS_STALE_VERSION' } })
  160. // The model-facing text names the remedy, not just the condition.
  161. expect(text(result)).toContain('file changed since it was read')
  162. expect(text(result)).toContain('re-read the file, then retry')
  163. })
  164. it('the stale remedy is actionable: re-reading the changed file unblocks the retried edit', async () => {
  165. await writeFile(join(dir, 'a.txt'), 'hello world')
  166. await call('read', { file_path: 'a.txt' })
  167. await writeFile(join(dir, 'a.txt'), 'hello brave world') // out-of-band change
  168. const stale = await call('edit', { file_path: 'a.txt', old_string: 'world', new_string: 'there' })
  169. expect(stale.isError).toBe(true)
  170. expect(stale.error).toMatchObject({ info: { code: 'FS_STALE_VERSION' } })
  171. // Follow the remedy: re-read (refreshes the observed version), then retry.
  172. expect((await call('read', { file_path: 'a.txt' })).isError).toBe(false)
  173. const retried = await call('edit', { file_path: 'a.txt', old_string: 'world', new_string: 'there' })
  174. expect(retried.isError).toBe(false)
  175. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('hello brave there')
  176. })
  177. it('rejects an ambiguous match without replace_all', async () => {
  178. await writeFile(join(dir, 'a.txt'), 'a a a')
  179. await call('read', { file_path: 'a.txt' })
  180. const result = await call('edit', { file_path: 'a.txt', old_string: 'a', new_string: 'b' })
  181. expect(result.isError).toBe(true)
  182. expect(result.error).toMatchObject({ info: { code: 'FS_AMBIGUOUS_EDIT' } })
  183. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('a a a')
  184. })
  185. it('replaces all matches with replace_all', async () => {
  186. await writeFile(join(dir, 'a.txt'), 'a a a')
  187. await call('read', { file_path: 'a.txt' })
  188. const result = await call('edit', { file_path: 'a.txt', old_string: 'a', new_string: 'b', replace_all: true })
  189. expect(result.isError).toBe(false)
  190. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('b b b')
  191. })
  192. it('supports a full write→edit cycle without an intervening read', async () => {
  193. await call('write', { file_path: 'a.txt', content: 'one two' })
  194. const result = await call('edit', { file_path: 'a.txt', old_string: 'two', new_string: 'three' })
  195. expect(result.isError).toBe(false)
  196. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('one three')
  197. })
  198. })
  199. describe('the gate records only through the events (no method coupling)', () => {
  200. it('a direct ctx.fs.readText records no observed-state, so a later edit rejects', async () => {
  201. await writeFile(join(dir, 'a.txt'), 'hello world')
  202. // Reach AROUND the tool — an explicit escape hatch for non-tool consumers.
  203. await ctx.fs.readText(await ctx.fs.resolve('a.txt'))
  204. // The model-facing edit still rejects: the read did not emit fs/observed.
  205. const result = await call('edit', { file_path: 'a.txt', old_string: 'world', new_string: 'there' })
  206. expect(result.isError).toBe(true)
  207. expect(result.error).toMatchObject({ info: { code: 'FS_NOT_OBSERVED' } })
  208. })
  209. })
  210. describe('deleted observed target', () => {
  211. it('a failed reread records absence so write can safely recreate the file', async () => {
  212. await writeFile(join(dir, 'a.txt'), 'original')
  213. await call('read', { file_path: 'a.txt' })
  214. await rm(join(dir, 'a.txt')) // out-of-band deletion
  215. // The original positive observation still protects the first mutation.
  216. const edit = await call('edit', { file_path: 'a.txt', old_string: 'original', new_string: 'x' })
  217. expect(edit.isError).toBe(true)
  218. expect(edit.error).toMatchObject({ info: { code: 'FS_STALE_VERSION' } })
  219. const write = await call('write', { file_path: 'a.txt', content: 'premature' })
  220. expect(write.isError).toBe(true)
  221. expect(write.error).toMatchObject({ info: { code: 'FS_STALE_VERSION' } })
  222. // A read-not-found is an authoritative negative observation for this
  223. // owner. It still fails as a read, but changes the next write guard.
  224. const reread = await call('read', { file_path: 'a.txt' })
  225. expect(reread.isError).toBe(true)
  226. expect(reread.error).toMatchObject({ info: { code: 'FS_NOT_FOUND' } })
  227. // Absence never authorizes edit: there is no content/version to edit.
  228. const retriedEdit = await call('edit', { file_path: 'a.txt', old_string: 'original', new_string: 'x' })
  229. expect(retriedEdit.isError).toBe(true)
  230. expect(retriedEdit.error).toMatchObject({ info: { code: 'FS_NOT_FOUND' } })
  231. // The retried write uses createIfAbsent; the provider remains responsible
  232. // for rejecting a concurrent creator at publication time.
  233. const recovered = await call('write', { file_path: 'a.txt', content: 'fresh' })
  234. expect(recovered.isError).toBe(false)
  235. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('fresh')
  236. })
  237. })
  238. describe('stat budget', () => {
  239. it('read stats once; write and edit never stat in the tool (the gate stats zero too)', async () => {
  240. await writeFile(join(dir, 'a.txt'), 'hello world')
  241. const statSpy = vi.spyOn(ctx.fs, 'stat')
  242. // read: exactly one stat (type + size routing + observed version).
  243. await call('read', { file_path: 'a.txt' })
  244. expect(statSpy).toHaveBeenCalledTimes(1)
  245. // edit (guarded, after the read): the gate supplies vObserved; the tool
  246. // does not stat to manufacture a basis. CAS happens in editText's lock.
  247. statSpy.mockClear()
  248. const edited = await call('edit', { file_path: 'a.txt', old_string: 'world', new_string: 'there' })
  249. expect(edited.isError).toBe(false)
  250. expect(statSpy).not.toHaveBeenCalled()
  251. // write (guarded replace, after the edit refreshed observed state): zero stat.
  252. statSpy.mockClear()
  253. const written = await call('write', { file_path: 'a.txt', content: 'fresh' })
  254. expect(written.isError).toBe(false)
  255. expect(statSpy).not.toHaveBeenCalled()
  256. statSpy.mockRestore()
  257. })
  258. it('a missing read still stats once and its recovery write stats zero times', async () => {
  259. const statSpy = vi.spyOn(ctx.fs, 'stat')
  260. const missing = await call('read', { file_path: 'missing.txt' })
  261. expect(missing.isError).toBe(true)
  262. expect(missing.error).toMatchObject({ info: { code: 'FS_NOT_FOUND' } })
  263. expect(statSpy).toHaveBeenCalledTimes(1)
  264. statSpy.mockClear()
  265. const created = await call('write', { file_path: 'missing.txt', content: 'fresh' })
  266. expect(created.isError).toBe(false)
  267. expect(statSpy).not.toHaveBeenCalled()
  268. statSpy.mockRestore()
  269. })
  270. })
  271. })
  272. // --------------------------------------------------------------------------
  273. // BARE deployment: the tool suite WITHOUT the policy gate.
  274. // --------------------------------------------------------------------------
  275. describe('bare provider (no dsh-fs-observation-policy)', () => {
  276. beforeEach(async () => {
  277. dir = await mkdtemp(join(tmpdir(), 'dsh-tool-fs-bare-'))
  278. ctx = new Context()
  279. await ctx.plugin(SystemPrompt)
  280. await ctx.plugin(ToolRuntime)
  281. await ctx.plugin(LocalFileSystem, { cwd: dir })
  282. fiber = await ctx.plugin(ToolFs)
  283. })
  284. it('read works (it never needed policy)', async () => {
  285. await writeFile(join(dir, 'a.txt'), 'alpha\nbeta')
  286. const result = await call('read', { file_path: 'a.txt' })
  287. expect(result.isError).toBe(false)
  288. expect(text(result)).toContain('1: alpha')
  289. })
  290. it('write unconditionally creates a new file', async () => {
  291. const result = await call('write', { file_path: 'new.txt', content: 'fresh' })
  292. expect(result.isError).toBe(false)
  293. expect(await readFile(join(dir, 'new.txt'), 'utf8')).toBe('fresh')
  294. })
  295. it('write unconditionally OVERWRITES an existing unread file', async () => {
  296. await writeFile(join(dir, 'a.txt'), 'original')
  297. const result = await call('write', { file_path: 'a.txt', content: 'clobbered' })
  298. expect(result.isError).toBe(false)
  299. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('clobbered')
  300. })
  301. it('edit unconditionally edits an UNREAD existing file', async () => {
  302. await writeFile(join(dir, 'a.txt'), 'hello world')
  303. const result = await call('edit', { file_path: 'a.txt', old_string: 'world', new_string: 'there' })
  304. expect(result.isError).toBe(false)
  305. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('hello there')
  306. })
  307. it('edit of a MISSING target reports FS_STALE_VERSION even on the unguarded path', async () => {
  308. const result = await call('edit', { file_path: 'missing.txt', old_string: 'a', new_string: 'b' })
  309. expect(result.isError).toBe(true)
  310. expect(result.error).toMatchObject({ info: { code: 'FS_STALE_VERSION' } })
  311. // Even without policy, the stale text carries the re-read remedy.
  312. expect(text(result)).toContain('file changed since it was read')
  313. expect(text(result)).toContain('re-read the file, then retry')
  314. })
  315. it('edit still enforces literal-match codes (FS_EDIT_NOT_FOUND), unrelated to freshness', async () => {
  316. await writeFile(join(dir, 'a.txt'), 'hello world')
  317. const result = await call('edit', { file_path: 'a.txt', old_string: 'absent', new_string: 'x' })
  318. expect(result.isError).toBe(true)
  319. expect(result.error).toMatchObject({ info: { code: 'FS_EDIT_NOT_FOUND' } })
  320. })
  321. it('neither write nor edit stats in the tool on the bare path', async () => {
  322. await writeFile(join(dir, 'a.txt'), 'hello world')
  323. const statSpy = vi.spyOn(ctx.fs, 'stat')
  324. expect((await call('write', { file_path: 'a.txt', content: 'x y' })).isError).toBe(false)
  325. expect((await call('edit', { file_path: 'a.txt', old_string: 'y', new_string: 'z' })).isError).toBe(false)
  326. expect(statSpy).not.toHaveBeenCalled()
  327. statSpy.mockRestore()
  328. })
  329. })
  330. // Per-session cwd: a relative file_path resolves against the calling session's workspace
  331. // (`exec.agent.session.header.cwd`), not the backend's config.cwd, so the
  332. // caller-selected session workspace wins, matching dsh-tool-bash.
  333. describe('per-session cwd', () => {
  334. let sessionDir: string
  335. beforeEach(async () => {
  336. dir = await mkdtemp(join(tmpdir(), 'dsh-tool-fs-cfg-'))
  337. sessionDir = await mkdtemp(join(tmpdir(), 'dsh-tool-fs-session-'))
  338. ctx = new Context()
  339. await ctx.plugin(SystemPrompt)
  340. await ctx.plugin(ToolRuntime)
  341. await ctx.plugin(LocalFileSystem, { cwd: dir }) // config.cwd = dir, NOT sessionDir
  342. await ctx.plugin(FsPolicy)
  343. fiber = await ctx.plugin(ToolFs)
  344. })
  345. afterEach(async () => { await rm(sessionDir, { recursive: true, force: true }) })
  346. const callIn = (sessionObj: object, name: string, args: unknown) =>
  347. ctx.tools.execute({
  348. signal: testToolSignal,
  349. callId: ToolCallId(`call-${++callCounter}`),
  350. name,
  351. arguments: args,
  352. agent: { session: sessionObj } as never,
  353. })
  354. it('writes a relative path into the SESSION cwd, not config.cwd', async () => {
  355. const result = await callIn({ header: { cwd: sessionDir } }, 'write', { file_path: 'note.txt', content: 'hi' })
  356. expect(result.isError).toBe(false)
  357. // Verify the WORLD: the file is in the session dir, and NOT in config.cwd.
  358. expect(await readFile(join(sessionDir, 'note.txt'), 'utf8')).toBe('hi')
  359. await expect(readFile(join(dir, 'note.txt'), 'utf8')).rejects.toMatchObject({ code: 'ENOENT' })
  360. })
  361. it('read + edit both resolve against the session cwd (end-to-end)', async () => {
  362. // ONE session object across both calls — observed-state keys by owner
  363. // identity, so read must record under the same owner the edit reads.
  364. const session = { header: { cwd: sessionDir } }
  365. await writeFile(join(sessionDir, 'code.txt'), 'alpha')
  366. expect((await callIn(session, 'read', { file_path: 'code.txt' })).isError).toBe(false)
  367. const edited = await callIn(session, 'edit', { file_path: 'code.txt', old_string: 'alpha', new_string: 'beta' })
  368. expect(edited.isError).toBe(false)
  369. expect(await readFile(join(sessionDir, 'code.txt'), 'utf8')).toBe('beta')
  370. })
  371. })
  372. // --------------------------------------------------------------------------
  373. // Abort-through-the-tool, tool-tier concurrency, and the fs/observed contract —
  374. // all through ctx.tools.execute() against the REAL backend + policy.
  375. // --------------------------------------------------------------------------
  376. describe('signal, concurrency, and the fs/observed contract', () => {
  377. beforeEach(async () => {
  378. dir = await mkdtemp(join(tmpdir(), 'dsh-tool-fs-'))
  379. ctx = new Context()
  380. await ctx.plugin(SystemPrompt)
  381. await ctx.plugin(ToolRuntime)
  382. await ctx.plugin(LocalFileSystem, { cwd: dir })
  383. await ctx.plugin(FsPolicy)
  384. fiber = await ctx.plugin(ToolFs)
  385. })
  386. const session = { header: {} }
  387. const callSig = (signal: AbortSignal, name: string, args: unknown) =>
  388. ctx.tools.execute({ callId: ToolCallId(`c-${++callCounter}`), name, arguments: args, agent: { session } as never, signal })
  389. const callOwned = (name: string, args: unknown) =>
  390. ctx.tools.execute({ signal: testToolSignal, callId: ToolCallId(`c-${++callCounter}`), name, arguments: args, agent: { session } as never })
  391. it('a pre-aborted registry call skips read/write/edit with ABORTED_BEFORE_DISPATCH', async () => {
  392. await writeFile(join(dir, 'a.txt'), 'hello')
  393. const read = await callSig(AbortSignal.abort(), 'read', { file_path: 'a.txt' })
  394. expect(read.isError).toBe(true)
  395. expect(read.error).toMatchObject({ info: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH } })
  396. const write = await callSig(AbortSignal.abort(), 'write', { file_path: 'new.txt', content: 'x' })
  397. expect(write.isError).toBe(true)
  398. expect(write.error).toMatchObject({ info: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH } })
  399. await expect(readFile(join(dir, 'new.txt'), 'utf8')).rejects.toMatchObject({ code: 'ENOENT' })
  400. // Read first (un-aborted, SAME session owner) so the edit clears the
  401. // observation gate; then the registry skips the aborted edit before its body.
  402. expect((await callOwned('read', { file_path: 'a.txt' })).isError).toBe(false)
  403. const edit = await callSig(AbortSignal.abort(), 'edit', { file_path: 'a.txt', old_string: 'hello', new_string: 'bye' })
  404. expect(edit.isError).toBe(true)
  405. expect(edit.error).toMatchObject({ info: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH } })
  406. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('hello') // unchanged
  407. })
  408. it('two concurrent edits of the same file, same session: one wins, one FS_STALE_VERSION', async () => {
  409. await writeFile(join(dir, 'a.txt'), 'base value here')
  410. // One read establishes the observed version both edits guard against; then
  411. // race two edits so both carry the SAME observed version (the barrier).
  412. expect((await callOwned('read', { file_path: 'a.txt' })).isError).toBe(false)
  413. const [one, two] = await Promise.all([
  414. callOwned('edit', { file_path: 'a.txt', old_string: 'base', new_string: 'ONE', replaceAll: false }),
  415. callOwned('edit', { file_path: 'a.txt', old_string: 'value', new_string: 'TWO', replaceAll: false }),
  416. ])
  417. const errors = [one, two].filter(r => r.isError)
  418. expect(errors).toHaveLength(1)
  419. expect(errors[0]?.error).toMatchObject({ info: { code: 'FS_STALE_VERSION' } })
  420. // The world is consistent: exactly one edit landed.
  421. const onDisk = await readFile(join(dir, 'a.txt'), 'utf8')
  422. expect(onDisk === 'ONE value here' || onDisk === 'base TWO here').toBe(true)
  423. })
  424. it('a stale observed version from an older read fails closed at edit CAS', async () => {
  425. await writeFile(join(dir, 'a.txt'), 'older content\n')
  426. const target = await ctx.fs.resolve('a.txt')
  427. const firstInfo = await ctx.fs.stat(target)
  428. if (!firstInfo) throw new Error('expected first stat')
  429. expect((await callOwned('read', { file_path: 'a.txt' })).isError).toBe(false)
  430. await writeFile(join(dir, 'a.txt'), 'newer current content\n')
  431. const secondInfo = await ctx.fs.stat(target)
  432. if (!secondInfo) throw new Error('expected second stat')
  433. expect(secondInfo.version).not.toBe(firstInfo.version)
  434. expect((await callOwned('read', { file_path: 'a.txt' })).isError).toBe(false)
  435. // Reproduce an older concurrent read winning the observation race.
  436. ctx.emit('fs/observed', target, { kind: 'present', version: firstInfo.version }, { agent: { session } })
  437. const edit = await callOwned('edit', {
  438. file_path: 'a.txt',
  439. old_string: 'newer',
  440. new_string: 'edited',
  441. })
  442. expect(edit.isError).toBe(true)
  443. expect(edit.error).toMatchObject({ info: { code: 'FS_STALE_VERSION' } })
  444. expect(await readFile(join(dir, 'a.txt'), 'utf8')).toBe('newer current content\n')
  445. })
  446. it('a throwing fs/observed listener surfaces as isError, but the mutation already hit disk', async () => {
  447. // fs/observed is a plain ctx.emit after the write succeeded; a throwing listener cannot
  448. // roll the write back — it only turns the tool result into isError.
  449. ctx.on('fs/observed', () => { throw new Error('recording bug') })
  450. const result = await callOwned('write', { file_path: 'w.txt', content: 'durable' })
  451. expect(result.isError).toBe(true)
  452. expect(await readFile(join(dir, 'w.txt'), 'utf8')).toBe('durable')
  453. })
  454. })