drain.spec.ts 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101
  1. import { afterEach, describe, expect, it, vi } from 'vitest'
  2. import { Context } from '@deepseek-ai/cordis'
  3. import { mkdtemp, rm } from 'node:fs/promises'
  4. import { tmpdir } from 'node:os'
  5. import { join } from 'node:path'
  6. import { credentialKey, credentialRef } from '@deepseek-ai/dsh-credentials'
  7. import { LocalCredentialProvider } from '../src/index.ts'
  8. // The atomic write is the gated asynchronous hold point inside a queued
  9. // write; gating it makes the dispose-versus-queued-write race fully
  10. // deterministic. The lock helper passes through so the gated operation still
  11. // runs inside its real acquire/release cycle.
  12. vi.mock('@deepseek-ai/dsh-atomic-write', async (importOriginal) => {
  13. const actual = await importOriginal<typeof import('@deepseek-ai/dsh-atomic-write')>()
  14. let gate: Promise<void> = Promise.resolve()
  15. return {
  16. ...actual,
  17. writeFileAtomic: vi.fn(() => gate),
  18. __setGate: (next: Promise<void>) => {
  19. gate = next
  20. },
  21. }
  22. })
  23. async function setGate(next: Promise<void>): Promise<void> {
  24. const mocked = await import('@deepseek-ai/dsh-atomic-write') as unknown as { __setGate: (next: Promise<void>) => void }
  25. mocked.__setGate(next)
  26. }
  27. const KEY = credentialRef('DSH_CRED_DRAIN_A')
  28. const OTHER = credentialRef('DSH_CRED_DRAIN_B')
  29. const RECORD = credentialKey('llm-drain', 'alpha')
  30. const OTHER_RECORD = credentialKey('llm-drain', 'beta')
  31. const cleanups: Array<() => Promise<void>> = []
  32. afterEach(async () => {
  33. await setGate(Promise.resolve())
  34. while (cleanups.length > 0) await cleanups.pop()!()
  35. })
  36. describe('write-drain teardown', () => {
  37. it('lets the in-flight write land and fails the queued one after disposal', async () => {
  38. const dir = await mkdtemp(join(tmpdir(), 'dsh-credentials-drain-'))
  39. cleanups.push(() => rm(dir, { recursive: true, force: true }))
  40. const ctx = new Context()
  41. const fiber = ctx.plugin(LocalCredentialProvider, { path: join(dir, '.credentials.yaml'), watch: false })
  42. await fiber
  43. const service = ctx.credentials
  44. let release!: () => void
  45. await setGate(new Promise<void>((resolveGate) => {
  46. release = resolveGate
  47. }))
  48. const first = service.set(KEY, 'one')
  49. // Let the first task pass its liveness checks and park on the gate, so it
  50. // is genuinely in-flight when disposal begins.
  51. await new Promise(resolvePause => setTimeout(resolvePause, 5))
  52. // Attach the rejection handler up front: the queued write fails while the
  53. // drain is still awaited, before any later `await expect` could run.
  54. const secondRejects = expect(service.set(OTHER, 'two')).rejects.toThrow(/disposed before the queued/)
  55. const disposal = fiber.dispose()
  56. // Give the drain disposer its first turn (set closed) before opening the gate.
  57. await new Promise(resolvePause => setTimeout(resolvePause, 10))
  58. release()
  59. await disposal
  60. await expect(first).resolves.toBeUndefined()
  61. await secondRejects
  62. expect(await service.resolve(KEY)).toEqual({ value: 'one', source: 'file' })
  63. expect(await service.resolve(OTHER)).toBeUndefined()
  64. })
  65. it('fails a queued record write after disposal on the same terms', async () => {
  66. const dir = await mkdtemp(join(tmpdir(), 'dsh-credentials-drain-record-'))
  67. cleanups.push(() => rm(dir, { recursive: true, force: true }))
  68. const ctx = new Context()
  69. const fiber = ctx.plugin(LocalCredentialProvider, { path: join(dir, '.credentials.yaml'), watch: false })
  70. await fiber
  71. const service = ctx.credentials
  72. let release!: () => void
  73. await setGate(new Promise<void>((resolveGate) => {
  74. release = resolveGate
  75. }))
  76. const first = service.modifyRecord(RECORD, () => Promise.resolve({ kind: 'grant', payload: { v: 1 } }))
  77. await new Promise(resolvePause => setTimeout(resolvePause, 5))
  78. const queuedModify = expect(service.modifyRecord(OTHER_RECORD, () => Promise.resolve({ kind: 'api-key' })))
  79. .rejects.toThrow(/disposed before the queued/)
  80. const queuedDelete = expect(service.deleteRecord(OTHER_RECORD)).rejects.toThrow(/disposed before the queued/)
  81. const disposal = fiber.dispose()
  82. await new Promise(resolvePause => setTimeout(resolvePause, 10))
  83. release()
  84. await disposal
  85. await expect(first).resolves.toEqual({ kind: 'grant', payload: { v: 1 } })
  86. await queuedModify
  87. await queuedDelete
  88. expect(await service.readRecord(OTHER_RECORD)).toBeUndefined()
  89. })
  90. })