job.test.ts 6.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196
  1. import { describe, expect } from "bun:test"
  2. import { Deferred, Effect } from "effect"
  3. import { BackgroundJob } from "@/background/job"
  4. import { testEffect } from "../lib/effect"
  5. const it = testEffect(BackgroundJob.defaultLayer)
  6. describe("background.job", () => {
  7. it.instance("tracks started jobs through completion", () =>
  8. Effect.gen(function* () {
  9. const jobs = yield* BackgroundJob.Service
  10. const latch = yield* Deferred.make<void>()
  11. const job = yield* jobs.start({
  12. type: "test",
  13. title: "test job",
  14. run: Deferred.await(latch).pipe(Effect.as("done")),
  15. })
  16. expect(job.id.startsWith("job_")).toBe(true)
  17. expect(job.status).toBe("running")
  18. expect(job.title).toBe("test job")
  19. yield* Deferred.succeed(latch, undefined)
  20. const done = yield* jobs.wait({ id: job.id })
  21. expect(done.timedOut).toBe(false)
  22. expect(done.info?.status).toBe("completed")
  23. expect(done.info?.output).toBe("done")
  24. expect((yield* jobs.list()).map((item) => item.id)).toEqual([job.id])
  25. }),
  26. )
  27. it.instance("returns a running snapshot when wait times out", () =>
  28. Effect.gen(function* () {
  29. const jobs = yield* BackgroundJob.Service
  30. const job = yield* jobs.start({
  31. type: "test",
  32. run: Effect.never,
  33. })
  34. const result = yield* jobs.wait({ id: job.id, timeout: 1 })
  35. expect(result.timedOut).toBe(true)
  36. expect(result.info?.status).toBe("running")
  37. }),
  38. )
  39. it.instance("deduplicates concurrent starts for a running id", () =>
  40. Effect.gen(function* () {
  41. const jobs = yield* BackgroundJob.Service
  42. const started = yield* Deferred.make<void>()
  43. const id = "job_test"
  44. const [first, second] = yield* Effect.all(
  45. [
  46. jobs.start({
  47. id,
  48. type: "test",
  49. run: Deferred.succeed(started, undefined).pipe(Effect.andThen(Effect.never)),
  50. }),
  51. jobs.start({
  52. id,
  53. type: "test",
  54. run: Effect.fail(new Error("duplicate started")),
  55. }),
  56. ],
  57. { concurrency: "unbounded" },
  58. )
  59. yield* Deferred.await(started)
  60. expect(first.id).toBe(id)
  61. expect(second.id).toBe(id)
  62. expect(first.status).toBe("running")
  63. expect(second.status).toBe("running")
  64. expect((yield* jobs.list()).map((item) => item.id)).toEqual([id])
  65. yield* jobs.cancel(id)
  66. }),
  67. )
  68. it.instance("waits for extensions before completing a running job", () =>
  69. Effect.gen(function* () {
  70. const jobs = yield* BackgroundJob.Service
  71. const first = yield* Deferred.make<void>()
  72. const second = yield* Deferred.make<void>()
  73. const job = yield* jobs.start({
  74. type: "test",
  75. run: Deferred.await(first).pipe(Effect.as("first")),
  76. })
  77. expect(yield* jobs.extend({ id: job.id, run: Deferred.await(second).pipe(Effect.as("second")) })).toBe(true)
  78. yield* Deferred.succeed(first, undefined)
  79. expect((yield* jobs.get(job.id))?.status).toBe("running")
  80. yield* Deferred.succeed(second, undefined)
  81. const done = yield* jobs.wait({ id: job.id })
  82. expect(done.info?.status).toBe("completed")
  83. expect(done.info?.output).toBe("second")
  84. }),
  85. )
  86. it.instance("rejects extensions after a job completes", () =>
  87. Effect.gen(function* () {
  88. const jobs = yield* BackgroundJob.Service
  89. const job = yield* jobs.start({ type: "test", run: Effect.succeed("done") })
  90. yield* jobs.wait({ id: job.id })
  91. expect(yield* jobs.extend({ id: job.id, run: Effect.succeed("late") })).toBe(false)
  92. expect((yield* jobs.get(job.id))?.output).toBe("done")
  93. }),
  94. )
  95. it.instance("records failed jobs", () =>
  96. Effect.gen(function* () {
  97. const jobs = yield* BackgroundJob.Service
  98. const job = yield* jobs.start({
  99. type: "test",
  100. run: Effect.fail(new Error("boom")),
  101. })
  102. const result = yield* jobs.wait({ id: job.id })
  103. expect(result.info?.status).toBe("error")
  104. expect(result.info?.error).toBe("boom")
  105. }),
  106. )
  107. it.instance("ignores stale settlements after restarting a failed job", () =>
  108. Effect.gen(function* () {
  109. const jobs = yield* BackgroundJob.Service
  110. const fail = yield* Deferred.make<void>()
  111. const interrupted = yield* Deferred.make<void>()
  112. const release = yield* Deferred.make<void>()
  113. const id = "job_test"
  114. yield* jobs.start({
  115. id,
  116. type: "test",
  117. run: Deferred.await(fail).pipe(Effect.andThen(Effect.fail(new Error("boom")))),
  118. })
  119. yield* jobs.extend({
  120. id,
  121. run: Effect.never.pipe(
  122. Effect.ensuring(Deferred.succeed(interrupted, undefined).pipe(Effect.andThen(Deferred.await(release)))),
  123. ),
  124. })
  125. yield* Deferred.succeed(fail, undefined)
  126. expect((yield* jobs.wait({ id })).info?.status).toBe("error")
  127. yield* Deferred.await(interrupted)
  128. yield* jobs.start({ id, type: "test", run: Effect.never })
  129. yield* Deferred.succeed(release, undefined)
  130. yield* Effect.yieldNow
  131. expect((yield* jobs.get(id))?.status).toBe("running")
  132. yield* jobs.cancel(id)
  133. }),
  134. )
  135. it.instance("can cancel running jobs", () =>
  136. Effect.gen(function* () {
  137. const jobs = yield* BackgroundJob.Service
  138. const interrupted = yield* Deferred.make<void>()
  139. const extendedInterrupted = yield* Deferred.make<void>()
  140. const job = yield* jobs.start({
  141. type: "test",
  142. run: Effect.never.pipe(Effect.ensuring(Deferred.succeed(interrupted, undefined))),
  143. })
  144. yield* jobs.extend({
  145. id: job.id,
  146. run: Effect.never.pipe(Effect.ensuring(Deferred.succeed(extendedInterrupted, undefined))),
  147. })
  148. const cancelled = yield* jobs.cancel(job.id)
  149. expect(cancelled?.status).toBe("cancelled")
  150. yield* Deferred.await(interrupted).pipe(Effect.timeout("1 second"))
  151. yield* Deferred.await(extendedInterrupted).pipe(Effect.timeout("1 second"))
  152. expect((yield* jobs.get(job.id))?.status).toBe("cancelled")
  153. }),
  154. )
  155. it.instance("returns immutable snapshots", () =>
  156. Effect.gen(function* () {
  157. const jobs = yield* BackgroundJob.Service
  158. const job = yield* jobs.start({
  159. type: "test",
  160. metadata: { value: "initial" },
  161. run: Effect.succeed("done"),
  162. })
  163. if (job.metadata) job.metadata.value = "changed"
  164. expect((yield* jobs.get(job.id))?.metadata?.value).toBe("initial")
  165. }),
  166. )
  167. })