skill-filesystem-watcher.spec.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452
  1. import { EventEmitter } from 'node:events'
  2. import type { Stats } from 'node:fs'
  3. import { mkdir, realpath, rm, symlink, writeFile } from 'node:fs/promises'
  4. import { join } from 'node:path'
  5. import { tmpdir } from 'node:os'
  6. import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
  7. import { Context } from '@deepseek-ai/cordis'
  8. import SkillRegistry from '@deepseek-ai/dsh-skill'
  9. interface FakeWatcherControl {
  10. emitter: EventEmitter
  11. closeCalls: number
  12. options: Record<string, unknown>
  13. path: string
  14. }
  15. interface FakeWatchFileControl {
  16. path: string
  17. listener(current: Stats, previous: Stats): void
  18. }
  19. interface FakeStatGate {
  20. started: PromiseWithResolvers<undefined>
  21. release: PromiseWithResolvers<undefined>
  22. }
  23. const watcherHarness = vi.hoisted(() => ({
  24. watchers: [] as FakeWatcherControl[],
  25. startupErrors: [] as Error[],
  26. closeErrors: 0,
  27. deferredReady: 0,
  28. watchFiles: [] as FakeWatchFileControl[],
  29. statGates: [] as FakeStatGate[],
  30. }))
  31. vi.mock('node:fs', async (importOriginal) => {
  32. const actual = await importOriginal<typeof import('node:fs')>()
  33. return {
  34. ...actual,
  35. watchFile(path: string, _options: unknown, listener: FakeWatchFileControl['listener']) {
  36. watcherHarness.watchFiles.push({ path, listener })
  37. },
  38. unwatchFile(path: string, listener: FakeWatchFileControl['listener']) {
  39. const index = watcherHarness.watchFiles.findIndex(control => control.path === path && control.listener === listener)
  40. if (index !== -1) watcherHarness.watchFiles.splice(index, 1)
  41. },
  42. }
  43. })
  44. vi.mock('node:fs/promises', async (importOriginal) => {
  45. const actual = await importOriginal<typeof import('node:fs/promises')>()
  46. return {
  47. ...actual,
  48. async stat(...args: Parameters<typeof actual.stat>) {
  49. const gate = watcherHarness.statGates.shift()
  50. if (gate !== undefined) {
  51. gate.started.resolve(undefined)
  52. await gate.release.promise
  53. }
  54. return await actual.stat(...args)
  55. },
  56. }
  57. })
  58. vi.mock('chokidar', () => ({
  59. default: {
  60. watch(path: unknown, options: Record<string, unknown>) {
  61. const emitter = new EventEmitter() as EventEmitter & { close(): Promise<void> }
  62. const control: FakeWatcherControl = { emitter, closeCalls: 0, options, path: String(path) }
  63. emitter.close = async () => {
  64. control.closeCalls += 1
  65. if (watcherHarness.closeErrors > 0) {
  66. watcherHarness.closeErrors -= 1
  67. throw new Error('close failed')
  68. }
  69. }
  70. watcherHarness.watchers.push(control)
  71. queueMicrotask(() => {
  72. if (watcherHarness.deferredReady > 0) {
  73. watcherHarness.deferredReady -= 1
  74. return
  75. }
  76. const error = watcherHarness.startupErrors.shift()
  77. if (error === undefined) emitter.emit('ready')
  78. else emitter.emit('error', error)
  79. })
  80. return emitter
  81. },
  82. },
  83. }))
  84. const SkillFileSystem = await import('../src/index.ts')
  85. /** Every temp dir created by this file, removed after each test. */
  86. const tempDirs: string[] = []
  87. afterEach(async () => {
  88. for (const dir of tempDirs.splice(0)) await rm(dir, { recursive: true, force: true })
  89. })
  90. async function tempDir(name: string): Promise<string> {
  91. const dir = await import('node:fs/promises').then(fs => fs.mkdtemp(join(tmpdir(), `dsh-${name}-`)))
  92. tempDirs.push(dir)
  93. return dir
  94. }
  95. async function writeSkill(root: string, name: string): Promise<void> {
  96. const directory = join(root, name)
  97. await mkdir(directory, { recursive: true })
  98. await writeFile(join(directory, 'SKILL.md'), `---\nname: ${name}\ndescription: ${name}\n---\n\nBody.\n`)
  99. }
  100. async function settle(): Promise<void> {
  101. await new Promise(resolve => setTimeout(resolve, 0))
  102. }
  103. beforeEach(() => {
  104. watcherHarness.watchers.length = 0
  105. watcherHarness.startupErrors.length = 0
  106. watcherHarness.closeErrors = 0
  107. watcherHarness.deferredReady = 0
  108. watcherHarness.watchFiles.length = 0
  109. watcherHarness.statGates.length = 0
  110. })
  111. describe('skill-filesystem watcher failures', () => {
  112. it('canonicalizes an existing root before opening its native watcher', async () => {
  113. const target = await tempDir('skill-watch-canonical-target')
  114. const aliasParent = await tempDir('skill-watch-canonical-alias')
  115. const alias = join(aliasParent, 'alias')
  116. await symlink(target, alias, process.platform === 'win32' ? 'junction' : 'dir')
  117. const root = join(alias, '.dsh/skills')
  118. await writeSkill(root, 'canonical-skill')
  119. const ctx = new Context()
  120. await ctx.plugin(SkillRegistry)
  121. const fiber = await ctx.plugin(SkillFileSystem, {
  122. dshHome: join(alias, '.dsh'),
  123. agentsHome: join(alias, '.agents'),
  124. watch: true,
  125. })
  126. expect((await ctx.skills.list()).map(skill => skill.name)).toEqual(['canonical-skill'])
  127. expect(watcherHarness.watchers[0]?.path).toBe(await realpath(root))
  128. expect(watcherHarness.watchers[0]?.options.persistent).toBe(true)
  129. await fiber.dispose()
  130. })
  131. it('preserves a symlink root when link following is disabled', async () => {
  132. const target = await tempDir('skill-watch-link-target')
  133. const aliasParent = await tempDir('skill-watch-link-alias')
  134. const alias = join(aliasParent, 'skills')
  135. await writeSkill(target, 'linked-skill')
  136. await symlink(target, alias, process.platform === 'win32' ? 'junction' : 'dir')
  137. const ctx = new Context()
  138. await ctx.plugin(SkillRegistry)
  139. const fiber = await ctx.plugin(SkillFileSystem, {
  140. includeDefaultRoots: false,
  141. customSkillDirs: [alias],
  142. watch: true,
  143. watchFollowSymlinks: false,
  144. })
  145. try {
  146. expect((await ctx.skills.list()).map(skill => skill.name)).toEqual(['linked-skill'])
  147. expect(watcherHarness.watchers[0]?.path).toBe(alias)
  148. expect(watcherHarness.watchers[0]?.options.followSymlinks).toBe(false)
  149. } finally {
  150. await fiber.dispose()
  151. await rm(aliasParent, { recursive: true, force: true })
  152. await rm(target, { recursive: true, force: true })
  153. }
  154. })
  155. it('ignores missing-path probes until the observed path actually changes', async () => {
  156. const home = await tempDir('skill-watch-missing-stable')
  157. const ctx = new Context()
  158. await ctx.plugin(SkillRegistry)
  159. const fiber = await ctx.plugin(SkillFileSystem, {
  160. dshHome: join(home, '.dsh'),
  161. agentsHome: join(home, '.agents'),
  162. watch: true,
  163. watchPollIntervalMs: 10,
  164. })
  165. expect(await ctx.skills.snapshot()).toEqual({ skills: [], complete: true })
  166. expect(watcherHarness.watchFiles).toHaveLength(2)
  167. let invalidations = 0
  168. ctx.on('skills/change', () => { invalidations += 1 })
  169. for (const control of watcherHarness.watchFiles) {
  170. control.listener({} as Stats, {} as Stats)
  171. }
  172. await settle()
  173. expect(invalidations).toBe(0)
  174. expect(watcherHarness.watchFiles).toHaveLength(2)
  175. await fiber.dispose()
  176. })
  177. it('keeps skills loadable across persistent watcher startup failures without caching them', async () => {
  178. const home = await tempDir('skill-watch-start-error')
  179. const root = join(home, '.dsh/skills')
  180. await writeSkill(root, 'retry-skill')
  181. watcherHarness.startupErrors.push(
  182. new Error('watch failed once'),
  183. new Error('watch failed twice'),
  184. new Error('watch failed three times'),
  185. )
  186. watcherHarness.closeErrors = 1
  187. const ctx = new Context()
  188. await ctx.plugin(SkillRegistry)
  189. const fiber = await ctx.plugin(SkillFileSystem, {
  190. dshHome: join(home, '.dsh'),
  191. agentsHome: join(home, '.agents'),
  192. watch: true,
  193. watchUsePolling: true,
  194. watchFollowSymlinks: false,
  195. watchPollIntervalMs: 10,
  196. watchStabilityThresholdMs: 20,
  197. })
  198. expect(await ctx.skills.snapshot()).toMatchObject({
  199. skills: [{ name: 'retry-skill' }],
  200. complete: false,
  201. })
  202. expect((await ctx.skills.get('retry-skill'))?.content).toBe('Body.')
  203. expect(await ctx.skills.snapshot()).toMatchObject({
  204. skills: [{ name: 'retry-skill' }],
  205. complete: false,
  206. })
  207. expect(watcherHarness.watchers).toHaveLength(3)
  208. expect(watcherHarness.watchers[0]?.options).toMatchObject({
  209. atomic: true,
  210. depth: 1,
  211. followSymlinks: false,
  212. usePolling: true,
  213. interval: 10,
  214. awaitWriteFinish: {
  215. stabilityThreshold: 20,
  216. pollInterval: 10,
  217. },
  218. })
  219. await fiber.dispose()
  220. })
  221. it('filters events, coalesces invalidation, recovers runtime errors, and contains late callbacks', async () => {
  222. const home = await tempDir('skill-watch-runtime-error')
  223. const root = join(home, '.dsh/skills')
  224. await writeSkill(root, 'watched-skill')
  225. const ctx = new Context()
  226. await ctx.plugin(SkillRegistry)
  227. const fiber = await ctx.plugin(SkillFileSystem, {
  228. dshHome: join(home, '.dsh'),
  229. agentsHome: join(home, '.agents'),
  230. watch: true,
  231. watchPollIntervalMs: 10,
  232. watchStabilityThresholdMs: 20,
  233. })
  234. expect((await ctx.skills.list()).map(skill => skill.name)).toEqual(['watched-skill'])
  235. let invalidations = 0
  236. ctx.on('skills/change', () => { invalidations += 1 })
  237. const first = watcherHarness.watchers[0]
  238. if (first === undefined) throw new Error('expected a root watcher')
  239. first.emitter.emit('change', join(first.path, 'notes.txt'))
  240. first.emitter.emit('change', join(home, 'outside.md'))
  241. first.emitter.emit('change', join(first.path, 'watched-skill/references.md'))
  242. first.emitter.emit('change', join(first.path, '.system/SKILL.md'))
  243. await settle()
  244. expect(invalidations).toBe(0)
  245. first.emitter.emit('change', join(first.path, 'watched-skill/SKILL.md'))
  246. first.emitter.emit('change', join(first.path, 'watched-skill/SKILL.md'))
  247. await settle()
  248. expect(invalidations).toBe(1)
  249. watcherHarness.closeErrors = 1
  250. watcherHarness.startupErrors.push(new Error('runtime rewatch failed'))
  251. first.emitter.emit('error', new Error('runtime watch failed'))
  252. await vi.waitFor(() => { expect(watcherHarness.watchers.length).toBeGreaterThanOrEqual(2) })
  253. expect(invalidations).toBeGreaterThanOrEqual(2)
  254. expect(await ctx.skills.snapshot()).toMatchObject({
  255. skills: [{ name: 'watched-skill' }],
  256. complete: true,
  257. })
  258. await fiber.dispose()
  259. first.emitter.emit('change', join(first.path, 'watched-skill/SKILL.md'))
  260. first.emitter.emit('error', new Error('late error'))
  261. await settle()
  262. })
  263. it('replaces a retained watcher when its root emits unlinkDir', async () => {
  264. const home = await tempDir('skill-watch-root-unlink')
  265. const root = join(home, '.dsh/skills')
  266. await writeSkill(root, 'removed-skill')
  267. const ctx = new Context()
  268. await ctx.plugin(SkillRegistry)
  269. const fiber = await ctx.plugin(SkillFileSystem, {
  270. dshHome: join(home, '.dsh'),
  271. agentsHome: join(home, '.agents'),
  272. watch: true,
  273. watchPollIntervalMs: 10,
  274. watchStabilityThresholdMs: 20,
  275. })
  276. expect((await ctx.skills.list()).map(skill => skill.name)).toEqual(['removed-skill'])
  277. const original = watcherHarness.watchers[0]
  278. if (original === undefined) throw new Error('expected a root watcher')
  279. await rm(root, { recursive: true })
  280. original.emitter.emit('unlinkDir', original.path)
  281. await vi.waitFor(() => { expect(original.closeCalls).toBeGreaterThan(0) })
  282. await vi.waitFor(() => {
  283. expect(watcherHarness.watchFiles.some(control => control.path === original.path)).toBe(true)
  284. })
  285. await fiber.dispose()
  286. })
  287. it('re-probes a retained root after child unlink and observes immediate recreation', async () => {
  288. const home = await tempDir('skill-watch-root-reprobe')
  289. const root = join(home, '.dsh/skills')
  290. await writeSkill(root, 'old-skill')
  291. const ctx = new Context()
  292. await ctx.plugin(SkillRegistry)
  293. const fiber = await ctx.plugin(SkillFileSystem, {
  294. dshHome: join(home, '.dsh'),
  295. agentsHome: join(home, '.agents'),
  296. watch: true,
  297. watchPollIntervalMs: 10,
  298. watchStabilityThresholdMs: 20,
  299. })
  300. expect((await ctx.skills.list()).map(skill => skill.name)).toEqual(['old-skill'])
  301. const original = watcherHarness.watchers[0]
  302. if (original === undefined) throw new Error('expected a root watcher')
  303. await rm(root, { recursive: true })
  304. original.emitter.emit('unlink', join(original.path, 'old-skill/SKILL.md'))
  305. await settle()
  306. expect(await ctx.skills.snapshot()).toEqual({ skills: [], complete: true })
  307. const missingRoot = watcherHarness.watchFiles.find(control => control.path === original.path)
  308. expect(missingRoot).toBeDefined()
  309. await writeSkill(root, 'recreated-skill')
  310. missingRoot!.listener({} as Stats, {} as Stats)
  311. await vi.waitFor(() => { expect(watcherHarness.watchers).toHaveLength(2) })
  312. await settle()
  313. expect((await ctx.skills.list()).map(skill => skill.name)).toEqual(['recreated-skill'])
  314. await fiber.dispose()
  315. })
  316. it('settles an opening watcher when plugin disposal races its ready event', async () => {
  317. const home = await tempDir('skill-watch-opening-dispose')
  318. const root = join(home, '.dsh/skills')
  319. await writeSkill(root, 'racing-skill')
  320. watcherHarness.deferredReady = 1
  321. const ctx = new Context()
  322. await ctx.plugin(SkillRegistry)
  323. let provider!: InstanceType<typeof SkillFileSystem.FileSystemSkillProvider>
  324. const disposeProvider = ctx.skills.registerProvider((control) => {
  325. provider = new SkillFileSystem.FileSystemSkillProvider(ctx, control, {
  326. dshHome: join(home, '.dsh'),
  327. agentsHome: join(home, '.agents'),
  328. watch: true,
  329. watchPollIntervalMs: 10,
  330. watchStabilityThresholdMs: 20,
  331. })
  332. return provider
  333. })
  334. const discovery = provider.list({})
  335. await vi.waitFor(() => { expect(watcherHarness.watchers).toHaveLength(1) })
  336. const first = watcherHarness.watchers[0]
  337. if (first === undefined) throw new Error('expected an opening root watcher')
  338. const disposal = provider.dispose()
  339. await expect(discovery).rejects.toThrow('skill-filesystem watcher disposed')
  340. await disposal
  341. disposeProvider()
  342. await settle()
  343. expect(first.closeCalls).toBeGreaterThan(0)
  344. })
  345. it('closes an opening watcher when disposal wins the mode probe', async () => {
  346. const home = await tempDir('skill-watch-probe-dispose')
  347. const root = join(home, '.dsh/skills')
  348. await writeSkill(root, 'racing-skill')
  349. watcherHarness.deferredReady = 1
  350. const statGate: FakeStatGate = {
  351. started: Promise.withResolvers<undefined>(),
  352. release: Promise.withResolvers<undefined>(),
  353. }
  354. watcherHarness.statGates.push(statGate)
  355. const ctx = new Context()
  356. await ctx.plugin(SkillRegistry)
  357. let provider!: InstanceType<typeof SkillFileSystem.FileSystemSkillProvider>
  358. const disposeProvider = ctx.skills.registerProvider((control) => {
  359. provider = new SkillFileSystem.FileSystemSkillProvider(ctx, control, {
  360. dshHome: join(home, '.dsh'),
  361. agentsHome: join(home, '.agents'),
  362. watch: true,
  363. watchPollIntervalMs: 10,
  364. watchStabilityThresholdMs: 20,
  365. })
  366. return provider
  367. })
  368. const discovery = provider.list({})
  369. await statGate.started.promise
  370. const disposal = provider.dispose()
  371. statGate.release.resolve(undefined)
  372. await expect(discovery).rejects.toThrow('skill-filesystem watcher disposed')
  373. await disposal
  374. expect(watcherHarness.watchers).toHaveLength(1)
  375. expect(watcherHarness.watchers[0]?.closeCalls).toBeGreaterThan(0)
  376. disposeProvider()
  377. })
  378. it('contains an opening watcher rejection during provider teardown', async () => {
  379. const home = await tempDir('skill-watch-opening-reject')
  380. const root = join(home, '.dsh/skills')
  381. await writeSkill(root, 'rejected-skill')
  382. watcherHarness.deferredReady = 1
  383. const ctx = new Context()
  384. await ctx.plugin(SkillRegistry)
  385. let provider!: InstanceType<typeof SkillFileSystem.FileSystemSkillProvider>
  386. const disposeProvider = ctx.skills.registerProvider((control) => {
  387. provider = new SkillFileSystem.FileSystemSkillProvider(ctx, control, {
  388. dshHome: join(home, '.dsh'),
  389. agentsHome: join(home, '.agents'),
  390. watch: true,
  391. watchPollIntervalMs: 10,
  392. watchStabilityThresholdMs: 20,
  393. })
  394. return provider
  395. })
  396. const discovery = provider.list({})
  397. await vi.waitFor(() => { expect(watcherHarness.watchers).toHaveLength(1) })
  398. const first = watcherHarness.watchers[0]
  399. if (first === undefined) throw new Error('expected an opening root watcher')
  400. first.emitter.emit('error', new Error('opening failed during disposal'))
  401. const disposal = provider.dispose()
  402. await expect(discovery).rejects.toThrow('opening failed during disposal')
  403. await disposal
  404. disposeProvider()
  405. })
  406. })