retention.spec.ts 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195
  1. /** Monotonic idle grace, independent window holds, and retryable terminal ownership. */
  2. import { afterEach, beforeEach, expect, it, vi } from 'vitest'
  3. import type { SubprocessTerminalActivity } from '@deepseek-ai/dsh-subprocess'
  4. import { TerminalRetention, type TerminalRetentionPolicy } from '../src/retention.ts'
  5. const owners: TerminalRetention[] = []
  6. beforeEach(() => { vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout', 'performance'] }) })
  7. afterEach(async () => { await Promise.allSettled(owners.splice(0).map(owner => owner.dispose())); vi.useRealTimers() })
  8. function fixture(policy: Partial<TerminalRetentionPolicy> = {}) {
  9. const inspect = vi.fn(async (): Promise<SubprocessTerminalActivity> => ({ state: 'idle', revision: 1 }))
  10. const terminate = vi.fn(async () => {})
  11. const failed = vi.fn()
  12. const owner = new TerminalRetention(
  13. { unattendedTimeoutMs: 100, activityPollIntervalMs: 10, cleanupRetryMs: 20, ...policy }, inspect, terminate, failed,
  14. )
  15. owners.push(owner)
  16. return { owner, inspect, terminate, failed }
  17. }
  18. it('starts the whole idle grace at the first positive observation and clears only at its deadline', async () => {
  19. const h = fixture()
  20. await vi.advanceTimersByTimeAsync(99)
  21. expect(h.terminate).not.toHaveBeenCalled()
  22. await vi.advanceTimersByTimeAsync(1)
  23. expect(h.terminate).toHaveBeenCalledOnce()
  24. await vi.advanceTimersByTimeAsync(1000)
  25. expect(h.terminate).toHaveBeenCalledOnce()
  26. })
  27. it('retains two independent windows and starts a fresh grace after the last disconnect', async () => {
  28. const h = fixture()
  29. const a = new AbortController()
  30. const b = new AbortController()
  31. const first = h.owner.retain(a.signal)[Symbol.asyncIterator]()
  32. const second = h.owner.retain(b.signal)[Symbol.asyncIterator]()
  33. expect(await first.next()).toMatchObject({ value: { type: 'retained' } })
  34. await second.next()
  35. a.abort()
  36. await vi.advanceTimersByTimeAsync(500)
  37. expect(h.inspect).not.toHaveBeenCalled()
  38. b.abort()
  39. await vi.advanceTimersByTimeAsync(99)
  40. expect(h.terminate).not.toHaveBeenCalled()
  41. await vi.advanceTimersByTimeAsync(1)
  42. expect(h.terminate).toHaveBeenCalledOnce()
  43. await first.return?.()
  44. await second.return?.()
  45. })
  46. it.each(['busy', 'unknown'] as const)('protects %s work for arbitrarily many deadlines, then gives it the full idle grace', async (state) => {
  47. const h = fixture()
  48. h.inspect.mockResolvedValue({ state, revision: 2 })
  49. await vi.advanceTimersByTimeAsync(1000)
  50. expect(h.terminate).not.toHaveBeenCalled()
  51. h.inspect.mockResolvedValue({ state: 'idle', revision: 3 })
  52. await vi.advanceTimersByTimeAsync(109)
  53. expect(h.terminate).not.toHaveBeenCalled()
  54. await vi.advanceTimersByTimeAsync(1)
  55. expect(h.terminate).toHaveBeenCalledOnce()
  56. })
  57. it('invalidates elapsed grace on failed inspection, input, and intervening shell revisions', async () => {
  58. const h = fixture()
  59. await vi.advanceTimersByTimeAsync(90)
  60. h.inspect.mockRejectedValueOnce(new Error('SSH unavailable'))
  61. await vi.advanceTimersByTimeAsync(90)
  62. h.owner.invalidate()
  63. await vi.advanceTimersByTimeAsync(90)
  64. h.inspect.mockResolvedValue({ state: 'idle', revision: 2 })
  65. await vi.advanceTimersByTimeAsync(100)
  66. expect(h.terminate).not.toHaveBeenCalled()
  67. await vi.advanceTimersByTimeAsync(10)
  68. expect(h.terminate).toHaveBeenCalledOnce()
  69. })
  70. it('discards an asynchronous observation invalidated by a reconnect and its subsequent disconnect', async () => {
  71. const h = fixture()
  72. await vi.advanceTimersByTimeAsync(90)
  73. const pending = Promise.withResolvers<SubprocessTerminalActivity>()
  74. h.inspect.mockReturnValueOnce(pending.promise)
  75. await vi.advanceTimersByTimeAsync(10)
  76. const abort = new AbortController()
  77. const stream = h.owner.retain(abort.signal)[Symbol.asyncIterator]()
  78. await stream.next()
  79. abort.abort()
  80. pending.resolve({ state: 'idle', revision: 1 })
  81. await vi.advanceTimersByTimeAsync(99)
  82. expect(h.terminate).not.toHaveBeenCalled()
  83. await vi.advanceTimersByTimeAsync(1)
  84. expect(h.terminate).toHaveBeenCalledOnce()
  85. await stream.return?.()
  86. })
  87. it('keeps failed cleanup closed to new holders and retries once without a new idle grace', async () => {
  88. const h = fixture()
  89. h.terminate.mockRejectedValueOnce(new Error('process still alive'))
  90. await vi.advanceTimersByTimeAsync(100)
  91. expect(h.failed).toHaveBeenCalledOnce()
  92. await expect(h.owner.retain(new AbortController().signal)[Symbol.asyncIterator]().next()).rejects.toThrow('unavailable')
  93. await vi.advanceTimersByTimeAsync(19)
  94. expect(h.terminate).toHaveBeenCalledOnce()
  95. await vi.advanceTimersByTimeAsync(1)
  96. expect(h.terminate).toHaveBeenCalledTimes(2)
  97. })
  98. it('joins explicit cleanup, closes held streams, and cancels timers on disposal', async () => {
  99. const h = fixture()
  100. const stream = h.owner.retain(new AbortController().signal)[Symbol.asyncIterator]()
  101. await stream.next()
  102. const ended = stream.next()
  103. const pending = Promise.withResolvers<undefined>()
  104. h.terminate.mockReturnValueOnce(pending.promise)
  105. const first = h.owner.close()
  106. expect(h.owner.close()).toBe(first)
  107. expect(await ended).toMatchObject({ done: true })
  108. pending.resolve(undefined)
  109. await h.owner.dispose()
  110. expect(vi.getTimerCount()).toBe(0)
  111. expect(h.terminate).toHaveBeenCalledOnce()
  112. })
  113. it('can disable automatic reclamation while retaining explicit cleanup and cancellation', async () => {
  114. const h = fixture({ unattendedTimeoutMs: 0 })
  115. await vi.advanceTimersByTimeAsync(1000)
  116. expect(h.inspect).not.toHaveBeenCalled()
  117. const signal = AbortSignal.abort(new Error('disconnected'))
  118. await expect(h.owner.retain(signal)[Symbol.asyncIterator]().next()).rejects.toThrow('disconnected')
  119. await h.owner.close()
  120. expect(h.terminate).toHaveBeenCalledOnce()
  121. })
  122. it('bounds native timers without shortening a configured cleanup retry', async () => {
  123. const h = fixture({ unattendedTimeoutMs: 0, cleanupRetryMs: 3_000_000_000 })
  124. h.terminate.mockRejectedValueOnce(new Error('retry'))
  125. await expect(h.owner.close()).rejects.toThrow('retry')
  126. await vi.advanceTimersByTimeAsync(2_999_999_999)
  127. expect(h.terminate).toHaveBeenCalledOnce()
  128. await vi.advanceTimersByTimeAsync(1)
  129. expect(h.terminate).toHaveBeenCalledTimes(2)
  130. })
  131. it('does not overlap activity queries when a disconnect schedules work during an outstanding observation', async () => {
  132. const h = fixture()
  133. const pending = Promise.withResolvers<SubprocessTerminalActivity>()
  134. h.inspect.mockReturnValueOnce(pending.promise)
  135. await vi.advanceTimersByTimeAsync(0)
  136. const abort = new AbortController()
  137. const stream = h.owner.retain(abort.signal)[Symbol.asyncIterator]()
  138. await stream.next()
  139. abort.abort()
  140. await vi.advanceTimersByTimeAsync(0)
  141. expect(h.inspect).toHaveBeenCalledOnce()
  142. pending.resolve({ state: 'idle', revision: 1 })
  143. await vi.advanceTimersByTimeAsync(10)
  144. expect(h.inspect).toHaveBeenCalledTimes(2)
  145. expect(h.terminate).not.toHaveBeenCalled()
  146. await stream.return?.()
  147. })
  148. it('waits for the in-flight observation before reporting a disposal cleanup failure', async () => {
  149. const h = fixture()
  150. const pending = Promise.withResolvers<SubprocessTerminalActivity>()
  151. h.inspect.mockReturnValueOnce(pending.promise)
  152. await vi.advanceTimersByTimeAsync(0)
  153. const error = new Error('cleanup failed')
  154. h.terminate.mockRejectedValueOnce(error)
  155. let settled = false
  156. const disposed = h.owner.dispose().catch((reason: unknown) => { settled = true; return reason })
  157. try {
  158. await vi.advanceTimersByTimeAsync(0)
  159. expect(settled).toBe(false)
  160. pending.resolve({ state: 'unknown', revision: 1 })
  161. expect(await disposed).toBe(error)
  162. } finally { pending.resolve({ state: 'unknown', revision: 1 }); await disposed }
  163. })
  164. it('waits for pending cleanup when the observation settles first during disposal', async () => {
  165. const h = fixture()
  166. const observed = Promise.withResolvers<SubprocessTerminalActivity>()
  167. const cleanup = Promise.withResolvers<undefined>()
  168. h.inspect.mockReturnValueOnce(observed.promise)
  169. h.terminate.mockReturnValueOnce(cleanup.promise)
  170. await vi.advanceTimersByTimeAsync(0)
  171. let settled = false
  172. const disposed = h.owner.dispose().then(() => { settled = true })
  173. try {
  174. observed.resolve({ state: 'unknown', revision: 1 })
  175. await vi.advanceTimersByTimeAsync(0)
  176. expect(settled).toBe(false)
  177. cleanup.resolve(undefined)
  178. await disposed
  179. expect(settled).toBe(true)
  180. } finally { observed.resolve({ state: 'unknown', revision: 1 }); cleanup.resolve(undefined); await disposed }
  181. })