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
9 changes: 8 additions & 1 deletion apps/server/src/environment/RemoteOpenTargets.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,10 @@ import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner";

import { makeStaleWhileRevalidate } from "./staleWhileRevalidate.ts";

const SSH_PORT = 22;
const TARGETS_CACHE_TTL = "60 seconds";

export class RemoteOpenTargets extends Context.Service<
RemoteOpenTargets,
Expand All @@ -31,7 +34,7 @@ export const make = Effect.gen(function* () {
const spawner = yield* ChildProcessSpawner.ChildProcessSpawner;
const net = yield* NetService.NetService;

const resolveTargets = Effect.gen(function* () {
const discoverTargets = Effect.gen(function* () {
// No local sshd means no name can work; advertise nothing so clients
// render a clear "no SSH route" state instead of links that hang.
// Check both loopback families: sshd can be bound IPv6-only.
Expand Down Expand Up @@ -67,6 +70,10 @@ export const make = Effect.gen(function* () {
return targets;
});

// Advertised on every client connect; probing sshd and tailscaled is slow
// on a loaded host, so serve the last result while it refreshes.
const resolveTargets = yield* makeStaleWhileRevalidate(discoverTargets, TARGETS_CACHE_TTL);

return RemoteOpenTargets.of({ resolveTargets: () => resolveTargets });
});

Expand Down
140 changes: 140 additions & 0 deletions apps/server/src/environment/staleWhileRevalidate.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,140 @@
import { assert, it } from "@effect/vitest";
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 Ref from "effect/Ref";
import * as Scope from "effect/Scope";
import * as TestClock from "effect/testing/TestClock";

import { makeStaleWhileRevalidate } from "./staleWhileRevalidate.ts";

it.effect("computes once, then answers from memory until the ttl lapses", () =>
Effect.gen(function* () {
const calls = yield* Ref.make(0);
const read = yield* makeStaleWhileRevalidate(
Ref.updateAndGet(calls, (count) => count + 1),
"60 seconds",
);

assert.equal(yield* read, 1);
yield* TestClock.adjust("59 seconds");
assert.equal(yield* read, 1);
assert.equal(yield* Ref.get(calls), 1);
}).pipe(Effect.scoped),
);

it.effect("serves the stale value while one background refresh replaces it", () =>
Effect.gen(function* () {
const calls = yield* Ref.make(0);
const release = yield* Deferred.make<void>();
const read = yield* makeStaleWhileRevalidate(
Effect.gen(function* () {
const call = yield* Ref.updateAndGet(calls, (count) => count + 1);
if (call > 1) yield* Deferred.await(release);
return call;
}),
"60 seconds",
);

assert.equal(yield* read, 1);
yield* TestClock.adjust("61 seconds");
// The refresh is blocked, yet callers are answered immediately and only
// one refresh starts however many callers arrive.
assert.equal(yield* read, 1);
assert.equal(yield* read, 1);
yield* TestClock.adjust("1 second");
assert.equal(yield* Ref.get(calls), 2);

yield* Deferred.succeed(release, undefined);
yield* TestClock.adjust("1 second");
assert.equal(yield* read, 2);
}),
);

it.effect("shares one first scan that an interrupted caller does not cancel", () =>
Effect.gen(function* () {
const calls = yield* Ref.make(0);
const release = yield* Deferred.make<void>();
const read = yield* makeStaleWhileRevalidate(
Ref.updateAndGet(calls, (count) => count + 1).pipe(Effect.tap(() => Deferred.await(release))),
"60 seconds",
);

// A connect that times out or disconnects mid-scan leaves the scan running.
const interrupted = yield* read.pipe(Effect.forkChild);
yield* Effect.yieldNow;
yield* Fiber.interrupt(interrupted);

const next = yield* read.pipe(Effect.forkChild);
yield* Effect.yieldNow;
yield* Deferred.succeed(release, undefined);
assert.equal(yield* Fiber.join(next), 1);
assert.equal(yield* Ref.get(calls), 1);
}).pipe(Effect.scoped),
);

it.effect("starts over after a failed first scan", () =>
Effect.gen(function* () {
const calls = yield* Ref.make(0);
const read = yield* makeStaleWhileRevalidate(
Ref.updateAndGet(calls, (count) => count + 1).pipe(
Effect.flatMap((call) => (call === 1 ? Effect.die("probe crashed") : Effect.succeed(call))),
),
"60 seconds",
);

assert.isTrue(Exit.isFailure(yield* Effect.exit(read)));
assert.equal(yield* read, 2);
}).pipe(Effect.scoped),
);

it.effect("abandons a hung refresh so a later call can start another", () =>
Effect.gen(function* () {
const calls = yield* Ref.make(0);
const read = yield* makeStaleWhileRevalidate(
Ref.updateAndGet(calls, (count) => count + 1).pipe(
Effect.flatMap((call) => (call === 2 ? Effect.never : Effect.succeed(call))),
),
"60 seconds",
);

assert.equal(yield* read, 1);
yield* TestClock.adjust("61 seconds");
assert.equal(yield* read, 1);
yield* TestClock.adjust("1 second");
assert.equal(yield* Ref.get(calls), 2);

// The hung refresh times out, releasing the claim for the next caller.
yield* TestClock.adjust("30 seconds");
assert.equal(yield* read, 1);
yield* TestClock.adjust("1 second");
assert.equal(yield* Ref.get(calls), 3);
assert.equal(yield* read, 3);
}).pipe(Effect.scoped),
);

it.effect("interrupts a background refresh when the owning scope closes", () =>
Effect.gen(function* () {
const interrupted = yield* Deferred.make<void>();
const scope = yield* Scope.make();
const calls = yield* Ref.make(0);
const read = yield* makeStaleWhileRevalidate(
Ref.updateAndGet(calls, (count) => count + 1).pipe(
Effect.flatMap((call) =>
call === 1
? Effect.succeed(call)
: Effect.never.pipe(Effect.onInterrupt(() => Deferred.succeed(interrupted, undefined))),
),
),
"60 seconds",
).pipe(Scope.provide(scope));

yield* read;
yield* TestClock.adjust("61 seconds");
yield* read;
yield* TestClock.adjust("1 second");
yield* Scope.close(scope, Exit.void);
yield* Deferred.await(interrupted);
}),
);
100 changes: 100 additions & 0 deletions apps/server/src/environment/staleWhileRevalidate.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
import * as Clock from "effect/Clock";
import * as Deferred from "effect/Deferred";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Ref from "effect/Ref";

const REFRESH_TIMEOUT = Duration.seconds(30);

type Entry<A> =
| { readonly _tag: "Empty" }
| { readonly _tag: "Scanning"; readonly scan: Deferred.Deferred<A> }
| {
readonly _tag: "Ready";
readonly value: A;
readonly expiresAtNanos: bigint;
readonly refreshing: boolean;
};

/**
* Memoizes a discovery that is cheap when warm and slow on a loaded host, for
* callers that sit on the client connect path. Until a value exists, callers
* wait on one shared scan. After that every call returns the last good value
* immediately; once the value is older than `ttl`, one refresh replaces it in
* the background.
*
* Scans run on their own fiber in the scope that built the cache, never on a
* caller's. Callers time out and disconnect mid-connect, and neither may cancel
* a scan other callers are waiting on, or throw away work a slow host needs
* more than one connect to finish. Only successes are stored: a failed scan
* leaves the previous state, so the next caller starts over. Expiry uses the
* monotonic clock so a wall-clock adjustment cannot keep a stale entry alive.
*
* Closing the scope interrupts any scan (and cleans up a scoped probe process).
* A refresh is capped at `REFRESH_TIMEOUT` so a probe hung on an unresponsive
* mount cannot leave the entry stale forever; a first scan is not, because
* callers that already timed out still pick up its result later.
*/
export const makeStaleWhileRevalidate = <A>(discover: Effect.Effect<A>, ttl: Duration.Input) => {
const ttlNanos = Duration.toNanosUnsafe(Duration.fromInputUnsafe(ttl));
return Effect.gen(function* () {
const scope = yield* Effect.scope;
const state = yield* Ref.make<Entry<A>>({ _tag: "Empty" });

const settle = (exit: Exit.Exit<A>) =>
Effect.gen(function* () {
const now = yield* Clock.monotonicTimeNanos;
yield* Ref.update(state, (entry): Entry<A> =>
Exit.isSuccess(exit)
? {
_tag: "Ready",
value: exit.value,
expiresAtNanos: now + ttlNanos,
refreshing: false,
}
: entry._tag === "Ready"
? { ...entry, refreshing: false }
: { _tag: "Empty" },
);
});
const fork = (scan: Effect.Effect<unknown>) =>
scan.pipe(Effect.ignoreCause({ log: true }), Effect.interruptible, Effect.forkIn(scope));
const firstScan = (scan: Deferred.Deferred<A>) =>
fork(
discover.pipe(
Effect.onExit((exit) => settle(exit).pipe(Effect.andThen(Deferred.done(scan, exit)))),
),
);
const refresh = fork(
discover.pipe(Effect.onExit(settle), Effect.timeoutOption(REFRESH_TIMEOUT)),
);

return Effect.gen(function* () {
const now = yield* Clock.monotonicTimeNanos;
// Claiming a scan and starting it are one uninterruptible step: a caller
// interrupted between them would leave a claimed scan that never runs.
const result = yield* Ref.modify(
state,
(entry): readonly [readonly [Effect.Effect<unknown>, Effect.Effect<A>], Entry<A>] => {
switch (entry._tag) {
case "Empty": {
const scan = Deferred.makeUnsafe<A>();
return [[firstScan(scan), Deferred.await(scan)], { _tag: "Scanning", scan }];
}
case "Scanning":
return [[Effect.void, Deferred.await(entry.scan)], entry];
case "Ready":
return entry.refreshing || entry.expiresAtNanos > now
? [[Effect.void, Effect.succeed(entry.value)], entry]
: [[refresh, Effect.succeed(entry.value)], { ...entry, refreshing: true }];
}
},
).pipe(
Effect.flatMap(([start, result]) => Effect.as(start, result)),
Effect.uninterruptible,
);
return yield* result;
});
});
};
112 changes: 59 additions & 53 deletions apps/server/src/process/externalLauncher.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1127,65 +1127,71 @@ it.effect.skipIf(windowsHost)("ignores unusable app bundles and keeps PATH launc
}).pipe(Effect.scoped, Effect.provide(NodeServices.layer)),
);

it.effect("memoizes editor discovery and refreshes after the cache window", () => {
let statCalls = 0;
const fileInfo = { type: "File" } as FileSystem.File.Info;
const launcherLayer = ExternalLauncher.layer.pipe(
Layer.provide(
Layer.mergeAll(
FileSystem.layerNoop({
stat: () =>
Effect.sync(() => {
statCalls += 1;
return fileInfo;
}),
}),
Path.layer,
Layer.succeed(
ChildProcessSpawner.ChildProcessSpawner,
ChildProcessSpawner.make(() => Effect.sync(() => makeMockDetachedHandle())),
it.effect(
"memoizes editor discovery and refreshes in the background after the cache window",
() => {
let statCalls = 0;
const fileInfo = { type: "File" } as FileSystem.File.Info;
const launcherLayer = ExternalLauncher.layer.pipe(
Layer.provide(
Layer.mergeAll(
FileSystem.layerNoop({
stat: () =>
Effect.sync(() => {
statCalls += 1;
return fileInfo;
}),
}),
Path.layer,
Layer.succeed(
ChildProcessSpawner.ChildProcessSpawner,
ChildProcessSpawner.make(() => Effect.sync(() => makeMockDetachedHandle())),
),
),
),
),
);
);

return Effect.gen(function* () {
const launcher = yield* ExternalLauncher.ExternalLauncher;
return Effect.gen(function* () {
const launcher = yield* ExternalLauncher.ExternalLauncher;

const first = yield* launcher.resolveAvailableEditors();
assert.equal(first.includes("vscode"), true);
const statCallsAfterFirstScan = statCalls;
assert.isAbove(statCallsAfterFirstScan, 0);

// Past the shared command-resolution cache TTL (30s) but within the
// discovery cache window: the memoized set is reused without any scan.
yield* TestClock.adjust("31 seconds");
const second = yield* launcher.resolveAvailableEditors();
assert.deepEqual([...second], [...first]);
assert.equal(statCalls, statCallsAfterFirstScan);

// Past the discovery cache window the next call rescans.
yield* TestClock.adjust("30 seconds");
yield* launcher.resolveAvailableEditors();
assert.isAbove(statCalls, statCallsAfterFirstScan);
}).pipe(
Effect.provide(
Layer.mergeAll(
launcherLayer,
Layer.succeed(HostProcessPlatform, "win32"),
ConfigProvider.layer(
ConfigProvider.fromEnv({
env: {
PATH: "C:\\t3-editor-discovery-cache-test",
PATHEXT: ".COM;.EXE;.BAT;.CMD",
},
}),
const first = yield* launcher.resolveAvailableEditors();
assert.equal(first.includes("vscode"), true);
const statCallsAfterFirstScan = statCalls;
assert.isAbove(statCallsAfterFirstScan, 0);

// Past the shared command-resolution cache TTL (30s) but within the
// discovery cache window: the memoized set is reused without any scan.
yield* TestClock.adjust("31 seconds");
const second = yield* launcher.resolveAvailableEditors();
assert.deepEqual([...second], [...first]);
assert.equal(statCalls, statCallsAfterFirstScan);

// Past the discovery cache window the next call still answers from the
// memoized set and rescans in the background.
yield* TestClock.adjust("30 seconds");
const third = yield* launcher.resolveAvailableEditors();
assert.deepEqual([...third], [...first]);
yield* Effect.yieldNow;
assert.isAbove(statCalls, statCallsAfterFirstScan);
}).pipe(
Effect.provide(
Layer.mergeAll(
launcherLayer,
Layer.succeed(HostProcessPlatform, "win32"),
ConfigProvider.layer(
ConfigProvider.fromEnv({
env: {
PATH: "C:\\t3-editor-discovery-cache-test",
PATHEXT: ".COM;.EXE;.BAT;.CMD",
},
}),
),
TestClock.layer(),
),
TestClock.layer(),
),
),
);
});
);
},
);

// 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
Expand Down
Loading
Loading