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
29 changes: 29 additions & 0 deletions apps/server/scripts/acp-mock-agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import * as NodeChildProcess from "node:child_process";
import * as NodeFS from "node:fs";

import * as Effect from "effect/Effect";
import * as Schema from "effect/Schema";
import * as Deferred from "effect/Deferred";

import * as NodeServices from "@effect/platform-node/NodeServices";
Expand Down Expand Up @@ -53,6 +54,12 @@ const waitForResumeRelease = process.env.T3_ACP_WAIT_FOR_RESUME_RELEASE === "1";
const completeFirstPromptOnCancel = process.env.T3_ACP_COMPLETE_FIRST_PROMPT_ON_CANCEL === "1";
const floodStderr = process.env.T3_ACP_FLOOD_STDERR === "1";
const hangPromptForever = process.env.T3_ACP_HANG_PROMPT_FOREVER === "1";
// Sends fs/write_text_file for this path, then fs/read_text_file, at the start of
// each prompt whatever the client advertised, and appends each outcome as a JSON
// line to T3_ACP_CLIENT_FS_PROBE_LOG_PATH.
const clientFsProbePath = process.env.T3_ACP_CLIENT_FS_PROBE_PATH;
const clientFsProbeLogPath = process.env.T3_ACP_CLIENT_FS_PROBE_LOG_PATH;
const encodeJson = Schema.encodeSync(Schema.fromJsonString(Schema.Unknown));
const hangAfterPermission = process.env.T3_ACP_HANG_AFTER_PERMISSION === "1";
const hangFirstPromptForever = process.env.T3_ACP_HANG_FIRST_PROMPT_FOREVER === "1";
const emitLateUpdateAfterCancel = process.env.T3_ACP_EMIT_LATE_UPDATE_AFTER_CANCEL === "1";
Expand Down Expand Up @@ -867,6 +874,28 @@ const program = Effect.gen(function* () {
beginAcpMockPrompt(cancelledSessions, requestedSessionId);
promptCount += 1;

if (clientFsProbePath !== undefined && clientFsProbeLogPath !== undefined) {
const probes = [
[
"fs/write_text_file",
{ sessionId: requestedSessionId, path: clientFsProbePath, content: "probe" },
],
["fs/read_text_file", { sessionId: requestedSessionId, path: clientFsProbePath }],
] as const;
for (const [method, params] of probes) {
const outcome = yield* agent.raw.request(method, params).pipe(
Effect.map((result) => ({ method, result })),
Effect.catch((error) =>
Effect.succeed({
method,
errorCode: error._tag === "AcpRequestError" ? error.code : error._tag,
}),
),
);
NodeFS.appendFileSync(clientFsProbeLogPath, `${encodeJson(outcome)}\n`, "utf8");
}
}

if (vibeRetryOutcome !== undefined) {
if (vibeRetryOutcome === "recovered") {
yield* Effect.sync(() =>
Expand Down
6 changes: 4 additions & 2 deletions apps/server/scripts/record-grok-acp-replay-fixture.ts
Original file line number Diff line number Diff line change
Expand Up @@ -235,12 +235,14 @@ function normalizeOutboundFrame(frame: Record<string, unknown>, runtimeInstructi
if (params === undefined) return frame;
switch (frame.method) {
case "initialize":
// Pin what T3 advertises (fs and terminal capabilities decide whether
// Grok routes file and shell work through T3); the rest is <any>.
return {
...frame,
params: Object.fromEntries(
Object.keys(params).map((key) => [
key,
key === "protocolVersion" ? params[key] : "<any>",
key === "protocolVersion" || key === "clientCapabilities" ? params[key] : "<any>",
]),
),
};
Expand Down Expand Up @@ -563,7 +565,7 @@ const recordScenario = Effect.fn("recordGrokScenario")(function* (fixtureName: s
generatedBy: "live-grok-recorder",
grokVersion: initializeMeta.agentVersion ?? "unknown",
normalization:
"Session ids are fixed UUIDs, the workspace is <workspace>, HOME is /home/grok-replay and the recording user is grok-replay. T3-owned prompt text, MCP servers and initialize params are <any>. Personal skills, machine identity, account settings and announcement broadcasts are removed, as are responses to Grok-internal request ids that T3's protocol drops. Timestamps are kept as recorded.",
"Session ids are fixed UUIDs, the workspace is <workspace>, HOME is /home/grok-replay and the recording user is grok-replay. T3-owned prompt text, MCP servers and initialize params other than clientCapabilities are <any>. Personal skills, machine identity, account settings and announcement broadcasts are removed, as are responses to Grok-internal request ids that T3's protocol drops. Timestamps are kept as recorded.",
droppedFrames,
},
entries: [
Expand Down
123 changes: 123 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -471,6 +471,28 @@ function makeMockRuntime(input: {
});
}

/** Serves opted-in client fs requests straight from disk, so tests see only the policy guard. */
function diskClientFileSystem(
fileSystem: FileSystem.FileSystem,
): NonNullable<AcpAdapterV2Flavor["clientFileSystem"]> {
const failed = () => EffectAcpErrors.AcpRequestError.internalError("test fs request failed");
return {
readTextFile: (request) =>
fileSystem.readFileString(request.path).pipe(
Effect.map((content) => ({ content })),
Effect.mapError(failed),
),
writeTextFile: (request) =>
fileSystem
.makeDirectory(NodePath.dirname(request.path), { recursive: true })
.pipe(
Effect.andThen(fileSystem.writeFileString(request.path, request.content)),
Effect.as({}),
Effect.mapError(failed),
),
};
}

function rawProtocolMethod(event: EffectAcpProtocol.AcpProtocolLogEvent): string | undefined {
if (event.stage !== "raw" || typeof event.payload !== "string") return undefined;
for (const line of event.payload.split("\n")) {
Expand Down Expand Up @@ -2111,6 +2133,105 @@ describe("AcpAdapterV2", () => {
}).pipe(Effect.provide(testLayer), Effect.scoped),
);

it.effect(
"answers fs requests method-not-found when the flavor does not opt into client fs",
() =>
Effect.gen(function* () {
const childProcessSpawner = yield* ChildProcessSpawner.ChildProcessSpawner;
const fileSystem = yield* FileSystem.FileSystem;
const idAllocator = yield* IdAllocatorV2;
const path = yield* Path.Path;
const serverConfig = yield* ServerConfig;
const selfInvocation = yield* resolveSelfInvocation();
const mockAgentPath = yield* path.fromFileUrl(
new URL("../../../scripts/acp-mock-agent.ts", import.meta.url),
);
const workspace = yield* fileSystem.makeTempDirectoryScoped({
prefix: "t3-acp-no-client-fs-",
});
const probePath = path.join(workspace, "planted.txt");
const probeLogPath = path.join(workspace, "fs-probe.jsonl");
const protocolEvents = yield* Queue.unbounded<EffectAcpProtocol.AcpProtocolLogEvent>();
const instanceId = ProviderInstanceId.make("acp-test-no-client-fs");
const adapter = makeAcpAdapterV2({
crypto: yield* Crypto.Crypto,
instanceId,
flavor: {
driver: ACP_TEST_DRIVER,
capabilities: AcpProviderCapabilitiesV2,
makeRuntime: makeMockRuntime({
childProcessSpawner,
mockAgentPath,
protocolEvents,
environment: {
T3_ACP_CLIENT_FS_PROBE_PATH: probePath,
T3_ACP_CLIENT_FS_PROBE_LOG_PATH: probeLogPath,
},
}),
},
fileSystem,
idAllocator,
serverConfig,
selfInvocation,
});
const threadId = ThreadId.make("thread-acp-no-client-fs");
const runtimePolicy = ProviderAdapterV2RuntimePolicy.make({
runtimeMode: "full-access",
interactionMode: "default",
cwd: workspace,
});
const modelSelection = { instanceId, model: "default" } as const;
const runtime = yield* adapter.openSession({
threadId,
providerSessionId: ProviderSessionId.make("provider-session-acp-no-client-fs"),
modelSelection,
runtimePolicy,
});
const providerThread = yield* runtime.ensureThread({
threadId,
modelSelection,
runtimePolicy,
});
yield* runtime.startTurn(
makeTurnInput({
threadId,
providerThread,
instanceId,
runtimePolicy,
now: yield* DateTime.now,
}),
);
yield* runtime.events.pipe(
Stream.takeUntil((event) => event.type === "turn.terminal"),
Stream.runDrain,
);

const initialize = Option.getOrThrow(
yield* Stream.fromQueue(protocolEvents).pipe(
Stream.filter(
(event) =>
event.direction === "outgoing" && rawProtocolMethod(event) === "initialize",
),
Stream.runHead,
),
);
assert.deepInclude(
(rawProtocolRequest(initialize)?.params as { clientCapabilities?: unknown })
?.clientCapabilities,
{ fs: { readTextFile: false, writeTextFile: false }, terminal: false },
);
const outcomes = (yield* fileSystem.readFileString(probeLogPath))
.trim()
.split("\n")
.map((line) => Option.getOrThrow(decodeUnknownJson(line)));
assert.deepEqual(outcomes, [
{ method: "fs/write_text_file", errorCode: -32601 },
{ method: "fs/read_text_file", errorCode: -32601 },
]);
assert.isFalse(yield* fileSystem.exists(probePath));
}).pipe(Effect.provide(testLayer), Effect.scoped),
);

it.effect("confines client-mediated writes under an explicit workspace-write sandbox", () =>
Effect.gen(function* () {
const childProcessSpawner = yield* ChildProcessSpawner.ChildProcessSpawner;
Expand Down Expand Up @@ -2143,6 +2264,7 @@ describe("AcpAdapterV2", () => {
driver: ACP_TEST_DRIVER,
capabilities: AcpProviderCapabilitiesV2,
makeRuntime,
clientFileSystem: diskClientFileSystem(fileSystem),
},
fileSystem,
idAllocator,
Expand Down Expand Up @@ -2222,6 +2344,7 @@ describe("AcpAdapterV2", () => {
}).pipe(Effect.andThen(runtime.handleReadTextFile(handler))),
}),
}),
clientFileSystem: diskClientFileSystem(fileSystem),
},
fileSystem,
idAllocator,
Expand Down
62 changes: 31 additions & 31 deletions apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,6 @@ import type {
AcpSessionRuntimeOptions,
AcpSessionRuntimeStartResult,
} from "../../provider/acp/AcpSessionRuntime.ts";
import { acpReadTextFile, acpWriteTextFile } from "../../provider/acp/AcpClientFs.ts";
import {
acpClientExecuteDisposition,
acpClientReadDisposition,
Expand Down Expand Up @@ -270,10 +269,12 @@ export interface AcpAdapterV2Flavor {
/** Native session mode to select for a runtime policy (e.g. Antigravity `yolo`). */
readonly sessionModeForPolicy?: (policy: ProviderAdapterV2RuntimePolicy) => string | undefined;
/**
* Serves the agent's `fs/read_text_file` and `fs/write_text_file` requests in
* place of the generic handlers, after the runtime policy guard. Receives the
* cwd of the policy active when the request arrives, which is null when the
* session has no workspace. Antigravity confines them to its workspace.
* Opts the session into the ACP client `fs` capability. Agents read and write
* files themselves under their own permission model unless a flavor sets
* this. Requests pass the runtime policy guard, then these handlers, which
* receive the cwd of the policy active when the request arrives (null when
* the session has no workspace). Antigravity sets it and confines requests
* to that workspace.
*/
readonly clientFileSystem?: {
readonly readTextFile: (
Expand Down Expand Up @@ -480,9 +481,10 @@ export interface AcpAdapterV2Options {
/** How agents spawn this install's `acp-mcp-bridge`; see `resolveSelfInvocation`. */
readonly selfInvocation: SelfInvocation;
/**
* Enables the ACP client `terminal` capability. Sessions advertise
* `terminal: true` and run agent-created terminals through this spawner
* with the provider instance's environment.
* Opts the session into the ACP client `terminal` capability. Agents run
* commands themselves unless an adapter sets this; with it, sessions
* advertise `terminal: true` and run agent-created terminals through this
* spawner with the provider instance's environment. Devin sets it.
*/
readonly clientTerminals?: {
readonly childProcessSpawner: ChildProcessSpawner.ChildProcessSpawner["Service"];
Expand Down Expand Up @@ -1973,7 +1975,10 @@ export function makeAcpAdapterV2(options: AcpAdapterV2Options): ProviderAdapterV
...(resumeSessionId === undefined ? {} : { resumeSessionId }),
interruptPromptOnCancel: flavor.interruptPromptOnCancel ?? false,
clientCapabilities: {
fs: { readTextFile: true, writeTextFile: true },
fs: {
readTextFile: flavor.clientFileSystem !== undefined,
writeTextFile: flavor.clientFileSystem !== undefined,
},
terminal: clientTerminals !== undefined,
elicitation: { form: {}, ...(flavor.onUrlElicitation ? { url: {} } : {}) },
...(flavor.clientCapabilitiesMeta ? { _meta: flavor.clientCapabilitiesMeta } : {}),
Expand Down Expand Up @@ -5395,31 +5400,26 @@ export function makeAcpAdapterV2(options: AcpAdapterV2Options): ProviderAdapterV
Effect.succeed(request),
requestContext.requestId,

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.

This changes the no-clientFileSystem path to leave both handlers unregistered, but the existing filesystem tests now explicitly opt in, so they no longer exercise that path. Could you add a focused adapter test that omits clientFileSystem and verifies read/write requests are rejected without touching disk?

Posted via Macroscope — Effect Service Conventions

);
// A flavor's own handlers replace the generic ones: effect-acp keeps
// only the last handler registered per method. They confine requests
// to the workspace of the policy the guard checks at request time,
// not the one the session opened with.
// Without the capability no fs handler is registered, so a stray
// request (OpenCode and Kilo send one after approved edits) gets
// method-not-found and cannot touch the disk. Opted-in handlers
// confine requests to the workspace of the policy the guard checks
// at request time, not the one the session opened with.
const clientFileSystem = flavor.clientFileSystem;
yield* targetRuntime.handleReadTextFile((request) =>
guardClientFsRead(request.path).pipe(
Effect.andThen(clientPolicyContext),
Effect.flatMap(({ policy }) =>
clientFileSystem === undefined
? acpReadTextFile(options.fileSystem, request)
: clientFileSystem.readTextFile(request, policy.cwd),
if (clientFileSystem !== undefined) {
yield* targetRuntime.handleReadTextFile((request) =>
guardClientFsRead(request.path).pipe(
Effect.andThen(clientPolicyContext),
Effect.flatMap(({ policy }) => clientFileSystem.readTextFile(request, policy.cwd)),
),
),
);
yield* targetRuntime.handleWriteTextFile((request) =>
guardClientFsWrite(request.path).pipe(
Effect.andThen(clientPolicyContext),
Effect.flatMap(({ policy }) =>
clientFileSystem === undefined
? acpWriteTextFile(options.fileSystem, request)
: clientFileSystem.writeTextFile(request, policy.cwd),
);
yield* targetRuntime.handleWriteTextFile((request) =>
guardClientFsWrite(request.path).pipe(
Effect.andThen(clientPolicyContext),
Effect.flatMap(({ policy }) => clientFileSystem.writeTextFile(request, policy.cwd)),
),
),
);
);
}
if (handlerOptions.mcp !== false) {
yield* wireAcpRuntimeMcpHandlers(targetRuntime, runtimeMcpBridge);
}
Expand Down
Loading
Loading