watcher.spec.ts 8.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223
  1. import { afterEach, describe, expect, it, vi } from 'vitest'
  2. import { Context } from 'cordis'
  3. import { chmod, mkdtemp, rm, writeFile } from 'node:fs/promises'
  4. import { tmpdir } from 'node:os'
  5. import { join } from 'node:path'
  6. import { credentialRef } from '@deepseek-ai/dsh-credentials'
  7. import { CredentialsLocal } from '../src/index.ts'
  8. // chokidar is the nondeterministic OS boundary: faking it lets these tests
  9. // drive the event pipeline (error events, races with unreadable files)
  10. // deterministically. Real end-to-end watching stays covered by local.spec.ts.
  11. vi.mock('chokidar', async () => {
  12. const { EventEmitter } = await import('node:events')
  13. class FakeWatcher extends EventEmitter {
  14. close = vi.fn(() => Promise.resolve())
  15. }
  16. const instances: Array<{ path: string; options: unknown; watcher: InstanceType<typeof FakeWatcher> }> = []
  17. return {
  18. watch: vi.fn((path: string, options: unknown) => {
  19. const watcher = new FakeWatcher()
  20. instances.push({ path, options, watcher })
  21. return watcher
  22. }),
  23. __instances: instances,
  24. }
  25. })
  26. interface FakeChokidar {
  27. __instances: Array<{
  28. path: string
  29. options: { awaitWriteFinish: { stabilityThreshold: number; pollInterval: number } }
  30. watcher: import('node:events').EventEmitter
  31. }>
  32. }
  33. async function fakeInstances(): Promise<FakeChokidar['__instances']> {
  34. const chokidar = await import('chokidar') as unknown as FakeChokidar
  35. return chokidar.__instances
  36. }
  37. const KEY = credentialRef('DSH_CRED_PIPE')
  38. const cleanups: Array<() => Promise<void>> = []
  39. afterEach(async () => {
  40. while (cleanups.length > 0) await cleanups.pop()!()
  41. ;(await fakeInstances()).length = 0
  42. })
  43. async function tempDir(): Promise<string> {
  44. const dir = await mkdtemp(join(tmpdir(), 'dsh-credentials-watch-'))
  45. cleanups.push(() => rm(dir, { recursive: true, force: true }))
  46. return dir
  47. }
  48. async function boot(config: ConstructorParameters<typeof CredentialsLocal>[1]): Promise<Context> {
  49. const ctx = new Context()
  50. const fiber = ctx.plugin(CredentialsLocal, config)
  51. cleanups.push(async () => {
  52. await fiber.dispose()
  53. })
  54. await fiber
  55. return ctx
  56. }
  57. describe('watcher pipeline', () => {
  58. it('clamps the write-settle poll interval for a zero debounce', async () => {
  59. const dir = await tempDir()
  60. await boot({ path: join(dir, '.env'), debounceMs: 0 })
  61. const [instance] = await fakeInstances()
  62. expect(instance!.options.awaitWriteFinish).toEqual({ stabilityThreshold: 0, pollInterval: 1 })
  63. })
  64. it('survives a watcher error and keeps publishing later edits', async () => {
  65. const dir = await tempDir()
  66. const path = join(dir, '.env')
  67. const ctx = await boot({ path, debounceMs: 5 })
  68. const [instance] = await fakeInstances()
  69. instance!.watcher.emit('error', new Error('watch backend failure'))
  70. expect(await ctx.credentials.resolve(KEY)).toBeUndefined()
  71. await writeFile(path, 'DSH_CRED_PIPE=arrived\n')
  72. instance!.watcher.emit('all', 'change', path)
  73. await vi.waitFor(async () => {
  74. expect(await ctx.credentials.resolve(KEY)).toEqual({ value: 'arrived', source: 'file' })
  75. })
  76. })
  77. it('keeps the last good snapshot when the file turns unreadable at runtime', async () => {
  78. const dir = await tempDir()
  79. const path = join(dir, '.env')
  80. await writeFile(path, 'DSH_CRED_PIPE=good\n')
  81. const ctx = await boot({ path, debounceMs: 5 })
  82. await chmod(path, 0o000)
  83. cleanups.push(() => chmod(path, 0o600))
  84. const [instance] = await fakeInstances()
  85. instance!.watcher.emit('all', 'change', path)
  86. // The warn-and-keep path is asynchronous; give the serialized refresh a turn.
  87. await new Promise(resolve => setTimeout(resolve, 50))
  88. expect(await ctx.credentials.resolve(KEY)).toEqual({ value: 'good', source: 'file' })
  89. })
  90. it('keeps the reload queue alive after an invariant violation escapes the fan-out', async () => {
  91. const dir = await tempDir()
  92. const path = join(dir, '.env')
  93. const ctx = await boot({ path, debounceMs: 5 })
  94. let arm = true
  95. ctx.on('credentials/updated', () => {
  96. if (!arm) return
  97. throw Object.assign(new Error('forged relation'), { code: 'INVARIANT' })
  98. })
  99. const [instance] = await fakeInstances()
  100. await writeFile(path, 'DSH_CRED_PIPE=first\n')
  101. instance!.watcher.emit('all', 'change', path)
  102. // The snapshot commits before the fan-out, so the value lands even though
  103. // the listener threw out of the refresh.
  104. await vi.waitFor(async () => {
  105. expect(await ctx.credentials.resolve(KEY)).toEqual({ value: 'first', source: 'file' })
  106. })
  107. arm = false
  108. await writeFile(path, 'DSH_CRED_PIPE=second\n')
  109. instance!.watcher.emit('all', 'change', path)
  110. await vi.waitFor(async () => {
  111. expect(await ctx.credentials.resolve(KEY)).toEqual({ value: 'second', source: 'file' })
  112. })
  113. })
  114. it('quiesces the refresh pipeline before dispose completes', async () => {
  115. const dir = await tempDir()
  116. const path = join(dir, '.env')
  117. await writeFile(path, 'DSH_CRED_PIPE=initial\n')
  118. const ctx = new Context()
  119. const fiber = ctx.plugin(CredentialsLocal, { path, debounceMs: 5 })
  120. await fiber
  121. let disposed = false
  122. let postDisposeCommits = 0
  123. ctx.on('credentials/updated', () => {
  124. if (disposed) postDisposeCommits += 1
  125. })
  126. await writeFile(path, 'DSH_CRED_PIPE=changed\n')
  127. const [instance] = await fakeInstances()
  128. // Two queued refreshes: dispose interrupts one mid-flight and the other
  129. // before it starts, so both closed guards must hold.
  130. instance!.watcher.emit('all', 'change', path)
  131. instance!.watcher.emit('all', 'change', path)
  132. await fiber.dispose()
  133. disposed = true
  134. instance!.watcher.emit('all', 'change', path)
  135. instance!.watcher.emit('ready')
  136. await new Promise(resolve => setTimeout(resolve, 100))
  137. expect(postDisposeCommits).toBe(0)
  138. })
  139. it('empties the snapshot when the document is deleted and emits the removals', async () => {
  140. const dir = await tempDir()
  141. const path = join(dir, '.env')
  142. await writeFile(path, 'DSH_CRED_PIPE=doomed\n')
  143. const ctx = await boot({ path, debounceMs: 5 })
  144. const seen: string[] = []
  145. ctx.on('credentials/updated', (ref) => {
  146. seen.push(ref)
  147. })
  148. await rm(path)
  149. const [instance] = await fakeInstances()
  150. instance!.watcher.emit('all', 'unlink', path)
  151. await vi.waitFor(async () => {
  152. expect(await ctx.credentials.resolve(KEY)).toBeUndefined()
  153. })
  154. expect(seen).toEqual([KEY])
  155. })
  156. it('publishes only seam-addressable keys and preserves the rest untouched', async () => {
  157. const dir = await tempDir()
  158. const path = join(dir, '.env')
  159. await writeFile(path, 'BAD-KEY=1\nDSH_CRED_PIPE=a\n')
  160. const ctx = await boot({ path, debounceMs: 5 })
  161. const seen: string[] = []
  162. ctx.on('credentials/updated', (ref) => {
  163. seen.push(ref)
  164. })
  165. await writeFile(path, 'BAD-KEY=2\nDSH_CRED_PIPE=b\n')
  166. const [instance] = await fakeInstances()
  167. instance!.watcher.emit('all', 'change', path)
  168. await vi.waitFor(async () => {
  169. expect(await ctx.credentials.resolve(KEY)).toEqual({ value: 'b', source: 'file' })
  170. })
  171. // The dash-named key is preserved file content the seam cannot address:
  172. // its change publishes nothing and breaks nothing.
  173. expect(seen).toEqual([KEY])
  174. })
  175. it('treats an event for a still-absent file as a no-op', async () => {
  176. const dir = await tempDir()
  177. const path = join(dir, '.env')
  178. const ctx = await boot({ path, debounceMs: 5 })
  179. const [instance] = await fakeInstances()
  180. instance!.watcher.emit('all', 'add', path)
  181. await new Promise(resolve => setTimeout(resolve, 50))
  182. expect(await ctx.credentials.resolve(KEY)).toBeUndefined()
  183. })
  184. it('reconciles at watcher ready so a change during setup is not missed', async () => {
  185. const dir = await tempDir()
  186. const path = join(dir, '.env')
  187. await writeFile(path, `${KEY}=a\n`)
  188. const ctx = await boot({ path, debounceMs: 5 })
  189. // Written after the initial load but before the watcher became active:
  190. // no 'all' event will ever fire for it.
  191. await writeFile(path, `${KEY}=written-before-ready\n`)
  192. const [instance] = await fakeInstances()
  193. instance!.watcher.emit('ready')
  194. await vi.waitFor(async () => {
  195. expect(await ctx.credentials.resolve(KEY)).toEqual({ value: 'written-before-ready', source: 'file' })
  196. })
  197. })
  198. })