flock.ts 9.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358
  1. import path from "path"
  2. import os from "os"
  3. import { randomBytes, randomUUID } from "crypto"
  4. import { mkdir, readFile, rm, stat, utimes, writeFile } from "fs/promises"
  5. import { Hash } from "./hash"
  6. import { Effect } from "effect"
  7. export type FlockGlobal = {
  8. state: string
  9. }
  10. export namespace Flock {
  11. let global: FlockGlobal | undefined
  12. export function setGlobal(g: FlockGlobal) {
  13. global = g
  14. }
  15. const root = () => {
  16. if (!global) throw new Error("Flock global not set")
  17. return path.join(global.state, "locks")
  18. }
  19. // Defaults for callers that do not provide timing options.
  20. const defaultOpts = {
  21. staleMs: 60_000,
  22. timeoutMs: 5 * 60_000,
  23. baseDelayMs: 100,
  24. maxDelayMs: 2_000,
  25. }
  26. export interface WaitEvent {
  27. key: string
  28. attempt: number
  29. delay: number
  30. waited: number
  31. }
  32. export type Wait = (input: WaitEvent) => void | Promise<void>
  33. export interface Options {
  34. dir?: string
  35. signal?: AbortSignal
  36. staleMs?: number
  37. timeoutMs?: number
  38. baseDelayMs?: number
  39. maxDelayMs?: number
  40. onWait?: Wait
  41. }
  42. type Opts = {
  43. staleMs: number
  44. timeoutMs: number
  45. baseDelayMs: number
  46. maxDelayMs: number
  47. }
  48. type Owned = {
  49. acquired: true
  50. startHeartbeat: (intervalMs?: number) => void
  51. release: () => Promise<void>
  52. }
  53. export interface Lease {
  54. release: () => Promise<void>
  55. [Symbol.asyncDispose]: () => Promise<void>
  56. }
  57. function code(err: unknown) {
  58. if (typeof err !== "object" || err === null || !("code" in err)) return
  59. const value = err.code
  60. if (typeof value !== "string") return
  61. return value
  62. }
  63. function sleep(ms: number, signal?: AbortSignal) {
  64. return new Promise<void>((resolve, reject) => {
  65. if (signal?.aborted) {
  66. reject(signal.reason ?? new Error("Aborted"))
  67. return
  68. }
  69. let timer: NodeJS.Timeout | undefined
  70. const done = () => {
  71. signal?.removeEventListener("abort", abort)
  72. resolve()
  73. }
  74. const abort = () => {
  75. if (timer) {
  76. clearTimeout(timer)
  77. }
  78. signal?.removeEventListener("abort", abort)
  79. reject(signal?.reason ?? new Error("Aborted"))
  80. }
  81. signal?.addEventListener("abort", abort, { once: true })
  82. timer = setTimeout(done, ms)
  83. })
  84. }
  85. function jitter(ms: number) {
  86. const j = Math.floor(ms * 0.3)
  87. const d = Math.floor(Math.random() * (2 * j + 1)) - j
  88. return Math.max(0, ms + d)
  89. }
  90. function mono() {
  91. return performance.now()
  92. }
  93. function wall() {
  94. return performance.timeOrigin + mono()
  95. }
  96. async function stats(file: string) {
  97. try {
  98. return await stat(file)
  99. } catch (err) {
  100. const errCode = code(err)
  101. if (errCode === "ENOENT" || errCode === "ENOTDIR") return
  102. throw err
  103. }
  104. }
  105. async function stale(lockDir: string, heartbeatPath: string, metaPath: string, staleMs: number) {
  106. // Stale detection allows automatic recovery after crashed owners.
  107. const now = wall()
  108. const heartbeat = await stats(heartbeatPath)
  109. if (heartbeat) {
  110. return now - heartbeat.mtimeMs > staleMs
  111. }
  112. const meta = await stats(metaPath)
  113. if (meta) {
  114. return now - meta.mtimeMs > staleMs
  115. }
  116. const dir = await stats(lockDir)
  117. if (!dir) {
  118. return false
  119. }
  120. return now - dir.mtimeMs > staleMs
  121. }
  122. async function tryAcquireLockDir(lockDir: string, opts: Opts): Promise<Owned | { acquired: false }> {
  123. const token = randomUUID?.() ?? randomBytes(16).toString("hex")
  124. const metaPath = path.join(lockDir, "meta.json")
  125. const heartbeatPath = path.join(lockDir, "heartbeat")
  126. try {
  127. await mkdir(lockDir, { mode: 0o700 })
  128. } catch (err) {
  129. if (code(err) !== "EEXIST") {
  130. throw err
  131. }
  132. if (!(await stale(lockDir, heartbeatPath, metaPath, opts.staleMs))) {
  133. return { acquired: false }
  134. }
  135. const breakerPath = lockDir + ".breaker"
  136. try {
  137. await mkdir(breakerPath, { mode: 0o700 })
  138. } catch (claimErr) {
  139. const errCode = code(claimErr)
  140. if (errCode === "EEXIST") {
  141. const breaker = await stats(breakerPath)
  142. if (breaker && wall() - breaker.mtimeMs > opts.staleMs) {
  143. await rm(breakerPath, { recursive: true, force: true }).catch(() => undefined)
  144. }
  145. return { acquired: false }
  146. }
  147. if (errCode === "ENOENT" || errCode === "ENOTDIR") {
  148. return { acquired: false }
  149. }
  150. throw claimErr
  151. }
  152. try {
  153. // Breaker ownership ensures only one contender performs stale cleanup.
  154. if (!(await stale(lockDir, heartbeatPath, metaPath, opts.staleMs))) {
  155. return { acquired: false }
  156. }
  157. await rm(lockDir, { recursive: true, force: true })
  158. try {
  159. await mkdir(lockDir, { mode: 0o700 })
  160. } catch (retryErr) {
  161. const errCode = code(retryErr)
  162. if (errCode === "EEXIST" || errCode === "ENOTEMPTY") {
  163. return { acquired: false }
  164. }
  165. throw retryErr
  166. }
  167. } finally {
  168. await rm(breakerPath, { recursive: true, force: true }).catch(() => undefined)
  169. }
  170. }
  171. const meta = {
  172. token,
  173. pid: process.pid,
  174. hostname: os.hostname(),
  175. createdAt: new Date().toISOString(),
  176. }
  177. await writeFile(heartbeatPath, "", { flag: "wx" }).catch(async () => {
  178. await rm(lockDir, { recursive: true, force: true })
  179. throw new Error("Lock acquired but heartbeat already existed (possible compromise).")
  180. })
  181. await writeFile(metaPath, JSON.stringify(meta, null, 2), { flag: "wx" }).catch(async () => {
  182. await rm(lockDir, { recursive: true, force: true })
  183. throw new Error("Lock acquired but meta.json already existed (possible compromise).")
  184. })
  185. let timer: NodeJS.Timeout | undefined
  186. const startHeartbeat = (intervalMs = Math.max(100, Math.floor(opts.staleMs / 3))) => {
  187. if (timer) return
  188. // Heartbeat prevents long critical sections from being evicted as stale.
  189. timer = setInterval(() => {
  190. const t = new Date()
  191. void utimes(heartbeatPath, t, t).catch(() => undefined)
  192. }, intervalMs)
  193. timer.unref?.()
  194. }
  195. const release = async () => {
  196. if (timer) {
  197. clearInterval(timer)
  198. timer = undefined
  199. }
  200. const current = await readFile(metaPath, "utf8")
  201. .then((raw) => {
  202. const parsed = JSON.parse(raw)
  203. if (!parsed || typeof parsed !== "object") return {}
  204. return {
  205. token: "token" in parsed && typeof parsed.token === "string" ? parsed.token : undefined,
  206. }
  207. })
  208. .catch((err) => {
  209. const errCode = code(err)
  210. if (errCode === "ENOENT" || errCode === "ENOTDIR") {
  211. throw new Error("Refusing to release: lock is compromised (metadata missing).")
  212. }
  213. if (err instanceof SyntaxError) {
  214. throw new Error("Refusing to release: lock is compromised (metadata invalid).")
  215. }
  216. throw err
  217. })
  218. // Token check prevents deleting a lock that was re-acquired by another process.
  219. if (current.token !== token) {
  220. throw new Error("Refusing to release: lock token mismatch (not the owner).")
  221. }
  222. await rm(lockDir, { recursive: true, force: true })
  223. }
  224. return {
  225. acquired: true,
  226. startHeartbeat,
  227. release,
  228. }
  229. }
  230. async function acquireLockDir(
  231. lockDir: string,
  232. input: { key: string; onWait?: Wait; signal?: AbortSignal },
  233. opts: Opts,
  234. ) {
  235. const stop = mono() + opts.timeoutMs
  236. let attempt = 0
  237. let waited = 0
  238. let delay = opts.baseDelayMs
  239. while (true) {
  240. input.signal?.throwIfAborted()
  241. const res = await tryAcquireLockDir(lockDir, opts)
  242. if (res.acquired) {
  243. return res
  244. }
  245. if (mono() > stop) {
  246. throw new Error(`Timed out waiting for lock: ${input.key}`)
  247. }
  248. attempt += 1
  249. const ms = jitter(delay)
  250. await input.onWait?.({
  251. key: input.key,
  252. attempt,
  253. delay: ms,
  254. waited,
  255. })
  256. await sleep(ms, input.signal)
  257. waited += ms
  258. delay = Math.min(opts.maxDelayMs, Math.floor(delay * 1.7))
  259. }
  260. }
  261. export async function acquire(key: string, input: Options = {}): Promise<Lease> {
  262. input.signal?.throwIfAborted()
  263. const cfg: Opts = {
  264. staleMs: input.staleMs ?? defaultOpts.staleMs,
  265. timeoutMs: input.timeoutMs ?? defaultOpts.timeoutMs,
  266. baseDelayMs: input.baseDelayMs ?? defaultOpts.baseDelayMs,
  267. maxDelayMs: input.maxDelayMs ?? defaultOpts.maxDelayMs,
  268. }
  269. const dir = input.dir ?? root()
  270. await mkdir(dir, { recursive: true })
  271. const lockfile = path.join(dir, Hash.fast(key) + ".lock")
  272. const lock = await acquireLockDir(
  273. lockfile,
  274. {
  275. key,
  276. onWait: input.onWait,
  277. signal: input.signal,
  278. },
  279. cfg,
  280. )
  281. lock.startHeartbeat()
  282. const release = () => lock.release()
  283. return {
  284. release,
  285. [Symbol.asyncDispose]() {
  286. return release()
  287. },
  288. }
  289. }
  290. export async function withLock<T>(key: string, fn: () => Promise<T>, input: Options = {}) {
  291. await using _ = await acquire(key, input)
  292. input.signal?.throwIfAborted()
  293. return await fn()
  294. }
  295. export const effect = Effect.fn("Flock.effect")(function* (key: string, input: Options = {}) {
  296. return yield* Effect.acquireRelease(
  297. Effect.promise((signal) => Flock.acquire(key, { ...input, signal })).pipe(
  298. Effect.withSpan("Flock.acquire", {
  299. attributes: { key },
  300. }),
  301. ),
  302. (lock) => Effect.promise(() => lock.release()).pipe(Effect.withSpan("Flock.release")),
  303. ).pipe(Effect.asVoid)
  304. })
  305. }