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
219 changes: 219 additions & 0 deletions apps/server/src/provider/Drivers/PiDriver.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,219 @@
import * as NodeServices from "@effect/platform-node/NodeServices";
import { assert, it } from "@effect/vitest";
import { ProviderInstanceId, type PiSettings } from "@t3tools/contracts";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as Layer from "effect/Layer";
import * as Queue from "effect/Queue";
import * as Schema from "effect/Schema";
import * as Sink from "effect/Sink";
import * as Stream from "effect/Stream";
import * as TestClock from "effect/testing/TestClock";
import { HttpClient } from "effect/http";
import { ChildProcess, ChildProcessSpawner } from "effect/process";

import * as BackgroundPolicy from "../../background/BackgroundPolicy.ts";
import * as ServerConfig from "../../config.ts";
import * as IdAllocator from "../../orchestration-v2/IdAllocator.ts";
import * as ServerSettings from "../../serverSettings.ts";
import { PiDriver } from "./PiDriver.ts";

const layerTest = Layer.mergeAll(
ServerConfig.layerTest("/machine", { prefix: "t3-pi-driver-" }),
IdAllocator.layer,
ServerSettings.layerTest({ enableProviderUpdateChecks: false }),
Layer.mock(BackgroundPolicy.BackgroundPolicy)({
shouldRunScopeWork: () => Effect.succeed(false),
}),
Layer.succeed(
HttpClient.HttpClient,
HttpClient.make(() => Effect.die("Unexpected HTTP")),
),
).pipe(Layer.provideMerge(NodeServices.layer));

const decodeRequest = Schema.decodeSync(
Schema.fromJsonString(Schema.Record(Schema.String, Schema.Unknown)),
);
const encoder = new TextEncoder();
const personalSkill = {
name: "skill:personal",
source: "skill",
sourceInfo: { scope: "user", path: "/home/.pi/agent/skills/personal/SKILL.md" },
};

// Respond through the real stdio transport, with a distinct command catalog for each cwd.
const makePiSpawner = Effect.gen(function* () {
const pendingCommand = yield* Deferred.make<void>();
const launches: Array<ChildProcess.StandardCommand> = [];
const spawner = ChildProcessSpawner.make((command) =>
Effect.gen(function* () {
assert.isTrue(ChildProcess.isStandardCommand(command));
if (!ChildProcess.isStandardCommand(command)) return yield* Effect.die("Unexpected pipeline");
launches.push(command);
const version = command.args.includes("--version");
const stdout = yield* Queue.unbounded<Uint8Array>();
const cwd = command.options.cwd;
return ChildProcessSpawner.makeHandle({
// Outside the valid PID range, so transport cleanup cannot signal a real process.
pid: ChildProcessSpawner.ProcessId(999_999_999),
exitCode: version ? Effect.succeed(ChildProcessSpawner.ExitCode(0)) : Effect.never,
isRunning: Effect.succeed(!version),
kill: () => Effect.void,
unref: Effect.succeed(Effect.void),
stdin: Sink.forEach((chunk: Uint8Array) => {
const request = decodeRequest(new TextDecoder().decode(chunk).trim());
if (cwd === "/pending" && request.type === "get_commands") {
return Deferred.succeed(pendingCommand, undefined).pipe(Effect.asVoid);
}
const failed = cwd === "/failed" && request.type === "get_commands";
const data =
request.type === "get_commands"
? {
commands: [
personalSkill,
...(cwd === "/machine"
? []
: [
{
name: `skill:${cwd?.slice(1)}`,
source: "skill",
sourceInfo: {
scope: "project",
path: `${cwd}/.agents/skills/SKILL.md`,
},
},
{ name: `prompt-${cwd?.slice(1)}`, source: "prompt" },
]),
],
}
: request.type === "get_available_models"
? { models: [{ provider: "test", id: "model" }] }
: {};
return Queue.offer(
stdout,
encoder.encode(
`${JSON.stringify({ type: "response", id: request.id, success: !failed, data, ...(failed ? { error: "commands unavailable" } : {}) })}\n`,
),
).pipe(Effect.asVoid);
}),
stdout: version ? Stream.succeed(encoder.encode("pi 1.0.2\n")) : Stream.fromQueue(stdout),
stderr: Stream.empty,
all: Stream.empty,
getInputFd: () => Sink.drain,
getOutputFd: () => Stream.empty,
});
}),
);
return { spawner, launches, pendingCommand };
});

const create = (config: Partial<PiSettings> = {}, enabled = true) =>
PiDriver.create({
instanceId: ProviderInstanceId.make("pi-workspace-test"),
displayName: "My Pi",
accentColor: "#abcdef",
environment: [{ name: "PI_CODING_AGENT_DIR", value: "/isolated-pi", sensitive: false }],
enabled,
config: { ...PiDriver.defaultConfig(), binaryPath: "custom-pi", ...config },
});

it.layer(layerTest)("PiDriver workspace discovery", (it) => {
it.effect("keeps each workspace's skills and commands separate from the machine catalog", () =>
Effect.gen(function* () {
const { spawner, launches } = yield* makePiSpawner;
const instance = yield* create({ launchArgs: '--approve --skill "extra skill"' }).pipe(
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
);
yield* instance.snapshot.refresh;
assert.isDefined(instance.snapshotForCwd);
const [first, second] = yield* Effect.all(
[instance.snapshotForCwd!("/first"), instance.snapshotForCwd!("/second")],
{ concurrency: "unbounded" },
);
assert.deepEqual(
first.skills.map((skill) => skill.name),
["personal", "first"],
);
assert.deepEqual(
second.skills.map((skill) => skill.name),
["personal", "second"],
);
assert.deepEqual(
first.slashCommands.map((command) => command.name),
["compact", "prompt-first"],
);
const machine = yield* instance.snapshot.getSnapshot;
assert.deepEqual(
machine.skills.map((skill) => skill.name),
["personal"],
);
assert.deepEqual(
machine.slashCommands.map((command) => command.name),
["compact"],
);
assert.equal(first.instanceId, instance.instanceId);
assert.equal(first.displayName, "My Pi");
assert.equal(first.accentColor, "#abcdef");
assert.deepEqual(first.models, machine.models);
const workspaceLaunch = launches.find((launch) => launch.options.cwd === "/first");
assert.isDefined(workspaceLaunch);
assert.equal(workspaceLaunch!.command, "custom-pi");
assert.includeMembers(
[...workspaceLaunch!.args],
["--approve", "--skill", "extra skill", "--no-session"],
);
assert.equal(workspaceLaunch!.options.env?.PI_CODING_AGENT_DIR, "/isolated-pi");
}).pipe(Effect.scoped),
);

it.effect("does not run a disabled provider's workspace probe", () =>
Effect.gen(function* () {
const instance = yield* create({}, false).pipe(
Effect.provideService(
ChildProcessSpawner.ChildProcessSpawner,
ChildProcessSpawner.make(() => Effect.die("Disabled Pi must not spawn")),
),
);
assert.isDefined(instance.snapshotForCwd);
const workspace = yield* instance.snapshotForCwd!("/first");
assert.isFalse(workspace.enabled);
assert.deepEqual(workspace.skills, []);
}).pipe(Effect.scoped),
);

it.effect(
"fails command discovery instead of returning an empty successful workspace catalog",
() =>
Effect.gen(function* () {
const { spawner } = yield* makePiSpawner;
const instance = yield* create().pipe(
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
);
yield* instance.snapshot.refresh;
assert.isDefined(instance.snapshotForCwd);
const error = yield* Effect.flip(instance.snapshotForCwd!("/failed"));
assert.equal(error._tag, "ProviderDriverError");
assert.equal(error.instanceId, instance.instanceId);
assert.deepEqual(
(yield* instance.snapshot.getSnapshot).skills.map((skill) => skill.name),
["personal"],
);
}).pipe(Effect.scoped),
);

it.effect("times out workspace discovery that needs interactive input", () =>
Effect.gen(function* () {
const { spawner, pendingCommand } = yield* makePiSpawner;
const instance = yield* create().pipe(
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
);
const probe = yield* instance.snapshotForCwd!("/pending").pipe(Effect.flip, Effect.forkChild);
yield* Deferred.await(pendingCommand);
yield* TestClock.adjust("15 seconds");
const error = yield* Fiber.join(probe);
assert.equal(error._tag, "ProviderDriverError");
assert.include(error.detail, "workspace commands");
}).pipe(Effect.scoped),
);
});
21 changes: 21 additions & 0 deletions apps/server/src/provider/Drivers/PiDriver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import { ProviderDriverError } from "../Errors.ts";
import {
buildInitialPiProviderSnapshot,
checkPiProviderStatus,
discoverPiCommandsForCwd,
enrichPiSnapshot,
} from "../PiProvider.ts";
import { makeManagedServerProvider } from "../makeManagedServerProvider.ts";
Expand Down Expand Up @@ -188,6 +189,26 @@ export const PiDriver: ProviderDriver<PiSettings, PiDriverEnv> = {
accentColor,
enabled,
snapshot,
snapshotForCwd: (workspaceCwd) =>
!effectiveConfig.enabled
? snapshot.getSnapshot
: Effect.all([
snapshot.getSnapshot,
discoverPiCommandsForCwd(effectiveConfig, processEnv, workspaceCwd).pipe(
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
Effect.mapError(
(cause) =>
new ProviderDriverError({
driver: DRIVER_KIND,
instanceId,
detail: "Failed to discover Pi workspace commands.",
cause,
}),
),
),
]).pipe(
Effect.map(([machineSnapshot, commands]) => ({ ...machineSnapshot, ...commands })),
),
orchestrationAdapter,
textGeneration,
} satisfies ProviderInstance;
Expand Down
76 changes: 58 additions & 18 deletions apps/server/src/provider/PiProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ import {
} from "../orchestration-v2/Adapters/piT3McpInjection.ts";
import {
makePiRpcConnection,
PiRpcError,
PiRpcTimeoutError,
piRecordField as recordField,
piRecordString as recordString,
} from "../orchestration-v2/Adapters/PiRpc.ts";
Expand Down Expand Up @@ -126,31 +128,41 @@ function parseDiscoveredModels(
return parsed;
}

const makePiDiscoveryConnection = Effect.fnUntraced(function* (
piSettings: PiSettings,
environment: NodeJS.ProcessEnv,
launchArgs: ReadonlyArray<string>,
cwd?: string,
) {
const launch = buildPiRpcLaunch({
launchArgs,
environment,
mcpSession: undefined,
extensionPath: undefined,
ephemeral: true,
});
const connection = yield* makePiRpcConnection({
command: piSettings.binaryPath || "pi",
args: launch.args,
cwd,
env: launch.env,
});
yield* Stream.fromQueue(connection.events).pipe(
Stream.runDrain,
Effect.ignore,
Effect.forkScoped,
);
return connection;
});

const discoverPiViaRpc = (
piSettings: PiSettings,
environment: NodeJS.ProcessEnv,
launchArgs: ReadonlyArray<string>,
cwd?: string,
) =>
Effect.gen(function* () {
const launch = buildPiRpcLaunch({
launchArgs,
environment,
mcpSession: undefined,
extensionPath: undefined,
ephemeral: true,
});
const connection = yield* makePiRpcConnection({
command: piSettings.binaryPath || "pi",
args: launch.args,
cwd,
env: launch.env,
});
yield* Stream.fromQueue(connection.events).pipe(
Stream.runDrain,
Effect.ignore,
Effect.forkScoped,
);
const connection = yield* makePiDiscoveryConnection(piSettings, environment, launchArgs, cwd);
const stateData = yield* connection.request({ type: "get_state" });
const modelsData = yield* connection.request({ type: "get_available_models" });
const commandsData = yield* connection
Expand All @@ -169,6 +181,34 @@ const discoverPiViaRpc = (
} satisfies PiDiscovery;
}).pipe(Effect.scoped);

/** Probe the full command catalog Pi exposes in a workspace without changing machine health. */
export const discoverPiCommandsForCwd = Effect.fn("discoverPiCommandsForCwd")(
function* (piSettings: PiSettings, environment: NodeJS.ProcessEnv, cwd: string) {
const launchArgs = resolvePiLaunchArgs(piSettings.launchArgs);
if (!launchArgs.ok) {
return yield* new PiRpcError({ operation: "launch", detail: launchArgs.message });
}
const connection = yield* makePiDiscoveryConnection(
piSettings,
environment,
launchArgs.args,
cwd,
);
// A failed read must not replace a previously usable workspace catalog with an empty one.
const commandsData = yield* connection.request({ type: "get_commands" });
const { slashCommands, skills } = parsePiDiscoveredCommands(commandsData);
return { slashCommands: withPiBuiltinSlashCommands(slashCommands), skills };
},
Effect.scoped,
Effect.timeoutOrElse({
duration: PI_RPC_DISCOVERY_TIMEOUT_MS,
orElse: () =>
Effect.fail(
new PiRpcTimeoutError({ operation: "discovery", timeoutMs: PI_RPC_DISCOVERY_TIMEOUT_MS }),
),
}),
);

const runPiVersionCommand = (piSettings: PiSettings, environment: NodeJS.ProcessEnv) =>
Effect.gen(function* () {
const command = piSettings.binaryPath || "pi";
Expand Down
Loading