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
52 changes: 51 additions & 1 deletion apps/desktop/src/electron/ElectronProtocol.test.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,11 @@
import { assert, describe, it } from "@effect/vitest";
import * as Cause from "effect/Cause";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as NodeServices from "@effect/platform-node/NodeServices";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as TestClock from "effect/testing/TestClock";
import { beforeEach, vi } from "vite-plus/test";

const { handleMock, netFetchMock, unhandleMock } = vi.hoisted(() => ({
Expand Down Expand Up @@ -174,7 +176,11 @@ describe("ElectronProtocol", () => {
targetOrigin: new URL("http://127.0.0.1:5733/"),
clerkFrontendApiHostname: undefined,
});
return yield* Effect.promise(() => handler!(new Request("t3code-dev://app/")));
const fiber = yield* Effect.forkChild(
Effect.promise(() => handler!(new Request("t3code-dev://app/"))),
);
yield* TestClock.adjust("50 millis");
return yield* Fiber.join(fiber);
}),
);

Expand All @@ -183,6 +189,50 @@ describe("ElectronProtocol", () => {
}).pipe(Effect.provide(layerProtocol)),
);

it.effect("rejects with the last renderer target failure after 50ms and 150ms retries", () =>
Effect.gen(function* () {
let handler: ((request: Request) => Promise<Response>) | undefined;
handleMock.mockImplementation((_scheme, nextHandler) => {
handler = nextHandler;
});
const lastFailure = new Error("connect ECONNREFUSED 127.0.0.1:5733 (3)");
netFetchMock
.mockRejectedValueOnce(new Error("connect ECONNREFUSED 127.0.0.1:5733 (1)"))
.mockRejectedValueOnce(new Error("connect ECONNREFUSED 127.0.0.1:5733 (2)"))
.mockRejectedValueOnce(lastFailure);

const rejection = yield* Effect.scoped(
Effect.gen(function* () {
const protocol = yield* ElectronProtocol.ElectronProtocol;
yield* protocol.registerDesktopProtocol({
scheme: "t3code-dev",
targetOrigin: new URL("http://127.0.0.1:5733/"),
clerkFrontendApiHostname: undefined,
});
const fiber = yield* Effect.forkChild(
Effect.promise(() =>
handler!(new Request("t3code-dev://app/")).then(
() => null,
(error: unknown) => error,
),
),
);
yield* TestClock.adjust("49 millis");
assert.equal(netFetchMock.mock.calls.length, 1);
yield* TestClock.adjust("1 millis");
assert.equal(netFetchMock.mock.calls.length, 2);
yield* TestClock.adjust("149 millis");
assert.equal(netFetchMock.mock.calls.length, 2);
yield* TestClock.adjust("1 millis");
return yield* Fiber.join(fiber);
}),
);

assert.strictEqual(rejection, lastFailure);
assert.equal(netFetchMock.mock.calls.length, 3);
}).pipe(Effect.provide(layerProtocol)),
);

it.effect("preserves protocol registration failures", () =>
Effect.gen(function* () {
const cause = new Error("protocol registration failed");
Expand Down
59 changes: 32 additions & 27 deletions apps/desktop/src/electron/ElectronProtocol.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,10 @@ 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 NodeTimersPromises from "node:timers/promises";
import * as Path from "effect/Path";
import * as Mime from "effect/http/Mime";
import * as Ref from "effect/Ref";
import * as Schedule from "effect/Schedule";
import * as Schema from "effect/Schema";
import * as Scope from "effect/Scope";

Expand Down Expand Up @@ -150,11 +150,29 @@ const registerDesktopSchemePrivileges = Effect.sync(registerDesktopSchemePrivile

export const layerSchemePrivileges = Layer.effectDiscard(registerDesktopSchemePrivileges);

async function proxyRequest(
class ElectronProtocolFetchError extends Schema.TaggedError<ElectronProtocolFetchError>()(
"ElectronProtocolFetchError",
{ cause: Schema.Defect() },
) {}

const netFetch = (url: string, init: RequestInit) =>
Effect.tryPromise({
try: () => Electron.net.fetch(url, init),
catch: (cause) => new ElectronProtocolFetchError({ cause }),
});

// The dev renderer target can briefly refuse connections while Vite restarts:
// retry idempotent requests after 50ms, then 150ms, and keep the last failure.
const fetchWithTransientRetry = (url: string, init: RequestInit) =>
netFetch(url, init).pipe(
Effect.retry({ schedule: Schedule.exponential("50 millis", 3), times: 2 }),
);

const proxyRequest = Effect.fn("desktop.protocol.proxyRequest")(function* (
request: Request,
targetOrigin: URL,
contentSecurityPolicy: string,
): Promise<Response> {
) {
const requestUrl = new URL(request.url);
if (requestUrl.host !== DESKTOP_HOST) {
return new Response(null, { status: 404 });
Expand Down Expand Up @@ -190,12 +208,10 @@ async function proxyRequest(
}
const response =
request.method === "GET" || request.method === "HEAD"
? await fetchWithTransientRetry(targetUrl.toString(), init)
: await Electron.net.fetch(targetUrl.toString(), init);
? yield* fetchWithTransientRetry(targetUrl.toString(), init)
: yield* netFetch(targetUrl.toString(), init);
return withContentSecurityPolicy(response, contentSecurityPolicy);
}

const TRANSIENT_FETCH_RETRY_DELAYS_MS = [0, 50, 150] as const;
});

// Serves the packaged web client without a backend: files resolve within the
// asset directory, and any other path falls back to index.html so the SPA
Expand Down Expand Up @@ -238,24 +254,6 @@ const serveDesktopAsset = Effect.fn("desktop.protocol.serveAsset")(function* (
});
});

async function fetchWithTransientRetry(url: string, init: RequestInit): Promise<Response> {
let lastError: unknown;

for (const delayMs of TRANSIENT_FETCH_RETRY_DELAYS_MS) {
if (delayMs > 0) {
await NodeTimersPromises.setTimeout(delayMs);
}

try {
return await Electron.net.fetch(url, init);
} catch (error) {
lastError = error;
}
}

throw lastError;
}

/** @public Service construction is part of the canonical Effect module API. */
export const make = Effect.gen(function* () {
const registered = yield* Ref.make(false);
Expand All @@ -278,7 +276,14 @@ export const make = Effect.gen(function* () {
contentSecurityPolicy,
);
}
return proxyRequest(request, input.targetOrigin, contentSecurityPolicy);
// Reject with net.fetch's own error, as an unproxied fetch would.
return runPromise(
proxyRequest(request, input.targetOrigin, contentSecurityPolicy).pipe(
Effect.catchTags({
ElectronProtocolFetchError: (error) => Effect.die(error.cause),
}),
),
);
});
},
catch: (cause) => new ElectronProtocolRegistrationError({ scheme: input.scheme, cause }),
Expand Down
2 changes: 2 additions & 0 deletions apps/desktop/src/ssh/DesktopSshEnvironment.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import {
} from "@t3tools/ssh/errors";
import * as SshTunnel from "@t3tools/ssh/tunnel";
import * as Context from "effect/Context";
import type * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
Expand All @@ -29,6 +30,7 @@ import * as DesktopSshPasswordPrompts from "./DesktopSshPasswordPrompts.ts";

export type DesktopSshEnvironmentRuntimeServices =
| ChildProcessSpawner.ChildProcessSpawner
| Crypto.Crypto
| FileSystem.FileSystem
| Path.Path
| HttpClient.HttpClient
Expand Down
1 change: 1 addition & 0 deletions apps/server/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
"@effect/platform-node-shared": "catalog:",
"@ff-labs/fff-node": "0.9.4",
"@napi-rs/keyring": "^1.3.0",
"@noble/hashes": "catalog:",
"@opencode-ai/sdk": "^1.3.15",
"@opencode/client": "2.0.23",
"@opencode/protocol": "2.0.23",
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/auth/ReusableDevAuth.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
// @effect-diagnostics-next-line nodeBuiltinImport:off -- Effect's Crypto has no timingSafeEqual.
import * as NodeCrypto from "node:crypto";
import { AuthSessionId } from "@t3tools/contracts";
import * as DateTime from "effect/DateTime";
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/auth/utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import type {
AuthClientPresentationMetadata,
} from "@t3tools/contracts";
import type * as HttpServerRequest from "effect/http/HttpServerRequest";
// @effect-diagnostics-next-line nodeBuiltinImport:off -- Effect's Crypto has no createHmac or timingSafeEqual.
import * as NodeCrypto from "node:crypto";
import * as Base64Url from "effect/encoding/Base64Url";
import * as Result from "effect/Result";
Expand Down
4 changes: 3 additions & 1 deletion apps/server/src/checkpointing/CheckpointDiffQuery.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import * as NodeCrypto from "@effect/platform-node/NodeCrypto";
import { assert, it, vi } from "@effect/vitest";
import { CheckpointRef, CheckpointScopeId, RunId, ThreadId } from "@t3tools/contracts";
import * as Effect from "effect/Effect";
Expand Down Expand Up @@ -59,6 +60,7 @@ function layerFor(input: {
}),
),
),
Layer.provideMerge(NodeCrypto.layer),
);
}

Expand All @@ -80,7 +82,7 @@ it.effect("computes V2 run diffs from projected checkpoint scopes", () => {
});
assert.deepEqual(diffCheckpoints.mock.calls[0]?.[0], {
cwd: "/repo",
fromCheckpointRef: checkpointRefForScopeOrdinal({
fromCheckpointRef: yield* checkpointRefForScopeOrdinal({
scopeId: firstScopeId,
ordinalWithinScope: 0,
}),
Expand Down
27 changes: 14 additions & 13 deletions apps/server/src/checkpointing/CheckpointDiffQuery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import {
type ThreadId,
} from "@t3tools/contracts";
import * as Context from "effect/Context";
import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Schema from "effect/Schema";
Expand Down Expand Up @@ -77,6 +78,7 @@ function buildTurnDiffResult(
export const make = Effect.gen(function* () {
const threads = yield* ThreadManagement.ThreadManagementService;
const checkpointStore = yield* CheckpointStore.CheckpointStore;
const crypto = yield* Crypto.Crypto;

const getTurnDiff: CheckpointDiffQuery["Service"]["getTurnDiff"] = Effect.fn("getTurnDiff")(
function* (input) {
Expand Down Expand Up @@ -155,21 +157,20 @@ export const make = Effect.gen(function* () {
});
}

// The root scope is shared by every run in this thread. Its runId
// tracks the latest owner, while ordinal zero stays the baseline.
const firstScope =
input.fromTurnCount === 0
? projection.checkpointScopes.find((scope) => scope.kind === "root_run")
: undefined;
const fromCheckpointRef =
input.fromTurnCount === 0
? (() => {
// The root scope is shared by every run in this thread. Its
// runId tracks the latest owner, while ordinal zero stays the baseline.
const firstScope = projection.checkpointScopes.find(
(scope) => scope.kind === "root_run",
);
return firstScope === undefined
? undefined
: checkpointRefForScopeOrdinal({
scopeId: firstScope.id,
ordinalWithinScope: 0,
});
})()
? firstScope === undefined
? undefined
: yield* checkpointRefForScopeOrdinal({
scopeId: firstScope.id,
ordinalWithinScope: 0,
}).pipe(Effect.provideService(Crypto.Crypto, crypto))
: readyCheckpoints.find((checkpoint) => checkpoint.appRunOrdinal === input.fromTurnCount)
?.ref;
if (fromCheckpointRef === undefined) {
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/cloud/CloudLink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
* and mint requests, and keeping the managed tunnel registered, recovered and
* released. HTTP handlers, server startup and shutdown all go through it.
*/
// @effect-diagnostics-next-line nodeBuiltinImport:off -- Effect's Crypto has no createPublicKey.
import * as NodeCrypto from "node:crypto";
import {
AuthStandardClientScopes,
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/cloud/environmentKeys.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
// @effect-diagnostics-next-line nodeBuiltinImport:off -- Effect's Crypto has no generateKeyPairSync.
import * as NodeCrypto from "node:crypto";
import * as Effect from "effect/Effect";
import * as Option from "effect/Option";
Expand Down
8 changes: 5 additions & 3 deletions apps/server/src/device/AgentDeviceTarget.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,9 @@ console.log(readFileSync(args[args.indexOf('--config') + 1], 'utf8'));
if (process.env.AGENT_DEVICE_DAEMON_BASE_URL) process.exit(2);`,
);
const shim = yield* ensureAgentDeviceShim({ entryPath, stateDir: dir });
const files = ["mini", "android"].map((host) => agentDeviceConfigPath(dir, host, path));
const files = yield* Effect.forEach(["mini", "android"], (host) =>
agentDeviceConfigPath(dir, host, path),
);
for (const [index, file] of files.entries())
yield* writeAgentDeviceConfig(file, {
baseUrl: `http://127.0.0.1:${1000 + index}`,
Expand Down Expand Up @@ -73,8 +75,8 @@ if (process.env.AGENT_DEVICE_DAEMON_BASE_URL) process.exit(2);`,
});
expect((yield* Effect.promise(() => invoke(files[0]!))).daemonAuthToken).toBe("new");
expect(yield* fs.readFileString(files[1]!)).toBe(second);
expect(agentDeviceSession("thread", "mini", "same-id")).not.toBe(
agentDeviceSession("thread", "android", "same-id"),
expect(yield* agentDeviceSession("thread", "mini", "same-id")).not.toBe(
yield* agentDeviceSession("thread", "android", "same-id"),
);
for (const args of [
["snapshot"],
Expand Down
18 changes: 12 additions & 6 deletions apps/server/src/device/AgentDeviceTarget.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import * as NodeCrypto from "node:crypto";
import * as Schema from "effect/Schema";
import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
import * as Hex from "effect/encoding/Hex";
import * as Schema from "effect/Schema";
import * as FileSystem from "effect/FileSystem";
import * as Path from "effect/Path";

Expand All @@ -12,15 +13,20 @@ const encodeEndpoint = Schema.encodeEffect(
),
);

const key = (value: string) =>
NodeCrypto.createHash("sha256").update(value).digest("hex").slice(0, 24);
const key = Effect.fn("AgentDeviceTarget.key")(function* (value: string) {
const crypto = yield* Crypto.Crypto;
const digest = yield* crypto
.digest("SHA-256", new TextEncoder().encode(value))
.pipe(Effect.orDie);
return Hex.encode(digest).slice(0, 24);
});

/** A stable file per host lets forwarded endpoints change without retargeting other commands. */
export const agentDeviceConfigPath = (stateDir: string, hostId: string, path: Path.Path) =>
path.join(stateDir, "device", "hosts", `${key(hostId)}.json`);
key(hostId).pipe(Effect.map((hash) => path.join(stateDir, "device", "hosts", `${hash}.json`)));

export const agentDeviceSession = (threadId: string, hostId: string, deviceId: string) =>
`t3-${key(JSON.stringify([threadId, hostId, deviceId]))}`;
key(JSON.stringify([threadId, hostId, deviceId])).pipe(Effect.map((hash) => `t3-${hash}`));

export const writeAgentDeviceConfig = Effect.fn("AgentDeviceTarget.writeConfig")(function* (
file: string,
Expand Down
3 changes: 2 additions & 1 deletion apps/server/src/device/DeviceMultiHost.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import { expect, it } from "@effect/vitest";
import { ThreadId } from "@t3tools/contracts";
import * as Deferred from "effect/Deferred";
import * as Fiber from "effect/Fiber";
import * as NodeCrypto from "@effect/platform-node/NodeCrypto";
import * as Effect from "effect/Effect";
import { HttpClient, HttpClientResponse } from "effect/http";
import * as ServerSettings from "../serverSettings.ts";
Expand Down Expand Up @@ -77,7 +78,7 @@ it.effect("keeps hosts independent when serials collide and another host fails",
order.push("write finished");
return "/host-config.json";
}),
).pipe(Effect.provideService(HttpClient.HttpClient, http));
).pipe(Effect.provide(NodeCrypto.layer), Effect.provideService(HttpClient.HttpClient, http));
expect(yield* service.agentReadinessIfSupported("b")).not.toBeNull();
const listed = yield* service.list;
expect(listed.devices.map((device) => device.hostId).sort()).toEqual(["a", "b"]);
Expand Down
Loading
Loading