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
41 changes: 21 additions & 20 deletions apps/server/src/process/externalLauncher.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import * as NodePath from "node:path";
import * as NodeServices from "@effect/platform-node/NodeServices";
import { assert, it } from "@effect/vitest";
import * as ConfigProvider from "effect/ConfigProvider";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as FileSystem from "effect/FileSystem";
Expand Down Expand Up @@ -1186,26 +1187,24 @@ it.effect("memoizes editor discovery and refreshes after the cache window", () =
);
});

// A client that disconnects mid-scan interrupts the shared discovery effect on
// the connection fiber. The cache must not retain that interrupt: doing so
// replayed it to every later connect for the whole TTL, so `server.getConfig`
// failed and no client could reconnect until the server restarted.
it.effect("rescans after an interrupted discovery instead of caching the interrupt", () => {
// Connects run discovery under a timeout and may disconnect mid-scan. Neither
// may cancel the scan: on a busy host every connect would time out partway
// through, cache nothing, and leave every client without editors.
it.effect("keeps scanning after the caller is interrupted and shares that scan", () => {
const fileInfo = { type: "File" } as FileSystem.File.Info;
let blockFirstScan = true;
let scans = 0;
const release = Deferred.makeUnsafe<void>();
let parkedStats = 0;
const launcherLayer = ExternalLauncher.layer.pipe(
Layer.provide(
Layer.mergeAll(
FileSystem.layerNoop({
// The first scan parks inside `stat` so the interrupt lands while
// discovery is in flight, which is what a client disconnecting
// mid-connect does to the shared effect.
// Scans park inside `stat` until released, so the interrupt lands
// while discovery is in flight.
stat: () =>
Effect.gen(function* () {
scans += 1;
if (blockFirstScan) {
return yield* Effect.never;
if (!Deferred.isDoneUnsafe(release)) {
parkedStats += 1;
yield* Deferred.await(release);
}
return fileInfo;
}),
Expand All @@ -1222,16 +1221,18 @@ it.effect("rescans after an interrupted discovery instead of caching the interru
return Effect.gen(function* () {
const launcher = yield* ExternalLauncher.ExternalLauncher;

const fiber = yield* Effect.forkChild(launcher.resolveAvailableEditors());
const interrupted = yield* Effect.forkChild(launcher.resolveAvailableEditors());
yield* Effect.yieldNow;
yield* Fiber.interrupt(fiber);
yield* Fiber.interrupt(interrupted);

// The next connect must still get a real answer well inside the TTL.
blockFirstScan = false;
scans = 0;
const editors = yield* launcher.resolveAvailableEditors();
// The next connect joins the running scan instead of starting its own.
const next = yield* Effect.forkChild(launcher.resolveAvailableEditors());
yield* Effect.yieldNow;
assert.equal(parkedStats, 1);

yield* Deferred.succeed(release, undefined);
const editors = yield* Fiber.join(next);
assert.equal(editors.includes("vscode"), true);
assert.isAbove(scans, 0);
}).pipe(
Effect.provide(
Layer.mergeAll(
Expand Down
79 changes: 56 additions & 23 deletions apps/server/src/process/externalLauncher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,13 +28,16 @@ import {
import * as Clock from "effect/Clock";
import * as Config from "effect/Config";
import * as Context from "effect/Context";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Encoding from "effect/Encoding";
import * as Exit from "effect/Exit";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Path from "effect/Path";
import * as Ref from "effect/Ref";
import * as Scope from "effect/Scope";
import * as Stream from "effect/Stream";
import * as ChildProcess from "effect/unstable/process/ChildProcess";
import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner";
Expand Down Expand Up @@ -462,21 +465,21 @@ const resolveFileManagerRevealKind = Effect.fn("externalLauncher.resolveFileMana
// the discovered set for a bounded window so repeat connects skip even the
// per-command cache lookups in @t3tools/shared/shell.
//
// This deliberately does not use `Effect.cachedWithTTL`: that memoizes the
// first caller's Exit whatever it is, including an interrupt. Callers run this
// on the connection fiber under a timeout (`resolveAvailableEditorsForConfig`),
// so one client disconnecting mid-scan would cache the interrupt and replay it
// to every later connect for the whole TTL, breaking `server.getConfig`
// permanently. Storing only on success means an interrupted scan leaves the
// cache untouched and the next connect simply rescans.
// The scan runs on its own fiber in the service scope, and every caller awaits
// that one scan. Callers apply a timeout (`resolveAvailableEditorsForConfig`)
// and disconnect mid-connect; neither may cancel a scan other connects are
// waiting on, or throw away work a slow host (a busy server at startup, a
// long PATH) needs more than one connect to finish. A failed scan clears the
// entry so the next caller starts over rather than replaying the failure.
// Expiry uses the monotonic clock (Clock.currentTimeNanos), matching the
// command-resolution cache in @t3tools/shared/shell, so a backward wall-clock
// adjustment cannot keep an expired entry alive.
const EDITOR_DISCOVERY_CACHE_TTL_NANOS = 60_000_000_000n;

interface EditorDiscoveryCacheEntry {
readonly editors: ReadonlyArray<EditorId>;
readonly expiresAtNanos: bigint;
readonly scan: Deferred.Deferred<ReadonlyArray<EditorId>>;
/** Undefined while the scan is still running. */
readonly expiresAtNanos: bigint | undefined;
}

/**
Expand Down Expand Up @@ -760,27 +763,57 @@ export const make = Effect.gen(function* () {
Effect.provideService(Path.Path, path),
);

const scope = yield* Scope.Scope;
const editorDiscoveryCache = yield* Ref.make<Option.Option<EditorDiscoveryCacheEntry>>(
Option.none(),
);
const cachedAvailableEditors = Effect.gen(function* () {
const nowNanos = yield* Clock.currentTimeNanos;
const entry = yield* Ref.get(editorDiscoveryCache);
if (Option.isSome(entry) && entry.value.expiresAtNanos > nowNanos) {
return entry.value.editors;
}
const editors = yield* provideCommandResolutionServices(resolveAvailableEditors()).pipe(
const runEditorDiscovery = (scan: Deferred.Deferred<ReadonlyArray<EditorId>>) =>
provideCommandResolutionServices(resolveAvailableEditors()).pipe(
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
Effect.onExit((exit) =>
Effect.gen(function* () {
const expiresAtNanos = (yield* Clock.currentTimeNanos) + EDITOR_DISCOVERY_CACHE_TTL_NANOS;
yield* Ref.update(editorDiscoveryCache, (current) =>
Option.isNone(current) || current.value.scan !== scan
? current
: Exit.isSuccess(exit)
? Option.some({ scan, expiresAtNanos })
: Option.none(),
);
yield* Deferred.done(scan, exit);
}),
),
Effect.interruptible,
Effect.forkIn(scope),
);
yield* Ref.set(
// Claiming the cache entry and starting its scan must not be split by an
// interrupt, or the entry would wait on a scan that never runs.
const acquireEditorDiscovery = Effect.gen(function* () {
const nowNanos = yield* Clock.currentTimeNanos;
const [scan, isNewScan] = yield* Ref.modify(
editorDiscoveryCache,
Option.some({
editors,
expiresAtNanos: nowNanos + EDITOR_DISCOVERY_CACHE_TTL_NANOS,
}),
(
current,
): [
[EditorDiscoveryCacheEntry["scan"], boolean],
Option.Option<EditorDiscoveryCacheEntry>,
] => {
if (
Option.isSome(current) &&
(current.value.expiresAtNanos === undefined || current.value.expiresAtNanos > nowNanos)
) {
return [[current.value.scan, false], current];
}
const scan = Deferred.makeUnsafe<ReadonlyArray<EditorId>>();
return [[scan, true], Option.some({ scan, expiresAtNanos: undefined })];

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Medium process/externalLauncher.ts:808

A discovery scan that never completes permanently blocks resolveAvailableEditors(): every later caller reuses the same unfinished Deferred and no replacement scan can start until restart. This happens because the pending entry stores expiresAtNanos: undefined and is cleared only by the scan's onExit; give pending scans a finite expiry so callers can start a replacement after the cache TTL.

Suggested change
return [[scan, true], Option.some({ scan, expiresAtNanos: undefined })];
return [[scan, true], Option.some({ scan, expiresAtNanos: nowNanos + EDITOR_DISCOVERY_CACHE_TTL_NANOS })];
🤖 Copy this AI Prompt to have your agent fix this:
In file @apps/server/src/process/externalLauncher.ts around line 808:

A discovery scan that never completes permanently blocks `resolveAvailableEditors()`: every later caller reuses the same unfinished `Deferred` and no replacement scan can start until restart. This happens because the pending entry stores `expiresAtNanos: undefined` and is cleared only by the scan's `onExit`; give pending scans a finite expiry so callers can start a replacement after the cache TTL.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not changing this. A scan that never settles can't block anything unboundedly: its only callers are server.getConfig/the config snapshot, which still bound it with the 5 s resolveAvailableEditorsForConfig timeout and degrade to no editors (the pre-PR behavior), and the late-snapshot wait, which ends with its subscription.

The suggested expiry would make things worse in the one case where a scan really hangs: a stat stuck on a dead network mount or PATH entry. A replacement scan every 60 s would hit the same entry and hang too, and each stuck fs.stat holds a libuv threadpool thread (4 by default), so piling up scans would starve every other filesystem call on the server. Keeping at most one scan in flight is deliberate. It also races with in-flight scans that take longer than the TTL, starting duplicates on a slow host, which is the situation this PR is about.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sorry, I'm unable to act on this request because you do not have permissions within this repository.

},
);
return editors;
});
if (isNewScan) {
yield* runEditorDiscovery(scan);
}
return scan;
}).pipe(Effect.uninterruptible);
const cachedAvailableEditors = Effect.flatMap(acquireEditorDiscovery, Deferred.await);

return ExternalLauncher.of({
resolveAvailableEditors: () => cachedAvailableEditors,
Expand Down
188 changes: 187 additions & 1 deletion apps/server/src/ws.test.ts
Original file line number Diff line number Diff line change
@@ -1,15 +1,29 @@
import { assert, it } from "@effect/vitest";
import { ORCHESTRATION_PROTOCOL_VERSION } from "@t3tools/contracts";
import {
ORCHESTRATION_PROTOCOL_VERSION,
type ServerConfig,
type ServerConfigStreamEvent,
} from "@t3tools/contracts";
import { HostProcessPlatform } from "@t3tools/shared/hostProcess";
import * as ConfigProvider from "effect/ConfigProvider";
import * as Deferred from "effect/Deferred";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Path from "effect/Path";
import * as Queue from "effect/Queue";
import * as Stream from "effect/Stream";
import * as TestClock from "effect/testing/TestClock";
import { ChildProcessSpawner } from "effect/unstable/process";

import * as ExternalLauncher from "./process/externalLauncher.ts";
import {
hasCompatibleOrchestrationProtocol,
resolveAvailableEditorsForConfig,
shouldUseBoundedThreadSnapshot,
withLateEditorConfig,
} from "./ws.ts";

it("accepts only the current orchestration protocol before websocket RPC setup", () => {
Expand Down Expand Up @@ -48,3 +62,175 @@ it.effect("does not block server config when editor discovery never resolves", (
assert.deepEqual(availableEditors, []);
}),
);

// Only the fields the late-editor fold reads or rewrites.
const snapshotConfig = (fields: Partial<ServerConfig>) =>
({ availableEditors: [], settings: {}, ...fields }) as unknown as ServerConfig;

const settingsUpdated = (settings: object): ServerConfigStreamEvent => ({
version: 1,
type: "settingsUpdated",
payload: { settings: settings as ServerConfig["settings"] },
});

it.effect("resends late editors without rolling back updates already sent", () =>
Effect.gen(function* () {
const settingsSent = yield* Deferred.make<void>();
const events = yield* withLateEditorConfig(
snapshotConfig({ settings: { enableProviderUpdateChecks: true } as never }),
Stream.make(settingsUpdated({ enableProviderUpdateChecks: false })),
{
resolveAvailableEditors: () => Effect.succeed(["file-manager"]),
// Holds the late snapshot until the settings change has gone out.
resolveFileManagerRevealKind: () =>
Deferred.await(settingsSent).pipe(Effect.as("file-explorer" as const)),
},
).pipe(
Stream.tap((event) =>
event.type === "settingsUpdated" ? Deferred.succeed(settingsSent, undefined) : Effect.void,
),
Stream.runCollect,
);

const [first, second] = Array.from(events);
assert.equal(events.length, 2);
assert.equal(first?.type, "settingsUpdated");
assert.equal(second?.type, "snapshot");
if (second?.type === "snapshot") {
assert.deepEqual(second.config.availableEditors, ["file-manager"]);
assert.equal(second.config.shellRevealInFileManagerKind, "file-explorer");
assert.deepEqual(second.config.settings, { enableProviderUpdateChecks: false } as never);
}
}),
);

it.effect("sends no late snapshot when the scan matches the snapshot", () =>
Effect.gen(function* () {
const events = yield* withLateEditorConfig(
snapshotConfig({ availableEditors: ["vscode"] }),
Stream.empty,
{
resolveAvailableEditors: () => Effect.succeed(["vscode"]),
resolveFileManagerRevealKind: () => Effect.succeed(undefined),
},
).pipe(Stream.runCollect);

assert.equal(events.length, 0);
}),
);

it.effect("resends a file manager reveal kind that missed the snapshot", () =>
Effect.gen(function* () {
const events = yield* withLateEditorConfig(
snapshotConfig({ availableEditors: ["file-manager"] }),
Stream.empty,
{
resolveAvailableEditors: () => Effect.succeed(["file-manager"]),
resolveFileManagerRevealKind: () => Effect.succeed("file-explorer"),
},
).pipe(Stream.runCollect);

const [late] = Array.from(events);
assert.equal(events.length, 1);
assert.equal(late?.type, "snapshot");
if (late?.type === "snapshot") {
assert.equal(late.config.shellRevealInFileManagerKind, "file-explorer");
}
}),
);

// The real launcher on Windows over a filesystem whose probes park until
// released, like a host too busy to finish discovery inside the snapshot timeout.
const makeParkedWindowsLauncher = Effect.gen(function* () {
const parkedProbes = yield* Queue.unbounded<void>();
const release = yield* Deferred.make<void>();
const launcher = yield* ExternalLauncher.make.pipe(
Effect.provide(
Layer.mergeAll(
FileSystem.layerNoop({
stat: () =>
Queue.offer(parkedProbes, undefined).pipe(
Effect.andThen(Deferred.await(release)),
Effect.as({ type: "File" } as FileSystem.File.Info),
),
}),
Path.layer,
Layer.succeed(
ChildProcessSpawner.ChildProcessSpawner,
ChildProcessSpawner.make(() => Effect.die("unexpected spawn")),
),
),
),
);
const onWindows = <A, E, R>(effect: Effect.Effect<A, E, R>) =>
effect.pipe(
Effect.provideService(HostProcessPlatform, "win32"),
Effect.provide(
ConfigProvider.layer(
ConfigProvider.fromEnv({
env: { PATH: "C:\\t3-late-editors-test", PATHEXT: ".EXE" },
}),
),
),
);
return {
editors: {
resolveAvailableEditors: () => onWindows(launcher.resolveAvailableEditors()),
resolveFileManagerRevealKind: () => onWindows(launcher.resolveFileManagerRevealKind()),
},
probeParked: Queue.take(parkedProbes),
releaseProbes: Deferred.succeed(release, undefined),
};
});

it.effect("recovers editors after a real scan outlasts the config timeout", () =>
Effect.gen(function* () {
const { editors, probeParked, releaseProbes } = yield* makeParkedWindowsLauncher;

const snapshotFiber = yield* resolveAvailableEditorsForConfig(
editors.resolveAvailableEditors(),
).pipe(Effect.forkChild);
yield* probeParked;
yield* TestClock.adjust(Duration.seconds(5));
const snapshotEditors = yield* Fiber.join(snapshotFiber);
assert.deepEqual(snapshotEditors, []);

const lateFiber = yield* withLateEditorConfig(
snapshotConfig({ availableEditors: snapshotEditors }),
Stream.empty,
editors,
).pipe(Stream.runCollect, Effect.forkChild);
yield* releaseProbes;

const [late] = Array.from(yield* Fiber.join(lateFiber));
assert.equal(late?.type, "snapshot");
if (late?.type === "snapshot") {
assert.equal(late.config.availableEditors.includes("vscode"), true);
}
}).pipe(Effect.scoped),
);

it.effect("recovers a reveal kind whose real probe outlasts the config timeout", () =>
Effect.gen(function* () {
const { editors, probeParked, releaseProbes } = yield* makeParkedWindowsLauncher;

// The snapshot's bounded probe timed out: file manager, but no reveal kind.
const lateFiber = yield* withLateEditorConfig(
snapshotConfig({ availableEditors: ["file-manager"] }),
Stream.empty,
{
resolveAvailableEditors: () => Effect.succeed(["file-manager"]),
resolveFileManagerRevealKind: editors.resolveFileManagerRevealKind,
},
).pipe(Stream.runCollect, Effect.forkChild);
yield* probeParked;
yield* TestClock.adjust(Duration.seconds(6));
yield* releaseProbes;

const [late] = Array.from(yield* Fiber.join(lateFiber));
assert.equal(late?.type, "snapshot");
if (late?.type === "snapshot") {
assert.equal(late.config.shellRevealInFileManagerKind, "file-explorer");
}
}).pipe(Effect.scoped),
);
Loading
Loading