Skip to content
Closed
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
10 changes: 7 additions & 3 deletions apps/server/scripts/record-grok-acp-replay-fixture.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import { ServerConfig } from "../src/config.ts";
import {
GROK_DEFAULT_INSTANCE_ID,
GROK_PROVIDER,
grokLaunchRuntimeMode,
makeGrokAdapterV2,
} from "../src/orchestration-v2/Adapters/GrokAdapterV2.ts";
import { ACP_PROTOCOL } from "../src/orchestration-v2/Adapters/AcpAdapterV2.ts";
Expand Down Expand Up @@ -188,12 +189,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 @@ -385,14 +388,15 @@ const recordScenario = Effect.fn("recordGrokScenario")(function* (fixtureName: s
serverConfig: yield* ServerConfig,
selfInvocation: yield* resolveSelfInvocation(),
// Production's runtime factory, with the protocol logger teeing raw lines.
makeRuntime: (input) =>
makeRuntime: ({ runtimePolicy, ...input }) =>
makeGrokAcpRuntime({
...input,
protocolLogging: tee.attachRuntime(),
interruptPromptOnCancel: input.interruptPromptOnCancel ?? false,
grokSettings: settings,
environment,
childProcessSpawner,
runtimeMode: grokLaunchRuntimeMode(runtimePolicy),
}),
}),
];
Expand Down Expand Up @@ -450,7 +454,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 @@ -470,6 +470,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 @@ -2110,6 +2132,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 @@ -2142,6 +2263,7 @@ describe("AcpAdapterV2", () => {
driver: ACP_TEST_DRIVER,
capabilities: AcpProviderCapabilitiesV2,
makeRuntime,
clientFileSystem: diskClientFileSystem(fileSystem),
},
fileSystem,
idAllocator,
Expand Down Expand Up @@ -2221,6 +2343,7 @@ describe("AcpAdapterV2", () => {
}).pipe(Effect.andThen(runtime.handleReadTextFile(handler))),
}),
}),
clientFileSystem: diskClientFileSystem(fileSystem),
},
fileSystem,
idAllocator,
Expand Down
Loading
Loading