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
91 changes: 90 additions & 1 deletion src/core/brain_cloud.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ import { decodeSse, type StreamFrame } from "./stream.js";
import { HttpError, StreamIncompleteError, StreamUnavailableError } from "./errors.js";
import { appendCustody } from "./custody.js";
import { hintFor } from "./error_hints.js";
import { fallbackCapabilities, type ResolvedCapabilities } from "./capabilities.js";

/** Version of the dev-session wire protocol this client speaks. */
export const DEV_PROTOCOL_VERSION = 1;
Expand All @@ -50,15 +51,98 @@ interface DevSessionCreated {
tools?: string[];
}

/**
* The dev-session protocol versions this build can decode, read from the
* active capability contract (`dev_session_protocol_versions`) and falling back
* to the single version the client declares on the wire.
*
* The contract has always carried this list; nothing read it. Sourcing the
* accepted set from the resolved contract is what makes the handshake real: a
* client shipped with a newer contract accepts the versions that contract
* names, without another release of this file.
*/
export function devProtocolVersions(contract?: Record<string, unknown>): number[] {
const raw = (contract ?? fallbackCapabilities().contract)["dev_session_protocol_versions"];
const versions = Array.isArray(raw)
? raw
.filter((v): v is number => typeof v === "number" && Number.isFinite(v))
.map((v) => Math.trunc(v))
: [];
return versions.length ? [...new Set(versions)] : [DEV_PROTOCOL_VERSION];
}

/**
* Validate the create response before a single byte of the stream is read.
* Returns null when the session is usable, or the message to fail the run with.
*
* Two distinct holes, both reachable from a 200:
*
* - No `session_id`. Nothing checked it, so `undefined` flowed straight into
* the path builders and the client issued requests against
* `/agent/dev/sessions/undefined/stream`.
* - A `protocol_version` this build does not speak. The client declares its
* version on create and the server echoes one back; the answer was never
* read. When the server ships a version with any changed frame shape, the
* tolerant `??` / `Number()` mapping in stream.ts turns fields this build no
* longer finds into zeros and empty strings rather than errors — a silently
* mis-decoded session instead of a refused one.
*
* This deliberately does NOT route into the legacy 404/403 downgrade: a
* version mismatch is not "this server has no dev route", and quietly dropping
* to the one-way chat stream would convert a loud, fixable incompatibility into
* an unexplained loss of the local tool round-trip.
*/
export function checkDevSession(created: unknown, speaks: readonly number[]): string | null {
const record = (created ?? {}) as Partial<DevSessionCreated>;
const spoken = speaks.length ? speaks : [DEV_PROTOCOL_VERSION];
const id = record.session_id;
if (typeof id !== "string" || id.trim() === "") {
return "cloud dev session was created without a session_id — refusing to open a stream against an unnamed session";
}
const version = record.protocol_version;
const list = spoken.map((v) => "v" + String(v)).join(", ");
if (typeof version !== "number" || !Number.isFinite(version)) {
return (
"cloud dev session did not declare a protocol_version; this build speaks " +
list +
" and will not attach to an unversioned session — upgrade the agent (npm i -g aether-agents@latest)"
);
}
if (!spoken.includes(Math.trunc(version))) {
return (
"cloud dev session speaks protocol v" +
String(version) +
" but this build speaks " +
list +
" — upgrade the agent (npm i -g aether-agents@latest)"
);
}
return null;
}

export class CloudBrain implements Brain {
private aborted = false;
private net: AbortController | null = null;
private sessionId: string | null = null;
private lastSeq = 0;
/** Serializes upstream result POSTs so they arrive in execution order. */
private upstream: Promise<void> = Promise.resolve();
/** Dev-session protocol versions this build will attach to. */
private readonly speaks: number[];

constructor(private readonly api: ApiClient) {}
/**
* `capabilities` is optional so existing call sites are unchanged; without it
* the packaged contract snapshot supplies the accepted version set. A caller
* that has already resolved the server contract should pass it, so a server
* that legitimately advertises a newer dev protocol is honored rather than
* refused on stale packaged data.
*/
constructor(
private readonly api: ApiClient,
capabilities?: ResolvedCapabilities,
) {
this.speaks = devProtocolVersions(capabilities?.contract);
}

run(task: TaskCommand): AsyncIterable<BrainEvent> {
const queue = new EventQueue();
Expand Down Expand Up @@ -87,6 +171,11 @@ export class CloudBrain implements Brain {
}
throw err;
}
// Deliberately outside the catch above: a malformed or version-mismatched
// create is NOT a legacy server, and must not degrade to the one-way chat
// stream. Fail the run loudly instead.
const refusal = checkDevSession(created, this.speaks);
if (refusal) throw new Error(refusal);
this.sessionId = created.session_id;
queue.push({ type: "stage", name: "execute", face: "⟨◉⟩" }); // uplink face
await this.devPump(queue);
Expand Down
98 changes: 90 additions & 8 deletions test/brain_cloud_dev.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@

import { test } from "node:test";
import assert from "node:assert/strict";
import { CloudBrain } from "../src/core/brain_cloud.js";
import { CloudBrain, checkDevSession, devProtocolVersions } from "../src/core/brain_cloud.js";
import { ApiClient } from "../src/core/transport.js";
import type { BrainEvent } from "../src/core/brain_protocol.js";
import type { TokenStore } from "../src/core/auth.js";
Expand All @@ -23,10 +23,16 @@ interface Call {

/** A dev-protocol server fake: create → JSON; stream attempts → scripted SSE
* bodies (one per reconnect); tool-results/control/DELETE → recorded JSON. */
function devServer(streams: string[][], opts?: { failToolResultTimes?: number }) {
function devServer(
streams: string[][],
opts?: { failToolResultTimes?: number; sessionId?: string | null; protocolVersion?: number | null },
) {
const calls: Call[] = [];
let attempt = 0;
let toolResultFailures = opts?.failToolResultTimes ?? 0;
// null means "omit the field entirely" — a 200 that answers neither question.
const sessionId = opts?.sessionId === undefined ? "devs_abc" : opts.sessionId;
const protocolVersion = opts?.protocolVersion === undefined ? 1 : opts.protocolVersion;
const fetchImpl = (async (input: RequestInfo | URL, init?: RequestInit) => {
const url = String(input);
const method = init?.method ?? "GET";
Expand All @@ -44,12 +50,10 @@ function devServer(streams: string[][], opts?: { failToolResultTimes?: number })
}) as unknown as Response;

if (url.endsWith("/agent/dev/sessions") && method === "POST") {
return json(200, {
session_id: "devs_abc",
protocol_version: 1,
model: "sonnet",
tools: ["read_file", "write_file"],
});
const created: Record<string, unknown> = { model: "sonnet", tools: ["read_file", "write_file"] };
if (sessionId !== null) created["session_id"] = sessionId;
if (protocolVersion !== null) created["protocol_version"] = protocolVersion;
return json(200, created);
}
if (url.includes("/stream") && method === "GET") {
const frames = streams[Math.min(attempt, streams.length - 1)] ?? [];
Expand Down Expand Up @@ -345,3 +349,81 @@ test("a non-404 create failure surfaces as an error, not a silent legacy downgra
assert.ok(!out.some((e) => e.type === "monologue"));
});
});

test("devProtocolVersions reads the accepted set from the capability contract", () => {
// The contract has always carried this list; until now nothing consumed it.
assert.deepEqual(devProtocolVersions({ dev_session_protocol_versions: [1, 2] }), [1, 2]);
// Packaged contract when no resolved contract is supplied.
assert.deepEqual(devProtocolVersions(), [1]);
// A contract with a missing or unusable list still leaves the client with the
// version it actually declares on the wire, never an empty accept-set.
assert.deepEqual(devProtocolVersions({}), [1]);
assert.deepEqual(devProtocolVersions({ dev_session_protocol_versions: ["two"] }), [1]);
});

test("checkDevSession names both versions so the message is actionable", () => {
assert.equal(checkDevSession({ session_id: "s", protocol_version: 1 }, [1]), null);
const mismatch = checkDevSession({ session_id: "s", protocol_version: 2 }, [1]);
assert.ok(mismatch);
assert.match(mismatch, /v2/);
assert.match(mismatch, /v1/);
assert.match(mismatch, /upgrade/i);
assert.match(String(checkDevSession({ protocol_version: 1 }, [1])), /session_id/);
assert.match(String(checkDevSession({ session_id: " ", protocol_version: 1 }, [1])), /session_id/);
});

test("a create response with no session_id fails the run instead of streaming /undefined/stream", async () => {
const { fetchImpl, calls } = devServer([[]], { sessionId: null });
await withFetch(fetchImpl, async () => {
const brain = new CloudBrain(new ApiClient("https://stub.test", tokens));
const out: BrainEvent[] = [];
for await (const ev of brain.run(TASK)) out.push(ev);
assert.ok(out.some((e) => e.type === "error"), "must surface an error");
assert.ok(
!calls.some((c) => c.url.includes("undefined")),
"must never build a request path out of an absent session id",
);
});
});

test("an unsupported dev protocol version fails the run and does NOT downgrade to legacy", async () => {
// A version mismatch is not "this server has no dev route". Folding it into
// the 404/403 downgrade would trade a loud, fixable incompatibility for a
// silent loss of the local tool round-trip.
const { fetchImpl, calls } = devServer([[]], { protocolVersion: 2 });
await withFetch(fetchImpl, async () => {
const brain = new CloudBrain(new ApiClient("https://stub.test", tokens));
const out: BrainEvent[] = [];
for await (const ev of brain.run(TASK)) out.push(ev);
const err = out.find((e) => e.type === "error");
assert.ok(err && err.type === "error");
assert.match(err.msg, /v2/);
assert.match(err.msg, /v1/);
assert.ok(!calls.some((c) => c.url.includes("/agent/chat/stream")), "must not fall back to the legacy stream");
assert.ok(!calls.some((c) => c.url.includes("/stream") && c.method === "GET"), "must not attach to the dev stream");
});
});

test("a build whose contract advertises v2 attaches to a v2 session", async () => {
// The accepted set comes from the resolved contract, so a client shipped with
// a newer contract negotiates upward without another edit to brain_cloud.ts.
const { fetchImpl, calls } = devServer(
[[frame({ type: "done", seq: 1, ok: true, uvt: 1, cents: 0 })]],
{ protocolVersion: 2 },
);
await withFetch(fetchImpl, async () => {
const brain = new CloudBrain(new ApiClient("https://stub.test", tokens), {
contract: { contract_version: 1, dev_session_protocol_versions: [1, 2] },
digest: "test",
source: "fallback",
overlay: null,
warnings: [],
});
const out: BrainEvent[] = [];
for await (const ev of brain.run(TASK)) out.push(ev);
assert.ok(!out.some((e) => e.type === "error"), JSON.stringify(out));
assert.ok(calls.some((c) => c.url.includes("/stream") && c.method === "GET"));
const done = out.find((e) => e.type === "done");
assert.ok(done && done.type === "done" && done.ok === true);
});
});