Harden Computer reconnect behavior

0663a852e75a · Devin AI · · parent e630419b28a2

Harden Computer reconnect behavior

Co-Authored-By: Christopher David <chris@openagents.com>
Co-Authored-By
Christopher David <chris@openagents.com>

Deploy story

What this commit did to the running system — joined from the forge receipt chain, the part a commit page elsewhere cannot show.

Not deployed through the forge lane

No push, promotion, build, or deploy receipt references this commit (receipts are scanned over a bounded recent window). Changes shipped by full node replacement carry their proof in the release gate receipt instead.

Changed files

  • modified packages/openagents-cli/src/cli.ts
  • modified packages/openagents-cli/src/computer-channel.ts
  • modified packages/openagents-cli/src/computer-up.ts
  • modified packages/openagents-cli/src/errors.ts
  • modified packages/openagents-cli/src/main.ts
  • modified packages/openagents-cli/src/runtime.ts
  • modified packages/openagents-cli/test/computer.test.ts

Diff

7 files changed, +615 -25

packages/openagents-cli/src/cli.ts modified +20 -2

@@ -33,7 +33,10 @@ import { formatAllowlist, resolveRoots, type Tier } from "./computer-policy.js";

33 33
import {
34 34
  ApiError,
35 35
  ComputerAlreadyPaired,
36
  ComputerMachineMismatch,
37
  ComputerMachineUnavailable,
36 38
  ComputerPairingInProgress,
39
  ComputerReconnectExhausted,
37 40
  InputError,
38 41
} from "./errors.js";
39 42
import { CredentialStore } from "./credential-store.js";

@@ -263,7 +266,22 @@ const computerUpCommand = Command.make("up", {}, () =>

263 266
    const flags = yield* rootCommand;
264 267
    const endpoint = yield* resolveApiEndpoint(endpointOverrides(flags));
265 268
    const up = yield* ComputerUp;
266
    const reason = yield* up.serve(endpoint.origin);
269
    const reason = yield* up.serve(endpoint.origin, VERSION);
270
    if (reason.includes("machine_unavailable")) {
271
      return yield* new ComputerMachineUnavailable({
272
        message: "The Computer machine is unavailable; the connection stopped.",
273
      });
274
    }
275
    if (reason.includes("machine_mismatch")) {
276
      return yield* new ComputerMachineMismatch({
277
        message: "The Computer machine does not match the paired identity; the connection stopped.",
278
      });
279
    }
280
    if (reason.includes("retry_exhausted")) {
281
      return yield* new ComputerReconnectExhausted({
282
        message: `The Computer connection stopped after bounded retries (${reason}).`,
283
      });
284
    }
267 285
    const output = yield* Output;
268 286
    yield* output.write(
269 287
      {

@@ -275,7 +293,7 @@ const computerUpCommand = Command.make("up", {}, () =>

275 293
  }),
276 294
).pipe(
277 295
  Command.withDescription(
278
    "Serve bounded Computer requests over an outbound connection until the server disconnects.",
296
    "Serve bounded Computer requests over an outbound connection. Transport loss and machine_reconnecting retry with bounded backoff; authorization refusals stop the command.",
279 297
  ),
280 298
);
281 299
packages/openagents-cli/src/computer-channel.ts modified +82 -14

@@ -10,8 +10,26 @@ export interface ComputerChannelOptions {

10 10
  readonly machineId: string;
11 11
  readonly hello: unknown;
12 12
  readonly heartbeatMillis?: number;
13
  readonly reconnectBackoffMillis?: number;
14
  readonly maximumReconnectAttempts?: number;
13 15
}
14 16
17
export interface ComputerSocket {
18
  readonly readyState: number;
19
  readonly send: (data: string) => void;
20
  readonly close: () => void;
21
  readonly on: (event: string, listener: (...args: ReadonlyArray<unknown>) => void) => void;
22
}
23
24
export interface ComputerSocketTransportInterface {
25
  readonly connect: (url: string) => ComputerSocket;
26
}
27
28
export class ComputerSocketTransport extends Context.Service<
29
  ComputerSocketTransport,
30
  ComputerSocketTransportInterface
31
>()("@openagentsinc/cli/ComputerSocketTransport") {}
32
15 33
export interface ComputerResponder {
16 34
  readonly chunk: (text: string) => void;
17 35
  readonly exit: (payload: Record<string, unknown>) => void;

@@ -56,9 +74,24 @@ const record = (value: unknown): Record<string, unknown> =>

56 74
const socketUrl = (origin: string, token: string): string =>
57 75
  `${origin.replace(/^http/u, "ws").replace(/\/$/u, "")}/controller/socket/websocket?vsn=2.0.0&token=${encodeURIComponent(token)}`;
58 76
59
const serveLive = (
77
const reconnectableTransportReason = (reason: string): boolean =>
78
  reason === "closed" ||
79
  reason === "phx_close" ||
80
  reason === "phx_error" ||
81
  reason === "heartbeat_timeout" ||
82
  reason === "socket_not_open" ||
83
  reason.startsWith("error:");
84
85
const wait = (milliseconds: number): Promise<void> =>
86
  new Promise((resolve) => {
87
    const timer = setTimeout(resolve, milliseconds);
88
    timer.unref();
89
  });
90
91
const serveConnection = (
60 92
  options: ComputerChannelOptions,
61 93
  handlers: ComputerChannelHandlers,
94
  transport: ComputerSocketTransportInterface,
62 95
): Promise<string> =>
63 96
  new Promise((resolve) => {
64 97
    const topic = `computer:${options.machineId}`;

@@ -68,21 +101,21 @@ const serveLive = (

68 101
    let heartbeat: NodeJS.Timeout | undefined;
69 102
    let heartbeatPending = false;
70 103
    let finished = false;
71
    const socket = new WebSocket(socketUrl(options.origin, Redacted.value(options.token)));
104
    const socket = transport.connect(socketUrl(options.origin, Redacted.value(options.token)));
72 105
73 106
    const finish = (reason: string): void => {
74 107
      if (finished) return;
75 108
      finished = true;
76 109
      if (heartbeat !== undefined) clearInterval(heartbeat);
77 110
      handlers.onClosed(reason);
78
      if (socket.readyState === WebSocket.OPEN || socket.readyState === WebSocket.CONNECTING) {
111
      if (socket.readyState === 1 || socket.readyState === 0) {
79 112
        socket.close();
80 113
      }
81 114
      resolve(reason);
82 115
    };
83 116
84 117
    const push = (event: string, payload: unknown, joined = true, forcedRef?: string): void => {
85
      if (socket.readyState !== WebSocket.OPEN) return;
118
      if (socket.readyState !== 1) return;
86 119
      reference += 1;
87 120
      const outgoing: Frame = [
88 121
        joined ? joinRef : null,

@@ -107,7 +140,7 @@ const serveLive = (

107 140
          finish("heartbeat_timeout");
108 141
          return;
109 142
        }
110
        if (socket.readyState !== WebSocket.OPEN) {
143
        if (socket.readyState !== 1) {
111 144
          finish("socket_not_open");
112 145
          return;
113 146
        }

@@ -119,10 +152,10 @@ const serveLive = (

119 152
      heartbeat.unref();
120 153
    });
121 154
122
    socket.on("message", (data: WebSocket.RawData) => {
155
    socket.on("message", (data: unknown) => {
123 156
      let parsed: unknown;
124 157
      try {
125
        parsed = JSON.parse(data.toString());
158
        parsed = JSON.parse(String(data));
126 159
      } catch {
127 160
        return;
128 161
      }

@@ -140,7 +173,9 @@ const serveLive = (

140 173
          handlers.onJoined();
141 174
          push("hello", options.hello);
142 175
        } else {
143
          finish(`join_refused:${JSON.stringify(payload.response ?? {})}`);
176
          const response = record(payload.response);
177
          const reason = response.reason;
178
          finish(`join_refused:${typeof reason === "string" ? reason : "unknown"}`);
144 179
        }
145 180
        return;
146 181
      }

@@ -165,15 +200,48 @@ const serveLive = (

165 200
        handlers.onCancel(requestId);
166 201
      }
167 202
    });
168
    socket.on("error", (cause: Error) => finish(`error:${cause.message}`));
203
    socket.on("error", (cause: unknown) =>
204
      finish(`error:${cause instanceof Error ? cause.message : String(cause)}`),
205
    );
169 206
    socket.on("close", () => finish("closed"));
170 207
  });
171 208
209
const serveLive = async (
210
  options: ComputerChannelOptions,
211
  handlers: ComputerChannelHandlers,
212
  transport: ComputerSocketTransportInterface,
213
): Promise<string> => {
214
  const maximumReconnectAttempts = options.maximumReconnectAttempts ?? 3;
215
  const backoffMillis = options.reconnectBackoffMillis ?? 250;
216
  let reconnectAttempts = 0;
217
  while (true) {
218
    const reason = await serveConnection(options, handlers, transport);
219
    const joinReason = reason.startsWith("join_refused:")
220
      ? reason.slice("join_refused:".length)
221
      : "";
222
    const retryable = joinReason === "machine_reconnecting" || reconnectableTransportReason(reason);
223
    if (!retryable || reconnectAttempts >= maximumReconnectAttempts) {
224
      return retryable && reconnectAttempts > 0
225
        ? `${joinReason === "machine_reconnecting" ? "machine_reconnecting" : "transport"}_retry_exhausted:${reason}`
226
        : reason;
227
    }
228
    reconnectAttempts += 1;
229
    const delay = Math.min(10_000, backoffMillis * 2 ** (reconnectAttempts - 1));
230
    handlers.onEvent(`reconnect:${reason}:${reconnectAttempts}`);
231
    await wait(delay);
232
  }
233
};
234
172 235
export const computerChannelNodeLayer = Layer.effect(
173 236
  ComputerChannel,
174
  Effect.succeed(
175
    ComputerChannel.of({
176
      serve: (options, handlers) => Effect.promise(() => serveLive(options, handlers)),
177
    }),
178
  ),
237
  Effect.gen(function* () {
238
    const transport = yield* ComputerSocketTransport;
239
    return ComputerChannel.of({
240
      serve: (options, handlers) => Effect.promise(() => serveLive(options, handlers, transport)),
241
    });
242
  }),
179 243
);
244
245
export const computerSocketNodeLayer = Layer.succeed(ComputerSocketTransport, {
246
  connect: (url) => new WebSocket(url) as unknown as ComputerSocket,
247
});
packages/openagents-cli/src/computer-up.ts modified +3 -3

@@ -17,7 +17,7 @@ import { CredentialStore } from "./credential-store.js";

17 17
import { InputError, type CliError } from "./errors.js";
18 18
19 19
export interface ComputerUpInterface {
20
  readonly serve: (origin: string) => Effect.Effect<string, CliError>;
20
  readonly serve: (origin: string, agentVersion: string) => Effect.Effect<string, CliError>;
21 21
}
22 22
23 23
export class ComputerUp extends Context.Service<ComputerUp, ComputerUpInterface>()(

@@ -89,7 +89,7 @@ export const computerUpLayer = Layer.effect(

89 89
    const probe = yield* ComputerProbe;
90 90
    const probeContext = yield* Effect.context<ComputerProbe>();
91 91
92
    const serve = Effect.fn("ComputerUp.serve")(function* (origin: string) {
92
    const serve = Effect.fn("ComputerUp.serve")(function* (origin: string, agentVersion: string) {
93 93
      const stored = yield* credentials.get(origin, "computer");
94 94
      if (Option.isNone(stored)) {
95 95
        return yield* new InputError({

@@ -264,7 +264,7 @@ export const computerUpLayer = Layer.effect(

264 264
          token: stored.value,
265 265
          machineId: status.value.machine_id,
266 266
          hello: {
267
            agent_version: "openagents-cli",
267
            agent_version: agentVersion,
268 268
            tier: config.tier,
269 269
            roots: config.roots,
270 270
            platform: `${process.platform}-${process.arch}`,
packages/openagents-cli/src/errors.ts modified +25 -1

@@ -156,6 +156,21 @@ export class ComputerStatusNetworkFailure extends Schema.TaggedErrorClass<Comput

156 156
  { message: Schema.String },
157 157
) {}
158 158
159
export class ComputerMachineUnavailable extends Schema.TaggedErrorClass<ComputerMachineUnavailable>()(
160
  "OpenAgentsCli.ComputerMachineUnavailable",
161
  { message: Schema.String },
162
) {}
163
164
export class ComputerMachineMismatch extends Schema.TaggedErrorClass<ComputerMachineMismatch>()(
165
  "OpenAgentsCli.ComputerMachineMismatch",
166
  { message: Schema.String },
167
) {}
168
169
export class ComputerReconnectExhausted extends Schema.TaggedErrorClass<ComputerReconnectExhausted>()(
170
  "OpenAgentsCli.ComputerReconnectExhausted",
171
  { message: Schema.String },
172
) {}
173
159 174
export type CliError =
160 175
  | InputError
161 176
  | ConfigurationError

@@ -178,7 +193,10 @@ export type CliError =

178 193
  | ComputerPairingExpired
179 194
  | ComputerPairingRefused
180 195
  | ComputerPairingNetworkFailure
181
  | ComputerStatusNetworkFailure;
196
  | ComputerStatusNetworkFailure
197
  | ComputerMachineUnavailable
198
  | ComputerMachineMismatch
199
  | ComputerReconnectExhausted;
182 200
183 201
export const exitCodeFor = (error: CliError): number => {
184 202
  switch (error._tag) {

@@ -198,6 +216,12 @@ export const exitCodeFor = (error: CliError): number => {

198 216
      return 11;
199 217
    case "OpenAgentsCli.ComputerStatusNetworkFailure":
200 218
      return 12;
219
    case "OpenAgentsCli.ComputerMachineUnavailable":
220
      return 13;
221
    case "OpenAgentsCli.ComputerMachineMismatch":
222
      return 14;
223
    case "OpenAgentsCli.ComputerReconnectExhausted":
224
      return 15;
201 225
    case "OpenAgentsCli.AuthenticationRequired":
202 226
    case "OpenAgentsCli.CredentialPersistenceUnavailable":
203 227
    case "OpenAgentsCli.CredentialStoreError":
packages/openagents-cli/src/main.ts modified +3

@@ -31,6 +31,9 @@ const cliErrorTags = new Set([

31 31
  "OpenAgentsCli.ComputerPairingRefused",
32 32
  "OpenAgentsCli.ComputerPairingNetworkFailure",
33 33
  "OpenAgentsCli.ComputerStatusNetworkFailure",
34
  "OpenAgentsCli.ComputerMachineUnavailable",
35
  "OpenAgentsCli.ComputerMachineMismatch",
36
  "OpenAgentsCli.ComputerReconnectExhausted",
34 37
]);
35 38
36 39
const isCliError = (value: unknown): value is CliError =>
packages/openagents-cli/src/runtime.ts modified +3 -2

@@ -6,7 +6,7 @@ import { apiTransportNodeLayer, networkPolicyLiveLayer } from "./api-transport.j

6 6
import { browserLauncherLayer } from "./browser-launcher.js";
7 7
import { computerConfigurationLayer } from "./computer-config.js";
8 8
import { computerClientLayer } from "./computer-client.js";
9
import { computerChannelNodeLayer } from "./computer-channel.js";
9
import { computerChannelNodeLayer, computerSocketNodeLayer } from "./computer-channel.js";
10 10
import { computerJournalLayer } from "./computer-journal.js";
11 11
import { computerProbeLayer } from "./computer-probe.js";
12 12
import { computerUpLayer } from "./computer-up.js";

@@ -42,10 +42,11 @@ const computerJournal = computerJournalLayer.pipe(Layer.provide(computerConfigur

42 42
const computerProbe = computerProbeLayer.pipe(
43 43
  Layer.provide(Layer.merge(computerConfiguration, NodeServices.layer)),
44 44
);
45
const computerChannel = computerChannelNodeLayer.pipe(Layer.provide(computerSocketNodeLayer));
45 46
const computerUp = computerUpLayer.pipe(
46 47
  Layer.provide(
47 48
    Layer.mergeAll(
48
      computerChannelNodeLayer,
49
      computerChannel,
49 50
      computerClient,
50 51
      computerConfiguration,
51 52
      computerJournal,
packages/openagents-cli/test/computer.test.ts modified +479 -3

@@ -1,11 +1,18 @@

1 1
import * as NodeServices from "@effect/platform-node/NodeServices";
2
import { Effect, Layer } from "effect";
2
import { Effect, Layer, Option, Redacted } from "effect";
3 3
import { mkdtemp, rm, stat } from "node:fs/promises";
4 4
import { homedir } from "node:os";
5 5
import { join } from "node:path";
6
import { describe, expect, it } from "vitest";
6
import { afterEach, describe, expect, it, vi } from "vitest";
7 7
8 8
import { runCliWith } from "../src/cli.js";
9
import {
10
  ComputerChannel,
11
  computerChannelNodeLayer,
12
  ComputerSocketTransport,
13
  type ComputerChannelHandlers,
14
  type ComputerSocket,
15
} from "../src/computer-channel.js";
9 16
import {
10 17
  ComputerConfiguration,
11 18
  computerConfigurationLayer,

@@ -34,8 +41,10 @@ import {

34 41
  toolchainCatalog,
35 42
} from "../src/computer-probe.js";
36 43
import { executeComputerCommand } from "../src/computer-executor.js";
44
import { ComputerClient, type ComputerStatus } from "../src/computer-client.js";
45
import { ComputerUp, computerUpLayer } from "../src/computer-up.js";
37 46
import { environmentLayerFromValues } from "../src/environment.js";
38
import { credentialStoreTestFileLayer } from "../src/credential-store.js";
47
import { CredentialStore, credentialStoreTestFileLayer } from "../src/credential-store.js";
39 48
import { pendingDeviceAuthorizationStoreTestLayer } from "../src/device-authorization-store.js";
40 49
import { persistedConfigurationTestLayer } from "../src/persisted-configuration.js";
41 50
import { outputTestLayer, type OutputDocument, type OutputMode } from "../src/output.js";

@@ -62,6 +71,50 @@ const computerJournalTestLayer = (path: string): Layer.Layer<ComputerJournal> =>

62 71
    ),
63 72
  );
64 73
74
class StubSocket implements ComputerSocket {
75
  readyState = 0;
76
  readonly sent: string[] = [];
77
  private readonly listeners = new Map<string, Array<(...args: ReadonlyArray<unknown>) => void>>();
78
79
  on(event: string, listener: (...args: ReadonlyArray<unknown>) => void): void {
80
    this.listeners.set(event, [...(this.listeners.get(event) ?? []), listener]);
81
  }
82
83
  send(data: string): void {
84
    this.sent.push(data);
85
  }
86
87
  close(): void {
88
    this.readyState = 3;
89
    this.emit("close");
90
  }
91
92
  open(): void {
93
    this.readyState = 1;
94
    this.emit("open");
95
  }
96
97
  message(value: unknown): void {
98
    this.emit("message", JSON.stringify(value));
99
  }
100
101
  emit(event: string, ...args: ReadonlyArray<unknown>): void {
102
    for (const listener of this.listeners.get(event) ?? []) listener(...args);
103
  }
104
}
105
106
const sentFrames = (socket: StubSocket): Array<ReadonlyArray<unknown>> =>
107
  socket.sent.map((value) => JSON.parse(value) as ReadonlyArray<unknown>);
108
109
const computerChannelTestLayer = (transport: {
110
  readonly connect: (url: string) => ComputerSocket;
111
}): Layer.Layer<ComputerChannel> =>
112
  computerChannelNodeLayer.pipe(Layer.provide(Layer.succeed(ComputerSocketTransport, transport)));
113
114
afterEach(() => {
115
  vi.useRealTimers();
116
});
117
65 118
describe("local Computer policy", () => {
66 119
  const root = "/workspace/project";
67 120
  const config = { tier: "probe" as const, roots: [root], preApproved: [] };

@@ -289,6 +342,415 @@ describe("local Computer journal", () => {

289 342
  });
290 343
});
291 344
345
describe("Computer channel", () => {
346
  const options = {
347
    origin: "https://openagents.example",
348
    token: Redacted.make("smct_test-secret"),
349
    machineId: "machine-1",
350
    hello: {
351
      agent_version: "0.2.1",
352
      tier: "curated",
353
      roots: ["/workspace"],
354
      probe: { schema: "openagents.computer_probe.v1" },
355
    },
356
    heartbeatMillis: 30_000,
357
    maximumReconnectAttempts: 0,
358
  };
359
360
  it("joins, sends hello, answers probes, and correlates run responses", async () => {
361
    const socket = new StubSocket();
362
    const events: string[] = [];
363
    const cancelled: string[] = [];
364
    let runResponder:
365
      | {
366
          readonly chunk: (text: string) => void;
367
          readonly exit: (payload: Record<string, unknown>) => void;
368
          readonly refused: (reason: string, detail: string) => void;
369
        }
370
      | undefined;
371
    const channelRun = Effect.runPromise(
372
      Effect.gen(function* () {
373
        const channel = yield* ComputerChannel;
374
        return yield* channel.serve(options, {
375
          onProbe: async () => ({ schema: "openagents.computer_probe.v1", roots: ["/workspace"] }),
376
          onRun: (_requestId, _payload, responder) => {
377
            runResponder = responder;
378
          },
379
          onCancel: (requestId) => cancelled.push(requestId),
380
          onJoined: () => events.push("joined"),
381
          onEvent: (event) => events.push(event),
382
          onClosed: (reason) => events.push(`closed:${reason}`),
383
        });
384
      }).pipe(Effect.provide(computerChannelTestLayer({ connect: () => socket }))),
385
    );
386
    await Promise.resolve();
387
    socket.open();
388
    expect(sentFrames(socket)[0]).toEqual(["1", "1", "computer:machine-1", "phx_join", {}]);
389
    socket.message(["1", "1", "computer:machine-1", "phx_reply", { status: "ok", response: {} }]);
390
    expect(sentFrames(socket)[1]).toEqual(["1", "2", "computer:machine-1", "hello", options.hello]);
391
392
    socket.message(["1", null, "computer:machine-1", "probe", { request_id: "probe-1" }]);
393
    await Promise.resolve();
394
    expect(sentFrames(socket)).toContainEqual([
395
      "1",
396
      "3",
397
      "computer:machine-1",
398
      "probe_result",
399
      {
400
        request_id: "probe-1",
401
        probe: { schema: "openagents.computer_probe.v1", roots: ["/workspace"] },
402
      },
403
    ]);
404
405
    socket.message([
406
      "1",
407
      null,
408
      "computer:machine-1",
409
      "run",
410
      { request_id: "run-1", argv: ["echo", "hello"], cwd: "/workspace" },
411
    ]);
412
    runResponder?.chunk("hello\n");
413
    runResponder?.exit({ status: "completed", exit_code: 0 });
414
    socket.message(["1", null, "computer:machine-1", "cancel", { request_id: "run-1" }]);
415
    expect(sentFrames(socket)).toContainEqual([
416
      "1",
417
      "4",
418
      "computer:machine-1",
419
      "chunk",
420
      { request_id: "run-1", text: "hello\n" },
421
    ]);
422
    expect(sentFrames(socket)).toContainEqual([
423
      "1",
424
      "5",
425
      "computer:machine-1",
426
      "exit",
427
      { request_id: "run-1", status: "completed", exit_code: 0 },
428
    ]);
429
    expect(cancelled).toEqual(["run-1"]);
430
    expect(events).toContain("joined");
431
    expect(JSON.stringify(sentFrames(socket))).not.toContain("smct_test-secret");
432
    socket.message(["1", null, "computer:machine-1", "phx_close", {}]);
433
    await expect(channelRun).resolves.toBe("phx_close");
434
  });
435
436
  it("retries machine_reconnecting and stops on authorization refusals", async () => {
437
    vi.useFakeTimers();
438
    const reconnectingSockets: StubSocket[] = [];
439
    const reconnectingEvents: string[] = [];
440
    const reconnectingRun = Effect.runPromise(
441
      Effect.gen(function* () {
442
        const channel = yield* ComputerChannel;
443
        return yield* channel.serve(
444
          { ...options, reconnectBackoffMillis: 0, maximumReconnectAttempts: 1 },
445
          {
446
            onProbe: async () => ({}),
447
            onRun: () => undefined,
448
            onCancel: () => undefined,
449
            onJoined: () => undefined,
450
            onEvent: (event) => reconnectingEvents.push(event),
451
            onClosed: () => undefined,
452
          },
453
        );
454
      }).pipe(
455
        Effect.provide(
456
          computerChannelTestLayer({
457
            connect: () => {
458
              const socket = new StubSocket();
459
              reconnectingSockets.push(socket);
460
              return socket;
461
            },
462
          }),
463
        ),
464
      ),
465
    );
466
    await Promise.resolve();
467
    reconnectingSockets[0]?.open();
468
    reconnectingSockets[0]?.message([
469
      "1",
470
      "1",
471
      "computer:machine-1",
472
      "phx_reply",
473
      { status: "error", response: { reason: "machine_reconnecting" } },
474
    ]);
475
    await vi.advanceTimersByTimeAsync(0);
476
    reconnectingSockets[1]?.open();
477
    reconnectingSockets[1]?.message([
478
      "1",
479
      "1",
480
      "computer:machine-1",
481
      "phx_reply",
482
      { status: "error", response: { reason: "machine_unavailable" } },
483
    ]);
484
    await expect(reconnectingRun).resolves.toBe("join_refused:machine_unavailable");
485
    expect(reconnectingEvents).toContain("reconnect:join_refused:machine_reconnecting:1");
486
487
    for (const reason of ["machine_unavailable", "machine_mismatch"] as const) {
488
      const socket = new StubSocket();
489
      const result = Effect.runPromise(
490
        Effect.gen(function* () {
491
          const channel = yield* ComputerChannel;
492
          return yield* channel.serve(options, {
493
            onProbe: async () => ({}),
494
            onRun: () => undefined,
495
            onCancel: () => undefined,
496
            onJoined: () => undefined,
497
            onEvent: () => undefined,
498
            onClosed: () => undefined,
499
          });
500
        }).pipe(Effect.provide(computerChannelTestLayer({ connect: () => socket }))),
501
      );
502
      await Promise.resolve();
503
      socket.open();
504
      socket.message([
505
        "1",
506
        "1",
507
        "computer:machine-1",
508
        "phx_reply",
509
        { status: "error", response: { reason } },
510
      ]);
511
      await expect(result).resolves.toBe(`join_refused:${reason}`);
512
    }
513
514
    const transportSockets: StubSocket[] = [];
515
    const transportRun = Effect.runPromise(
516
      Effect.gen(function* () {
517
        const channel = yield* ComputerChannel;
518
        return yield* channel.serve(
519
          { ...options, reconnectBackoffMillis: 0, maximumReconnectAttempts: 1 },
520
          {
521
            onProbe: async () => ({}),
522
            onRun: () => undefined,
523
            onCancel: () => undefined,
524
            onJoined: () => undefined,
525
            onEvent: (event) => reconnectingEvents.push(event),
526
            onClosed: () => undefined,
527
          },
528
        );
529
      }).pipe(
530
        Effect.provide(
531
          computerChannelTestLayer({
532
            connect: () => {
533
              const socket = new StubSocket();
534
              transportSockets.push(socket);
535
              return socket;
536
            },
537
          }),
538
        ),
539
      ),
540
    );
541
    await Promise.resolve();
542
    transportSockets[0]?.open();
543
    transportSockets[0]?.close();
544
    await vi.advanceTimersByTimeAsync(0);
545
    transportSockets[1]?.open();
546
    transportSockets[1]?.close();
547
    await expect(transportRun).resolves.toBe("transport_retry_exhausted:closed");
548
  });
549
550
  it("ends after an unacknowledged heartbeat", async () => {
551
    vi.useFakeTimers();
552
    const socket = new StubSocket();
553
    const result = Effect.runPromise(
554
      Effect.gen(function* () {
555
        const channel = yield* ComputerChannel;
556
        return yield* channel.serve(
557
          { ...options, heartbeatMillis: 10 },
558
          {
559
            onProbe: async () => ({}),
560
            onRun: () => undefined,
561
            onCancel: () => undefined,
562
            onJoined: () => undefined,
563
            onEvent: () => undefined,
564
            onClosed: () => undefined,
565
          },
566
        );
567
      }).pipe(Effect.provide(computerChannelTestLayer({ connect: () => socket }))),
568
    );
569
    socket.open();
570
    await vi.advanceTimersByTimeAsync(20);
571
    expect(sentFrames(socket)).toContainEqual([null, "2", "phoenix", "heartbeat", {}]);
572
    await expect(result).resolves.toBe("heartbeat_timeout");
573
  });
574
});
575
576
describe("Computer up service", () => {
577
  const probeReport = {
578
    schema: "openagents.computer_probe.v1",
579
    host: {
580
      platform: "linux",
581
      release: "test",
582
      architecture: "x64",
583
      hostname: "test",
584
      shell: "",
585
      cpuCount: 1,
586
      totalMemoryBytes: 1,
587
      uptimeSeconds: 1,
588
    },
589
    codingAgents: [],
590
    toolchains: [],
591
    roots: ["/workspace"],
592
    worktrees: [],
593
  } as const;
594
  const status: ComputerStatus = {
595
    machine_id: "machine-1",
596
    name: "test-computer",
597
    status: "active",
598
    token_expires_at: "2099-01-01T00:00:00.000Z",
599
  };
600
601
  const fixture = (tier: "probe" | "curated" | "shell", roots = ["/workspace"]) => {
602
    let channelHandlers: ComputerChannelHandlers | undefined;
603
    let channelOptions:
604
      | {
605
          readonly hello: unknown;
606
          readonly token: Redacted.Redacted<string>;
607
        }
608
      | undefined;
609
    const entries: Array<Record<string, unknown>> = [];
610
    const output: string[] = [];
611
    const terminal: Array<Record<string, unknown>> = [];
612
    const journal = Layer.succeed(
613
      ComputerJournal,
614
      ComputerJournal.of({
615
        append: (entry) => Effect.sync(() => entries.push(entry)),
616
        read: () => Effect.succeed([]),
617
      }),
618
    );
619
    const channel = Layer.succeed(
620
      ComputerChannel,
621
      ComputerChannel.of({
622
        serve: (options, handlers) => {
623
          channelOptions = options;
624
          channelHandlers = handlers;
625
          handlers.onJoined();
626
          return Effect.succeed("phx_close");
627
        },
628
      }),
629
    );
630
    const client = Layer.succeed(
631
      ComputerClient,
632
      ComputerClient.of({
633
        start: () => Effect.die("unused"),
634
        wait: () => Effect.die("unused"),
635
        status: () => Effect.succeed(Option.some(status)),
636
      }),
637
    );
638
    const credentials = Layer.succeed(
639
      CredentialStore,
640
      CredentialStore.of({
641
        get: () => Effect.succeed(Option.some(Redacted.make("smct_test-secret"))),
642
        set: () => Effect.void,
643
        remove: () => Effect.void,
644
      }),
645
    );
646
    const probe = Layer.succeed(
647
      ComputerProbe,
648
      ComputerProbe.of({ probe: () => Effect.succeed(probeReport) }),
649
    );
650
    const config = computerConfigurationTestLayer({
651
      tier,
652
      roots,
653
      preApproved: tier === "shell" ? ["node"] : [],
654
    });
655
    const layer = computerUpLayer.pipe(
656
      Layer.provide(Layer.mergeAll(channel, client, credentials, journal, probe, config)),
657
    );
658
    return {
659
      layer,
660
      entries,
661
      output,
662
      terminal,
663
      handlers: () => channelHandlers,
664
      hello: () => channelOptions?.hello,
665
      token: () => channelOptions?.token,
666
      emitRun: (payload: Record<string, unknown>) => {
667
        const handlers = channelHandlers;
668
        if (handlers === undefined) throw new Error("channel was not started");
669
        handlers.onRun("request-1", payload, {
670
          chunk: (text) => output.push(text),
671
          exit: (value) => terminal.push(value),
672
          refused: (reason, detail) => terminal.push({ reason, detail }),
673
        });
674
      },
675
    };
676
  };
677
678
  it("reports the CLI version, clamps limits, streams output, and journals outcomes", async () => {
679
    const value = fixture("shell", [process.cwd()]);
680
    await Effect.runPromise(
681
      Effect.gen(function* () {
682
        const up = yield* ComputerUp;
683
        return yield* up.serve("https://openagents.example", "0.1.7");
684
      }).pipe(Effect.provide(value.layer)),
685
    );
686
    expect(value.hello()).toMatchObject({ agent_version: "0.1.7", tier: "shell" });
687
    expect(value.token()).toBeDefined();
688
    value.emitRun({
689
      argv: [process.execPath, "-e", "process.stdout.write('abcdef')"],
690
      cwd: process.cwd(),
691
      tier: "shell",
692
      timeout_ms: 999_999,
693
      maximum_output_bytes: 3,
694
    });
695
    for (let attempt = 0; attempt < 50 && value.terminal.length === 0; attempt += 1) {
696
      await new Promise((resolve) => setTimeout(resolve, 10));
697
    }
698
    expect(value.output.join("")).toBe("abc");
699
    expect(value.terminal[0]).toMatchObject({
700
      status: "completed",
701
      exit_code: 0,
702
      truncated: true,
703
    });
704
    expect(value.entries.map((entry) => entry.outcome)).toEqual(
705
      expect.arrayContaining(["pending", "running", "completed"]),
706
    );
707
    expect(value.entries.map((entry) => entry.decision)).toContain("allowed");
708
    expect(JSON.stringify(value.entries)).not.toContain("smct_test-secret");
709
  });
710
711
  it("refuses requests for tiers, roots, allowlists, and denied patterns", async () => {
712
    const value = fixture("probe");
713
    await Effect.runPromise(
714
      Effect.gen(function* () {
715
        const up = yield* ComputerUp;
716
        return yield* up.serve("https://openagents.example", "0.1.7");
717
      }).pipe(Effect.provide(value.layer)),
718
    );
719
    const cases = [
720
      {
721
        argv: ["echo", "hello"],
722
        cwd: "/workspace",
723
        tier: "shell",
724
        reason: "tier_insufficient",
725
      },
726
      {
727
        argv: ["node", "--version"],
728
        cwd: "/outside",
729
        tier: "probe",
730
        reason: "root_not_declared",
731
      },
732
      {
733
        argv: ["not-allowlisted"],
734
        cwd: "/workspace",
735
        tier: "curated",
736
        reason: "tier_insufficient",
737
      },
738
      {
739
        argv: ["cat", ".env"],
740
        cwd: "/workspace",
741
        tier: "probe",
742
        reason: "denied_argument",
743
      },
744
    ] as const;
745
    for (const request of cases) {
746
      value.emitRun(request);
747
      expect(value.terminal.at(-1)).toMatchObject({ reason: request.reason });
748
    }
749
    expect(value.entries.filter((entry) => entry.outcome === "refused")).toHaveLength(4);
750
    expect(JSON.stringify(value.entries)).not.toContain("smct_test-secret");
751
  });
752
});
753
292 754
describe("local Computer execution", () => {
293 755
  it("executes argv directly with scrubbed environment and bounded output", async () => {
294 756
    const chunks: string[] = [];

@@ -322,6 +784,20 @@ describe("local Computer execution", () => {

322 784
    expect(outcome.timedOut).toBe(false);
323 785
    expect(outcome.exitCode).toBe(null);
324 786
  });
787
788
  it("keeps a natural exit truthful when cancellation arrives afterward", async () => {
789
    const execution = executeComputerCommand(
790
      [process.execPath, "-e", "process.exit(0)"],
791
      process.cwd(),
792
      { timeoutMillis: 5_000, maximumOutputBytes: 64 },
793
      () => undefined,
794
    );
795
    const outcome = await execution.done;
796
    execution.cancel();
797
    expect(outcome.cancelled).toBe(false);
798
    expect(outcome.timedOut).toBe(false);
799
    expect(outcome.exitCode).toBe(0);
800
  });
325 801
});
326 802
327 803
describe("Computer CLI output", () => {

This page updates live while a promote is in flight · changelog