diff --git a/apps/mobile/src/components/ProviderIcon.tsx b/apps/mobile/src/components/ProviderIcon.tsx index 374738d0aeca..219fd344c943 100644 --- a/apps/mobile/src/components/ProviderIcon.tsx +++ b/apps/mobile/src/components/ProviderIcon.tsx @@ -26,6 +26,20 @@ export function ProviderIcon(props: ProviderIconProps) { ); } + if (props.provider === "ohMyPi") { + return ( + + + + ); + } + if (props.provider === "claudeAgent") { return ( diff --git a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts new file mode 100644 index 000000000000..b118ecb44de0 --- /dev/null +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts @@ -0,0 +1,441 @@ +// @effect-diagnostics nodeBuiltinImport:off +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { expect, it } from "@effect/vitest"; +import { + ApprovalRequestId, + OH_MY_PI_DEFAULT_MODEL, + ProviderInstanceId, + ThreadId, + type ProviderRuntimeEvent, +} from "@t3tools/contracts"; +import * as Deferred from "effect/Deferred"; +import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as FileSystem from "effect/FileSystem"; +import * as Layer from "effect/Layer"; +import * as Path from "effect/Path"; +import * as Queue from "effect/Queue"; +import * as Stream from "effect/Stream"; +import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner"; +import * as NodeURL from "node:url"; +import * as BackgroundPolicy from "../../background/BackgroundPolicy.ts"; +import { ServerConfig } from "../../config.ts"; +import { ServerSettingsService } from "../../serverSettings.ts"; +import { execScriptSource, writeFakeCli } from "../../testUtils/fakeCli.ts"; +import { NoOpProviderEventLoggers, ProviderEventLoggers } from "../Layers/ProviderEventLoggers.ts"; +import { OhMyPiDriver } from "./OhMyPiDriver.ts"; + +const testLayer = ServerConfig.layerTest(process.cwd(), { prefix: "t3-omp-driver-" }).pipe( + Layer.provideMerge(NodeServices.layer), + Layer.provideMerge(ServerSettingsService.layerTest()), + Layer.provideMerge( + Layer.mock(BackgroundPolicy.BackgroundPolicy)({ + shouldRunScopeWork: () => Effect.succeed(false), + }), + ), + Layer.provideMerge(Layer.succeed(ProviderEventLoggers, NoOpProviderEventLoggers)), +); +const instanceId = ProviderInstanceId.make("omp-test"); +const threadId = ThreadId.make("omp-thread"); + +it.layer(testLayer)("OhMyPi driver", (it) => { + it.effect("does not start a disabled CLI", () => + Effect.gen(function* () { + const instance = yield* OhMyPiDriver.create({ + instanceId, + displayName: undefined, + enabled: false, + environment: [], + config: OhMyPiDriver.defaultConfig(), + }); + expect((yield* instance.snapshot.refresh).status).toBe("disabled"); + expect((yield* instance.snapshot.getSnapshot).supportsConversationRollback).toBe(false); + }).pipe( + Effect.provideService( + ChildProcessSpawner.ChildProcessSpawner, + ChildProcessSpawner.make(() => Effect.die("Disabled provider spawned a process")), + ), + Effect.scoped, + ), + ); + + it.effect( + "refreshes the CLI catalog, excludes custom models, and retains models on discovery errors", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const { cwd } = yield* ServerConfig; + const directory = yield* fs.makeTempDirectoryScoped(); + const catalogFile = path.join(directory, "models.json"); + const binaryPath = yield* Effect.sync(() => + writeFakeCli({ + directory, + name: "omp-models-mock", + source: ` + if (process.cwd() !== process.env.T3_OMP_EXPECTED_CWD) process.exit(3); + if (process.argv[2] === "--version") { process.stdout.write("omp/18.1.14"); process.exit(0); } + if (process.argv[2] !== "models" || process.argv[3] !== "--json") process.exit(4); + const fs = await import("node:fs"); + const output = fs.readFileSync(process.env.T3_OMP_MODEL_FILE, "utf8"); + if (output === "exit") process.exit(5); + process.stdout.write(output); + `, + }), + ); + const instance = yield* OhMyPiDriver.create({ + instanceId, + displayName: undefined, + enabled: true, + environment: [ + { name: "T3_OMP_MODEL_FILE", value: catalogFile, sensitive: false }, + { name: "T3_OMP_EXPECTED_CWD", value: cwd, sensitive: false }, + ], + config: { ...OhMyPiDriver.defaultConfig(), binaryPath, customModels: ["custom/model"] }, + }); + yield* fs.writeFileString( + catalogFile, + '{"models":[{"provider":"openai","id":"gpt","name":"GPT","thinking":["low","high"]},{"provider":"openai","id":"gpt","name":"duplicate"}]}', + ); + const first = yield* instance.snapshot.refresh; + expect(first.status).toBe("ready"); + expect(first.models.map((model) => model.slug)).toEqual([ + OH_MY_PI_DEFAULT_MODEL, + "openai/gpt", + ]); + expect(first.models[1]?.capabilities?.optionDescriptors).toMatchObject([ + { + id: "thinking", + options: [{ id: "off" }, { id: "auto" }, { id: "low" }, { id: "high" }], + }, + ]); + for (const output of ["not json", "exit", '{"models":[{"name":"missing ID"}]}']) { + yield* fs.writeFileString(catalogFile, output); + const failed = yield* instance.snapshot.refresh; + expect(failed.status).toBe("error"); + expect(failed.models).toEqual(first.models); + expect(failed.message).toContain("omp models --json"); + } + yield* fs.writeFileString(catalogFile, '{"models":[]}'); + const empty = yield* instance.snapshot.refresh; + expect(empty.status).toBe("ready"); + expect(empty.models.map((model) => model.slug)).toEqual([OH_MY_PI_DEFAULT_MODEL]); + }).pipe(Effect.scoped), + ); + + for (const sendDuringPreparation of [false, true]) { + it.effect(`steers within one turn (send during preparation: ${sendDuringPreparation})`, () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const directory = yield* fs.makeTempDirectoryScoped(); + const logPath = path.join(directory, "steering.jsonl"); + const binaryPath = yield* Effect.sync(() => + writeFakeCli({ + directory, + name: "omp-steering-mock", + env: { + T3_ACP_REQUEST_LOG_PATH: logPath, + T3_ACP_COMPLETE_FIRST_PROMPT_ON_CANCEL: "1", + }, + source: execScriptSource({ + scriptPath: NodeURL.fileURLToPath( + new URL("../../../scripts/acp-mock-agent.ts", import.meta.url), + ), + }), + }), + ); + const instance = yield* OhMyPiDriver.create({ + instanceId, + displayName: undefined, + enabled: true, + environment: [], + config: { ...OhMyPiDriver.defaultConfig(), binaryPath }, + }); + const events = yield* Queue.unbounded(); + yield* instance.adapter.streamEvents.pipe( + Stream.runForEach((event) => Queue.offer(events, event)), + Effect.forkChild, + ); + yield* instance.adapter.startSession({ + threadId, + cwd: directory, + runtimeMode: "full-access", + }); + const first = yield* instance.adapter + .sendTurn({ + threadId, + input: "long task", + // Forces asynchronous configuration before the first prompt dispatch. + modelSelection: { instanceId, model: "composer-2" }, + }) + .pipe(Effect.forkChild); + const seen: ProviderRuntimeEvent[] = []; + if (!sendDuringPreparation) { + // Wait for a running tool; the other case submits concurrently while + // the first call is still configuring the session. + while (true) { + const event = yield* Queue.take(events); + seen.push(event); + if (event.type === "item.updated" && event.payload.itemType === "command_execution") + break; + } + } + const steered = yield* instance.adapter.sendTurn({ threadId, input: "do this instead" }); + const original = yield* Fiber.join(first); + while (true) { + const event = yield* Queue.take(events); + seen.push(event); + if (event.type === "turn.completed") break; + } + expect(steered.turnId).toBe(original.turnId); + expect( + seen.some( + (event) => + event.type === "content.delta" && + event.payload.streamKind === "reasoning_text" && + event.payload.delta === "native-cancel-received", + ), + ).toBe(true); + expect(seen.filter((event) => event.type === "turn.started")).toHaveLength(1); + expect(seen.filter((event) => event.type === "turn.completed")).toMatchObject([ + { turnId: original.turnId, payload: { state: "completed" } }, + ]); + const requests = (yield* fs.readFileString(logPath)) + .trim() + .split("\n") + .map((line) => JSON.parse(line) as { method: string }); + expect( + requests + .filter( + (request) => + request.method === "session/prompt" || request.method === "session/cancel", + ) + .map((request) => request.method), + ).toEqual(["session/prompt", "session/cancel", "session/prompt"]); + yield* instance.adapter.stopAll(); + }).pipe(Effect.scoped), + ); + } + + it.effect("stops during attachment preparation without dispatching a prompt", () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const directory = yield* fs.makeTempDirectoryScoped(); + const logPath = path.join(directory, "stop-preparing.jsonl"); + const readStarted = yield* Deferred.make(); + const releaseRead = yield* Deferred.make(); + const binaryPath = yield* Effect.sync(() => + writeFakeCli({ + directory, + name: "omp-stop-mock", + env: { T3_ACP_REQUEST_LOG_PATH: logPath }, + source: execScriptSource({ + scriptPath: NodeURL.fileURLToPath( + new URL("../../../scripts/acp-mock-agent.ts", import.meta.url), + ), + }), + }), + ); + const instance = yield* OhMyPiDriver.create({ + instanceId, + displayName: undefined, + enabled: true, + environment: [], + config: { ...OhMyPiDriver.defaultConfig(), binaryPath }, + }).pipe( + Effect.provideService(FileSystem.FileSystem, { + ...fs, + readFile: () => + Effect.gen(function* () { + yield* Deferred.succeed(readStarted, undefined); + yield* Deferred.await(releaseRead); + return new Uint8Array([1]); + }), + }), + ); + const events = yield* Queue.unbounded(); + yield* instance.adapter.streamEvents.pipe( + Stream.runForEach((event) => Queue.offer(events, event)), + Effect.forkChild, + ); + yield* instance.adapter.startSession({ + threadId, + cwd: directory, + runtimeMode: "full-access", + }); + const sending = yield* instance.adapter + .sendTurn({ + threadId, + input: "read this image", + attachments: [ + { + type: "image", + id: "omp-thread-00000000-0000-4000-8000-000000000000", + name: "image.png", + mimeType: "image/png", + sizeBytes: 1, + }, + ], + }) + .pipe(Effect.forkChild); + yield* Deferred.await(readStarted); + yield* instance.adapter.interruptTurn(threadId); + yield* Deferred.succeed(releaseRead, undefined); + const stopped = yield* Fiber.join(sending); + while (true) { + const event = yield* Queue.take(events); + if (event.type !== "turn.completed") continue; + expect(event.turnId).toBe(stopped.turnId); + expect(event.payload.state).toBe("cancelled"); + break; + } + expect(yield* fs.readFileString(logPath)).not.toContain('"method":"session/prompt"'); + const next = yield* instance.adapter.sendTurn({ threadId, input: "new task" }); + expect(next.turnId).not.toBe(stopped.turnId); + expect(yield* fs.readFileString(logPath)).toContain('"method":"session/prompt"'); + yield* instance.adapter.stopAll(); + }).pipe(Effect.scoped), + ); + + it.effect.skipIf(process.env.T3_OH_MY_PI_MODELS_PROBE !== "1")( + "discovers models from the installed OhMyPi CLI", + () => + Effect.gen(function* () { + const instance = yield* OhMyPiDriver.create({ + instanceId, + displayName: undefined, + enabled: true, + environment: [], + config: { + ...OhMyPiDriver.defaultConfig(), + binaryPath: process.env.T3_OH_MY_PI_BINARY ?? "omp", + }, + }); + const snapshot = yield* instance.snapshot.refresh; + expect(snapshot.status).toBe("ready"); + expect( + snapshot.models.filter((model) => model.slug !== OH_MY_PI_DEFAULT_MODEL).length, + ).toBeGreaterThan(0); + expect(yield* instance.adapter.listSessions()).toEqual([]); + }).pipe(Effect.scoped), + ); + + it.effect( + "discovers models without starting ACP, streams a turn, handles approvals and resumes through ACP", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const directory = yield* fs.makeTempDirectoryScoped(); + const logPath = path.join(directory, "requests.jsonl"); + const argvPath = path.join(directory, "argv.txt"); + const binaryPath = yield* Effect.sync(() => + writeFakeCli({ + directory, + name: "omp-mock", + env: { + T3_ACP_REQUEST_LOG_PATH: logPath, + T3_ACP_EMIT_TOOL_CALLS: "1", + T3_ACP_ALLOW_ONCE_OPTION_ID: "omp-allow-42", + }, + source: + ` + if (process.argv[2] === "--version") { + process.stdout.write("omp/18.1.14"); + process.exit(0); + } + if (process.argv[2] === "models") { + if (process.argv[3] !== "--json") process.exit(2); + process.stdout.write(JSON.stringify({ models: [ + { provider: "anthropic", id: "sonnet", name: "Sonnet", thinking: ["low", "high"] }, + { provider: "openai", id: "gpt", name: "GPT", thinking: ["high", "xhigh"] }, + ] })); + process.exit(0); + } + ` + + execScriptSource({ + scriptPath: NodeURL.fileURLToPath( + new URL("../../../scripts/acp-mock-agent.ts", import.meta.url), + ), + argvLogPath: argvPath, + }), + }), + ); + const instance = yield* OhMyPiDriver.create({ + instanceId, + displayName: "My OMP", + enabled: true, + environment: [], + config: { ...OhMyPiDriver.defaultConfig(), binaryPath }, + }); + const refreshed = yield* instance.snapshot.refresh; + expect(refreshed.status).toBe("ready"); + expect(refreshed.models.map((model) => model.slug)).toContain("anthropic/sonnet"); + expect(refreshed.models.map((model) => model.slug)).toContain("openai/gpt"); + expect(refreshed.version).toBe("18.1.14"); + expect(yield* fs.exists(logPath)).toBe(false); + const events = yield* Queue.unbounded(); + yield* instance.adapter.streamEvents.pipe( + Stream.runForEach((event) => Queue.offer(events, event)), + Effect.forkChild, + ); + const session = yield* instance.adapter.startSession({ + threadId, + cwd: directory, + runtimeMode: "approval-required", + modelSelection: { instanceId, model: OH_MY_PI_DEFAULT_MODEL }, + }); + expect(session.provider).toBe("ohMyPi"); + expect((yield* instance.snapshot.getSnapshot).models).toEqual(refreshed.models); + const turn = yield* instance.adapter + .sendTurn({ + threadId, + input: "hello", + attachments: [], + modelSelection: { instanceId, model: "composer-2" }, + }) + .pipe(Effect.forkChild); + const seen: ProviderRuntimeEvent[] = []; + while (true) { + const event = yield* Queue.take(events); + seen.push(event); + if (event.type === "request.opened") { + yield* instance.adapter.respondToRequest( + threadId, + ApprovalRequestId.make(event.requestId!), + "accept", + ); + } + if (event.type === "turn.completed") break; + } + yield* Fiber.join(turn); + // Session model/config changes must not rewrite the shared catalog or default. + expect((yield* instance.snapshot.getSnapshot).models).toEqual(refreshed.models); + expect(seen.some((event) => event.type === "content.delta")).toBe(true); + expect(seen.some((event) => event.type === "request.resolved")).toBe(true); + yield* instance.adapter.stopSession(threadId); + yield* instance.adapter.startSession({ + threadId, + cwd: directory, + runtimeMode: "approval-required", + resumeCursor: session.resumeCursor, + }); + const interruptedTurn = yield* instance.adapter + .sendTurn({ threadId, input: "wait for approval", attachments: [] }) + .pipe(Effect.forkChild); + while ((yield* Queue.take(events)).type !== "request.opened") {} + yield* instance.adapter.interruptTurn(threadId); + yield* Fiber.join(interruptedTurn).pipe(Effect.exit); + expect(yield* instance.adapter.hasSession(threadId)).toBe(true); + const requests = yield* fs.readFileString(logPath); + expect(requests).toContain('"methodId":"agent"'); + expect(requests).toContain('"method":"session/load"'); + expect(requests).not.toContain('"value":"oh-my-pi-default"'); + expect(yield* fs.readFileString(argvPath)).toContain("acp\t--approval-mode\talways-ask"); + yield* instance.adapter.stopAll(); + expect(yield* instance.adapter.listSessions()).toEqual([]); + }).pipe(Effect.scoped), + ); +}); diff --git a/apps/server/src/provider/Drivers/OhMyPiDriver.ts b/apps/server/src/provider/Drivers/OhMyPiDriver.ts new file mode 100644 index 000000000000..fc810e196a1a --- /dev/null +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.ts @@ -0,0 +1,226 @@ +import { + OH_MY_PI_DEFAULT_MODEL, + OhMyPiSettings, + ProviderDriverKind, + TextGenerationError, +} from "@t3tools/contracts"; +import { createModelCapabilities } from "@t3tools/shared/model"; +import * as Crypto from "effect/Crypto"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import * as Path from "effect/Path"; +import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; +import * as SubscriptionRef from "effect/SubscriptionRef"; +import { ChildProcessSpawner } from "effect/unstable/process"; +import * as BackgroundPolicy from "../../background/BackgroundPolicy.ts"; +import { ServerConfig } from "../../config.ts"; +import { ServerSettingsService } from "../../serverSettings.ts"; +import { ProviderDriverError } from "../Errors.ts"; +import { makeOhMyPiAdapter } from "../Layers/OhMyPiAdapter.ts"; +import { ProviderEventLoggers } from "../Layers/ProviderEventLoggers.ts"; +import { makeManagedServerProvider } from "../makeManagedServerProvider.ts"; +import { + defaultProviderContinuationIdentity, + type ProviderDriver, + type ProviderInstance, +} from "../ProviderDriver.ts"; +import { mergeProviderInstanceEnvironment } from "../ProviderInstanceEnvironment.ts"; +import { makeManualOnlyProviderMaintenanceCapabilities } from "../providerMaintenance.ts"; +import { + buildServerProvider, + isCommandMissingCause, + type ServerProviderDraft, +} from "../providerSnapshot.ts"; +import { probeOhMyPiModels } from "./OhMyPiModels.ts"; +import { withInstanceIdentity } from "./instanceIdentity.ts"; + +const DRIVER = ProviderDriverKind.make("ohMyPi"); +const decodeSettings = Schema.decodeSync(OhMyPiSettings); +const capabilities = createModelCapabilities({ optionDescriptors: [] }); + +export type OhMyPiDriverEnv = + | BackgroundPolicy.BackgroundPolicy + | ChildProcessSpawner.ChildProcessSpawner + | Crypto.Crypto + | FileSystem.FileSystem + | Path.Path + | ProviderEventLoggers + | ServerConfig + | ServerSettingsService; + +export const OhMyPiDriver: ProviderDriver = { + driverKind: DRIVER, + metadata: { displayName: "OhMyPi", supportsMultipleInstances: true }, + configSchema: OhMyPiSettings, + defaultConfig: () => decodeSettings({}), + create: ({ instanceId, displayName, accentColor, environment, enabled, config }) => + Effect.gen(function* () { + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const serverConfig = yield* ServerConfig; + const eventLoggers = yield* ProviderEventLoggers; + const effectiveConfig = { ...config, enabled }; + const processEnv = mergeProviderInstanceEnvironment(environment); + const continuationIdentity = defaultProviderContinuationIdentity({ + driverKind: DRIVER, + instanceId, + }); + const stampIdentity = withInstanceIdentity({ + instanceId, + driverKind: DRIVER, + displayName, + accentColor, + continuationGroupKey: continuationIdentity.continuationKey, + }); + const initial = { + ...buildServerProvider({ + presentation: { displayName: "OhMyPi", showInteractionModeToggle: true }, + enabled, + checkedAt: DateTime.formatIso(yield* DateTime.now), + models: [ + { + slug: OH_MY_PI_DEFAULT_MODEL, + name: "OhMyPi default", + isDefault: true, + isCustom: false, + capabilities, + }, + ], + probe: { + installed: false, + version: null, + status: "warning", + auth: { status: "unknown" }, + message: enabled + ? "Checking OhMyPi availability." + : "OhMyPi is disabled in T3 Code settings.", + }, + }), + supportsConversationRollback: false, + supportsTextGeneration: false, + } satisfies ServerProviderDraft; + const metadata = yield* SubscriptionRef.make(initial); + const getSnapshot = SubscriptionRef.get(metadata).pipe(Effect.map(stampIdentity)); + const checkProvider = Effect.gen(function* () { + if (!enabled) return yield* getSnapshot; + const result = yield* probeOhMyPiModels(effectiveConfig, processEnv, serverConfig.cwd).pipe( + Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner), + Effect.result, + ); + const checkedAt = DateTime.formatIso(yield* DateTime.now); + yield* SubscriptionRef.update(metadata, (draft): ServerProviderDraft => { + if (result._tag === "Success") { + return { + ...draft, + installed: true, + version: result.success.version, + models: [...initial.models, ...result.success.models], + status: "ready", + checkedAt, + message: + "Uses OhMyPi's local credentials. Run omp on the server to configure a model and sign in.", + }; + } + const cause = result.failure; + const missing = isCommandMissingCause(cause); + return { + ...draft, + installed: !missing, + status: "error", + checkedAt, + message: missing + ? "OhMyPi CLI (omp) is not installed or not on PATH." + : "Could not load OhMyPi models. Check the binary path and run omp models --json on the server. The previous model list is unchanged.", + }; + }); + return yield* getSnapshot; + }); + const snapshot = yield* makeManagedServerProvider({ + resolveMaintenance: () => + Effect.succeed( + makeManualOnlyProviderMaintenanceCapabilities({ + provider: DRIVER, + packageName: "@oh-my-pi/pi-coding-agent", + }), + ), + getSettings: Effect.succeed(effectiveConfig), + streamSettings: Stream.empty, + haveSettingsChanged: () => false, + refreshOnInterval: false, + initialSnapshot: () => getSnapshot, + checkProvider, + enrichSnapshot: ({ publishSnapshot }) => + SubscriptionRef.changes(metadata).pipe( + Stream.runForEach((draft) => publishSnapshot(stampIdentity(draft))), + ), + }).pipe( + Effect.mapError( + (cause) => + new ProviderDriverError({ + driver: DRIVER, + instanceId, + detail: "Failed to build OhMyPi snapshot.", + cause, + }), + ), + ); + const adapter = yield* makeOhMyPiAdapter(effectiveConfig, { + instanceId, + environment: processEnv, + ...(eventLoggers.native ? { nativeEventLogger: eventLoggers.native } : {}), + onAvailableCommands: (commands, cwd) => + SubscriptionRef.update(metadata, (draft) => ({ + ...draft, + workspaceSnapshots: [ + ...(draft.workspaceSnapshots ?? []).filter((workspace) => workspace.cwd !== cwd), + { + cwd, + checkedAt: draft.checkedAt, + skills: [], + slashCommands: commands + .filter((command) => command.name.trim()) + .map((command) => ({ + name: command.name, + description: command.description, + ...(command.input ? { input: command.input } : {}), + })), + }, + ].slice(-32), + })), + }); + const unsupported = (operation: string) => + Effect.fail( + new TextGenerationError({ + operation, + detail: + "OhMyPi does not support background text generation in T3 Code. Select another provider for this action.", + }), + ); + return { + instanceId, + driverKind: DRIVER, + continuationIdentity, + displayName, + accentColor, + enabled, + snapshot: { ...snapshot, getSnapshot }, + snapshotForCwd: (cwd) => + getSnapshot.pipe( + Effect.map((snapshot) => ({ + ...snapshot, + slashCommands: + snapshot.workspaceSnapshots?.find((workspace) => workspace.cwd === cwd) + ?.slashCommands ?? [], + })), + ), + adapter, + textGeneration: { + generateCommitMessage: () => unsupported("generateCommitMessage"), + generatePrContent: () => unsupported("generatePrContent"), + generateBranchName: () => unsupported("generateBranchName"), + generateThreadTitle: () => unsupported("generateThreadTitle"), + }, + } satisfies ProviderInstance; + }), +}; diff --git a/apps/server/src/provider/Drivers/OhMyPiModels.ts b/apps/server/src/provider/Drivers/OhMyPiModels.ts new file mode 100644 index 000000000000..f5adedcd25ce --- /dev/null +++ b/apps/server/src/provider/Drivers/OhMyPiModels.ts @@ -0,0 +1,94 @@ +import type { OhMyPiSettings, ServerProviderModel } from "@t3tools/contracts"; +import { createModelCapabilities } from "@t3tools/shared/model"; +import { resolveSpawnCommand } from "@t3tools/shared/shell"; +import * as Effect from "effect/Effect"; +import * as Schema from "effect/Schema"; +import { ChildProcess } from "effect/unstable/process"; +import { parseGenericCliVersion, spawnAndCollect } from "../providerSnapshot.ts"; + +const ModelsOutput = Schema.Struct({ + models: Schema.Array( + Schema.Struct({ + provider: Schema.NonEmptyString, + id: Schema.NonEmptyString, + name: Schema.String, + thinking: Schema.optionalKey(Schema.NullOr(Schema.Array(Schema.NonEmptyString))), + }), + ), +}); +const decodeModels = Schema.decodeEffect(Schema.fromJsonString(ModelsOutput)); + +class OhMyPiModelsError extends Schema.TaggedError()("OhMyPiModelsError", { + detail: Schema.String, +}) {} + +/** Discover the installed CLI's catalog without authenticating or starting an ACP session. */ +export const probeOhMyPiModels = Effect.fn("probeOhMyPiModels")(function* ( + settings: Pick, + environment: NodeJS.ProcessEnv, + cwd: string, +) { + const command = settings.binaryPath || "omp"; + const run = Effect.fn("OhMyPiModels.run")(function* (args: ReadonlyArray) { + const spawn = yield* resolveSpawnCommand(command, args, { env: environment }); + return yield* spawnAndCollect( + command, + ChildProcess.make(spawn.command, spawn.args, { + cwd, + env: environment, + shell: spawn.shell, + }), + ); + }); + const [catalog, version] = yield* Effect.all( + [ + Effect.gen(function* () { + const result = yield* run(["models", "--json"]); + if (result.code !== 0) { + return yield* new OhMyPiModelsError({ + detail: "omp models --json exited unsuccessfully.", + }); + } + return yield* decodeModels(result.stdout); + }), + run(["--version"]).pipe( + Effect.timeout("4 seconds"), + Effect.map((result) => (result.code === 0 ? parseGenericCliVersion(result.stdout) : null)), + Effect.catch(() => Effect.succeed(null)), + ), + ], + { concurrency: "unbounded" }, + ); + const seen = new Set(); + const models = catalog.models.flatMap((model): ServerProviderModel[] => { + const slug = `${model.provider}/${model.id}`; + if (seen.has(slug)) return []; + seen.add(slug); + const thinking = [...new Set(model.thinking ?? [])]; + return [ + { + slug, + name: model.name.trim() || model.id, + subProvider: model.provider, + isCustom: false, + capabilities: createModelCapabilities({ + optionDescriptors: + thinking.length > 0 + ? [ + { + id: "thinking", + label: "Thinking", + type: "select", + options: [...new Set(["off", "auto", ...thinking])].map((id) => ({ + id, + label: id === "off" ? "Off" : id === "auto" ? "Auto" : id, + })), + }, + ] + : [], + }), + }, + ]; + }); + return { models, version }; +}, Effect.timeout("30 seconds")); diff --git a/apps/server/src/provider/Layers/OhMyPiAdapter.ts b/apps/server/src/provider/Layers/OhMyPiAdapter.ts new file mode 100644 index 000000000000..7065ed074772 --- /dev/null +++ b/apps/server/src/provider/Layers/OhMyPiAdapter.ts @@ -0,0 +1,1045 @@ +/** + * OhMyPiAdapterLive — OhMyPi CLI (`omp acp`) via ACP. + * + * @module OhMyPiAdapterLive + */ + +import { + ApprovalRequestId, + type OhMyPiSettings, + type ProviderOptionSelection, + EventId, + type ProviderApprovalDecision, + type ProviderInteractionMode, + type ProviderRuntimeEvent, + type ProviderSession, + ProviderDriverKind, + ProviderInstanceId, + RuntimeRequestId, + type RuntimeMode, + type ThreadId, + TurnId, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Crypto from "effect/Crypto"; +import * as Deferred from "effect/Deferred"; +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; +import * as FileSystem from "effect/FileSystem"; +import * as Option from "effect/Option"; +import * as Path from "effect/Path"; +import * as PubSub from "effect/PubSub"; +import * as Schema from "effect/Schema"; +import * as Scope from "effect/Scope"; +import * as Semaphore from "effect/Semaphore"; +import * as Stream from "effect/Stream"; +import * as SynchronizedRef from "effect/SynchronizedRef"; +import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner"; +import * as EffectAcpErrors from "effect-acp/errors"; +import type * as EffectAcpSchema from "effect-acp/schema"; + +import { resolveAttachmentPath } from "../../attachmentStore.ts"; +import { ServerConfig } from "../../config.ts"; +import { buildRuntimeInstructions } from "../RuntimeInstructions.ts"; +import * as McpProviderSession from "../../mcp/McpProviderSession.ts"; +import { + ProviderAdapterProcessError, + ProviderAdapterRequestError, + ProviderAdapterSessionNotFoundError, + ProviderAdapterValidationError, +} from "../Errors.ts"; +import { mapAcpToAdapterError } from "../acp/AcpAdapterSupport.ts"; +import type * as AcpSessionRuntime from "../acp/AcpSessionRuntime.ts"; +import { + makeAcpAssistantItemEvent, + makeAcpContentDeltaEvent, + makeAcpPlanUpdatedEvent, + makeAcpRequestOpenedEvent, + makeAcpRequestResolvedEvent, + makeAcpToolCallEvent, +} from "../acp/AcpCoreRuntimeEvents.ts"; +import { parsePermissionRequest } from "../acp/AcpRuntimeModel.ts"; +import { makeAcpNativeLoggerFactory } from "../acp/AcpNativeLogging.ts"; +import { + applyOhMyPiAcpModelSelection, + makeOhMyPiAcpRuntime, + selectOhMyPiPermissionOption, +} from "../acp/OhMyPiAcpSupport.ts"; +import type { ProviderAdapterShape } from "../Services/ProviderAdapter.ts"; +import type { ProviderAdapterError } from "../Errors.ts"; +type OhMyPiAdapterShape = ProviderAdapterShape; +import { type EventNdjsonLogger, makeEventNdjsonLogger } from "./EventNdjsonLogger.ts"; +const encodeUnknownJsonStringExit = Schema.encodeUnknownExit(Schema.fromJsonString(Schema.Unknown)); + +const PROVIDER = ProviderDriverKind.make("ohMyPi"); +const OH_MY_PI_RESUME_VERSION = 1 as const; +function encodeJsonStringForDiagnostics(input: unknown): string | undefined { + const result = encodeUnknownJsonStringExit(input); + return Exit.isSuccess(result) ? result.value : undefined; +} + +export interface OhMyPiAdapterLiveOptions { + readonly environment?: NodeJS.ProcessEnv; + readonly nativeEventLogPath?: string; + readonly nativeEventLogger?: EventNdjsonLogger; + /** + * Selections are honored when `modelSelection.instanceId` matches this value. + * Defaults to the legacy built-in instance id (`ohMyPi`). + */ + readonly instanceId?: ProviderInstanceId; + readonly onAvailableCommands?: ( + commands: ReadonlyArray, + cwd: string, + ) => Effect.Effect; +} + +interface PendingApproval { + readonly decision: Deferred.Deferred; + readonly kind: string | "unknown"; +} + +interface OhMyPiSessionContext { + readonly threadId: ThreadId; + session: ProviderSession; + readonly scope: Scope.Closeable; + readonly acp: AcpSessionRuntime.AcpSessionRuntime["Service"]; + notificationFiber: Fiber.Fiber | undefined; + readonly pendingApprovals: Map; + readonly turns: Array<{ id: TurnId; items: Array }>; + lastPlanFingerprint: string | undefined; + activeTurnId: TurnId | undefined; + /** Number of sendTurn prompts currently in flight or being prepared. + * >0 means a turn is actively running, so a new sendTurn is a steer that + * continues it, and only the last remaining prompt settles the turn. */ + promptsInFlight: number; + interruptionVersion: number; + stopped: boolean; +} + +function settlePendingApprovalsAsCancelled( + pendingApprovals: ReadonlyMap, +): Effect.Effect { + const pendingEntries = Array.from(pendingApprovals.values()); + return Effect.forEach( + pendingEntries, + (pending) => Deferred.succeed(pending.decision, "cancel").pipe(Effect.ignore), + { + discard: true, + }, + ); +} + +const ResumeCursor = Schema.Struct({ + schemaVersion: Schema.Literal(1), + sessionId: Schema.NonEmptyString, +}); +const decodeResumeCursor = Schema.decodeUnknownOption(ResumeCursor); + +function parseOhMyPiResume(raw: unknown): { sessionId: string } | undefined { + return Option.getOrUndefined(decodeResumeCursor(raw)); +} + +function applyRequestedSessionConfiguration(input: { + readonly runtime: AcpSessionRuntime.AcpSessionRuntime["Service"]; + readonly runtimeMode: RuntimeMode; + readonly interactionMode: ProviderInteractionMode | undefined; + readonly modelSelection: + | { + readonly model: string; + readonly options?: ReadonlyArray | null | undefined; + } + | undefined; + readonly mapError: (context: { + readonly cause: import("effect-acp/errors").AcpError; + readonly method: "session/set_config_option" | "session/set_mode"; + }) => E; +}): Effect.Effect { + return Effect.gen(function* () { + if (input.modelSelection) { + yield* applyOhMyPiAcpModelSelection({ + runtime: input.runtime, + model: input.modelSelection.model, + selections: input.modelSelection.options, + mapError: ({ cause }) => + input.mapError({ + cause, + method: "session/set_config_option", + }), + }); + } + + const requestedModeId = input.interactionMode === "plan" ? "plan" : "default"; + const modes = yield* input.runtime.getModeState; + if (!modes?.availableModes.some((mode) => mode.id === requestedModeId)) return; + + yield* input.runtime.setMode(requestedModeId).pipe( + Effect.mapError((cause) => + input.mapError({ + cause, + method: "session/set_mode", + }), + ), + ); + }); +} + +function selectAutoApprovedPermissionOption( + request: EffectAcpSchema.RequestPermissionRequest, +): string | undefined { + const allowAlwaysOption = request.options.find((option) => option.kind === "allow_always"); + if (typeof allowAlwaysOption?.optionId === "string" && allowAlwaysOption.optionId.trim()) { + return allowAlwaysOption.optionId.trim(); + } + + const allowOnceOption = request.options.find((option) => option.kind === "allow_once"); + if (typeof allowOnceOption?.optionId === "string" && allowOnceOption.optionId.trim()) { + return allowOnceOption.optionId.trim(); + } + + return undefined; +} + +export function makeOhMyPiAdapter( + ohMyPiSettings: OhMyPiSettings, + options?: OhMyPiAdapterLiveOptions, +) { + return Effect.gen(function* () { + const boundInstanceId = options?.instanceId ?? ProviderInstanceId.make("ohMyPi"); + const fileSystem = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const childProcessSpawner = yield* ChildProcessSpawner.ChildProcessSpawner; + const serverConfig = yield* Effect.service(ServerConfig); + const crypto = yield* Crypto.Crypto; + const nativeEventLogger = + options?.nativeEventLogger ?? + (options?.nativeEventLogPath !== undefined + ? yield* makeEventNdjsonLogger(options.nativeEventLogPath, { + stream: "native", + }) + : undefined); + const managedNativeEventLogger = + options?.nativeEventLogger === undefined ? nativeEventLogger : undefined; + const makeAcpNativeLoggers = yield* makeAcpNativeLoggerFactory(); + + const sessions = new Map(); + const threadLocksRef = yield* SynchronizedRef.make(new Map()); + const runtimeEventPubSub = yield* PubSub.unbounded(); + + const nowIso = Effect.map(DateTime.now, DateTime.formatIso); + const randomUUIDv4 = crypto.randomUUIDv4.pipe( + Effect.mapError( + (cause) => + new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "crypto/randomUUIDv4", + detail: "Failed to generate OhMyPi runtime identifier.", + cause, + }), + ), + ); + const nextEventId = Effect.map(randomUUIDv4, (id) => EventId.make(id)); + const makeEventStamp = () => Effect.all({ eventId: nextEventId, createdAt: nowIso }); + const mapExtensionFailure = (effect: Effect.Effect) => + effect.pipe( + Effect.mapError( + (cause) => + new EffectAcpErrors.AcpTransportError({ + detail: "Failed to process OhMyPi ACP extension event.", + cause, + }), + ), + ); + + const offerRuntimeEvent = (event: ProviderRuntimeEvent) => + PubSub.publish(runtimeEventPubSub, event).pipe(Effect.asVoid); + + const getThreadSemaphore = (threadId: string) => + SynchronizedRef.modifyEffect(threadLocksRef, (current) => { + const existing: Option.Option = Option.fromNullishOr( + current.get(threadId), + ); + return Option.match(existing, { + onNone: () => + Semaphore.make(1).pipe( + Effect.map((semaphore) => { + const next = new Map(current); + next.set(threadId, semaphore); + return [semaphore, next] as const; + }), + ), + onSome: (semaphore) => Effect.succeed([semaphore, current] as const), + }); + }); + + const withThreadLock = (threadId: string, effect: Effect.Effect) => + Effect.flatMap(getThreadSemaphore(threadId), (semaphore) => semaphore.withPermit(effect)); + + const logNative = ( + threadId: ThreadId, + method: string, + payload: unknown, + _source: "acp.jsonrpc", + ) => + Effect.gen(function* () { + if (!nativeEventLogger) return; + const observedAt = yield* nowIso; + yield* nativeEventLogger.write( + { + observedAt, + event: { + id: yield* randomUUIDv4, + kind: "notification", + provider: PROVIDER, + createdAt: observedAt, + method, + threadId, + payload, + }, + }, + threadId, + ); + }); + + const emitPlanUpdate = ( + ctx: OhMyPiSessionContext, + payload: { + readonly explanation?: string | null; + readonly plan: ReadonlyArray<{ + readonly step: string; + readonly status: "pending" | "inProgress" | "completed"; + }>; + }, + rawPayload: unknown, + source: "acp.jsonrpc", + method: string, + ) => + Effect.gen(function* () { + const fingerprint = `${ctx.activeTurnId ?? "no-turn"}:${encodeJsonStringForDiagnostics(payload) ?? "[unserializable payload]"}`; + if (ctx.lastPlanFingerprint === fingerprint) { + return; + } + ctx.lastPlanFingerprint = fingerprint; + yield* offerRuntimeEvent( + makeAcpPlanUpdatedEvent({ + stamp: yield* makeEventStamp(), + provider: PROVIDER, + threadId: ctx.threadId, + turnId: ctx.activeTurnId, + payload, + source, + method, + rawPayload, + }), + ); + }); + + const requireSession = ( + threadId: ThreadId, + ): Effect.Effect => { + const ctx = sessions.get(threadId); + if (!ctx || ctx.stopped || ctx.session.status === "error") { + return Effect.fail( + new ProviderAdapterSessionNotFoundError({ provider: PROVIDER, threadId }), + ); + } + return Effect.succeed(ctx); + }; + + const stopSessionInternal = (ctx: OhMyPiSessionContext) => + Effect.gen(function* () { + if (ctx.stopped) return; + ctx.stopped = true; + yield* settlePendingApprovalsAsCancelled(ctx.pendingApprovals); + if (ctx.notificationFiber) { + yield* Fiber.interrupt(ctx.notificationFiber); + } + yield* Effect.ignore(Scope.close(ctx.scope, Exit.void)); + sessions.delete(ctx.threadId); + yield* offerRuntimeEvent({ + type: "session.exited", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId: ctx.threadId, + payload: { exitKind: "graceful" }, + }); + }); + + const startSession: OhMyPiAdapterShape["startSession"] = (input) => + withThreadLock( + input.threadId, + Effect.gen(function* () { + if (input.provider !== undefined && input.provider !== PROVIDER) { + return yield* new ProviderAdapterValidationError({ + provider: PROVIDER, + operation: "startSession", + issue: `Expected provider '${PROVIDER}' but received '${input.provider}'.`, + }); + } + if (!input.cwd?.trim()) { + return yield* new ProviderAdapterValidationError({ + provider: PROVIDER, + operation: "startSession", + issue: "cwd is required and must be non-empty.", + }); + } + + const cwd = path.resolve(input.cwd.trim()); + const ohMyPiModelSelection = + input.modelSelection?.instanceId === boundInstanceId ? input.modelSelection : undefined; + const existing = sessions.get(input.threadId); + if (existing && !existing.stopped) { + yield* stopSessionInternal(existing); + } + + const pendingApprovals = new Map(); + const sessionScope = yield* Scope.make("sequential"); + let sessionScopeTransferred = false; + yield* Effect.addFinalizer(() => + sessionScopeTransferred ? Effect.void : Scope.close(sessionScope, Exit.void), + ); + let ctx!: OhMyPiSessionContext; + + const resumeSessionId = parseOhMyPiResume(input.resumeCursor)?.sessionId; + const acpNativeLoggers = makeAcpNativeLoggers({ + nativeEventLogger, + provider: PROVIDER, + threadId: input.threadId, + }); + + const mcpSession = McpProviderSession.readMcpProviderSession(input.threadId); + const acp = yield* makeOhMyPiAcpRuntime({ + ohMyPiSettings, + ...(options?.environment || mcpSession?.agentDeviceEnvironment + ? { + environment: McpProviderSession.withAgentDeviceEnvironment( + options?.environment ?? process.env, + mcpSession, + ), + } + : {}), + childProcessSpawner, + cwd, + runtimeMode: input.runtimeMode, + ...(resumeSessionId ? { resumeSessionId } : {}), + clientInfo: { name: "t3-code", version: "0.0.0" }, + ...(mcpSession + ? { + mcpServers: [ + { + type: "http" as const, + name: "t3-code", + url: mcpSession.endpoint, + headers: [ + { + name: "Authorization", + value: mcpSession.authorizationHeader, + }, + ], + }, + ], + } + : {}), + ...acpNativeLoggers, + }).pipe( + Effect.provideService(Crypto.Crypto, crypto), + Effect.provideService(Scope.Scope, sessionScope), + Effect.mapError( + (cause) => + new ProviderAdapterProcessError({ + provider: PROVIDER, + threadId: input.threadId, + detail: cause.message, + cause, + }), + ), + ); + const started = yield* Effect.gen(function* () { + yield* acp.handleRequestPermission((params) => + mapExtensionFailure( + Effect.gen(function* () { + yield* logNative( + input.threadId, + "session/request_permission", + params, + "acp.jsonrpc", + ); + if (input.runtimeMode === "full-access") { + const autoApprovedOptionId = selectAutoApprovedPermissionOption(params); + if (autoApprovedOptionId !== undefined) { + return { + outcome: { + outcome: "selected" as const, + optionId: autoApprovedOptionId, + }, + }; + } + } + const permissionRequest = parsePermissionRequest(params); + const requestId = ApprovalRequestId.make(yield* randomUUIDv4); + const runtimeRequestId = RuntimeRequestId.make(requestId); + const decision = yield* Deferred.make(); + pendingApprovals.set(requestId, { + decision, + kind: permissionRequest.kind, + }); + yield* offerRuntimeEvent( + makeAcpRequestOpenedEvent({ + stamp: yield* makeEventStamp(), + provider: PROVIDER, + threadId: input.threadId, + turnId: ctx?.activeTurnId, + requestId: runtimeRequestId, + permissionRequest, + detail: + permissionRequest.detail ?? + encodeJsonStringForDiagnostics(params)?.slice(0, 2000) ?? + "[unserializable params]", + args: params, + source: "acp.jsonrpc", + method: "session/request_permission", + rawPayload: params, + }), + ); + const resolved = yield* Deferred.await(decision); + pendingApprovals.delete(requestId); + yield* offerRuntimeEvent( + makeAcpRequestResolvedEvent({ + stamp: yield* makeEventStamp(), + provider: PROVIDER, + threadId: input.threadId, + turnId: ctx?.activeTurnId, + requestId: runtimeRequestId, + permissionRequest, + decision: resolved, + }), + ); + const optionId = selectOhMyPiPermissionOption(params, resolved); + return { + outcome: + optionId === undefined + ? ({ outcome: "cancelled" } as const) + : { + outcome: "selected" as const, + optionId, + }, + }; + }), + ), + ); + return yield* acp.start(); + }).pipe( + Effect.mapError((error) => + mapAcpToAdapterError(PROVIDER, input.threadId, "session/start", error), + ), + ); + + yield* applyRequestedSessionConfiguration({ + runtime: acp, + runtimeMode: input.runtimeMode, + interactionMode: undefined, + modelSelection: ohMyPiModelSelection, + mapError: ({ cause, method }) => + mapAcpToAdapterError(PROVIDER, input.threadId, method, cause), + }); + + const now = yield* nowIso; + const session: ProviderSession = { + provider: PROVIDER, + providerInstanceId: boundInstanceId, + status: "ready", + runtimeMode: input.runtimeMode, + cwd, + model: ohMyPiModelSelection?.model, + threadId: input.threadId, + resumeCursor: { + schemaVersion: OH_MY_PI_RESUME_VERSION, + sessionId: started.sessionId, + }, + createdAt: now, + updatedAt: now, + }; + + ctx = { + threadId: input.threadId, + session, + scope: sessionScope, + acp, + notificationFiber: undefined, + pendingApprovals, + turns: [], + lastPlanFingerprint: undefined, + activeTurnId: undefined, + promptsInFlight: 0, + interruptionVersion: 0, + stopped: false, + }; + + const nf = yield* Stream.runDrain( + Stream.mapEffect(acp.getEvents(), (event) => + Effect.gen(function* () { + switch (event._tag) { + case "EventStreamBarrier": + yield* Deferred.succeed(event.acknowledge, undefined); + return; + case "ConfigOptionsUpdated": + return; + case "AvailableCommandsUpdated": + yield* ( + options?.onAvailableCommands?.(event.availableCommands, cwd) ?? Effect.void + ); + return; + case "ConnectionTerminated": + ctx.session = { ...ctx.session, status: "error", updatedAt: yield* nowIso }; + yield* settlePendingApprovalsAsCancelled(ctx.pendingApprovals); + yield* offerRuntimeEvent({ + type: "session.exited", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId: ctx.threadId, + payload: { + exitKind: "error", + recoverable: true, + reason: event.error.message, + }, + }); + return; + case "ModeChanged": + return; + case "AssistantItemStarted": + yield* offerRuntimeEvent( + makeAcpAssistantItemEvent({ + stamp: yield* makeEventStamp(), + provider: PROVIDER, + threadId: ctx.threadId, + turnId: ctx.activeTurnId, + itemId: event.itemId, + lifecycle: "item.started", + }), + ); + return; + case "AssistantItemCompleted": + yield* offerRuntimeEvent( + makeAcpAssistantItemEvent({ + stamp: yield* makeEventStamp(), + provider: PROVIDER, + threadId: ctx.threadId, + turnId: ctx.activeTurnId, + itemId: event.itemId, + lifecycle: "item.completed", + }), + ); + return; + case "PlanUpdated": + yield* logNative( + ctx.threadId, + "session/update", + event.rawPayload, + "acp.jsonrpc", + ); + yield* emitPlanUpdate( + ctx, + event.payload, + event.rawPayload, + "acp.jsonrpc", + "session/update", + ); + return; + case "ToolCallUpdated": + yield* logNative( + ctx.threadId, + "session/update", + event.rawPayload, + "acp.jsonrpc", + ); + yield* offerRuntimeEvent( + makeAcpToolCallEvent({ + stamp: yield* makeEventStamp(), + provider: PROVIDER, + threadId: ctx.threadId, + turnId: ctx.activeTurnId, + toolCall: event.toolCall, + rawPayload: event.rawPayload, + }), + ); + return; + case "ThoughtDelta": + case "ContentDelta": + yield* logNative( + ctx.threadId, + "session/update", + event.rawPayload, + "acp.jsonrpc", + ); + yield* offerRuntimeEvent( + makeAcpContentDeltaEvent({ + stamp: yield* makeEventStamp(), + provider: PROVIDER, + threadId: ctx.threadId, + turnId: ctx.activeTurnId, + ...(event._tag === "ContentDelta" && event.itemId + ? { itemId: event.itemId } + : {}), + ...(event._tag === "ThoughtDelta" ? { streamKind: "reasoning_text" } : {}), + text: event.text, + rawPayload: event.rawPayload, + }), + ); + return; + } + }), + ), + ).pipe( + Effect.catch((cause) => + Effect.logError("Failed to process OhMyPi runtime notification.", { cause }), + ), + // Fork into the session scope, not the calling fiber. `forkChild` + // makes this a child of `startSession`, and Effect interrupts a + // fiber's children when it completes, so the consumer died as soon + // as `startSession` returned and every later notification was + // dropped. The scope is created, stored on the context and closed + // on teardown already; only the fork target was wrong. + Effect.forkIn(ctx.scope), + ); + + ctx.notificationFiber = nf; + sessions.set(input.threadId, ctx); + sessionScopeTransferred = true; + + yield* offerRuntimeEvent({ + type: "session.started", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId: input.threadId, + payload: { resume: started.initializeResult }, + }); + yield* offerRuntimeEvent({ + type: "session.state.changed", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId: input.threadId, + payload: { state: "ready", reason: "OhMyPi ACP session ready" }, + }); + yield* offerRuntimeEvent({ + type: "thread.started", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId: input.threadId, + payload: { providerThreadId: started.sessionId }, + }); + + return session; + }).pipe(Effect.scoped), + ); + + const sendTurn: OhMyPiAdapterShape["sendTurn"] = (input) => + Effect.gen(function* () { + const scope = yield* Scope.Scope; + // Serialize preparation through dispatch so a steer always sees the + // reserved turn and can cancel a prompt that has reached the runtime. + const sending = yield* withThreadLock( + input.threadId, + Effect.gen(function* () { + const ctx = yield* requireSession(input.threadId); + const interruptionVersion = ctx.interruptionVersion; + // Keep a steering prompt under the active T3 turn until both RPCs settle. + const steeringTurnId = ctx.promptsInFlight > 0 ? ctx.activeTurnId : undefined; + const turnId = steeringTurnId ?? TurnId.make(yield* randomUUIDv4); + // Count this prompt immediately so a superseded in-flight prompt + // resolving from here on does not settle the turn; the matching + // decrement is the `ensuring` below. + ctx.activeTurnId = turnId; + ctx.promptsInFlight += 1; + const dispatched = yield* Deferred.make(); + + const sending = yield* Effect.gen(function* () { + const turnModelSelection = + input.modelSelection?.instanceId === boundInstanceId + ? input.modelSelection + : undefined; + const model = turnModelSelection?.model ?? ctx.session.model; + yield* applyRequestedSessionConfiguration({ + runtime: ctx.acp, + runtimeMode: ctx.session.runtimeMode, + interactionMode: input.interactionMode, + modelSelection: + model === undefined + ? undefined + : { + model, + options: turnModelSelection?.options, + }, + mapError: ({ cause, method }) => + mapAcpToAdapterError(PROVIDER, input.threadId, method, cause), + }); + const modelConfig = (yield* ctx.acp.getConfigOptions).find( + (option) => option.category === "model", + ); + const resolvedModel = + modelConfig?.type === "select" ? modelConfig.currentValue : model; + ctx.activeTurnId = turnId; + if (steeringTurnId === undefined) { + ctx.lastPlanFingerprint = undefined; + } + ctx.session = { + ...ctx.session, + activeTurnId: turnId, + updatedAt: yield* nowIso, + }; + + if (steeringTurnId === undefined) { + yield* offerRuntimeEvent({ + type: "turn.started", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId: input.threadId, + turnId, + payload: { model: resolvedModel }, + }); + } + + const promptParts: Array = []; + const rawPrompt = input.input?.trim() ?? ""; + if (rawPrompt) { + promptParts.push({ type: "text", text: rawPrompt }); + } + if (input.attachments && input.attachments.length > 0) { + for (const attachment of input.attachments) { + // Send images inline. Generic files reach the agent + // through the path line ProviderService puts in the prompt. + if (attachment.type !== "image") { + continue; + } + const attachmentPath = resolveAttachmentPath({ + attachmentsDir: serverConfig.attachmentsDir, + attachment, + }); + if (!attachmentPath) { + return yield* new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "session/prompt", + detail: `Invalid attachment id '${attachment.id}'.`, + }); + } + const bytes = yield* fileSystem.readFile(attachmentPath).pipe( + Effect.mapError( + (cause) => + new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "session/prompt", + detail: cause.message, + cause, + }), + ), + ); + promptParts.push({ + type: "image", + data: Buffer.from(bytes).toString("base64"), + mimeType: attachment.mimeType, + }); + } + } + + if (promptParts.length === 0) { + return yield* new ProviderAdapterValidationError({ + provider: PROVIDER, + operation: "sendTurn", + issue: "Turn requires non-empty text or attachments.", + }); + } + + if (steeringTurnId !== undefined) { + yield* settlePendingApprovalsAsCancelled(ctx.pendingApprovals); + yield* ctx.acp.cancel.pipe( + Effect.mapError((error) => + mapAcpToAdapterError(PROVIDER, input.threadId, "session/cancel", error), + ), + ); + } + + // ACP has no system-message field; keep runtime context separate from the user's text. + const result = + interruptionVersion !== ctx.interruptionVersion + ? { stopReason: "cancelled" as const } + : yield* ctx.acp + .prompt( + { + prompt: [ + ...promptParts, + { + type: "text", + text: buildRuntimeInstructions({ + harness: "OhMyPi", + model: resolvedModel, + }), + }, + ], + }, + { dispatched }, + ) + .pipe( + Effect.mapError((error) => + mapAcpToAdapterError(PROVIDER, input.threadId, "session/prompt", error), + ), + ); + + yield* ctx.acp.drainEvents; + const turnRecord = ctx.turns.find((turn) => turn.id === turnId); + if (turnRecord) { + turnRecord.items.push({ prompt: promptParts, result }); + } else { + ctx.turns.push({ id: turnId, items: [{ prompt: promptParts, result }] }); + } + ctx.session = { + ...ctx.session, + activeTurnId: turnId, + updatedAt: yield* nowIso, + model: resolvedModel, + }; + + // Only the last remaining prompt settles the turn — a steer- + // superseded prompt resolving (usually cancelled) while another is + // in flight or pending must leave the merged turn running. + if (ctx.promptsInFlight === 1) { + yield* offerRuntimeEvent({ + type: "turn.completed", + ...(yield* makeEventStamp()), + provider: PROVIDER, + threadId: input.threadId, + turnId, + payload: { + state: result.stopReason === "cancelled" ? "cancelled" : "completed", + stopReason: result.stopReason ?? null, + }, + }); + } + + return { + threadId: input.threadId, + turnId, + resumeCursor: ctx.session.resumeCursor, + }; + }).pipe( + Effect.ensuring( + Effect.sync(() => { + ctx.promptsInFlight = Math.max(0, ctx.promptsInFlight - 1); + }), + ), + Effect.forkIn(scope), + ); + yield* Effect.raceFirst(Deferred.await(dispatched), Fiber.join(sending)); + return sending; + }), + ); + return yield* Fiber.join(sending); + }).pipe(Effect.scoped); + + const interruptTurn: OhMyPiAdapterShape["interruptTurn"] = (threadId) => + Effect.gen(function* () { + const ctx = yield* requireSession(threadId); + // Invalidate preparation even when ACP has no prompt to cancel yet. + ctx.interruptionVersion += 1; + yield* settlePendingApprovalsAsCancelled(ctx.pendingApprovals); + yield* Effect.ignore( + ctx.acp.cancel.pipe( + Effect.mapError((error) => + mapAcpToAdapterError(PROVIDER, threadId, "session/cancel", error), + ), + ), + ); + }); + + const respondToRequest: OhMyPiAdapterShape["respondToRequest"] = ( + threadId, + requestId, + decision, + ) => + Effect.gen(function* () { + const ctx = yield* requireSession(threadId); + const pending = ctx.pendingApprovals.get(requestId); + if (!pending) { + return yield* new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "session/request_permission", + detail: `Unknown pending approval request: ${requestId}`, + }); + } + yield* Deferred.succeed(pending.decision, decision); + }); + + const respondToUserInput: OhMyPiAdapterShape["respondToUserInput"] = () => + Effect.fail( + new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "session/elicitation", + detail: "No pending OhMyPi user-input request.", + }), + ); + + const readThread: OhMyPiAdapterShape["readThread"] = (threadId) => + Effect.gen(function* () { + const ctx = yield* requireSession(threadId); + return { threadId, turns: ctx.turns }; + }); + + const rollbackThread: OhMyPiAdapterShape["rollbackThread"] = () => + Effect.fail( + new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "rollbackThread", + detail: "OhMyPi does not support conversation rollback through ACP.", + }), + ); + + const stopSession: OhMyPiAdapterShape["stopSession"] = (threadId) => + withThreadLock( + threadId, + Effect.gen(function* () { + const ctx = yield* requireSession(threadId); + yield* stopSessionInternal(ctx); + }), + ); + + const listSessions: OhMyPiAdapterShape["listSessions"] = () => + Effect.sync(() => Array.from(sessions.values(), (c) => ({ ...c.session }))); + + const hasSession: OhMyPiAdapterShape["hasSession"] = (threadId) => + Effect.sync(() => { + const c = sessions.get(threadId); + return c !== undefined && !c.stopped && c.session.status !== "error"; + }); + + const stopAll: OhMyPiAdapterShape["stopAll"] = () => + Effect.forEach(sessions.values(), stopSessionInternal, { discard: true }); + + yield* Effect.addFinalizer(() => + Effect.forEach(sessions.values(), stopSessionInternal, { discard: true }).pipe( + Effect.catch((cause) => + Effect.logError("Failed to emit OhMyPi session shutdown event.", { cause }), + ), + Effect.tap(() => PubSub.shutdown(runtimeEventPubSub)), + Effect.tap(() => managedNativeEventLogger?.close() ?? Effect.void), + ), + ); + + const streamEvents = Stream.fromPubSub(runtimeEventPubSub); + + return { + provider: PROVIDER, + capabilities: { sessionModelSwitch: "in-session", supportsConversationRollback: false }, + compaction: { type: "slash-command", command: "/compact" }, + startSession, + sendTurn, + interruptTurn, + readThread, + rollbackThread, + respondToRequest, + respondToUserInput, + stopSession, + listSessions, + hasSession, + stopAll, + streamEvents, + } satisfies OhMyPiAdapterShape; + }); +} diff --git a/apps/server/src/provider/acp/OhMyPiAcpCliProbe.test.ts b/apps/server/src/provider/acp/OhMyPiAcpCliProbe.test.ts new file mode 100644 index 000000000000..0b2e95604d48 --- /dev/null +++ b/apps/server/src/provider/acp/OhMyPiAcpCliProbe.test.ts @@ -0,0 +1,29 @@ +/** Opt in with T3_OH_MY_PI_ACP_PROBE=1; initializes the installed CLI without sending a prompt. */ +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { expect, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as FileSystem from "effect/FileSystem"; +import { ChildProcessSpawner } from "effect/unstable/process"; +import { makeOhMyPiAcpRuntime } from "./OhMyPiAcpSupport.ts"; + +it.effect.skipIf(process.env.T3_OH_MY_PI_ACP_PROBE !== "1")( + "initializes the installed OhMyPi ACP CLI", + () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const cwd = yield* fs.makeTempDirectoryScoped(); + const agentDir = yield* fs.makeTempDirectoryScoped(); + const runtime = yield* makeOhMyPiAcpRuntime({ + ohMyPiSettings: { binaryPath: process.env.T3_OH_MY_PI_BINARY ?? "omp" }, + childProcessSpawner: yield* ChildProcessSpawner.ChildProcessSpawner, + environment: { ...process.env, PI_CODING_AGENT_DIR: agentDir }, + cwd, + runtimeMode: "approval-required", + clientInfo: { name: "t3-omp-probe", version: "0.0.0" }, + }); + const initialized = yield* runtime.initialize(); + expect(initialized.agentInfo?.name).toBe("oh-my-pi"); + expect(initialized.authMethods?.some((method) => method.id === "agent")).toBe(true); + expect(initialized.agentCapabilities?.loadSession).toBe(true); + }).pipe(Effect.scoped, Effect.provide(NodeServices.layer)), +); diff --git a/apps/server/src/provider/acp/OhMyPiAcpSupport.test.ts b/apps/server/src/provider/acp/OhMyPiAcpSupport.test.ts new file mode 100644 index 000000000000..d9fad5973c7c --- /dev/null +++ b/apps/server/src/provider/acp/OhMyPiAcpSupport.test.ts @@ -0,0 +1,139 @@ +import { describe, expect, it } from "@effect/vitest"; +import { OH_MY_PI_DEFAULT_MODEL } from "@t3tools/contracts"; +import * as Effect from "effect/Effect"; +import type * as AcpSchema from "effect-acp/schema"; +import { + applyOhMyPiAcpModelSelection, + ohMyPiApprovalMode, + selectOhMyPiPermissionOption, +} from "./OhMyPiAcpSupport.ts"; + +const config: ReadonlyArray = [ + { + id: "model", + category: "model", + name: "Model", + type: "select", + currentValue: "anthropic/sonnet", + options: [ + { value: "anthropic/sonnet", name: "Sonnet" }, + { value: "openai/gpt", name: "GPT" }, + ], + }, + { + id: "thinking", + category: "thought_level", + name: "Thinking", + type: "select", + currentValue: "high", + options: [ + { value: "off", name: "Off" }, + { value: "high", name: "High" }, + ], + }, +]; + +describe("OhMyPi ACP", () => { + it("overrides OMP's default yolo policy with the selected permission mode", () => { + expect(ohMyPiApprovalMode("approval-required")).toBe("always-ask"); + expect(ohMyPiApprovalMode("auto")).toBe("always-ask"); + expect(ohMyPiApprovalMode("auto-accept-edits")).toBe("write"); + expect(ohMyPiApprovalMode("full-access")).toBe("yolo"); + }); + it.effect("keeps the configured default and applies thinking after changing models", () => + Effect.gen(function* () { + const calls: unknown[] = []; + const runtime = { + getConfigOptions: Effect.succeed(config), + setModel: (model: string) => + Effect.sync(() => { + calls.push(["model", model]); + }), + setConfigOption: (id: string, value: string | boolean) => + Effect.sync(() => { + calls.push([id, value]); + return { configOptions: config }; + }), + }; + yield* applyOhMyPiAcpModelSelection({ + runtime, + model: OH_MY_PI_DEFAULT_MODEL, + selections: [], + mapError: ({ cause }) => cause, + }); + expect(calls).toEqual([]); + yield* applyOhMyPiAcpModelSelection({ + runtime, + model: "openai/gpt", + selections: [ + { id: "thinking", value: "high" }, + { id: "mode", value: "plan" }, + ], + mapError: ({ cause }) => cause, + }); + expect(calls).toEqual([ + ["model", "openai/gpt"], + ["thinking", "high"], + ]); + }), + ); + + it.effect("keeps the new model's default when a saved thinking level is unavailable", () => + Effect.gen(function* () { + let currentConfig = config; + const applied: unknown[] = []; + const runtime = { + getConfigOptions: Effect.sync(() => currentConfig), + setModel: (_model: string) => + Effect.sync(() => { + currentConfig = [ + { + id: "thinking", + category: "thought_level", + name: "Thinking", + type: "select", + currentValue: "low", + options: [{ value: "low", name: "Low" }], + }, + ]; + }), + setConfigOption: (id: string, value: string | boolean) => + Effect.sync(() => { + applied.push([id, value]); + return { configOptions: currentConfig }; + }), + }; + yield* applyOhMyPiAcpModelSelection({ + runtime, + model: "openai/gpt", + selections: [{ id: "thinking", value: "high" }], + mapError: ({ cause }) => cause, + }); + expect(applied).toEqual([]); + expect(currentConfig[0]?.currentValue).toBe("low"); + yield* applyOhMyPiAcpModelSelection({ + runtime, + model: "openai/gpt", + selections: [{ id: "thinking", value: "low" }], + mapError: ({ cause }) => cause, + }); + expect(applied).toEqual([["thinking", "low"]]); + }), + ); + + it("uses opaque permission IDs and cancels unavailable choices", () => { + const request: AcpSchema.RequestPermissionRequest = { + sessionId: "session", + toolCall: { toolCallId: "tool" }, + options: [ + { optionId: "yes-42", name: "Allow", kind: "allow_once" }, + { optionId: "no-7", name: "Deny", kind: "reject_once" }, + ], + }; + expect(selectOhMyPiPermissionOption(request, "accept")).toBe("yes-42"); + expect(selectOhMyPiPermissionOption(request, "acceptForSession")).toBe("yes-42"); + expect(selectOhMyPiPermissionOption(request, "decline")).toBe("no-7"); + expect(selectOhMyPiPermissionOption(request, "cancel")).toBeUndefined(); + expect(selectOhMyPiPermissionOption({ ...request, options: [] }, "accept")).toBeUndefined(); + }); +}); diff --git a/apps/server/src/provider/acp/OhMyPiAcpSupport.ts b/apps/server/src/provider/acp/OhMyPiAcpSupport.ts new file mode 100644 index 000000000000..f47f2421a436 --- /dev/null +++ b/apps/server/src/provider/acp/OhMyPiAcpSupport.ts @@ -0,0 +1,114 @@ +import { + OH_MY_PI_DEFAULT_MODEL, + type OhMyPiSettings, + type ProviderApprovalDecision, + type ProviderOptionSelection, + type RuntimeMode, +} from "@t3tools/contracts"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner"; +import type * as AcpErrors from "effect-acp/errors"; +import type * as AcpSchema from "effect-acp/schema"; +import * as AcpSessionRuntime from "./AcpSessionRuntime.ts"; + +interface OhMyPiAcpRuntimeInput extends Omit< + AcpSessionRuntime.AcpSessionRuntimeOptions, + "authMethodId" | "clientCapabilities" | "spawn" +> { + readonly childProcessSpawner: ChildProcessSpawner.ChildProcessSpawner["Service"]; + readonly ohMyPiSettings: Pick; + readonly environment?: NodeJS.ProcessEnv; + readonly runtimeMode?: RuntimeMode; +} + +export const makeOhMyPiAcpRuntime = Effect.fn("makeOhMyPiAcpRuntime")(function* ( + input: OhMyPiAcpRuntimeInput, +) { + const context = yield* Layer.build( + AcpSessionRuntime.layer({ + ...input, + spawn: { + command: input.ohMyPiSettings.binaryPath || "omp", + args: [ + "acp", + ...(input.runtimeMode ? ["--approval-mode", ohMyPiApprovalMode(input.runtimeMode)] : []), + ], + cwd: input.cwd, + ...(input.environment ? { env: input.environment } : {}), + }, + authMethodId: "agent", + // OMP runs its own tools when these capabilities are absent. + clientCapabilities: { fs: { readTextFile: false, writeTextFile: false }, terminal: false }, + }).pipe( + Layer.provide( + Layer.succeed(ChildProcessSpawner.ChildProcessSpawner, input.childProcessSpawner), + ), + ), + ); + return yield* Effect.service(AcpSessionRuntime.AcpSessionRuntime).pipe(Effect.provide(context)); +}); + +/** Permission IDs are opaque; select by the ACP kind supplied by the agent. */ +export function selectOhMyPiPermissionOption( + request: AcpSchema.RequestPermissionRequest, + decision: ProviderApprovalDecision, +): string | undefined { + if (decision === "cancel") return undefined; + const kind = + decision === "acceptForSession" + ? "allow_always" + : decision === "accept" + ? "allow_once" + : "reject_once"; + return ( + request.options.find((option) => option.kind === kind)?.optionId ?? + (decision === "acceptForSession" + ? request.options.find((option) => option.kind === "allow_once")?.optionId + : undefined) + ); +} + +export const applyOhMyPiAcpModelSelection = Effect.fn("applyOhMyPiAcpModelSelection")(function* < + E, +>(input: { + readonly runtime: Pick< + AcpSessionRuntime.AcpSessionRuntime["Service"], + "getConfigOptions" | "setModel" | "setConfigOption" + >; + readonly model: string | null | undefined; + readonly selections: ReadonlyArray | null | undefined; + readonly mapError: (context: { readonly cause: AcpErrors.AcpError }) => E; +}) { + if (input.model && input.model !== OH_MY_PI_DEFAULT_MODEL) { + yield* input.runtime + .setModel(input.model) + .pipe(Effect.mapError((cause) => input.mapError({ cause }))); + } + // The model change can alter the available thinking levels. + const config = yield* input.runtime.getConfigOptions; + for (const selection of input.selections ?? []) { + const option = config.find( + (entry) => entry.id === selection.id && entry.category === "thought_level", + ); + if (!option || option.type !== "select") continue; + const choices = option.options.flatMap((entry) => ("value" in entry ? [entry] : entry.options)); + // A model switch can leave a saved thinking level that the new model lacks. + if (!choices.some((choice) => choice.value === selection.value)) continue; + yield* input.runtime + .setConfigOption(option.id, selection.value) + .pipe(Effect.mapError((cause) => input.mapError({ cause }))); + } +}); + +export function ohMyPiApprovalMode(mode: RuntimeMode): string { + switch (mode) { + case "full-access": + return "yolo"; + case "auto-accept-edits": + return "write"; + case "auto": + case "approval-required": + return "always-ask"; + } +} diff --git a/apps/server/src/provider/builtInDrivers.ts b/apps/server/src/provider/builtInDrivers.ts index 60e3402eed42..57cd4579ee53 100644 --- a/apps/server/src/provider/builtInDrivers.ts +++ b/apps/server/src/provider/builtInDrivers.ts @@ -26,6 +26,7 @@ import { CursorDriver, type CursorDriverEnv } from "./Drivers/CursorDriver.ts"; import { GrokDriver, type GrokDriverEnv } from "./Drivers/GrokDriver.ts"; import { OpenCodeDriver, type OpenCodeDriverEnv } from "./Drivers/OpenCodeDriver.ts"; import { AntigravityDriver, type AntigravityDriverEnv } from "./Drivers/AntigravityDriver.ts"; +import { OhMyPiDriver, type OhMyPiDriverEnv } from "./Drivers/OhMyPiDriver.ts"; import type { AnyProviderDriver } from "./ProviderDriver.ts"; /** @@ -39,7 +40,8 @@ export type BuiltInDriversEnv = | CursorDriverEnv | GrokDriverEnv | OpenCodeDriverEnv - | AntigravityDriverEnv; + | AntigravityDriverEnv + | OhMyPiDriverEnv; /** * Ordered list of built-in drivers. Order matters only for tie-breaking in @@ -53,4 +55,5 @@ export const BUILT_IN_DRIVERS: ReadonlyArray { }).pipe(Effect.provide(makeServerSettingsLayer())), ); + for (const useInstances of [false, true]) { + it.effect(`excludes OhMyPi from text generation fallback (instances: ${useInstances})`, () => + Effect.gen(function* () { + const serverConfig = yield* ServerConfig.ServerConfig; + const fileSystem = yield* FileSystem.FileSystem; + const serverSettings = yield* ServerSettingsModule.ServerSettingsService; + yield* fileSystem.writeFileString( + serverConfig.settingsPath, + useInstances + ? '{"providerInstances":{"codex":{"driver":"codex","enabled":false,"config":{}},"claudeAgent":{"driver":"claudeAgent","enabled":false,"config":{}},"ohMyPi":{"driver":"ohMyPi","enabled":true,"config":{}},"opencode":{"driver":"opencode","enabled":true,"config":{}}}}' + : '{"providers":{"codex":{"enabled":false},"claudeAgent":{"enabled":false},"ohMyPi":{"enabled":true},"opencode":{"enabled":true}}}', + ); + + const settings = yield* serverSettings.getSettings; + assert.equal(settings.textGenerationModelSelection.instanceId, "opencode"); + }).pipe(Effect.provide(makeServerSettingsLayer())), + ); + } + it.effect("keeps unused providers disabled in existing sparse settings files", () => Effect.gen(function* () { const serverConfig = yield* ServerConfig.ServerConfig; diff --git a/apps/server/src/serverSettings.ts b/apps/server/src/serverSettings.ts index 0b64d445adf8..3bb37ca5a381 100644 --- a/apps/server/src/serverSettings.ts +++ b/apps/server/src/serverSettings.ts @@ -327,6 +327,8 @@ function fallbackTextGenerationProvider(settings: ServerSettings): ServerSetting // (codex enabled) when the Providers UI has only written providerInstances. const fallbackEntry = Object.entries(settings.providers).find(([driver, provider]) => { const instance = settings.providerInstances[ProviderInstanceId.make(driver)]; + // OhMyPi supports conversation turns only, not background text generation. + if ((instance?.driver ?? driver) === "ohMyPi") return false; return instance === undefined ? provider.enabled : resolveProviderInstanceEnabled(instance); }); const fallback = fallbackEntry ? ProviderDriverKind.make(fallbackEntry[0]) : undefined; diff --git a/apps/web/src/components/Icons.tsx b/apps/web/src/components/Icons.tsx index b13040152e18..e40ab0d643c9 100644 --- a/apps/web/src/components/Icons.tsx +++ b/apps/web/src/components/Icons.tsx @@ -740,3 +740,15 @@ export const PiAgentIcon: Icon = ({ className, ...props }) => ( ); + +export const OhMyPiIcon: Icon = (props) => ( + + + +); diff --git a/apps/web/src/components/chat/providerIconUtils.ts b/apps/web/src/components/chat/providerIconUtils.ts index db0e5ca222f3..3727d795933a 100644 --- a/apps/web/src/components/chat/providerIconUtils.ts +++ b/apps/web/src/components/chat/providerIconUtils.ts @@ -7,6 +7,7 @@ import { Icon, OpenAI, OpenCodeIcon, + OhMyPiIcon, } from "../Icons"; export const PROVIDER_ICON_BY_PROVIDER: Partial> = { @@ -16,6 +17,7 @@ export const PROVIDER_ICON_BY_PROVIDER: Partial [ProviderDriverKind.make("cursor")]: CursorIcon, [ProviderDriverKind.make("grok")]: GrokIcon, [ProviderDriverKind.make("antigravity")]: AntigravityIcon, + [ProviderDriverKind.make("ohMyPi")]: OhMyPiIcon, }; export type ModelEsque = { diff --git a/apps/web/src/components/settings/ProviderInstanceCard.tsx b/apps/web/src/components/settings/ProviderInstanceCard.tsx index 7b386e062c18..7024dce0128f 100644 --- a/apps/web/src/components/settings/ProviderInstanceCard.tsx +++ b/apps/web/src/components/settings/ProviderInstanceCard.tsx @@ -469,7 +469,9 @@ export function ProviderInstanceCard({ ? instance.driver : null; const customModels = - instance.driver === "antigravity" ? [] : readConfigCustomModels(instance.config); + instance.driver === "antigravity" || instance.driver === "ohMyPi" + ? [] + : readConfigCustomModels(instance.config); // Server-returned models may lag behind settings writes. Treat probe // models as the source for built-ins only; custom rows come directly // from the current instance config so add/remove reflects immediately. diff --git a/apps/web/src/components/settings/ProviderModelsSection.tsx b/apps/web/src/components/settings/ProviderModelsSection.tsx index 505a86ff0eac..d20fef92e7d0 100644 --- a/apps/web/src/components/settings/ProviderModelsSection.tsx +++ b/apps/web/src/components/settings/ProviderModelsSection.tsx @@ -171,6 +171,7 @@ export function ProviderModelsSection({ onFavoriteModelsChange, onModelOrderChange, }: ProviderModelsSectionProps) { + const supportsCustomModels = driverKind !== "antigravity" && driverKind !== "ohMyPi"; const [input, setInput] = useState(""); const [isAdding, setIsAdding] = useState(false); const [filter, setFilter] = useState(""); @@ -223,7 +224,7 @@ export function ProviderModelsSection({ }, [displayModels]); const handleAdd = () => { - if (driverKind === "antigravity") return; + if (!supportsCustomModels) return; const normalized = normalizeCustomModelSlug(input); if (!normalized) { setError("Enter a model slug."); @@ -586,7 +587,7 @@ export function ProviderModelsSection({ })} - {driverKind === "antigravity" ? null : isAdding ? ( + {!supportsCustomModels ? null : isAdding ? (
)} - {driverKind !== "antigravity" && error ? ( + {supportsCustomModels && error ? (

{error}

) : null}
diff --git a/apps/web/src/components/settings/providerDriverMeta.ts b/apps/web/src/components/settings/providerDriverMeta.ts index 4bf4da3919ba..537a585a4d6f 100644 --- a/apps/web/src/components/settings/providerDriverMeta.ts +++ b/apps/web/src/components/settings/providerDriverMeta.ts @@ -5,6 +5,7 @@ import { CursorSettings, GrokSettings, OpenCodeSettings, + OhMyPiSettings, ProviderDriverKind, } from "@t3tools/contracts"; import type * as Schema from "effect/Schema"; @@ -16,6 +17,7 @@ import { type Icon, OpenAI, OpenCodeIcon, + OhMyPiIcon, } from "../Icons"; type ProviderSettingsSchema = { @@ -44,6 +46,12 @@ export interface ProviderClientDefinition { } const PROVIDER_CLIENT_DEFINITIONS: readonly ProviderClientDefinition[] = [ + { + value: ProviderDriverKind.make("ohMyPi"), + label: "OhMyPi", + icon: OhMyPiIcon, + settingsSchema: OhMyPiSettings, + }, { value: ProviderDriverKind.make("codex"), label: "Codex", diff --git a/apps/web/src/components/settings/settingsSearch.ts b/apps/web/src/components/settings/settingsSearch.ts index f56579990714..14c518058b60 100644 --- a/apps/web/src/components/settings/settingsSearch.ts +++ b/apps/web/src/components/settings/settingsSearch.ts @@ -377,7 +377,7 @@ export const SETTINGS_SEARCH_ITEMS = [ title: "Providers", to: "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/settings/providers", searchTerms: [ - "agents cli codex claude cursor grok opencode antigravity google sign in sign out install subscription instances authentication api key models configuration binary path config directory endpoint arguments environment variables display name accent color custom favorite hidden auto compact", + "agents cli codex claude cursor grok opencode ohmypi omp antigravity google sign in sign out install subscription instances authentication api key models configuration binary path config directory endpoint arguments environment variables display name accent color custom favorite hidden auto compact", ], }, { diff --git a/apps/web/src/modelSelection.test.ts b/apps/web/src/modelSelection.test.ts index 959a5eb8c481..8fd29e3ed8ee 100644 --- a/apps/web/src/modelSelection.test.ts +++ b/apps/web/src/modelSelection.test.ts @@ -657,6 +657,35 @@ describe("instance-scoped model selection", () => { } }); + it("offers only account catalog models for OhMyPi despite custom model settings", () => { + const driver = ProviderDriverKind.make("ohMyPi"); + const customId = ProviderInstanceId.make("ohMyPi_work"); + const nativeModel = "openai/gpt"; + const settings: UnifiedSettings = { + ...DEFAULT_UNIFIED_SETTINGS, + providers: { + ...DEFAULT_UNIFIED_SETTINGS.providers, + ohMyPi: { + ...DEFAULT_UNIFIED_SETTINGS.providers.ohMyPi, + customModels: ["api-only-model"], + }, + }, + providerInstances: { + [customId]: { driver, config: { customModels: ["unknown-model"] } }, + }, + }; + const entries = deriveProviderInstanceEntries([ + provider({ provider: driver, instanceId: "ohMyPi", models: [nativeModel] }), + provider({ provider: driver, instanceId: customId, models: [nativeModel] }), + ]); + + for (const entry of entries) { + expect(getAppModelOptionsForInstance(settings, entry).map((model) => model.slug)).toEqual([ + nativeModel, + ]); + } + }); + it("resolves the Antigravity default marker without creating an unavailable model", () => { const instanceId = ProviderInstanceId.make("antigravity_work"); const nativeModel = "gemini-3.1-pro"; diff --git a/apps/web/src/modelSelection.ts b/apps/web/src/modelSelection.ts index 1a39f06098ac..f167db7acddd 100644 --- a/apps/web/src/modelSelection.ts +++ b/apps/web/src/modelSelection.ts @@ -58,7 +58,7 @@ function readInstanceCustomModels( instanceId: ProviderInstanceId, driverKind: ProviderDriverKind, ): ReadonlyArray { - if (driverKind === "antigravity") return []; + if (driverKind === "antigravity" || driverKind === "ohMyPi") return []; const instance = settings.providerInstances?.[instanceId]; const config = instance?.config; if (config !== null && typeof config === "object") { diff --git a/docs/user/composer.md b/docs/user/composer.md index 4a8df5333664..ce60df749e36 100644 --- a/docs/user/composer.md +++ b/docs/user/composer.md @@ -33,8 +33,8 @@ device until you sign back into the same account. ## Custom models On web and desktop, use Settings → Providers → **Models** to add an unlisted model with a custom -name and options. Only options supported by the provider integration affect turns. Antigravity -uses its account catalog and does not support custom models. +name and options. Only options supported by the provider integration affect turns. Antigravity and +OhMyPi use their discovered catalogs and do not support custom models. ## Model defaults diff --git a/docs/user/install.md b/docs/user/install.md index 17e9291bf1a0..810b90a4c183 100644 --- a/docs/user/install.md +++ b/docs/user/install.md @@ -67,14 +67,15 @@ and enable the provider you want. Installation, login, and configuration belong to that environment's machine, even when you connect from a phone or another computer. -| Provider | Install and authenticate | -| ----------- | -------------------------------------------------------------------------------------------- | -| Codex | Install [Codex CLI](https://developers.openai.com/codex/cli), then run `codex login`. | -| Claude | Install [Claude Code](https://claude.com/product/claude-code), then run `claude auth login`. | -| Cursor | Install [Cursor CLI](https://cursor.com/cli), then run `agent login`. | -| Grok Build | Install [Grok Build CLI](https://x.ai/cli), then run `grok login`. | -| OpenCode | Install [OpenCode](https://opencode.ai), then run `opencode auth login`. | -| Antigravity | Install and sign in with Google from T3 Code's provider settings. | +| Provider | Install and authenticate | +| ----------- | --------------------------------------------------------------------------------------------- | +| Codex | Install [Codex CLI](https://developers.openai.com/codex/cli), then run `codex login`. | +| Claude | Install [Claude Code](https://claude.com/product/claude-code), then run `claude auth login`. | +| Cursor | Install [Cursor CLI](https://cursor.com/cli), then run `agent login`. | +| Grok Build | Install [Grok Build CLI](https://x.ai/cli), then run `grok login`. | +| OpenCode | Install [OpenCode](https://opencode.ai), then run `opencode auth login`. | +| OhMyPi | Install [OhMyPi](https://omp.sh), then run `omp` to select a model and configure credentials. | +| Antigravity | Install and sign in with Google from T3 Code's provider settings. | Provider CLIs must be on the server's `PATH`. If T3 Code cannot find one, set its **Binary path** in provider settings, especially when using a version manager. @@ -97,6 +98,25 @@ For provider-specific setup and accounts, see [Codex](./providers-codex.md), [Claude](./providers-claude.md), [OpenCode](./providers-opencode.md), and [Antigravity](./providers-antigravity.md). +### OhMyPi + +Enable OhMyPi in **Settings → Providers** after configuring `omp` on the environment's +machine. Set **Binary path** if it is not on the server's `PATH`. T3 Code launches +`omp acp` using the CLI's existing credentials; remote and mobile clients use that +same environment. + +Available models and thinking levels load from OhMyPi before you start a thread. +Refresh provider status after changing your OhMyPi credentials or model configuration. +Choose **OhMyPi default** to use the model configured in OhMyPi, or select a model +from the picker. You can switch models, stop turns, and continue a saved session +after reconnecting. +OhMyPi's native slash commands appear once a session has started in that workspace. + +**Auto** uses the same approval policy as **Supervised**. **Auto-accept edits** +allows workspace writes, and **Full access** allows all tool tiers. OhMyPi's +structured question dialogs, conversation rollback, and background title, branch, +commit, and PR generation are not currently supported in T3 Code. + ## Next steps - [Working with threads](./thread-sidebar.md): start tasks and organize parallel work. diff --git a/packages/contracts/src/model.ts b/packages/contracts/src/model.ts index bce1a766bc9b..3f5f36820ff7 100644 --- a/packages/contracts/src/model.ts +++ b/packages/contracts/src/model.ts @@ -161,6 +161,8 @@ export const PREFERRED_DEFAULT_CODEX_MODELS: ReadonlyArray = [ "gpt-5.6-terra", ]; export const DEFAULT_TEXT_GENERATION_MODEL = "gpt-5.6-luna"; +/** Keep OhMyPi's configured model. Never send this alias to ACP. */ +export const OH_MY_PI_DEFAULT_MODEL = "oh-my-pi-default"; /** Keep the official Antigravity session's current model. Never send this ID to ACP. */ export const ANTIGRAVITY_DEFAULT_MODEL = "antigravity-default"; export const DEFAULT_TEXT_GENERATION_REASONING_EFFORT = "low"; @@ -171,6 +173,7 @@ export const DEFAULT_MODEL_BY_PROVIDER: Partial> = { + [ProviderDriverKind.make("ohMyPi")]: "OhMyPi", [ProviderDriverKind.make("antigravity")]: "Antigravity", [CODEX_DRIVER_KIND]: "Codex", [CLAUDE_DRIVER_KIND]: "Claude", diff --git a/packages/contracts/src/settings.ts b/packages/contracts/src/settings.ts index dd6136461fc1..5e3a96324f63 100644 --- a/packages/contracts/src/settings.ts +++ b/packages/contracts/src/settings.ts @@ -712,6 +712,31 @@ export const GrokSettings = makeProviderSettingsSchema( ); export type GrokSettings = typeof GrokSettings.Type; +export const OhMyPiSettings = makeProviderSettingsSchema( + { + // Opt in from Settings before starting the CLI. + enabled: Schema.Boolean.pipe( + Schema.withDecodingDefault(Effect.succeed(false)), + Schema.annotateKey({ providerSettingsForm: { hidden: true } }), + ), + binaryPath: makeBinaryPathSetting("omp").pipe( + Schema.annotateKey({ + title: "Binary path", + description: "Path to the OhMyPi CLI binary.", + providerSettingsForm: { placeholder: "omp", clearWhenEmpty: "omit" }, + }), + ), + customModels: Schema.Array(CustomModelSetting).pipe( + Schema.withDecodingDefault(Effect.succeed([])), + Schema.annotateKey({ providerSettingsForm: { hidden: true } }), + ), + }, + { + order: ["binaryPath"], + }, +); +export type OhMyPiSettings = typeof OhMyPiSettings.Type; + /** * Antigravity ACP auth methods. Personal and Enterprise open a Google sign-in * in the browser. The API key and Agent Platform methods take credentials from @@ -1064,6 +1089,7 @@ export const ServerSettings = Schema.Struct({ codex: CodexSettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), claudeAgent: ClaudeSettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), cursor: CursorSettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), + ohMyPi: OhMyPiSettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), grok: GrokSettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), opencode: OpenCodeSettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), antigravity: AntigravitySettings.pipe(Schema.withDecodingDefault(Effect.succeed({}))), @@ -1221,6 +1247,12 @@ const GrokSettingsPatch = Schema.Struct({ customModels: Schema.optionalKey(Schema.Array(CustomModelSetting)), }); +const OhMyPiSettingsPatch = Schema.Struct({ + enabled: Schema.optionalKey(Schema.Boolean), + binaryPath: Schema.optionalKey(TrimmedString), + customModels: Schema.optionalKey(Schema.Array(CustomModelSetting)), +}); + const AntigravitySettingsPatch = Schema.Struct({ enabled: Schema.optionalKey(Schema.Boolean), authMethod: Schema.optionalKey(AntigravityAuthMethod), @@ -1298,6 +1330,7 @@ export const ServerSettingsPatch = Schema.Struct({ codex: Schema.optionalKey(CodexSettingsPatch), claudeAgent: Schema.optionalKey(ClaudeSettingsPatch), cursor: Schema.optionalKey(CursorSettingsPatch), + ohMyPi: Schema.optionalKey(OhMyPiSettingsPatch), grok: Schema.optionalKey(GrokSettingsPatch), opencode: Schema.optionalKey(OpenCodeSettingsPatch), antigravity: Schema.optionalKey(AntigravitySettingsPatch),