effect-flock-worker.ts 1.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960
  1. import fs from "fs/promises"
  2. import os from "os"
  3. import { Effect } from "effect"
  4. import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
  5. import { EffectFlock } from "@opencode-ai/util/effect-flock"
  6. import { Global } from "@opencode-ai/util/global"
  7. type Msg = {
  8. key: string
  9. dir: string
  10. holdMs?: number
  11. ready?: string
  12. active?: string
  13. done?: string
  14. }
  15. function sleep(ms: number) {
  16. return new Promise<void>((resolve) => setTimeout(resolve, ms))
  17. }
  18. const msg: Msg = JSON.parse(process.argv[2])
  19. const testGlobal = Global.layerWith({
  20. home: os.homedir(),
  21. data: os.tmpdir(),
  22. cache: os.tmpdir(),
  23. config: os.tmpdir(),
  24. state: os.tmpdir(),
  25. bin: os.tmpdir(),
  26. log: os.tmpdir(),
  27. })
  28. const testLayer = AppNodeBuilder.build(EffectFlock.node, [[Global.node, testGlobal]])
  29. async function job() {
  30. if (msg.ready) await fs.writeFile(msg.ready, String(process.pid))
  31. if (msg.active) await fs.writeFile(msg.active, String(process.pid), { flag: "wx" })
  32. try {
  33. if (msg.holdMs && msg.holdMs > 0) await sleep(msg.holdMs)
  34. if (msg.done) await fs.appendFile(msg.done, "1\n")
  35. } finally {
  36. if (msg.active) await fs.rm(msg.active, { force: true })
  37. }
  38. }
  39. await Effect.runPromise(
  40. Effect.gen(function* () {
  41. const flock = yield* EffectFlock.Service
  42. yield* flock.withLock(
  43. Effect.promise(() => job()),
  44. msg.key,
  45. msg.dir,
  46. )
  47. }).pipe(Effect.provide(testLayer)),
  48. ).catch((err) => {
  49. const text = err instanceof Error ? (err.stack ?? err.message) : String(err)
  50. process.stderr.write(text)
  51. process.exit(1)
  52. })