|
|
@@ -25,6 +25,9 @@ type Active = {
|
|
|
pending: number
|
|
|
next: number
|
|
|
output?: { sequence: number; text: string }
|
|
|
+ tail: Deferred.Deferred<void>
|
|
|
+ promoted: Deferred.Deferred<Info>
|
|
|
+ onPromote?: Effect.Effect<void>
|
|
|
}
|
|
|
|
|
|
type State = {
|
|
|
@@ -38,15 +41,31 @@ type FinishResult = {
|
|
|
scope?: Scope.Closeable
|
|
|
}
|
|
|
|
|
|
+type PromoteResult = {
|
|
|
+ info?: Info
|
|
|
+ promoted?: Deferred.Deferred<Info>
|
|
|
+ onPromote?: Effect.Effect<void>
|
|
|
+}
|
|
|
+
|
|
|
type StartResult = { info: Info } | { info: Info; scope: Scope.Closeable; token: object }
|
|
|
|
|
|
-type ExtendResult = { extended: false } | { extended: true; scope: Scope.Closeable; token: object; sequence: number }
|
|
|
+type ExtendResult =
|
|
|
+ | { extended: false }
|
|
|
+ | {
|
|
|
+ extended: true
|
|
|
+ previous: Deferred.Deferred<void>
|
|
|
+ scope: Scope.Closeable
|
|
|
+ tail: Deferred.Deferred<void>
|
|
|
+ token: object
|
|
|
+ sequence: number
|
|
|
+ }
|
|
|
|
|
|
export type StartInput = {
|
|
|
id?: string
|
|
|
type: string
|
|
|
title?: string
|
|
|
metadata?: Record<string, unknown>
|
|
|
+ onPromote?: Effect.Effect<void>
|
|
|
run: Effect.Effect<string, unknown>
|
|
|
}
|
|
|
|
|
|
@@ -71,6 +90,8 @@ export interface Interface {
|
|
|
readonly start: (input: StartInput) => Effect.Effect<Info>
|
|
|
readonly extend: (input: ExtendInput) => Effect.Effect<boolean>
|
|
|
readonly wait: (input: WaitInput) => Effect.Effect<WaitResult>
|
|
|
+ readonly waitForPromotion: (id: string) => Effect.Effect<Info>
|
|
|
+ readonly promote: (id: string) => Effect.Effect<Info | undefined>
|
|
|
readonly cancel: (id: string) => Effect.Effect<Info | undefined>
|
|
|
}
|
|
|
|
|
|
@@ -128,6 +149,7 @@ export const make = Effect.gen(function* () {
|
|
|
: "error"
|
|
|
const next = {
|
|
|
...job,
|
|
|
+ onPromote: undefined,
|
|
|
pending: 0,
|
|
|
output,
|
|
|
info: {
|
|
|
@@ -182,6 +204,8 @@ export const make = Effect.gen(function* () {
|
|
|
const id = input.id ?? Identifier.ascending("job")
|
|
|
const started_at = yield* Clock.currentTimeMillis
|
|
|
const done = yield* Deferred.make<Info>()
|
|
|
+ const promoted = yield* Deferred.make<Info>()
|
|
|
+ const tail = yield* Deferred.make<void>()
|
|
|
const result = yield* SynchronizedRef.modifyEffect(
|
|
|
state.jobs,
|
|
|
Effect.fnUntraced(function* (jobs) {
|
|
|
@@ -205,6 +229,9 @@ export const make = Effect.gen(function* () {
|
|
|
token,
|
|
|
pending: 1,
|
|
|
next: 1,
|
|
|
+ tail,
|
|
|
+ promoted,
|
|
|
+ onPromote: input.onPromote,
|
|
|
}
|
|
|
return [{ info: snapshot(job), scope, token }, new Map(jobs).set(id, job)] as readonly [
|
|
|
StartResult,
|
|
|
@@ -212,7 +239,14 @@ export const make = Effect.gen(function* () {
|
|
|
]
|
|
|
}),
|
|
|
)
|
|
|
- if ("scope" in result) yield* fork(result.scope, id, result.token, 0, restore(input.run))
|
|
|
+ if ("scope" in result)
|
|
|
+ yield* fork(
|
|
|
+ result.scope,
|
|
|
+ id,
|
|
|
+ result.token,
|
|
|
+ 0,
|
|
|
+ restore(input.run).pipe(Effect.ensuring(Deferred.succeed(tail, undefined))),
|
|
|
+ )
|
|
|
return result.info
|
|
|
}),
|
|
|
)
|
|
|
@@ -221,23 +255,34 @@ export const make = Effect.gen(function* () {
|
|
|
const extend: Interface["extend"] = Effect.fn("BackgroundJob.extend")(function* (input) {
|
|
|
return yield* Effect.uninterruptibleMask((restore) =>
|
|
|
Effect.gen(function* () {
|
|
|
+ const tail = yield* Deferred.make<void>()
|
|
|
const result = yield* SynchronizedRef.modify(
|
|
|
state.jobs,
|
|
|
(jobs): readonly [ExtendResult, Map<string, Active>] => {
|
|
|
const job = jobs.get(input.id)
|
|
|
if (!job || job.info.status !== "running") return [{ extended: false }, jobs]
|
|
|
return [
|
|
|
- { extended: true, scope: job.scope, token: job.token, sequence: job.next },
|
|
|
+ { extended: true, previous: job.tail, scope: job.scope, tail, token: job.token, sequence: job.next },
|
|
|
new Map(jobs).set(input.id, {
|
|
|
...job,
|
|
|
pending: job.pending + 1,
|
|
|
next: job.next + 1,
|
|
|
+ tail,
|
|
|
}),
|
|
|
]
|
|
|
},
|
|
|
)
|
|
|
if (!result.extended) return false
|
|
|
- yield* fork(result.scope, input.id, result.token, result.sequence, restore(input.run))
|
|
|
+ yield* fork(
|
|
|
+ result.scope,
|
|
|
+ input.id,
|
|
|
+ result.token,
|
|
|
+ result.sequence,
|
|
|
+ Deferred.await(result.previous).pipe(
|
|
|
+ Effect.andThen(restore(input.run)),
|
|
|
+ Effect.ensuring(Deferred.succeed(result.tail, undefined)),
|
|
|
+ ),
|
|
|
+ )
|
|
|
return true
|
|
|
}),
|
|
|
)
|
|
|
@@ -254,6 +299,40 @@ export const make = Effect.gen(function* () {
|
|
|
return { info: snapshot(job), timedOut: true }
|
|
|
})
|
|
|
|
|
|
+ const waitForPromotion: Interface["waitForPromotion"] = Effect.fn("BackgroundJob.waitForPromotion")(function* (id) {
|
|
|
+ const job = (yield* SynchronizedRef.get(state.jobs)).get(id)
|
|
|
+ if (!job || job.info.status !== "running") return yield* Effect.never
|
|
|
+ if (job.info.metadata?.background === true) return snapshot(job)
|
|
|
+ return yield* Deferred.await(job.promoted)
|
|
|
+ })
|
|
|
+
|
|
|
+ const promote: Interface["promote"] = Effect.fn("BackgroundJob.promote")(function* (id) {
|
|
|
+ const result = yield* SynchronizedRef.modifyEffect(
|
|
|
+ state.jobs,
|
|
|
+ Effect.fnUntraced(function* (jobs) {
|
|
|
+ const job = jobs.get(id)
|
|
|
+ if (!job || job.info.status !== "running") return [{}, jobs] as readonly [PromoteResult, Map<string, Active>]
|
|
|
+ if (job.info.metadata?.background === true)
|
|
|
+ return [{ info: snapshot(job) }, jobs] as readonly [PromoteResult, Map<string, Active>]
|
|
|
+ const next = {
|
|
|
+ ...job,
|
|
|
+ onPromote: undefined,
|
|
|
+ info: {
|
|
|
+ ...job.info,
|
|
|
+ metadata: { ...job.info.metadata, background: true },
|
|
|
+ },
|
|
|
+ }
|
|
|
+ return [
|
|
|
+ { info: snapshot(next), onPromote: job.onPromote, promoted: job.promoted },
|
|
|
+ new Map(jobs).set(id, next),
|
|
|
+ ] as readonly [PromoteResult, Map<string, Active>]
|
|
|
+ }),
|
|
|
+ )
|
|
|
+ if (result.info && result.promoted) yield* Deferred.succeed(result.promoted, result.info).pipe(Effect.ignore)
|
|
|
+ if (result.onPromote) yield* result.onPromote.pipe(Effect.ignore)
|
|
|
+ return result.info
|
|
|
+ })
|
|
|
+
|
|
|
const cancel: Interface["cancel"] = Effect.fn("BackgroundJob.cancel")(function* (id) {
|
|
|
const completed_at = yield* Clock.currentTimeMillis
|
|
|
const result = yield* SynchronizedRef.modify(state.jobs, (jobs): readonly [FinishResult, Map<string, Active>] => {
|
|
|
@@ -262,6 +341,7 @@ export const make = Effect.gen(function* () {
|
|
|
if (job.info.status !== "running") return [{ info: snapshot(job) }, jobs]
|
|
|
const next = {
|
|
|
...job,
|
|
|
+ onPromote: undefined,
|
|
|
pending: 0,
|
|
|
info: {
|
|
|
...job.info,
|
|
|
@@ -276,7 +356,7 @@ export const make = Effect.gen(function* () {
|
|
|
return result.info
|
|
|
})
|
|
|
|
|
|
- return Service.of({ list, get, start, extend, wait, cancel })
|
|
|
+ return Service.of({ list, get, start, extend, wait, waitForPromotion, promote, cancel })
|
|
|
})
|
|
|
|
|
|
export const layer = Layer.effect(Service, make)
|