Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
176 changes: 130 additions & 46 deletions apps/server/src/provider/providerMaintenanceRunner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,14 +7,18 @@ import {
} from "@t3tools/contracts";
import { ServerProviderUpdateError } from "@t3tools/contracts";
import * as Cause from "effect/Cause";
import * as Context from "effect/Context";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Fiber from "effect/Fiber";
import * as Layer from "effect/Layer";
import * as Ref from "effect/Ref";
import * as Schema from "effect/Schema";
import * as Scope from "effect/Scope";
import * as Sink from "effect/Sink";
import * as Stream from "effect/Stream";
import * as TestClock from "effect/testing/TestClock";
import { HttpClient, HttpClientResponse } from "effect/unstable/http";
import { ChildProcessSpawner } from "effect/unstable/process";
import { HostProcessEnvironment, HostProcessPlatform } from "@t3tools/shared/hostProcess";
Expand Down Expand Up @@ -209,18 +213,16 @@ function makeRegistry(
}

const makeTestRunner = (registry: ProviderRegistryShape) =>
Effect.service(ProviderMaintenanceRunner.ProviderMaintenanceRunner).pipe(
Effect.provide(
ProviderMaintenanceRunner.layer.pipe(
Layer.provide(
Layer.mergeAll(
Layer.succeed(ProviderRegistry, registry),
// Fresh per runner so a version cached by one test cannot leak into another.
Layer.sync(ProviderVersionCache, () => new Map()),
),
),
ProviderMaintenanceRunner.layer.pipe(
Layer.provide(
Layer.mergeAll(
Layer.succeed(ProviderRegistry, registry),
// Fresh per runner so a version cached by one test cannot leak into another.
Layer.sync(ProviderVersionCache, () => new Map()),
),
),
Layer.build,
Effect.map(Context.get(ProviderMaintenanceRunner.ProviderMaintenanceRunner)),
);

describe("providerMaintenanceRunner", () => {
Expand Down Expand Up @@ -797,58 +799,140 @@ describe("providerMaintenanceRunner", () => {
);
});

it.effect(
"releases the running-provider marker when interrupted after queuing but before the lock run starts",
() =>
it.effect.each(["queued", "running"] as const)(
"finishes a %s update after its client disconnects and keeps duplicate updates blocked",
(disconnectAt) =>
Effect.gen(function* () {
const { registry } = yield* makeRegistry(baseProvider);
let blockQueuedState = true;
const queuedStateWrittenLatch: { resolve: () => void } = { resolve: () => {} };
const releaseQueuedStateLatch: { resolve: () => void } = { resolve: () => {} };
const queuedStateWritten = new Promise<void>((resolve) => {
queuedStateWrittenLatch.resolve = resolve;
});
const releaseQueuedState = new Promise<void>((resolve) => {
releaseQueuedStateLatch.resolve = resolve;
});

const { registry } = yield* makeRegistry();
const reachedState = yield* Deferred.make<void>();
const releaseUpdate = yield* Deferred.make<void>();
const settled = yield* Deferred.make<ServerProviderUpdateState>();
const updater = yield* makeTestRunner({
...registry,
setProviderMaintenanceActionState: Effect.fn(
"providerMaintenanceRunner.test.blockQueuedState",
)(function* (input) {
const providers = yield* registry.setProviderMaintenanceActionState(input);
if (input.state?.status === "queued" && blockQueuedState) {
queuedStateWrittenLatch.resolve();
yield* Effect.promise(() => releaseQueuedState);
}
return providers;
}),
});

const first = yield* updater.updateProvider(CODEX_DRIVER).pipe(Effect.forkScoped);
yield* Effect.promise(() => queuedStateWritten);
blockQueuedState = false;
setProviderMaintenanceActionState: (input) =>
Effect.gen(function* () {
const providers = yield* registry.setProviderMaintenanceActionState(input);
if (input.state?.status === disconnectAt) {
yield* Deferred.succeed(reachedState, undefined);
// A queued update hasn't spawned its installer yet.
if (disconnectAt === "queued") yield* Deferred.await(releaseUpdate);
}
if (input.state?.finishedAt) yield* Deferred.succeed(settled, input.state);
return providers;
}),
}).pipe(
Effect.provide(
mockSpawnerLayer(() => ({
exitCode: Deferred.await(releaseUpdate).pipe(
Effect.as(ChildProcessSpawner.ExitCode(0)),
),
})),
),
);

yield* Fiber.interrupt(first);
releaseQueuedStateLatch.resolve();
const client = yield* updater.updateProvider(CODEX_DRIVER).pipe(Effect.forkScoped);
yield* Deferred.await(reachedState);
yield* Fiber.interrupt(client);
assert.strictEqual((yield* registry.getProviders)[0]?.updateState?.status, disconnectAt);

const second = yield* updater.updateProvider(CODEX_DRIVER).pipe(Effect.exit);
assert.strictEqual(Exit.isSuccess(second), true);
if (Exit.isSuccess(second)) {
assert.strictEqual(second.value.providers[0]?.updateState?.status, "succeeded");
const duplicate = yield* updater.updateProvider(CODEX_DRIVER).pipe(Effect.exit);
assert.strictEqual(Exit.isFailure(duplicate), true);
if (Exit.isFailure(duplicate)) {
assert.include(String(Cause.squash(duplicate.cause)), "already running");
}

yield* Deferred.succeed(releaseUpdate, undefined);
const terminal = yield* Deferred.await(settled);
assert.strictEqual(terminal.status, "succeeded");
assert.isNotNull(terminal.finishedAt);
assert.strictEqual((yield* registry.getProviders)[0]?.updateState?.status, "succeeded");
}).pipe(Effect.provide(Layer.mergeAll(NonWindowsPlatform, latestVersionHttpClient("0.0.0")))),
);

it.effect.each(["queued", "running"] as const)(
"records failure when server shutdown interrupts a %s update",
(interruptAt) =>
Effect.gen(function* () {
const { registry } = yield* makeRegistry();
const serverScope = yield* Scope.make();
const reachedState = yield* Deferred.make<void>();
const updater = yield* makeTestRunner({
...registry,
setProviderMaintenanceActionState: (input) =>
Effect.gen(function* () {
const providers = yield* registry.setProviderMaintenanceActionState(input);
if (input.state?.status === interruptAt) {
yield* Deferred.succeed(reachedState, undefined);
if (interruptAt === "queued") return yield* Effect.never;
}
return providers;
}),
}).pipe(Scope.provide(serverScope));

const client = yield* updater.updateProvider(CODEX_DRIVER).pipe(Effect.forkScoped);
yield* Deferred.await(reachedState);
yield* Scope.close(serverScope, Exit.void);
yield* Fiber.await(client);

const state = (yield* registry.getProviders)[0]?.updateState;
assert.strictEqual(state?.status, "failed");
assert.strictEqual(state?.message, "Provider update was interrupted. Try again.");
assert.isNotNull(state?.finishedAt);
if (interruptAt === "queued") assert.isNull(state?.startedAt);
else assert.isNotNull(state?.startedAt);
}).pipe(
Effect.provide(
Layer.mergeAll(
NonWindowsPlatform,
latestVersionHttpClient("0.0.0"),
mockSpawnerLayer(() => ({ stdout: "updated" })),
mockSpawnerLayer(() => ({ exitCode: Effect.never })),
),
),
),
);

it.effect("times out a detached installer and allows a retry", () =>
Effect.gen(function* () {
const { registry } = yield* makeRegistry();
const commandStarted = yield* Deferred.make<void>();
const settled = yield* Deferred.make<ServerProviderUpdateState>();
let commands = 0;
const updater = yield* makeTestRunner({
...registry,
setProviderMaintenanceActionState: (input) =>
registry
.setProviderMaintenanceActionState(input)
.pipe(
Effect.tap(() =>
input.state?.finishedAt ? Deferred.succeed(settled, input.state) : Effect.void,
),
),
}).pipe(
Effect.provide(
mockSpawnerLayer(() => ({
exitCode:
commands++ === 0
? Deferred.succeed(commandStarted, undefined).pipe(Effect.andThen(Effect.never))
: Effect.succeed(ChildProcessSpawner.ExitCode(0)),
})),
),
);

const client = yield* updater.updateProvider(CODEX_DRIVER).pipe(Effect.forkScoped);
yield* Deferred.await(commandStarted);
yield* Fiber.interrupt(client);
yield* TestClock.adjust("5 minutes");
const state = yield* Deferred.await(settled);
assert.strictEqual(state.status, "failed");
assert.strictEqual(state.message, "Update timed out.");
assert.isNotNull(state.finishedAt);

const retried = yield* updater.updateProvider(CODEX_DRIVER);
assert.strictEqual(retried.providers[0]?.updateState?.status, "succeeded");
assert.strictEqual(commands, 2);
}).pipe(Effect.provide(Layer.mergeAll(NonWindowsPlatform, latestVersionHttpClient("0.0.0")))),
);

it.effect("resolves npm to a .cmd shim and routes through the shell on win32", () => {
const captured: Array<{
readonly command: string;
Expand Down
44 changes: 36 additions & 8 deletions apps/server/src/provider/providerMaintenanceRunner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import * as Data from "effect/Data";
import * as DateTime from "effect/DateTime";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Ref from "effect/Ref";
Expand Down Expand Up @@ -213,6 +214,7 @@ function makeUpdateState(input: {

/** @public Service construction is part of the canonical Effect module API. */
export const make = Effect.fn("ProviderMaintenanceRunner.make")(function* () {
const scope = yield* Effect.scope;
const providerRegistry = yield* ProviderRegistry;
const spawner = yield* ChildProcessSpawner.ChildProcessSpawner;
const httpClient = yield* HttpClient.HttpClient;
Expand Down Expand Up @@ -302,8 +304,8 @@ export const make = Effect.fn("ProviderMaintenanceRunner.make")(function* () {
}),
);

const updateProvider: ProviderMaintenanceRunnerShape["updateProvider"] = Effect.fn(
"ProviderMaintenanceRunner.updateProvider",
const runUpdate: ProviderMaintenanceRunnerShape["updateProvider"] = Effect.fn(
"ProviderMaintenanceRunner.runUpdate",
)(function* (target) {
const provider = typeof target === "string" ? target : target.provider;
const instanceId =
Expand All @@ -323,12 +325,17 @@ export const make = Effect.fn("ProviderMaintenanceRunner.make")(function* () {
});
}

const setUpdateState = (state: ServerProviderUpdateState | null) =>
providerRegistry.setProviderMaintenanceActionState({
instanceId,
action: "update",
state,
});
const updateStateRef = yield* Ref.make<ServerProviderUpdateState | null>(null);
const setUpdateState = (state: ServerProviderUpdateState) =>
Ref.set(updateStateRef, state).pipe(
Effect.andThen(
providerRegistry.setProviderMaintenanceActionState({
instanceId,
action: "update",
state,
}),
),
);
const setQueuedState = setUpdateState(
makeUpdateState({
status: "queued",
Expand Down Expand Up @@ -455,6 +462,22 @@ export const make = Effect.fn("ProviderMaintenanceRunner.make")(function* () {
run: runProviderUpdate(),
})
.pipe(
// Interruption skips catchCause. Record a terminal state even when
// shutdown interrupts an update waiting for the package-manager lock.
Effect.onInterrupt(() =>
Effect.gen(function* () {
const state = yield* Ref.get(updateStateRef);
if (state?.status !== "queued" && state?.status !== "running") return;
yield* setUpdateState(
makeUpdateState({
status: "failed",
startedAt: state.startedAt,
finishedAt: yield* nowIso,
message: "Provider update was interrupted. Try again.",
}),
);
}),
),
Effect.mapError((error) =>
isServerProviderUpdateError(error)
? new ServerProviderUpdateError({
Expand All @@ -466,6 +489,11 @@ export const make = Effect.fn("ProviderMaintenanceRunner.make")(function* () {
);
});

// The server owns the installer; a disconnected client only stops waiting
// for its result. Closing the service scope still stops and settles updates.
const updateProvider: ProviderMaintenanceRunnerShape["updateProvider"] = (target) =>
runUpdate(target).pipe(Effect.forkIn(scope), Effect.flatMap(Fiber.join));

return ProviderMaintenanceRunner.of({
updateProvider,
});
Expand Down
Loading
Loading