lease.spec.ts 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434
  1. /**
  2. * Cross-process write-lock behavior, exercised through fresh backend
  3. * instances over one shared root: kernel `flock` locks conflict between two
  4. * descriptors even inside one process, so a second instance behaves exactly
  5. * like a second process. Exclusion while a holder is live, immediate
  6. * admission after close, lock-file residue rules, and the inode verification
  7. * that defeats an unlinked-and-recreated lock path. Filesystem and flock
  8. * refusals are injected through the module mocks below: POSIX modes cannot
  9. * express them on Windows, and an injected error is the only deterministic
  10. * cross-platform refusal. Real cross-process exclusion and crash release are
  11. * pinned by lease.two-process.e2e.ts.
  12. */
  13. import { existsSync } from 'node:fs'
  14. import { mkdtemp, readdir, readFile, rm, writeFile } from 'node:fs/promises'
  15. import { tmpdir } from 'node:os'
  16. import { join } from 'node:path'
  17. import { afterEach, describe, expect, it, vi } from 'vitest'
  18. import { Context } from '@deepseek-ai/cordis'
  19. import { SESSION_FORMAT_VERSION, SessionId, SessionSeq } from '@deepseek-ai/dsh-session'
  20. import type { SessionHeader } from '@deepseek-ai/dsh-session'
  21. import {
  22. SessionAlreadyExistsError,
  23. SessionAlreadyOwnedError,
  24. SessionPersistenceNotFoundError,
  25. } from '@deepseek-ai/dsh-session-persistence'
  26. import type { SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
  27. import JsonlSessionPersistence from '../src/index.ts'
  28. import { LEASE_FILENAME, SessionWriteLease } from '../src/lease.ts'
  29. import type { JsonlSessionHandle } from '../src/storage.ts'
  30. import { sessionDir } from '../src/format.ts'
  31. // The lock's base name, duplicated for the hoisted mock factories: they run
  32. // while `../src/lease.ts` is still evaluating, before LEASE_FILENAME exists.
  33. const LOCK = vi.hoisted(() => 'session.lock')
  34. const refuse = vi.hoisted(() => ({
  35. /** Next open of a lock file fails EACCES (read-only directory). */
  36. lockOpen: false,
  37. /** Next flock call fails EACCES (a non-contention kernel refusal). */
  38. flock: false,
  39. /** Next flock call fails EWOULDBLOCK. */
  40. flockBusy: false,
  41. /** Next stat of a lock file fails EACCES (unreadable path). */
  42. lockStat: false,
  43. /** For N further lock-path stats: unlink and recreate the file first, so the locked inode is orphaned. */
  44. swapLockOnStat: 0,
  45. /** Next lock-path stat: unlink the file first, so the verify read finds nothing. */
  46. dropLockOnStat: false,
  47. }))
  48. vi.mock('node:fs/promises', async (importOriginal) => {
  49. const actual = await importOriginal<typeof import('node:fs/promises')>()
  50. const denied = (syscall: string): never => {
  51. throw Object.assign(new Error(`EACCES: injected ${syscall} refusal`), { code: 'EACCES' })
  52. }
  53. return {
  54. ...actual,
  55. open: (async (path: unknown, ...rest: never[]) => {
  56. if (refuse.lockOpen && String(path).endsWith(LOCK)) {
  57. refuse.lockOpen = false
  58. denied('open')
  59. }
  60. return (actual.open as (path: unknown, ...args: never[]) => Promise<unknown>)(path, ...rest)
  61. }) as typeof actual.open,
  62. stat: (async (path: unknown, ...rest: never[]) => {
  63. const at = String(path)
  64. if (at.endsWith(LOCK)) {
  65. if (refuse.lockStat) {
  66. refuse.lockStat = false
  67. denied('stat')
  68. }
  69. if (refuse.dropLockOnStat) {
  70. refuse.dropLockOnStat = false
  71. await actual.unlink(at)
  72. } else if (refuse.swapLockOnStat > 0) {
  73. refuse.swapLockOnStat -= 1
  74. await actual.unlink(at)
  75. await actual.writeFile(at, '')
  76. }
  77. }
  78. return (actual.stat as (path: unknown, ...args: never[]) => Promise<unknown>)(path, ...rest)
  79. }) as typeof actual.stat,
  80. }
  81. })
  82. vi.mock('@deepseek-ai/node-addon-system/flock', async (importOriginal) => {
  83. const actual = await importOriginal<typeof import('@deepseek-ai/node-addon-system/flock')>()
  84. return {
  85. tryLockExclusive: async (fd: number): Promise<void> => {
  86. if (refuse.flock) {
  87. refuse.flock = false
  88. throw Object.assign(new Error('EACCES: injected flock refusal'), { code: 'EACCES' })
  89. }
  90. if (refuse.flockBusy) {
  91. refuse.flockBusy = false
  92. throw Object.assign(new Error('EWOULDBLOCK: injected contention'), { code: 'EWOULDBLOCK' })
  93. }
  94. return actual.tryLockExclusive(fd)
  95. },
  96. }
  97. })
  98. const dirs: string[] = []
  99. const contexts: Context[] = []
  100. afterEach(async () => {
  101. refuse.lockOpen = false
  102. refuse.flock = false
  103. refuse.flockBusy = false
  104. refuse.lockStat = false
  105. refuse.swapLockOnStat = 0
  106. refuse.dropLockOnStat = false
  107. for (const ctx of contexts.splice(0)) await ctx.fiber.dispose()
  108. for (const dir of dirs.splice(0)) await rm(dir, { recursive: true, force: true })
  109. })
  110. function meta(id: string, cwd = '/work'): SessionHeader {
  111. return { version: SESSION_FORMAT_VERSION, id: SessionId(id), createdAt: 1_000, cwd, isSeeded: false }
  112. }
  113. async function freshRoot(): Promise<string> {
  114. const root = await mkdtemp(join(tmpdir(), 'dsh-jsonl-lease-'))
  115. dirs.push(root)
  116. return root
  117. }
  118. async function mount(root: string): Promise<SessionPersistence> {
  119. const ctx = new Context()
  120. contexts.push(ctx)
  121. await ctx.plugin(JsonlSessionPersistence, { root, compression: 'none' })
  122. return ctx.sessionPersistence
  123. }
  124. function lockPath(root: string, id: string, cwd = '/work'): string {
  125. return join(sessionDir(root, cwd, SessionId(id)), LEASE_FILENAME)
  126. }
  127. /**
  128. * Make the next lock release do its real work, then report failure — as a
  129. * close(2) that freed the descriptor but returned EIO would.
  130. */
  131. function failReleaseOnce(): void {
  132. const spy = vi.spyOn(SessionWriteLease.prototype, 'release')
  133. spy.mockImplementationOnce(async function (this: SessionWriteLease) {
  134. spy.mockRestore()
  135. await this.release()
  136. throw Object.assign(new Error('EIO: injected release failure'), { code: 'EIO' })
  137. })
  138. }
  139. const EVENTS = [
  140. { type: 'turn/start', seq: SessionSeq(0), time: 1, data: { turn: 1 } },
  141. { type: 'turn/end', seq: SessionSeq(1), time: 2, data: { turn: 1, reason: { kind: 'completed' } } },
  142. ] as const
  143. describe('cross-process write lock', () => {
  144. it('excludes a second instance while the holder is live, and admits it after close', async () => {
  145. const root = await freshRoot()
  146. const first = await mount(root)
  147. const second = await mount(root)
  148. const holder = await first.create(meta('excluded'))
  149. await holder.append([...EVENTS])
  150. // Another instance over the same root cannot create or write-open the id:
  151. // the materialized duplicate is an existence fact, the write open an
  152. // ownership one.
  153. await expect(second.create(meta('excluded'))).rejects.toBeInstanceOf(SessionAlreadyExistsError)
  154. await expect(second.open(SessionId('excluded'), 'write')).rejects.toBeInstanceOf(SessionAlreadyOwnedError)
  155. // Unmaterialized creates hold no lock and leave no artifact, so a rival
  156. // instance's create succeeds; the collision surfaces at the loser's first
  157. // materializing write, where the winner already holds the lock.
  158. const pendingWinner = await first.create(meta('excluded-pending'))
  159. const pendingLoser = await second.create(meta('excluded-pending'))
  160. await pendingWinner.append([...EVENTS])
  161. await expect(pendingLoser.append([...EVENTS])).rejects.toBeInstanceOf(SessionAlreadyOwnedError)
  162. await pendingLoser.close()
  163. await pendingWinner.close()
  164. // Reads never touch the lock.
  165. const reader = await second.open(SessionId('excluded'), 'read')
  166. expect((await reader.read()).events.map(event => event.seq)).toEqual([0, 1])
  167. await reader.close()
  168. await holder.close()
  169. // POSIX keeps the materialized session's lock file (Windows locks a kernel
  170. // object with no filesystem footprint); the kernel lock itself is gone.
  171. if (process.platform !== 'win32') expect(existsSync(lockPath(root, 'excluded'))).toBe(true)
  172. const reopened = await second.open(SessionId('excluded'), 'write')
  173. await reopened.append([{ type: 'turn/start', seq: SessionSeq(2), time: 3, data: { turn: 2 } }])
  174. await reopened.close()
  175. })
  176. it.skipIf(process.platform === 'win32')('removing the lock file forfeits a wedged holder: a fresh inode admits a successor', async () => {
  177. const root = await freshRoot()
  178. const first = await mount(root)
  179. const second = await mount(root)
  180. const wedged = await first.create(meta('wedged'))
  181. await wedged.append([...EVENTS])
  182. // The documented escape hatch for a live-but-stuck holder: deleting the
  183. // lock file orphans the held inode, and a successor locks the fresh one.
  184. await rm(lockPath(root, 'wedged'))
  185. const successor = await second.open(SessionId('wedged'), 'write')
  186. await successor.append([{ type: 'turn/start', seq: SessionSeq(2), time: 3, data: { turn: 2 } }])
  187. await successor.close()
  188. await wedged.close()
  189. })
  190. it('write-opening an absent session leaves no lock residue', async () => {
  191. const root = await freshRoot()
  192. const backend = await mount(root)
  193. await expect(backend.open(SessionId('absent'), 'write')).rejects.toBeInstanceOf(SessionPersistenceNotFoundError)
  194. expect(existsSync(join(root, LEASE_FILENAME))).toBe(false)
  195. })
  196. it('write-opening an absent id under an existing project directory reports not-found', async () => {
  197. const root = await freshRoot()
  198. const backend = await mount(root)
  199. const writer = await backend.create(meta('present-sibling'))
  200. await writer.append([...EVENTS])
  201. await writer.close()
  202. // The project directory exists but the id's session directory does not:
  203. // the generation scan reports absence rather than misreading a sibling.
  204. await expect(backend.open(SessionId('absent-sibling'), 'write')).rejects.toBeInstanceOf(SessionPersistenceNotFoundError)
  205. })
  206. it('a never-materialized create leaves no filesystem footprint at all', async () => {
  207. const root = await freshRoot()
  208. const backend = await mount(root)
  209. const handle = await backend.create(meta('erased'))
  210. // The lock is taken only at the first materializing write, so an
  211. // unmaterialized session creates neither its directory nor a lock file.
  212. expect(existsSync(join(lockPath(root, 'erased'), '..'))).toBe(false)
  213. await handle.close()
  214. expect(existsSync(join(lockPath(root, 'erased'), '..'))).toBe(false)
  215. await expect(backend.stat(SessionId('erased'))).resolves.toBeUndefined()
  216. })
  217. it('materialization publishes the lock before the first log bytes and keeps it on the handle', async () => {
  218. const root = await freshRoot()
  219. const first = await mount(root)
  220. const second = await mount(root)
  221. const creator = await first.create(meta('lazy-lock'))
  222. await creator.append([...EVENTS])
  223. // The materializing append acquired and retained the lock.
  224. if (process.platform !== 'win32') expect(existsSync(lockPath(root, 'lazy-lock'))).toBe(true)
  225. await expect(second.open(SessionId('lazy-lock'), 'write')).rejects.toBeInstanceOf(SessionAlreadyOwnedError)
  226. // A later append reuses the held lock rather than re-acquiring.
  227. await creator.append([{ type: 'turn/start', seq: SessionSeq(2), time: 3, data: { turn: 2 } }])
  228. await creator.close()
  229. const reopened = await second.open(SessionId('lazy-lock'), 'write')
  230. await reopened.close()
  231. })
  232. it('an explicitly flushed empty session takes the lock with its header', async () => {
  233. const root = await freshRoot()
  234. const first = await mount(root)
  235. const second = await mount(root)
  236. const creator = await first.create(meta('flush-lock'))
  237. await creator.flush()
  238. await expect(second.open(SessionId('flush-lock'), 'write')).rejects.toBeInstanceOf(SessionAlreadyOwnedError)
  239. await creator.close()
  240. })
  241. it.skipIf(process.platform === 'win32')('surfaces a filesystem refusal opening the lock file', async () => {
  242. const root = await freshRoot()
  243. const backend = await mount(root)
  244. const writer = await backend.create(meta('open-blocked'))
  245. await writer.append([...EVENTS])
  246. await writer.close()
  247. refuse.lockOpen = true
  248. await expect(backend.open(SessionId('open-blocked'), 'write')).rejects.toThrow(/EACCES/)
  249. })
  250. it.skipIf(process.platform === 'win32')('surfaces a non-contention flock failure', async () => {
  251. const root = await freshRoot()
  252. const backend = await mount(root)
  253. const writer = await backend.create(meta('flock-blocked'))
  254. await writer.append([...EVENTS])
  255. await writer.close()
  256. refuse.flock = true
  257. await expect(backend.open(SessionId('flock-blocked'), 'write')).rejects.toThrow(/EACCES/)
  258. })
  259. it.skipIf(process.platform === 'win32')('maps the EWOULDBLOCK contention spelling to already-owned', async () => {
  260. const root = await freshRoot()
  261. const backend = await mount(root)
  262. const writer = await backend.create(meta('win-contended'))
  263. await writer.append([...EVENTS])
  264. await writer.close()
  265. // Some libcs spell flock(2) contention EWOULDBLOCK rather than EAGAIN.
  266. refuse.flockBusy = true
  267. await expect(backend.open(SessionId('win-contended'), 'write')).rejects.toBeInstanceOf(SessionAlreadyOwnedError)
  268. })
  269. it.skipIf(process.platform === 'win32')('surfaces a lock-path stat refusal from the inode verification', async () => {
  270. const root = await freshRoot()
  271. const backend = await mount(root)
  272. const writer = await backend.create(meta('stat-blocked'))
  273. await writer.append([...EVENTS])
  274. await writer.close()
  275. refuse.lockStat = true
  276. await expect(backend.open(SessionId('stat-blocked'), 'write')).rejects.toThrow(/EACCES/)
  277. })
  278. it.skipIf(process.platform === 'win32')('retries when the locked inode is no longer the lock path, and wins on a stable pass', async () => {
  279. const root = await freshRoot()
  280. const backend = await mount(root)
  281. const writer = await backend.create(meta('churned'))
  282. await writer.append([...EVENTS])
  283. await writer.close()
  284. // One churn (unlink+recreate under the verify stat) orphans the first
  285. // locked inode; the retry locks the fresh file and verifies clean.
  286. refuse.swapLockOnStat = 1
  287. const reopened = await backend.open(SessionId('churned'), 'write')
  288. await reopened.append([{ type: 'turn/start', seq: SessionSeq(2), time: 3, data: { turn: 2 } }])
  289. await reopened.close()
  290. })
  291. it.skipIf(process.platform === 'win32')('retries when the lock path vanishes under the verify read', async () => {
  292. const root = await freshRoot()
  293. const backend = await mount(root)
  294. const writer = await backend.create(meta('vanished'))
  295. await writer.append([...EVENTS])
  296. await writer.close()
  297. refuse.dropLockOnStat = true
  298. const reopened = await backend.open(SessionId('vanished'), 'write')
  299. await reopened.close()
  300. })
  301. it.skipIf(process.platform === 'win32')('gives up as already-owned when the lock path never stabilizes', async () => {
  302. const root = await freshRoot()
  303. const backend = await mount(root)
  304. const writer = await backend.create(meta('unstable'))
  305. await writer.append([...EVENTS])
  306. await writer.close()
  307. // Churn on every attempt: the bounded retry refuses rather than spinning.
  308. refuse.swapLockOnStat = 3
  309. await expect(backend.open(SessionId('unstable'), 'write')).rejects.toBeInstanceOf(SessionAlreadyOwnedError)
  310. })
  311. it('a failing lock release still frees the in-process claim on close', async () => {
  312. const root = await freshRoot()
  313. const backend = await mount(root)
  314. const holder = await backend.create(meta('release-fails'))
  315. await holder.append([...EVENTS])
  316. failReleaseOnce()
  317. await expect(holder.close()).rejects.toThrow(/injected release failure/)
  318. // The claim is freed despite the failed release: the id is not wedged.
  319. const reopened = await backend.open(SessionId('release-fails'), 'write')
  320. await reopened.append([{ type: 'turn/start', seq: SessionSeq(2), time: 3, data: { turn: 2 } }])
  321. await reopened.close()
  322. })
  323. it('a write-open failure with a failing release aggregates both and frees the claim', async () => {
  324. const root = await freshRoot()
  325. const backend = await mount(root)
  326. const writer = await backend.create(meta('open-and-release-fail'))
  327. await writer.append([...EVENTS])
  328. await writer.close()
  329. // Corrupt the stored header line so the open fails after the lock is
  330. // acquired (a garbled tail would be recovered as torn, not refused).
  331. const dir = join(lockPath(root, 'open-and-release-fail'), '..')
  332. const log = (await readdir(dir)).find(name => name.endsWith('.jsonl'))
  333. const stored = await readFile(join(dir, String(log)), 'utf8')
  334. await writeFile(join(dir, String(log)), `#${stored.slice(1)}`)
  335. failReleaseOnce()
  336. const outcome = await backend.open(SessionId('open-and-release-fail'), 'write').then(() => undefined, (error: unknown) => error)
  337. expect(outcome).toBeInstanceOf(AggregateError)
  338. const errors = (outcome as AggregateError).errors as Error[]
  339. expect(errors).toHaveLength(2)
  340. expect(String(errors[0])).toMatch(/corrupt/i)
  341. expect(String(errors[1])).toMatch(/injected release failure/)
  342. // The original diagnostic survives, and the claim is freed: the next
  343. // attempt reports the corruption again rather than a phantom owner.
  344. await expect(backend.open(SessionId('open-and-release-fail'), 'write')).rejects.toThrow(/corrupt/i)
  345. })
  346. it('a drain failure and a release failure reject close as one AggregateError', async () => {
  347. const root = await freshRoot()
  348. const backend = await mount(root)
  349. const holder = await backend.create(meta('drain-and-release-fail')) as unknown as JsonlSessionHandle
  350. await holder.append([...EVENTS])
  351. vi.spyOn(backend as unknown as { persistBatch: () => Promise<void> }, 'persistBatch')
  352. .mockRejectedValueOnce(new Error('injected drain refusal'))
  353. holder.enqueueLive({ type: 'turn/start', seq: SessionSeq(2), time: 3, data: { turn: 2 } }, () => {})
  354. failReleaseOnce()
  355. const outcome = await holder.close().then(() => undefined, (error: unknown) => error)
  356. expect(outcome).toBeInstanceOf(AggregateError)
  357. expect((outcome as AggregateError).errors.map(String).join('\n')).toMatch(/drain refusal[\s\S]*release failure/)
  358. // Both failures reported, and the id is still not wedged.
  359. const reopened = await backend.open(SessionId('drain-and-release-fail'), 'write')
  360. await reopened.close()
  361. })
  362. it.skipIf(process.platform === 'win32')('release is idempotent and never removes the lock file', async () => {
  363. const root = await freshRoot()
  364. const dir = join(root, 'solo')
  365. const lease = await SessionWriteLease.acquire(dir, SessionId('solo'))
  366. await lease.release()
  367. await lease.release()
  368. // The file survives every release, keeping the stable inode later
  369. // lockers verify against; the kernel lock died with the descriptor.
  370. expect(existsSync(join(dir, LOCK))).toBe(true)
  371. const successor = await SessionWriteLease.acquire(dir, SessionId('solo'))
  372. await successor.release()
  373. expect(existsSync(join(dir, LOCK))).toBe(true)
  374. })
  375. it('keeps distinct sessions independently lockable', async () => {
  376. const root = await freshRoot()
  377. const backend = await mount(root)
  378. const a = await backend.create(meta('indep-a'))
  379. const b = await backend.create(meta('indep-b'))
  380. await a.append([...EVENTS])
  381. await b.append([...EVENTS])
  382. if (process.platform !== 'win32') {
  383. expect((await readdir(join(lockPath(root, 'indep-a'), '..'))).filter(name => name === LOCK)).toHaveLength(1)
  384. }
  385. await a.close()
  386. await b.close()
  387. })
  388. })