Skip to content
Open
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
166 changes: 158 additions & 8 deletions apps/desktop/src/backend/DesktopBackendPool.test.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,16 @@
import { assert, describe, it } from "@effect/vitest";
import * as Deferred from "effect/Deferred";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Queue from "effect/Queue";
import * as Sink from "effect/Sink";
import * as TestClock from "effect/testing/TestClock";
import * as Ref from "effect/Ref";
import * as Stream from "effect/Stream";
import { HttpClient } from "effect/unstable/http";
import { HttpClient, HttpClientResponse } from "effect/unstable/http";
import { ChildProcessSpawner } from "effect/unstable/process";

import * as DesktopObservability from "../app/DesktopObservability.ts";
Expand Down Expand Up @@ -43,18 +47,20 @@ function makeStubInstance(

function makePoolLayer(
labelRef: Ref.Ref<string>,
spawner = ChildProcessSpawner.make(() => Effect.die("unexpected child process spawn")),
): Layer.Layer<DesktopBackendPool.DesktopBackendPool> {
return DesktopBackendPool.layer.pipe(
Layer.provideMerge(
Layer.mergeAll(
FileSystem.layerNoop({}),
Layer.succeed(
ChildProcessSpawner.ChildProcessSpawner,
ChildProcessSpawner.make(() => Effect.die("unexpected child process spawn")),
),
FileSystem.layerNoop({ exists: () => Effect.succeed(true) }),
Layer.succeed(ChildProcessSpawner.ChildProcessSpawner, spawner),
Layer.succeed(
HttpClient.HttpClient,
HttpClient.make(() => Effect.die("unexpected HTTP request")),
HttpClient.make((request) =>
Effect.succeed(
HttpClientResponse.fromWeb(request, new Response(null, { status: 200 })),
),
),
),
Layer.succeed(DesktopObservability.DesktopBackendOutputLogFactory, {
forInstance: () =>
Expand All @@ -78,7 +84,7 @@ function makePoolLayer(
updateCancellations: Stream.empty,
}),
Layer.succeed(DesktopBackendConfiguration.DesktopBackendConfiguration, {
resolvePrimary: Effect.die("unexpected primary config resolve"),
resolvePrimary: Effect.succeed(backendConfig),
resolvePrimaryLabel: Ref.get(labelRef),
resolveWsl: () => Effect.die("unexpected WSL config resolve"),
} satisfies DesktopBackendConfiguration.DesktopBackendConfiguration["Service"]),
Expand Down Expand Up @@ -106,7 +112,151 @@ function makePoolLayer(
);
}

const backendConfig: DesktopBackendStartConfig = {
executablePath: "/backend",
args: [],
entryPath: "/backend",
cwd: "/",
env: {},
bootstrap: {
mode: "desktop",
noBrowser: true,
port: 3773,
t3Home: "/tmp/t3-test",
host: "127.0.0.1",
desktopBootstrapToken: "test-token",
tailscaleServeEnabled: false,
tailscaleServePort: 443,
},
bootstrapDelivery: "stdin",
extendEnv: false,
httpBaseUrl: new URL("http://127.0.0.1:3773"),
captureOutput: false,
preflightFailure: Option.none(),
};

const secondarySpec: DesktopBackendPool.BackendInstanceSpec = {
id: DesktopBackendPool.BackendInstanceId("wsl:Ubuntu"),
label: Effect.succeed("WSL (Ubuntu)"),
configResolve: Effect.succeed(backendConfig),
};

const makeProcessHarness = Effect.gen(function* () {
const started = yield* Queue.unbounded<{
readonly pid: number;
readonly exit: Deferred.Deferred<ChildProcessSpawner.ExitCode>;
}>();
const alive = new Set<number>();
let nextPid = 100;
const spawner = ChildProcessSpawner.make(() =>
Effect.gen(function* () {
const pid = nextPid++;
const exit = yield* Deferred.make<ChildProcessSpawner.ExitCode>();
yield* Effect.acquireRelease(
Effect.sync(() => alive.add(pid)),
() =>
Effect.gen(function* () {
alive.delete(pid);
yield* Deferred.succeed(exit, ChildProcessSpawner.ExitCode(0));
}),
);
yield* Queue.offer(started, { pid, exit });
return ChildProcessSpawner.makeHandle({
pid: ChildProcessSpawner.ProcessId(pid),
stdin: Sink.drain,
stdout: Stream.empty,
stderr: Stream.empty,
all: Stream.empty,
exitCode: Deferred.await(exit),
isRunning: Effect.sync(() => alive.has(pid)),
kill: () => Deferred.succeed(exit, ChildProcessSpawner.ExitCode(0)).pipe(Effect.asVoid),
getInputFd: () => Sink.drain,
getOutputFd: () => Stream.empty,
unref: Effect.succeed(Effect.void),
});
}),
);
const label = yield* Ref.make("Windows");
return { started, alive, layer: makePoolLayer(label, spawner) };
});

describe("DesktopBackendPool", () => {
it.effect("unregister stops only the secondary and allows it to be enabled again", () =>
Effect.gen(function* () {
const harness = yield* makeProcessHarness;
yield* Effect.gen(function* () {
const pool = yield* DesktopBackendPool.DesktopBackendPool;
const primary = yield* pool.primary;
yield* primary.start;
const primaryProcess = yield* Queue.take(harness.started);
const secondary = yield* pool.register(secondarySpec);
yield* secondary.start;
const secondaryProcess = yield* Queue.take(harness.started);

yield* pool.unregister(secondary.id);

assert.isFalse((yield* secondary.snapshot).desiredRunning);
assert.isTrue(Option.isNone((yield* secondary.snapshot).activePid));
assert.isFalse((yield* secondary.snapshot).restartScheduled);
assert.isTrue(Option.isNone(yield* pool.get(secondary.id)));
assert.isFalse(harness.alive.has(secondaryProcess.pid));
assert.deepEqual([...harness.alive], [primaryProcess.pid]);
assert.isTrue((yield* primary.snapshot).desiredRunning);

const replacement = yield* pool.register(secondarySpec);
yield* replacement.start;
const replacementProcess = yield* Queue.take(harness.started);
assert.notEqual(replacementProcess.pid, secondaryProcess.pid);
assert.sameMembers([...harness.alive], [primaryProcess.pid, replacementProcess.pid]);
yield* pool.unregister(replacement.id);
assert.deepEqual([...harness.alive], [primaryProcess.pid]);
}).pipe(Effect.provide(harness.layer));
assert.equal(harness.alive.size, 0);
}),
);

it.effect("unregister cancels a pending secondary restart", () =>
Effect.gen(function* () {
const label = yield* Ref.make("Windows");
let resolves = 0;
yield* Effect.gen(function* () {
const pool = yield* DesktopBackendPool.DesktopBackendPool;
const secondary = yield* pool.register({
...secondarySpec,
configResolve: Effect.sync(() => {
resolves += 1;
return {
...backendConfig,
preflightFailure: Option.some({ reason: "WSL is starting", fatal: false }),
};
}),
});
yield* secondary.start;
assert.isTrue((yield* secondary.snapshot).restartScheduled);
yield* pool.unregister(secondary.id);
assert.isFalse((yield* secondary.snapshot).desiredRunning);
assert.isFalse((yield* secondary.snapshot).restartScheduled);
yield* TestClock.adjust(Duration.minutes(1));
assert.equal(resolves, 1);
}).pipe(Effect.provide(makePoolLayer(label)));
}),
);

it.effect("pool shutdown releases running primary and secondary processes", () =>
Effect.gen(function* () {
const harness = yield* makeProcessHarness;
yield* Effect.gen(function* () {
const pool = yield* DesktopBackendPool.DesktopBackendPool;
yield* (yield* pool.primary).start;
yield* Queue.take(harness.started);
yield* (yield* pool.register(secondarySpec)).start;
yield* Queue.take(harness.started);
assert.equal(harness.alive.size, 2);
}).pipe(Effect.provide(harness.layer));
assert.equal(harness.alive.size, 0);
}),
);

it.effect("layerTest exposes registered instances by id", () =>
Effect.gen(function* () {
const pool = yield* DesktopBackendPool.DesktopBackendPool;
Expand Down
7 changes: 4 additions & 3 deletions apps/desktop/src/backend/DesktopBackendPool.ts
Original file line number Diff line number Diff line change
Expand Up @@ -337,12 +337,13 @@ export const layer = Layer.effect(
] as const);
}
return Effect.gen(function* () {
// Provide the captured factory services first, then the child scope
// last so instance finalizers are owned by the unregisterable scope.
// The captured runtime context also contains the pool's Scope, despite
// the narrower factory requirements type. Override it inside that
// context so unregister owns the instance's process and restart loop.
const instanceScope = yield* Scope.fork(layerScope, "sequential");
const instance = yield* DesktopBackendManager.makeBackendInstance(spec).pipe(
Effect.provide(factoryContext),
Scope.provide(instanceScope),
Effect.provide(factoryContext),
);
const next = new Map(current);
next.set(spec.id, {
Expand Down
Loading