Uh oh!
There was an error while loading. Please reload this page.
- Notifications
You must be signed in to change notification settings - Fork 516
refactor(stack): centralize service activation policy#6042
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Uh oh!
There was an error while loading. Please reload this page.
Changes from all commits
990dd0f9feb6ff6689f3b5838abeFile filter
Filter by extension
Conversations
Uh oh!
There was an error while loading. Please reload this page.
Jump to
Uh oh!
There was an error while loading. Please reload this page.
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -7,6 +7,7 @@ import { | ||
| FiberMap, | ||
| Layer, | ||
| Context, | ||
| Semaphore, | ||
| Stream, | ||
| SubscriptionRef, | ||
| } from "effect"; | ||
| @@ -100,7 +101,10 @@ export class Orchestrator extends Context.Service< | ||
| healthy: Deferred.Deferred<void>; | ||
| completed: Deferred.Deferred<number>; | ||
| stopped: Deferred.Deferred<void>; | ||
| readonly readinessWaiters: Set<Deferred.Deferred<void>>; | ||
| readonly completionWaiters: Set<Deferred.Deferred<number>>; | ||
| stoppedByUser: boolean; | ||
| requested: boolean; | ||
| } | ||
| const services = new Map<string, ServiceSignals>(); | ||
| @@ -113,12 +117,16 @@ export class Orchestrator extends Context.Service< | ||
| healthy: Deferred.makeUnsafe<void>(), | ||
| completed: Deferred.makeUnsafe<number>(), | ||
| stopped: Deferred.makeUnsafe<void>(), | ||
| readinessWaiters: new Set(), | ||
| completionWaiters: new Set(), | ||
| stoppedByUser: false, | ||
| requested: false, | ||
| }); | ||
| } | ||
| // FiberMap to track running service fibers — auto-interrupted on scope close | ||
| const fibers = yield* FiberMap.make<string>(); | ||
| const startServiceLock = Semaphore.makeUnsafe(1); | ||
| // Helper: send a validated FSM event — only does the state transition | ||
| const sendEvent = ( | ||
| @@ -130,6 +138,13 @@ export class Orchestrator extends Context.Service< | ||
| return transition(svc.state, event); | ||
| }; | ||
| const signalHealthy = (svc: ServiceSignals): Effect.Effect<void> => | ||
| Effect.forEach( | ||
| [svc.healthy, ...svc.readinessWaiters], | ||
| (waiter) => Deferred.succeed(waiter, void 0), | ||
| { discard: true }, | ||
| ); | ||
| // Helper: run all hooks for a given trigger in sequence | ||
| const runHooks = (def: ServiceDef, trigger: HookTrigger): Effect.Effect<void> => | ||
| Effect.gen(function* () { | ||
| @@ -349,7 +364,7 @@ export class Orchestrator extends Context.Service< | ||
| yield* runHooks(def, "healthy"); | ||
| const current = SubscriptionRef.getUnsafe(svcSig.state); | ||
| if (current.status !== "Failed") { | ||
| yield* Deferred.succeed(svcSig.healthy, void 0); | ||
| yield* signalHealthy(svcSig); | ||
| } | ||
| } | ||
| } | ||
| @@ -382,7 +397,7 @@ export class Orchestrator extends Context.Service< | ||
| if (svcSig) { | ||
| const current = SubscriptionRef.getUnsafe(svcSig.state); | ||
| if (current.status !== "Failed") { | ||
| yield* Deferred.succeed(svcSig.healthy, void 0); | ||
| yield* signalHealthy(svcSig); | ||
| } | ||
| } | ||
| } | ||
| @@ -441,8 +456,14 @@ export class Orchestrator extends Context.Service< | ||
| const handleResult = (r: SpawnResult) => | ||
| Effect.gen(function* () { | ||
| if (r._tag === "Exited") { | ||
| const completeSig = services.get(def.name)?.completed; | ||
| if (completeSig) yield* Deferred.succeed(completeSig, r.exitCode); | ||
| const service = services.get(def.name); | ||
| if (service !== undefined) { | ||
| yield* Effect.forEach( | ||
| [service.completed, ...service.completionWaiters], | ||
| (waiter) => Deferred.succeed(waiter, r.exitCode), | ||
| { discard: true }, | ||
| ); | ||
| } | ||
| if (r.exitCode !== 0 && r.exitCode !== 143) { | ||
| yield* appendRecentServiceLogs( | ||
| def.name, | ||
| @@ -538,7 +559,15 @@ export class Orchestrator extends Context.Service< | ||
| } | ||
| }; | ||
| collectDependents(name); | ||
| return graph.startOrder.filter((def) => names.has(def.name)); | ||
| return graph.startOrder.filter((def) => { | ||
| const service = services.get(def.name); | ||
| return ( | ||
| names.has(def.name) && | ||
| (def.name === name || | ||
| FiberMap.hasUnsafe(fibers, def.name) || | ||
| (service?.requested === true && service.stoppedByUser !== true)) | ||
| ); | ||
| }); | ||
| }; | ||
| const waitReadySingle = (def: ServiceDef): Effect.Effect<void, ServiceReadyError> => | ||
| @@ -547,7 +576,6 @@ export class Orchestrator extends Context.Service< | ||
| if (!svc) return Effect.void; | ||
| const restartPolicy = def.restart ?? defaults.restart; | ||
| // Check if already failed | ||
| const current = SubscriptionRef.getUnsafe(svc.state); | ||
| if (current.status === "Failed") { | ||
| return Effect.fail( | ||
| @@ -559,8 +587,17 @@ export class Orchestrator extends Context.Service< | ||
| } | ||
| if (restartPolicy === "no") { | ||
| // One-shot: wait for completed, check exit code | ||
| return Deferred.await(svc.completed).pipe( | ||
| const completion = Deferred.makeUnsafe<number>(); | ||
| svc.completionWaiters.add(completion); | ||
| return Deferred.isDone(svc.completed).pipe( | ||
| Effect.flatMap((alreadyCompleted) => | ||
| alreadyCompleted | ||
| ? Deferred.await(svc.completed).pipe( | ||
| Effect.flatMap((exitCode) => Deferred.succeed(completion, exitCode)), | ||
| ) | ||
| : Effect.void, | ||
| ), | ||
| Effect.andThen(Deferred.await(completion)), | ||
| Effect.flatMap((exitCode) => | ||
| exitCode === 0 | ||
| ? Effect.void | ||
| @@ -572,35 +609,67 @@ export class Orchestrator extends Context.Service< | ||
| }), | ||
| ), | ||
| ), | ||
| Effect.ensuring(Effect.sync(() => svc.completionWaiters.delete(completion))), | ||
| ); | ||
| } | ||
| // Long-running: race healthy vs failure | ||
| return Effect.race( | ||
| Deferred.await(svc.healthy), | ||
| SubscriptionRef.changes(svc.state).pipe( | ||
| Stream.filter((s) => s.status === "Failed"), | ||
| Stream.take(1), | ||
| Stream.runDrain, | ||
| Effect.andThen( | ||
| Effect.gen(function* () { | ||
| const current = SubscriptionRef.getUnsafe(svc.state); | ||
| return yield* Effect.fail( | ||
| new ServiceReadyError({ | ||
| name: def.name, | ||
| reason: current.error ?? "Service entered Failed state", | ||
| }), | ||
| ); | ||
| }), | ||
| if (current.status === "Stopped") { | ||
| return Effect.fail( | ||
| new ServiceReadyError({ | ||
| name: def.name, | ||
| reason: "Service stopped before becoming ready", | ||
| }), | ||
| ); | ||
| } | ||
| // A waiter belongs to the readiness request, not to one service signal | ||
| // generation. Dependency retries may replace svc.healthy while this | ||
| // request is pending, so signal the stable waiter from any generation. | ||
| const readiness = Deferred.makeUnsafe<void>(); | ||
| svc.readinessWaiters.add(readiness); | ||
| return Deferred.isDone(svc.healthy).pipe( | ||
| Effect.flatMap((alreadyHealthy) => | ||
| alreadyHealthy ? Deferred.succeed(readiness, void 0) : Effect.void, | ||
| ), | ||
| Effect.andThen( | ||
| Effect.race( | ||
| Deferred.await(readiness), | ||
| SubscriptionRef.changes(svc.state).pipe( | ||
| Stream.filter( | ||
| (state) => state.status === "Failed" || state.status === "Stopped", | ||
| ), | ||
| Stream.take(1), | ||
| Stream.runDrain, | ||
| Effect.andThen( | ||
| Deferred.isDone(readiness).pipe( | ||
| Effect.flatMap((ready) => { | ||
| if (ready) return Effect.void; | ||
| const terminal = SubscriptionRef.getUnsafe(svc.state); | ||
| return Effect.fail( | ||
| new ServiceReadyError({ | ||
| name: def.name, | ||
| reason: | ||
| terminal.status === "Stopped" | ||
| ? "Service stopped before becoming ready" | ||
| : (terminal.error ?? "Service entered Failed state"), | ||
| }), | ||
| ); | ||
| }), | ||
| ), | ||
| ), | ||
| ), | ||
| ), | ||
| ), | ||
| Effect.ensuring(Effect.sync(() => svc.readinessWaiters.delete(readiness))), | ||
| ); | ||
| }); | ||
| return { | ||
| start: () => | ||
| Effect.gen(function* () { | ||
| for (const def of graph.startOrder) { | ||
| const service = services.get(def.name); | ||
| if (service !== undefined) service.requested = true; | ||
| yield* FiberMap.run(fibers, def.name, runServiceSafe(def)); | ||
| } | ||
| }), | ||
| @@ -612,10 +681,74 @@ export class Orchestrator extends Context.Service< | ||
| return yield* Effect.fail(new ServiceNotFoundError({ name })); | ||
| } | ||
| const order = graph.startOrderFor(name); | ||
| const orderNames = new Set(order.map((service) => service.name)); | ||
| const resetNames = new Set<string>(); | ||
| let dependencySignalsReset = false; | ||
| for (const d of order) { | ||
| const service = services.get(d.name); | ||
| const state = service?.state; | ||
| const status = | ||
| state === undefined ? undefined : SubscriptionRef.getUnsafe(state).status; | ||
| const restartPolicy = d.restart ?? defaults.restart; | ||
| // A successful one-shot remains a satisfied dependency after its | ||
| // process exits. Naturally stopped long-running services can be | ||
| // started again even when their restart policy did not relaunch them. | ||
| if ( | ||
| status === "Stopped" && | ||
| d.name !== name && | ||
| service?.stoppedByUser !== true && | ||
| restartPolicy === "no" | ||
| ) { | ||
| continue; | ||
jgoux marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| } | ||
| if (status === "Failed") { | ||
| // The failed fiber may still be unwinding its process scope. | ||
| // Remove it before replacing the signals used by the retry. | ||
| yield* FiberMap.remove(fibers, d.name); | ||
| yield* resetService(d.name); | ||
jgoux marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. jgoux marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| resetNames.add(d.name); | ||
| dependencySignalsReset = true; | ||
| } else if (status === "Stopped") { | ||
| yield* FiberMap.remove(fibers, d.name); | ||
| yield* resetService(d.name); | ||
| resetNames.add(d.name); | ||
| dependencySignalsReset = true; | ||
jgoux marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| } else if (dependencySignalsReset && status === "Pending") { | ||
| // This fiber may be waiting on Deferreds replaced by a retried dependency. | ||
| // Relaunch it so it observes the dependency's new signal generation. | ||
| yield* FiberMap.remove(fibers, d.name); | ||
| yield* resetService(d.name); | ||
| resetNames.add(d.name); | ||
| } | ||
| if (service !== undefined) service.requested = true; | ||
| yield* FiberMap.run(fibers, d.name, runServiceSafe(d), { onlyIfMissing: true }); | ||
| } | ||
| }), | ||
| // A caller may retry a dependency directly while dependents started by an | ||
| // earlier request are still waiting on its old Deferred generation. | ||
| for (const d of graph.startOrder) { | ||
| if (orderNames.has(d.name)) continue; | ||
| const dependsOnResetService = graph | ||
| .startOrderFor(d.name) | ||
| .some((dependency) => resetNames.has(dependency.name)); | ||
| const service = services.get(d.name); | ||
| if ( | ||
| dependsOnResetService && | ||
| service !== undefined && | ||
| FiberMap.hasUnsafe(fibers, d.name) && | ||
| SubscriptionRef.getUnsafe(service.state).status === "Pending" | ||
| ) { | ||
| yield* FiberMap.remove(fibers, d.name); | ||
| yield* resetService(d.name); | ||
jgoux marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. jgoux marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| service.requested = true; | ||
| yield* FiberMap.run(fibers, d.name, runServiceSafe(d), { | ||
jgoux marked this conversation as resolved.
Uh oh!There was an error while loading. Please reload this page. | ||
| onlyIfMissing: true, | ||
| }); | ||
| } | ||
| } | ||
| }).pipe(startServiceLock.withPermit), | ||
| stop: () => | ||
| Effect.gen(function* () { | ||
| @@ -698,6 +831,8 @@ export class Orchestrator extends Context.Service< | ||
| yield* resetService(affectedDef.name); | ||
| } | ||
| for (const affectedDef of affected) { | ||
| const service = services.get(affectedDef.name); | ||
| if (service !== undefined) service.requested = true; | ||
| yield* FiberMap.run(fibers, affectedDef.name, runServiceSafe(affectedDef)); | ||
| } | ||
| }), | ||
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.