job.test.ts 7.9 KB

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