Compare commits
6 commits
dev
...
fix/v2-ses
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ab3850d1b4 | ||
|
|
6ec6593e9e | ||
|
|
93ca32f41a | ||
|
|
fd07e58384 | ||
|
|
ecf6c532d5 | ||
|
|
911664db67 |
2 changed files with 62 additions and 11 deletions
|
|
@ -39,13 +39,14 @@ export interface Coordinator<Key, A, E> {
|
||||||
/** One Session's process-local execution lane: one active demand and at most one coalesced follow-up. */
|
/** One Session's process-local execution lane: one active demand and at most one coalesced follow-up. */
|
||||||
type Entry<A, E> = {
|
type Entry<A, E> = {
|
||||||
readonly done: Deferred.Deferred<A, E>
|
readonly done: Deferred.Deferred<A, E>
|
||||||
readonly settled: Deferred.Deferred<Exit.Exit<A, E>>
|
readonly settled: Deferred.Deferred<Exit.Exit<A, E> | undefined>
|
||||||
current: Demand
|
current: Demand
|
||||||
pending?: Demand
|
pending?: Demand
|
||||||
explicitWaiter?: Deferred.Deferred<A, E>
|
explicitWaiter?: Deferred.Deferred<A, E>
|
||||||
interruptSeq?: number
|
interruptSeq?: number
|
||||||
owner?: Fiber.Fiber<void, never>
|
owner?: Fiber.Fiber<void, never>
|
||||||
stopping: boolean
|
stopping: boolean
|
||||||
|
advisoryRetryAvailable: boolean
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Combines follow-up demand: runs dominate, while wakes retain the newest durable admission sequence. */
|
/** Combines follow-up demand: runs dominate, while wakes retain the newest durable admission sequence. */
|
||||||
|
|
@ -81,12 +82,16 @@ export const make = <Key, A, E>(options: {
|
||||||
}),
|
}),
|
||||||
)
|
)
|
||||||
|
|
||||||
const makeEntry = (current: Demand, explicitWaiter?: Deferred.Deferred<A, E>): Entry<A, E> => ({
|
const makeEntry = (current: Demand, options?: {
|
||||||
|
readonly explicitWaiter?: Deferred.Deferred<A, E>
|
||||||
|
readonly advisoryRetryAvailable?: boolean
|
||||||
|
}): Entry<A, E> => ({
|
||||||
done: Deferred.makeUnsafe<A, E>(),
|
done: Deferred.makeUnsafe<A, E>(),
|
||||||
settled: Deferred.makeUnsafe<Exit.Exit<A, E>>(),
|
settled: Deferred.makeUnsafe<Exit.Exit<A, E> | undefined>(),
|
||||||
current,
|
current,
|
||||||
explicitWaiter,
|
explicitWaiter: options?.explicitWaiter,
|
||||||
stopping: false,
|
stopping: false,
|
||||||
|
advisoryRetryAvailable: options?.advisoryRetryAvailable ?? current._tag === "wake",
|
||||||
})
|
})
|
||||||
|
|
||||||
const start = (key: Key, entry: Entry<A, E>, demand: Demand, successor = false) => {
|
const start = (key: Key, entry: Entry<A, E>, demand: Demand, successor = false) => {
|
||||||
|
|
@ -132,6 +137,7 @@ export const make = <Key, A, E>(options: {
|
||||||
const pending = entry.pending
|
const pending = entry.pending
|
||||||
entry.pending = undefined
|
entry.pending = undefined
|
||||||
entry.current = pending
|
entry.current = pending
|
||||||
|
entry.advisoryRetryAvailable = pending._tag === "wake"
|
||||||
start(key, entry, pending, true)
|
start(key, entry, pending, true)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -141,16 +147,23 @@ export const make = <Key, A, E>(options: {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
const successor = entry.pending !== undefined ? makeEntry(entry.pending, entry.explicitWaiter) : undefined
|
const successor =
|
||||||
|
entry.pending !== undefined
|
||||||
|
? makeEntry(entry.pending, { explicitWaiter: entry.explicitWaiter })
|
||||||
|
: exit._tag === "Failure" && demand._tag === "wake" && !entry.stopping && entry.advisoryRetryAvailable
|
||||||
|
? makeEntry(demand, { explicitWaiter: entry.explicitWaiter, advisoryRetryAvailable: false })
|
||||||
|
: undefined
|
||||||
|
const retrying = successor !== undefined && entry.pending === undefined
|
||||||
if (successor === undefined) active.delete(key)
|
if (successor === undefined) active.delete(key)
|
||||||
else active.set(key, successor)
|
else active.set(key, successor)
|
||||||
if (successor !== undefined) start(key, successor, successor.current, true)
|
if (successor !== undefined) start(key, successor, successor.current, true)
|
||||||
Deferred.doneUnsafe(entry.done, exit)
|
Deferred.doneUnsafe(entry.done, exit)
|
||||||
Deferred.doneUnsafe(entry.settled, Effect.succeed(exit))
|
Deferred.doneUnsafe(entry.settled, Effect.succeed(retrying ? undefined : exit))
|
||||||
if (
|
if (
|
||||||
exit._tag === "Failure" &&
|
exit._tag === "Failure" &&
|
||||||
!(entry.stopping && Cause.hasInterruptsOnly(exit.cause)) &&
|
!(entry.stopping && Cause.hasInterruptsOnly(exit.cause)) &&
|
||||||
demand._tag === "wake" &&
|
demand._tag === "wake" &&
|
||||||
|
successor === undefined &&
|
||||||
options.onFailure !== undefined
|
options.onFailure !== undefined
|
||||||
) {
|
) {
|
||||||
report(Effect.suspend(() => options.onFailure!(key, exit.cause)))
|
report(Effect.suspend(() => options.onFailure!(key, exit.cause)))
|
||||||
|
|
@ -184,7 +197,7 @@ export const make = <Key, A, E>(options: {
|
||||||
Deferred.await(shutdown).pipe(Effect.as(Exit.void)),
|
Deferred.await(shutdown).pipe(Effect.as(Exit.void)),
|
||||||
)
|
)
|
||||||
if (closed) break
|
if (closed) break
|
||||||
if (exit._tag === "Failure" && firstFailure === undefined) firstFailure = exit.cause
|
if (exit?._tag === "Failure" && firstFailure === undefined) firstFailure = exit.cause
|
||||||
}
|
}
|
||||||
if (firstFailure !== undefined) return yield* Effect.failCause(firstFailure)
|
if (firstFailure !== undefined) return yield* Effect.failCause(firstFailure)
|
||||||
})
|
})
|
||||||
|
|
|
||||||
|
|
@ -33,7 +33,9 @@ describe("SessionRunCoordinator", () => {
|
||||||
Effect.scoped(
|
Effect.scoped(
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const drained = yield* Deferred.make<void>()
|
const drained = yield* Deferred.make<void>()
|
||||||
const coordinator = yield* SessionRunCoordinator.make({ drain: () => Deferred.succeed(drained, undefined) })
|
const coordinator = yield* SessionRunCoordinator.make({
|
||||||
|
drain: () => Deferred.succeed(drained, undefined),
|
||||||
|
})
|
||||||
|
|
||||||
yield* coordinator.wake("session")
|
yield* coordinator.wake("session")
|
||||||
yield* Deferred.await(drained)
|
yield* Deferred.await(drained)
|
||||||
|
|
@ -44,7 +46,9 @@ describe("SessionRunCoordinator", () => {
|
||||||
it.effect("does nothing when interrupted while idle", () =>
|
it.effect("does nothing when interrupted while idle", () =>
|
||||||
Effect.scoped(
|
Effect.scoped(
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const coordinator = yield* SessionRunCoordinator.make({ drain: () => Effect.void })
|
const coordinator = yield* SessionRunCoordinator.make({
|
||||||
|
drain: () => Effect.void,
|
||||||
|
})
|
||||||
|
|
||||||
yield* coordinator.interrupt("session")
|
yield* coordinator.interrupt("session")
|
||||||
}),
|
}),
|
||||||
|
|
@ -55,7 +59,9 @@ describe("SessionRunCoordinator", () => {
|
||||||
Effect.scoped(
|
Effect.scoped(
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
let runs = 0
|
let runs = 0
|
||||||
const coordinator = yield* SessionRunCoordinator.make({ drain: () => Effect.sync(() => runs++) })
|
const coordinator = yield* SessionRunCoordinator.make({
|
||||||
|
drain: () => Effect.sync(() => runs++),
|
||||||
|
})
|
||||||
|
|
||||||
yield* coordinator.interrupt("session", 2)
|
yield* coordinator.interrupt("session", 2)
|
||||||
yield* coordinator.wake("session", 1)
|
yield* coordinator.wake("session", 1)
|
||||||
|
|
@ -722,6 +728,36 @@ describe("SessionRunCoordinator", () => {
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
it.effect("retries a failed advisory successor once without another wake", () =>
|
||||||
|
Effect.scoped(
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const firstGate = yield* Deferred.make<void>()
|
||||||
|
let runs = 0
|
||||||
|
const failure = new Error("transient wake failure")
|
||||||
|
const coordinator = yield* SessionRunCoordinator.make<string, void, Error>({
|
||||||
|
drain: () =>
|
||||||
|
Effect.sync(() => ++runs).pipe(
|
||||||
|
Effect.flatMap((run) =>
|
||||||
|
run === 1
|
||||||
|
? Deferred.await(firstGate).pipe(Effect.andThen(Effect.fail(failure)))
|
||||||
|
: run === 2
|
||||||
|
? Effect.fail(failure)
|
||||||
|
: Effect.void,
|
||||||
|
),
|
||||||
|
),
|
||||||
|
})
|
||||||
|
|
||||||
|
yield* coordinator.wake("session", 1)
|
||||||
|
yield* coordinator.wake("session", 2)
|
||||||
|
yield* Deferred.succeed(firstGate, undefined)
|
||||||
|
const idle = yield* coordinator.awaitIdle("session").pipe(Effect.exit)
|
||||||
|
|
||||||
|
expect(runs).toBe(3)
|
||||||
|
expect(Exit.isSuccess(idle)).toBeTrue()
|
||||||
|
}),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
it.effect("upgrades an active wake when an explicit run joins it", () =>
|
it.effect("upgrades an active wake when an explicit run joins it", () =>
|
||||||
Effect.scoped(
|
Effect.scoped(
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
|
|
@ -916,8 +952,9 @@ describe("SessionRunCoordinator", () => {
|
||||||
const failure = new Error("wake failed")
|
const failure = new Error("wake failed")
|
||||||
const reported: Cause.Cause<Error>[] = []
|
const reported: Cause.Cause<Error>[] = []
|
||||||
const reportedOnce = yield* Deferred.make<void>()
|
const reportedOnce = yield* Deferred.make<void>()
|
||||||
|
let runs = 0
|
||||||
const coordinator = yield* SessionRunCoordinator.make<string, void, Error>({
|
const coordinator = yield* SessionRunCoordinator.make<string, void, Error>({
|
||||||
drain: () => Effect.fail(failure),
|
drain: () => Effect.sync(() => runs++).pipe(Effect.andThen(Effect.fail(failure))),
|
||||||
onFailure: (_key, cause) =>
|
onFailure: (_key, cause) =>
|
||||||
Effect.sync(() => reported.push(cause)).pipe(Effect.andThen(Deferred.succeed(reportedOnce, undefined))),
|
Effect.sync(() => reported.push(cause)).pipe(Effect.andThen(Deferred.succeed(reportedOnce, undefined))),
|
||||||
})
|
})
|
||||||
|
|
@ -927,6 +964,7 @@ describe("SessionRunCoordinator", () => {
|
||||||
yield* Effect.yieldNow
|
yield* Effect.yieldNow
|
||||||
|
|
||||||
expect(reported).toHaveLength(1)
|
expect(reported).toHaveLength(1)
|
||||||
|
expect(runs).toBe(2)
|
||||||
expect(Cause.squash(reported[0]!)).toBe(failure)
|
expect(Cause.squash(reported[0]!)).toBe(failure)
|
||||||
}),
|
}),
|
||||||
),
|
),
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue