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