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
3 changes: 2 additions & 1 deletion apps/server/src/device/DeviceHubProxy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -196,7 +196,7 @@ const handler = Effect.gen(function* () {
(!readOnly && /\/api\/stream-(mode|settings)$/.test(hubPath));
yield* authenticate(controlsDevice ? AuthOrchestrationOperateScope : AuthOrchestrationReadScope);
const devices = yield* DeviceService.DeviceService;
const ready = yield* devices.currentReadiness();
const ready = yield* devices.currentReadiness(url.value.searchParams.get("hostId") ?? undefined);
if (!ready) {
return HttpServerResponse.text("Device hub is not running", { status: 503 });
}
Expand All @@ -206,6 +206,7 @@ const handler = Effect.gen(function* () {
// The ticket authenticates here and must not travel on to the hub.
const upstreamSearch = new URLSearchParams(url.value.search);
upstreamSearch.delete("wsTicket");
upstreamSearch.delete("hostId");
const search = upstreamSearch.size > 0 ? `?${upstreamSearch.toString()}` : "";
const upstreamPath = `${hubPath}${search}`;
if (upgrade) {
Expand Down
86 changes: 86 additions & 0 deletions apps/server/src/device/DeviceMultiHost.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
import { expect, it } from "@effect/vitest";
import { ThreadId } from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import { HttpClient, HttpClientResponse } from "effect/unstable/http";
import { ServerSettingsService } from "../serverSettings.ts";
import { DeviceHostError, DeviceHost } from "./DeviceHost.ts";
import { makeWithHosts } from "./DeviceService.ts";

it.effect("keeps hosts independent when serials collide and another host fails", () =>
Effect.gen(function* () {
const host = (id: string, failed = false): DeviceHost["Service"] => {
const ready = {
hub: { origin: `http://${id}` },
agentDevice: { baseUrl: `http://${id}`, token: "test", entryPath: "/agent-device" },
run: () => Effect.succeed({ stdout: "", stderr: "", code: 0 }),
helpers: { serveSimAxSettings: null, serveSimCli: null },
};
return {
id,
summary: Effect.succeed({
id,
label: id,
kind: "local",
hubInstalled: true,
agentDeviceInstalled: true,
platforms: [{ platform: "android", available: true }],
}),
platformAvailability: (platform) => Effect.succeed({ platform, available: true }),
ensureReady: () =>
failed
? Effect.fail(
new DeviceHostError({ hostId: id, step: "connect", cause: new Error("offline") }),
)
: Effect.succeed(ready),
ensureAgentReady: () => Effect.succeed(ready),
current: Effect.succeed(ready),
stopAgent: Effect.void,
stop: Effect.void,
};
};
const http = HttpClient.make((request) =>
Effect.succeed(
HttpClientResponse.fromWeb(
request,
Response.json({
simulators: [],
emulators: [
{
id: "emulator-5554",
name: "Pixel",
version: "36",
platform: "android",
booted: true,
physical: false,
},
],
}),
),
),
);
const hosts = new Map(["a", "b", "offline"].map((id) => [id, host(id, id === "offline")]));
const service = yield* makeWithHosts(hosts).pipe(
Effect.provideService(HttpClient.HttpClient, http),
);
const listed = yield* service.list;
expect(listed.devices.map((device) => device.hostId).sort()).toEqual(["a", "b"]);
expect(listed.hostStatuses.offline?.status).toBe("failed");
const threadId = ThreadId.make("thread");
for (const hostId of ["a", "b"])
yield* service.open({ threadId, hostId, deviceId: "emulator-5554", platform: "android" });
yield* service.close({ threadId, hostId: "a", deviceId: "emulator-5554" });
const state = yield* service.state;
expect(state.devices).toHaveLength(2);
expect(state.sessions.map((session) => session.hostId)).toEqual(["b"]);
expect(state.hostStatuses.a?.status).toBe("ready");
expect(state.hostStatuses.offline?.status).toBe("failed");
yield* service.agentReadinessIfSupported("b");
expect((yield* service.state).hostStatuses.b?.status).toBe("ready");
yield* service.configure({ enabled: false });
expect((yield* service.state).hostStatuses).toEqual({});
}).pipe(
Effect.provide(
ServerSettingsService.layerTest({ enableDeviceSupport: true, enableAgentDeviceAccess: true }),
),
),
);
1 change: 1 addition & 0 deletions apps/server/src/device/DeviceService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import { type DeviceService, make, stateStream } from "./DeviceService.ts";
const baseState: DeviceServiceState = {
hosts: [],
hostStatus: "idle",
hostStatuses: {},
devices: [],
sessions: [],
onboardingCompleted: false,
Expand Down
109 changes: 58 additions & 51 deletions apps/server/src/device/DeviceService.ts
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -136,8 +136,9 @@ interface ServiceState {
const vendorPrefix = (platform: DevicePlatform) =>
platform === "ios" ? "/vendor/serve-sim" : "/vendor/serve-emu";

export const make = Effect.gen(function* () {
const localHost = yield* DeviceHost.DeviceHost;
export const makeWithHosts = Effect.fn("DeviceService.makeWithHosts")(function* (
hosts: ReadonlyMap<DeviceHostId, DeviceHost.DeviceHost["Service"]>,
Comment thread
juliusmarminge marked this conversation as resolved.
) {
const settings = yield* ServerSettings.ServerSettingsService;
const lifecycleLock = yield* Semaphore.make(1);
const readDeviceSettings = settings.getSettings.pipe(
Expand All @@ -152,16 +153,15 @@ export const make = Effect.gen(function* () {
),
);
const initialSettings = yield* readDeviceSettings;
const hosts: ReadonlyMap<DeviceHostId, DeviceHost.DeviceHost["Service"]> = new Map([
[localHost.id, localHost],
]);

const httpClient = (yield* HttpClient.HttpClient).pipe(HttpClient.withScope);
const statePubSub = yield* PubSub.unbounded<DeviceServiceState>();
const initialHosts = yield* Effect.forEach(hosts.values(), (host) => host.summary);
const stateRef = yield* SynchronizedRef.make<ServiceState>({
state: {
hosts: initialHosts,
hostStatus: initialSettings.enabled ? "idle" : "disabled",
hostStatuses: {},
devices: [],
sessions: [],
onboardingCompleted: initialSettings.onboardingCompleted,
Expand All @@ -187,6 +187,18 @@ export const make = Effect.gen(function* () {
return host;
});

const setHostStatus = (
hostId: DeviceHostId,
status: DeviceServiceState["hostStatuses"][string],
) =>
publish((state) => ({
...state,
...(hostId === LOCAL_DEVICE_HOST_ID
? { hostStatus: status.status, hostStatusDetail: status.detail }
: {}),
hostStatuses: { ...state.hostStatuses, [hostId]: status },
}));

const readiness: DeviceService["Service"]["readiness"] = Effect.fn("DeviceService.readiness")(
function* (hostId) {
const host = yield* resolveHost(hostId);
Expand All @@ -198,34 +210,19 @@ export const make = Effect.gen(function* () {
});
}
const ready = yield* host
.ensureReady((phase) =>
publish((state) => ({ ...state, hostStatus: phase, hostStatusDetail: undefined })).pipe(
Effect.asVoid,
),
)
.ensureReady((status) => setHostStatus(host.id, { status }).pipe(Effect.asVoid))
.pipe(
Effect.tapError((error) =>
publish((state) => ({
...state,
hostStatus: "failed",
hostStatusDetail: error.message,
})),
setHostStatus(host.id, { status: "failed", detail: error.message }),
),
Effect.mapError(
(error) => new DeviceHostUnavailableError({ hostId: host.id, reason: error.message }),
),
);
yield* SynchronizedRef.get(stateRef).pipe(
Effect.flatMap(({ state }) =>
state.hostStatus === "ready"
? Effect.void
: publish((current) => ({
...current,
hostStatus: "ready",
hostStatusDetail: undefined,
})),
),
);
const { state } = yield* SynchronizedRef.get(stateRef);
if (state.hostStatuses[host.id]?.status !== "ready") {
yield* setHostStatus(host.id, { status: "ready" });
}
return { hostId: host.id, ...ready };
},
lifecycleLock.withPermit,
Expand All @@ -249,25 +246,18 @@ export const make = Effect.gen(function* () {
const summary = yield* host.summary;
if (!summary.platforms.some((platform) => platform.available)) return null;
const ready = yield* host
.ensureAgentReady((phase) =>
publish((state) => ({ ...state, hostStatus: phase, hostStatusDetail: undefined })).pipe(
Effect.asVoid,
),
)
.ensureAgentReady((phase) => setHostStatus(host.id, { status: phase }).pipe(Effect.asVoid))
.pipe(
Effect.tapError((error) =>
publish((state) => ({
...state,
hostStatus: "failed",
hostStatusDetail: error.message,
})),
setHostStatus(host.id, { status: "failed", detail: error.message }),
),
Effect.mapError(
(error) => new DeviceHostUnavailableError({ hostId: host.id, reason: error.message }),
),
);
const hostSummaries = yield* Effect.forEach(hosts.values(), (candidate) => candidate.summary);
yield* publish((state) => ({ ...state, hosts: hostSummaries, hostStatus: "ready" }));
yield* publish((state) => ({ ...state, hosts: hostSummaries }));
yield* setHostStatus(host.id, { status: "ready" });
return { hostId: host.id, ...ready };
}, lifecycleLock.withPermit);

Expand Down Expand Up @@ -357,27 +347,37 @@ export const make = Effect.gen(function* () {
return yield* publish((state) => ({
...state,
hosts: hostSummaries,
devices,
hostStatusDetail: detail,
...(ready.hostId === LOCAL_DEVICE_HOST_ID ? { hostStatusDetail: detail } : {}),
devices: [
...state.devices.filter((device) => device.hostId !== ready.hostId),
...devices,
],
hostStatuses: {
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
...state.hostStatuses,
[ready.hostId]: { status: "ready", ...(detail ? { detail } : {}) },
},
}));
}),
);
});

const list: DeviceService["Service"]["list"] = Effect.gen(function* () {
if (!(yield* readDeviceSettings).enabled) return (yield* SynchronizedRef.get(stateRef)).state;
const ready = yield* readiness();
return yield* refresh(ready);
}).pipe(
Effect.tapError((error) =>
publish((state) =>
state.hostStatus === "disabled"
? state
: { ...state, hostStatus: "failed", hostStatusDetail: error.message },
),
),
Effect.withSpan("DeviceService.list"),
);
yield* Effect.forEach(
hosts.values(),
(host) =>
Effect.gen(function* () {
const ready = yield* readinessIfSupported(host.id);
if (ready) yield* refresh(ready);
}).pipe(
Effect.catch((error) =>
setHostStatus(host.id, { status: "failed", detail: error.message }),
),
Comment on lines +373 to +375

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 device/DeviceService.ts:373

A concurrent list can mark a host { status: "failed" } after configure({ enabled: false }) has disabled device support, leaving disabled host-aware clients with a stale failure status indefinitely. The catch unconditionally calls setHostStatus, and disabled lists return without clearing hostStatuses; skip this update when state.hostStatus is already "disabled".

-          Effect.catch((error) =>
-            setHostStatus(host.id, { status: "failed", detail: error.message }),
-          ),
+          Effect.catch((error) =>
+            SynchronizedRef.get(stateRef).pipe(
+              Effect.flatMap(({ state }) =>
+                state.hostStatus === "disabled"
+                  ? Effect.void
+                  : setHostStatus(host.id, { status: "failed", detail: error.message }),
+              ),
+            ),
+          ),
🚀 Reply "fix it for me" or copy this AI Prompt for your agent:
In file @apps/server/src/device/DeviceService.ts around lines 373-375:

A concurrent `list` can mark a host `{ status: "failed" }` after `configure({ enabled: false })` has disabled device support, leaving disabled host-aware clients with a stale failure status indefinitely. The catch unconditionally calls `setHostStatus`, and disabled lists return without clearing `hostStatuses`; skip this update when `state.hostStatus` is already `"disabled"`.

),
{ concurrency: 4 },
);
return (yield* SynchronizedRef.get(stateRef)).state;
}).pipe(Effect.withSpan("DeviceService.list"));

const configure: DeviceService["Service"]["configure"] = Effect.fn("DeviceService.configure")(
function* (input) {
Expand Down Expand Up @@ -415,6 +415,7 @@ export const make = Effect.gen(function* () {
...state,
hostStatus: nextEnabled ? "idle" : "disabled",
hostStatusDetail: undefined,
hostStatuses: {},
devices: nextEnabled ? state.devices : [],
sessions: nextEnabled ? state.sessions : [],
bootingDevices: nextEnabled ? state.bootingDevices : [],
Expand Down Expand Up @@ -636,6 +637,7 @@ export const make = Effect.gen(function* () {
const closing = state.sessions.filter(
(session) =>
session.threadId === input.threadId &&
(input.hostId === undefined || session.hostId === input.hostId) &&
(input.deviceId === undefined || session.deviceId === input.deviceId),
);
if (closing.length === 0) return;
Expand Down Expand Up @@ -749,6 +751,11 @@ export const make = Effect.gen(function* () {
});
});

export const make = Effect.gen(function* () {
const host = yield* DeviceHost.DeviceHost;
return yield* makeWithHosts(new Map([[host.id, host]]));
});

export const layer = Layer.effect(DeviceService, make).pipe(Layer.provide(LocalDeviceHost.layer));

/** State stream for WS subscribers: current snapshot first, then every change. */
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/mcp/McpDeviceToolkit.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ const state = {
},
],
hostStatus: "ready" as const,
hostStatuses: { local: { status: "ready" as const } },
devices: [device],
sessions: [],
onboardingCompleted: true,
Expand Down
5 changes: 4 additions & 1 deletion apps/server/src/mcp/toolkits/device/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,9 @@ const handlers = {
const target =
input.deviceId !== undefined
? { hostId: input.hostId ?? LOCAL_DEVICE_HOST_ID, deviceId: input.deviceId }
: sessions.at(-1);
: sessions
.filter((session) => input.hostId === undefined || session.hostId === input.hostId)
.at(-1);
if (!target) {
return yield* new DeviceToolUnavailableError({
reason: "No device is open in this thread. Call device_open first.",
Expand All @@ -183,6 +185,7 @@ const handlers = {
const devices = yield* DeviceService.DeviceService;
yield* devices.close({
threadId: scope.threadId,
...(input.hostId === undefined ? {} : { hostId: input.hostId }),
...(input.deviceId === undefined ? {} : { deviceId: input.deviceId }),
...(input.shutdown === undefined ? {} : { shutdown: input.shutdown }),
});
Expand Down
3 changes: 2 additions & 1 deletion apps/server/src/server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1669,7 +1669,8 @@ const NodeHttpServerTestWithWsDeflate = HttpServer.layerTestClient.pipe(

const EMPTY_DEVICE_STATE: DeviceServiceState = {
hosts: [],
hostStatus: "idle",
hostStatus: "disabled",
hostStatuses: {},
devices: [],
sessions: [],
onboardingCompleted: false,
Expand Down
Loading
Loading