From 28316f8016d317836f3453024795247ee9fa86fd Mon Sep 17 00:00:00 2001 From: Luiz Ferraz Date: Fri, 11 Sep 2026 19:33:19 +0000 Subject: [PATCH 1/6] feat(providers): add OhMyPi provider via ACP - Run the omp CLI over ACP with approvals, resume, and model selection - Add OhMyPi icons and settings entries across web, desktop, and mobile - Document omp installation in the user guide --- apps/mobile/src/components/ProviderIcon.tsx | 14 + .../src/provider/Drivers/OhMyPiDriver.test.ts | 152 +++ .../src/provider/Drivers/OhMyPiDriver.ts | 255 +++++ .../src/provider/Layers/OhMyPiAdapter.ts | 1010 +++++++++++++++++ .../provider/acp/OhMyPiAcpCliProbe.test.ts | 29 + .../src/provider/acp/OhMyPiAcpSupport.test.ts | 108 ++ .../src/provider/acp/OhMyPiAcpSupport.ts | 156 +++ apps/server/src/provider/builtInDrivers.ts | 5 +- apps/web/src/components/Icons.tsx | 12 + .../src/components/chat/providerIconUtils.ts | 2 + .../components/settings/providerDriverMeta.ts | 8 + .../src/components/settings/settingsSearch.ts | 2 +- docs/user/install.md | 34 +- packages/contracts/src/model.ts | 4 + packages/contracts/src/settings.ts | 33 + 15 files changed, 1814 insertions(+), 10 deletions(-) create mode 100644 apps/server/src/provider/Drivers/OhMyPiDriver.test.ts create mode 100644 apps/server/src/provider/Drivers/OhMyPiDriver.ts create mode 100644 apps/server/src/provider/Layers/OhMyPiAdapter.ts create mode 100644 apps/server/src/provider/acp/OhMyPiAcpCliProbe.test.ts create mode 100644 apps/server/src/provider/acp/OhMyPiAcpSupport.test.ts create mode 100644 apps/server/src/provider/acp/OhMyPiAcpSupport.ts 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..1221eaa7a0d5 --- /dev/null +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts @@ -0,0 +1,152 @@ +// @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 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( + "probes without a session, 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: 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 }, + }); + expect((yield* instance.snapshot.refresh).status).toBe("ready"); + const probeLog = yield* fs.readFileString(logPath); + expect(probeLog).toContain('"method":"initialize"'); + expect(probeLog).not.toContain('"method":"session/new"'); + 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"); + const turn = yield* instance.adapter + .sendTurn({ threadId, input: "hello", attachments: [] }) + .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); + 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..77594452fe8a --- /dev/null +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.ts @@ -0,0 +1,255 @@ +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 { makeOhMyPiAcpRuntime, ohMyPiModelsFromConfig } from "../acp/OhMyPiAcpSupport.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, + providerModelsFromSettings, + type ServerProviderDraft, +} from "../providerSnapshot.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 crypto = yield* Crypto.Crypto; + 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: providerModelsFromSettings( + [ + { + slug: OH_MY_PI_DEFAULT_MODEL, + name: "OhMyPi default", + isDefault: true, + isCustom: false, + capabilities, + }, + ], + config.customModels, + 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 onConfigOptionsUpdated = (options: Parameters[0]) => + SubscriptionRef.update(metadata, (draft) => { + const models = ohMyPiModelsFromConfig(options); + return models.length > 0 + ? { + ...draft, + models: providerModelsFromSettings(models, config.customModels, capabilities), + } + : draft; + }); + const checkProvider = Effect.gen(function* () { + if (!enabled) return yield* getSnapshot; + const result = yield* Effect.gen(function* () { + const runtime = yield* makeOhMyPiAcpRuntime({ + ohMyPiSettings: effectiveConfig, + childProcessSpawner: spawner, + environment: processEnv, + cwd: serverConfig.cwd, + clientInfo: { name: "t3-code", version: "0.0.0" }, + }); + return yield* runtime.initialize(); + }).pipe( + Effect.provideService(Crypto.Crypto, crypto), + Effect.scoped, + Effect.timeout("30 seconds"), + 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.agentInfo?.version ?? null, + 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 = cause._tag === "AcpSpawnError" && isCommandMissingCause(cause.cause); + return { + ...draft, + installed: !missing, + status: "error", + checkedAt, + message: missing + ? "OhMyPi CLI (omp) is not installed or not on PATH." + : "OhMyPi's ACP health check failed. Check the binary path and run omp acp on the server.", + }; + }); + 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 } : {}), + onSessionStarted: (started) => + onConfigOptionsUpdated(started.sessionSetupResult.configOptions ?? []), + onConfigOptionsUpdated, + 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/Layers/OhMyPiAdapter.ts b/apps/server/src/provider/Layers/OhMyPiAdapter.ts new file mode 100644 index 000000000000..0a35bd15a994 --- /dev/null +++ b/apps/server/src/provider/Layers/OhMyPiAdapter.ts @@ -0,0 +1,1010 @@ +/** + * 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 onSessionStarted?: ( + started: AcpSessionRuntime.AcpSessionRuntimeStartResult, + ) => Effect.Effect; + readonly onConfigOptionsUpdated?: ( + options: ReadonlyArray, + ) => Effect.Effect; + 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; + 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* options?.onSessionStarted?.(started) ?? Effect.void; + 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, + 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": + yield* options?.onConfigOptionsUpdated?.(event.configOptions) ?? Effect.void; + 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 "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.itemId ? { itemId: event.itemId } : {}), + 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 ctx = yield* requireSession(input.threadId); + // OMP cancels an in-flight prompt before queuing a steer. Keep both + // prompt RPCs under the same T3 turn until the last one settles. + 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.promptsInFlight += 1; + + return 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.", + }); + } + + // ACP has no system-message field; keep runtime context separate from the user's text. + const result = yield* ctx.acp + .prompt({ + prompt: [ + ...promptParts, + { + type: "text", + text: buildRuntimeInstructions({ harness: "OhMyPi", model: resolvedModel }), + }, + ], + }) + .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); + }), + ), + ); + }); + + const interruptTurn: OhMyPiAdapterShape["interruptTurn"] = (threadId) => + Effect.gen(function* () { + const ctx = yield* requireSession(threadId); + 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..04029dbabd2c --- /dev/null +++ b/apps/server/src/provider/acp/OhMyPiAcpSupport.test.ts @@ -0,0 +1,108 @@ +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, + ohMyPiModelsFromConfig, + 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("publishes upstream model IDs, a default alias, and native thinking choices", () => { + const models = ohMyPiModelsFromConfig(config); + expect(models.map((model) => model.slug)).toEqual(["anthropic/sonnet", "openai/gpt"]); + expect(models[0]).toMatchObject({ + isDefault: true, + aliases: [OH_MY_PI_DEFAULT_MODEL], + subProvider: "anthropic", + capabilities: { optionDescriptors: [{ id: "thinking", currentValue: "high" }] }, + }); + }); + + 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("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..9cbb1f3ac37c --- /dev/null +++ b/apps/server/src/provider/acp/OhMyPiAcpSupport.ts @@ -0,0 +1,156 @@ +import { + OH_MY_PI_DEFAULT_MODEL, + type OhMyPiSettings, + type ProviderApprovalDecision, + type ProviderOptionSelection, + type RuntimeMode, + type ServerProviderModel, +} from "@t3tools/contracts"; +import { createModelCapabilities } from "@t3tools/shared/model"; +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 function ohMyPiModelsFromConfig( + options: ReadonlyArray, +): ReadonlyArray { + const model = options.find((option) => option.category === "model" || option.id === "model"); + if (model?.type !== "select") return []; + const thinking = options.find((option) => option.category === "thought_level"); + const capabilities = createModelCapabilities({ + optionDescriptors: + thinking?.type === "select" + ? [ + { + id: thinking.id, + label: thinking.name, + type: "select", + currentValue: thinking.currentValue, + options: thinking.options + .flatMap((entry) => ("value" in entry ? [entry] : entry.options)) + .map((entry) => ({ id: entry.value, label: entry.name })), + }, + ] + : [], + }); + const seen = new Set(); + return model.options + .flatMap((entry) => ("value" in entry ? [entry] : entry.options)) + .flatMap((entry): ServerProviderModel[] => { + if (!entry.value.trim() || seen.has(entry.value)) return []; + seen.add(entry.value); + return [ + { + slug: entry.value, + name: entry.name || entry.value, + subProvider: entry.value.split("/")[0], + isCustom: false, + ...(entry.value === model.currentValue + ? { isDefault: true, aliases: [OH_MY_PI_DEFAULT_MODEL] } + : {}), + capabilities, + }, + ]; + }); +} + +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; + 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 ( ); + +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/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/docs/user/install.md b/docs/user/install.md index 17e9291bf1a0..f5642f341e46 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,23 @@ 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. + +Start a thread with **OhMyPi default** to use the model configured in OhMyPi. The +session supplies available models and thinking levels to the model 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), From 858f07144c3fab70f2a51ffa906192f5d34ba071 Mon Sep 17 00:00:00 2001 From: Luiz Ferraz Date: Fri, 11 Sep 2026 20:29:17 +0000 Subject: [PATCH 2/6] feat(providers): discover OhMyPi models without starting ACP - Probe status via `omp models --json` instead of an ACP initialize handshake - Load models and thinking options before a thread starts; discovery errors keep the previous model list - Report the CLI version from `omp --version` --- .../src/provider/Drivers/OhMyPiDriver.test.ts | 130 ++++++++++++++++-- .../src/provider/Drivers/OhMyPiDriver.ts | 30 ++-- .../src/provider/Drivers/OhMyPiModels.ts | 94 +++++++++++++ docs/user/install.md | 8 +- 4 files changed, 230 insertions(+), 32 deletions(-) create mode 100644 apps/server/src/provider/Drivers/OhMyPiModels.ts diff --git a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts index 1221eaa7a0d5..b90b3aa40ce2 100644 --- a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts @@ -59,7 +59,98 @@ it.layer(testLayer)("OhMyPi driver", (it) => { ); it.effect( - "probes without a session, streams a turn, handles approvals and resumes through ACP", + "refreshes the CLI catalog, preserves 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", + "custom/model", + ]); + 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, + "custom/model", + ]); + }).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; @@ -76,12 +167,27 @@ it.layer(testLayer)("OhMyPi driver", (it) => { T3_ACP_EMIT_TOOL_CALLS: "1", T3_ACP_ALLOW_ONCE_OPTION_ID: "omp-allow-42", }, - source: execScriptSource({ - scriptPath: NodeURL.fileURLToPath( - new URL("../../../scripts/acp-mock-agent.ts", import.meta.url), - ), - argvLogPath: argvPath, - }), + 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", reasoning: true }, + { provider: "openai", id: "gpt", name: "GPT", reasoning: false }, + ] })); + process.exit(0); + } + ` + + execScriptSource({ + scriptPath: NodeURL.fileURLToPath( + new URL("../../../scripts/acp-mock-agent.ts", import.meta.url), + ), + argvLogPath: argvPath, + }), }), ); const instance = yield* OhMyPiDriver.create({ @@ -91,10 +197,12 @@ it.layer(testLayer)("OhMyPi driver", (it) => { environment: [], config: { ...OhMyPiDriver.defaultConfig(), binaryPath }, }); - expect((yield* instance.snapshot.refresh).status).toBe("ready"); - const probeLog = yield* fs.readFileString(logPath); - expect(probeLog).toContain('"method":"initialize"'); - expect(probeLog).not.toContain('"method":"session/new"'); + 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)), diff --git a/apps/server/src/provider/Drivers/OhMyPiDriver.ts b/apps/server/src/provider/Drivers/OhMyPiDriver.ts index 77594452fe8a..06cf422a41fd 100644 --- a/apps/server/src/provider/Drivers/OhMyPiDriver.ts +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.ts @@ -20,7 +20,7 @@ import { ServerSettingsService } from "../../serverSettings.ts"; import { ProviderDriverError } from "../Errors.ts"; import { makeOhMyPiAdapter } from "../Layers/OhMyPiAdapter.ts"; import { ProviderEventLoggers } from "../Layers/ProviderEventLoggers.ts"; -import { makeOhMyPiAcpRuntime, ohMyPiModelsFromConfig } from "../acp/OhMyPiAcpSupport.ts"; +import { ohMyPiModelsFromConfig } from "../acp/OhMyPiAcpSupport.ts"; import { makeManagedServerProvider } from "../makeManagedServerProvider.ts"; import { defaultProviderContinuationIdentity, @@ -35,6 +35,7 @@ import { providerModelsFromSettings, type ServerProviderDraft, } from "../providerSnapshot.ts"; +import { probeOhMyPiModels } from "./OhMyPiModels.ts"; import { withInstanceIdentity } from "./instanceIdentity.ts"; const DRIVER = ProviderDriverKind.make("ohMyPi"); @@ -59,7 +60,6 @@ export const OhMyPiDriver: ProviderDriver = { create: ({ instanceId, displayName, accentColor, environment, enabled, config }) => Effect.gen(function* () { const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; - const crypto = yield* Crypto.Crypto; const serverConfig = yield* ServerConfig; const eventLoggers = yield* ProviderEventLoggers; const effectiveConfig = { ...config, enabled }; @@ -120,19 +120,8 @@ export const OhMyPiDriver: ProviderDriver = { }); const checkProvider = Effect.gen(function* () { if (!enabled) return yield* getSnapshot; - const result = yield* Effect.gen(function* () { - const runtime = yield* makeOhMyPiAcpRuntime({ - ohMyPiSettings: effectiveConfig, - childProcessSpawner: spawner, - environment: processEnv, - cwd: serverConfig.cwd, - clientInfo: { name: "t3-code", version: "0.0.0" }, - }); - return yield* runtime.initialize(); - }).pipe( - Effect.provideService(Crypto.Crypto, crypto), - Effect.scoped, - Effect.timeout("30 seconds"), + const result = yield* probeOhMyPiModels(effectiveConfig, processEnv, serverConfig.cwd).pipe( + Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner), Effect.result, ); const checkedAt = DateTime.formatIso(yield* DateTime.now); @@ -141,7 +130,12 @@ export const OhMyPiDriver: ProviderDriver = { return { ...draft, installed: true, - version: result.success.agentInfo?.version ?? null, + version: result.success.version, + models: providerModelsFromSettings( + [...initial.models.filter((model) => !model.isCustom), ...result.success.models], + config.customModels, + capabilities, + ), status: "ready", checkedAt, message: @@ -149,7 +143,7 @@ export const OhMyPiDriver: ProviderDriver = { }; } const cause = result.failure; - const missing = cause._tag === "AcpSpawnError" && isCommandMissingCause(cause.cause); + const missing = isCommandMissingCause(cause); return { ...draft, installed: !missing, @@ -157,7 +151,7 @@ export const OhMyPiDriver: ProviderDriver = { checkedAt, message: missing ? "OhMyPi CLI (omp) is not installed or not on PATH." - : "OhMyPi's ACP health check failed. Check the binary path and run omp acp on the server.", + : "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; 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/docs/user/install.md b/docs/user/install.md index f5642f341e46..810b90a4c183 100644 --- a/docs/user/install.md +++ b/docs/user/install.md @@ -105,9 +105,11 @@ machine. Set **Binary path** if it is not on the server's `PATH`. T3 Code launch `omp acp` using the CLI's existing credentials; remote and mobile clients use that same environment. -Start a thread with **OhMyPi default** to use the model configured in OhMyPi. The -session supplies available models and thinking levels to the model picker. You -can switch models, stop turns, and continue a saved session after reconnecting. +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** From 12a532c04015d29f827b259fde74bccc7fa8795c Mon Sep 17 00:00:00 2001 From: Luiz Ferraz Date: Sat, 12 Sep 2026 23:58:11 +0000 Subject: [PATCH 3/6] fix(providers): address OhMyPi review feedback --- .../src/provider/Drivers/OhMyPiDriver.test.ts | 84 +++++++++++++++++-- .../src/provider/Drivers/OhMyPiDriver.ts | 31 +++---- .../src/provider/Layers/OhMyPiAdapter.ts | 12 ++- apps/server/src/serverSettings.test.ts | 19 +++++ apps/server/src/serverSettings.ts | 2 + .../settings/ProviderInstanceCard.tsx | 4 +- .../settings/ProviderModelsSection.tsx | 7 +- apps/web/src/modelSelection.test.ts | 29 +++++++ apps/web/src/modelSelection.ts | 2 +- docs/user/composer.md | 4 +- 10 files changed, 159 insertions(+), 35 deletions(-) diff --git a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts index b90b3aa40ce2..90d4f308525a 100644 --- a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts @@ -59,7 +59,7 @@ it.layer(testLayer)("OhMyPi driver", (it) => { ); it.effect( - "refreshes the CLI catalog, preserves custom models, and retains models on discovery errors", + "refreshes the CLI catalog, excludes custom models, and retains models on discovery errors", () => Effect.gen(function* () { const fs = yield* FileSystem.FileSystem; @@ -101,7 +101,6 @@ it.layer(testLayer)("OhMyPi driver", (it) => { expect(first.models.map((model) => model.slug)).toEqual([ OH_MY_PI_DEFAULT_MODEL, "openai/gpt", - "custom/model", ]); expect(first.models[1]?.capabilities?.optionDescriptors).toMatchObject([ { @@ -119,13 +118,86 @@ it.layer(testLayer)("OhMyPi driver", (it) => { 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, - "custom/model", - ]); + expect(empty.models.map((model) => model.slug)).toEqual([OH_MY_PI_DEFAULT_MODEL]); }).pipe(Effect.scoped), ); + it.effect("cancels an active ACP prompt before steering within the same turn", () => + 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" }) + .pipe(Effect.forkChild); + const seen: ProviderRuntimeEvent[] = []; + // The mock emits a tool event only after receiving the first prompt, then + // waits for cancellation. No timer or premature approval releases it. + 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.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.skipIf(process.env.T3_OH_MY_PI_MODELS_PROBE !== "1")( "discovers models from the installed OhMyPi CLI", () => diff --git a/apps/server/src/provider/Drivers/OhMyPiDriver.ts b/apps/server/src/provider/Drivers/OhMyPiDriver.ts index 06cf422a41fd..28b4147e2af2 100644 --- a/apps/server/src/provider/Drivers/OhMyPiDriver.ts +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.ts @@ -32,7 +32,6 @@ import { makeManualOnlyProviderMaintenanceCapabilities } from "../providerMainte import { buildServerProvider, isCommandMissingCause, - providerModelsFromSettings, type ServerProviderDraft, } from "../providerSnapshot.ts"; import { probeOhMyPiModels } from "./OhMyPiModels.ts"; @@ -80,19 +79,15 @@ export const OhMyPiDriver: ProviderDriver = { presentation: { displayName: "OhMyPi", showInteractionModeToggle: true }, enabled, checkedAt: DateTime.formatIso(yield* DateTime.now), - models: providerModelsFromSettings( - [ - { - slug: OH_MY_PI_DEFAULT_MODEL, - name: "OhMyPi default", - isDefault: true, - isCustom: false, - capabilities, - }, - ], - config.customModels, - capabilities, - ), + models: [ + { + slug: OH_MY_PI_DEFAULT_MODEL, + name: "OhMyPi default", + isDefault: true, + isCustom: false, + capabilities, + }, + ], probe: { installed: false, version: null, @@ -114,7 +109,7 @@ export const OhMyPiDriver: ProviderDriver = { return models.length > 0 ? { ...draft, - models: providerModelsFromSettings(models, config.customModels, capabilities), + models, } : draft; }); @@ -131,11 +126,7 @@ export const OhMyPiDriver: ProviderDriver = { ...draft, installed: true, version: result.success.version, - models: providerModelsFromSettings( - [...initial.models.filter((model) => !model.isCustom), ...result.success.models], - config.customModels, - capabilities, - ), + models: [...initial.models, ...result.success.models], status: "ready", checkedAt, message: diff --git a/apps/server/src/provider/Layers/OhMyPiAdapter.ts b/apps/server/src/provider/Layers/OhMyPiAdapter.ts index 0a35bd15a994..2c75a07c1587 100644 --- a/apps/server/src/provider/Layers/OhMyPiAdapter.ts +++ b/apps/server/src/provider/Layers/OhMyPiAdapter.ts @@ -737,8 +737,7 @@ export function makeOhMyPiAdapter( const sendTurn: OhMyPiAdapterShape["sendTurn"] = (input) => Effect.gen(function* () { const ctx = yield* requireSession(input.threadId); - // OMP cancels an in-flight prompt before queuing a steer. Keep both - // prompt RPCs under the same T3 turn until the last one settles. + // 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 @@ -839,6 +838,15 @@ export function makeOhMyPiAdapter( }); } + 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 = yield* ctx.acp .prompt({ diff --git a/apps/server/src/serverSettings.test.ts b/apps/server/src/serverSettings.test.ts index 5d2571e72dc3..c2a88b97b0ae 100644 --- a/apps/server/src/serverSettings.test.ts +++ b/apps/server/src/serverSettings.test.ts @@ -674,6 +674,25 @@ it.layer(NodeServices.layer)("server settings", (it) => { }).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/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/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 From 95542bd221894b9962bb0cd05af0b1170e08b2a7 Mon Sep 17 00:00:00 2001 From: Luiz Ferraz Date: Sun, 13 Sep 2026 00:08:23 +0000 Subject: [PATCH 4/6] fix(providers): keep OhMyPi model catalog independent of sessions --- .../src/provider/Drivers/OhMyPiDriver.test.ts | 14 ++++- .../src/provider/Drivers/OhMyPiDriver.ts | 14 ----- .../src/provider/Layers/OhMyPiAdapter.ts | 8 --- .../src/provider/acp/OhMyPiAcpSupport.test.ts | 55 +++++++++++++++---- .../src/provider/acp/OhMyPiAcpSupport.ts | 48 +--------------- 5 files changed, 57 insertions(+), 82 deletions(-) diff --git a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts index 90d4f308525a..5cef8495b368 100644 --- a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts @@ -248,8 +248,8 @@ it.layer(testLayer)("OhMyPi driver", (it) => { 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", reasoning: true }, - { provider: "openai", id: "gpt", name: "GPT", reasoning: false }, + { provider: "anthropic", id: "sonnet", name: "Sonnet", thinking: ["low", "high"] }, + { provider: "openai", id: "gpt", name: "GPT", thinking: ["high", "xhigh"] }, ] })); process.exit(0); } @@ -287,8 +287,14 @@ it.layer(testLayer)("OhMyPi driver", (it) => { 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: [] }) + .sendTurn({ + threadId, + input: "hello", + attachments: [], + modelSelection: { instanceId, model: "composer-2" }, + }) .pipe(Effect.forkChild); const seen: ProviderRuntimeEvent[] = []; while (true) { @@ -304,6 +310,8 @@ it.layer(testLayer)("OhMyPi driver", (it) => { 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); diff --git a/apps/server/src/provider/Drivers/OhMyPiDriver.ts b/apps/server/src/provider/Drivers/OhMyPiDriver.ts index 28b4147e2af2..fc810e196a1a 100644 --- a/apps/server/src/provider/Drivers/OhMyPiDriver.ts +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.ts @@ -20,7 +20,6 @@ import { ServerSettingsService } from "../../serverSettings.ts"; import { ProviderDriverError } from "../Errors.ts"; import { makeOhMyPiAdapter } from "../Layers/OhMyPiAdapter.ts"; import { ProviderEventLoggers } from "../Layers/ProviderEventLoggers.ts"; -import { ohMyPiModelsFromConfig } from "../acp/OhMyPiAcpSupport.ts"; import { makeManagedServerProvider } from "../makeManagedServerProvider.ts"; import { defaultProviderContinuationIdentity, @@ -103,16 +102,6 @@ export const OhMyPiDriver: ProviderDriver = { } satisfies ServerProviderDraft; const metadata = yield* SubscriptionRef.make(initial); const getSnapshot = SubscriptionRef.get(metadata).pipe(Effect.map(stampIdentity)); - const onConfigOptionsUpdated = (options: Parameters[0]) => - SubscriptionRef.update(metadata, (draft) => { - const models = ohMyPiModelsFromConfig(options); - return models.length > 0 - ? { - ...draft, - models, - } - : draft; - }); const checkProvider = Effect.gen(function* () { if (!enabled) return yield* getSnapshot; const result = yield* probeOhMyPiModels(effectiveConfig, processEnv, serverConfig.cwd).pipe( @@ -180,9 +169,6 @@ export const OhMyPiDriver: ProviderDriver = { instanceId, environment: processEnv, ...(eventLoggers.native ? { nativeEventLogger: eventLoggers.native } : {}), - onSessionStarted: (started) => - onConfigOptionsUpdated(started.sessionSetupResult.configOptions ?? []), - onConfigOptionsUpdated, onAvailableCommands: (commands, cwd) => SubscriptionRef.update(metadata, (draft) => ({ ...draft, diff --git a/apps/server/src/provider/Layers/OhMyPiAdapter.ts b/apps/server/src/provider/Layers/OhMyPiAdapter.ts index 2c75a07c1587..30d6b889d5fe 100644 --- a/apps/server/src/provider/Layers/OhMyPiAdapter.ts +++ b/apps/server/src/provider/Layers/OhMyPiAdapter.ts @@ -88,12 +88,6 @@ export interface OhMyPiAdapterLiveOptions { * Defaults to the legacy built-in instance id (`ohMyPi`). */ readonly instanceId?: ProviderInstanceId; - readonly onSessionStarted?: ( - started: AcpSessionRuntime.AcpSessionRuntimeStartResult, - ) => Effect.Effect; - readonly onConfigOptionsUpdated?: ( - options: ReadonlyArray, - ) => Effect.Effect; readonly onAvailableCommands?: ( commands: ReadonlyArray, cwd: string, @@ -539,7 +533,6 @@ export function makeOhMyPiAdapter( ), ); - yield* options?.onSessionStarted?.(started) ?? Effect.void; yield* applyRequestedSessionConfiguration({ runtime: acp, runtimeMode: input.runtimeMode, @@ -588,7 +581,6 @@ export function makeOhMyPiAdapter( yield* Deferred.succeed(event.acknowledge, undefined); return; case "ConfigOptionsUpdated": - yield* options?.onConfigOptionsUpdated?.(event.configOptions) ?? Effect.void; return; case "AvailableCommandsUpdated": yield* ( diff --git a/apps/server/src/provider/acp/OhMyPiAcpSupport.test.ts b/apps/server/src/provider/acp/OhMyPiAcpSupport.test.ts index 04029dbabd2c..d9fad5973c7c 100644 --- a/apps/server/src/provider/acp/OhMyPiAcpSupport.test.ts +++ b/apps/server/src/provider/acp/OhMyPiAcpSupport.test.ts @@ -4,7 +4,6 @@ import * as Effect from "effect/Effect"; import type * as AcpSchema from "effect-acp/schema"; import { applyOhMyPiAcpModelSelection, - ohMyPiModelsFromConfig, ohMyPiApprovalMode, selectOhMyPiPermissionOption, } from "./OhMyPiAcpSupport.ts"; @@ -41,17 +40,6 @@ describe("OhMyPi ACP", () => { expect(ohMyPiApprovalMode("auto-accept-edits")).toBe("write"); expect(ohMyPiApprovalMode("full-access")).toBe("yolo"); }); - it("publishes upstream model IDs, a default alias, and native thinking choices", () => { - const models = ohMyPiModelsFromConfig(config); - expect(models.map((model) => model.slug)).toEqual(["anthropic/sonnet", "openai/gpt"]); - expect(models[0]).toMatchObject({ - isDefault: true, - aliases: [OH_MY_PI_DEFAULT_MODEL], - subProvider: "anthropic", - capabilities: { optionDescriptors: [{ id: "thinking", currentValue: "high" }] }, - }); - }); - it.effect("keeps the configured default and applies thinking after changing models", () => Effect.gen(function* () { const calls: unknown[] = []; @@ -90,6 +78,49 @@ describe("OhMyPi ACP", () => { }), ); + 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", diff --git a/apps/server/src/provider/acp/OhMyPiAcpSupport.ts b/apps/server/src/provider/acp/OhMyPiAcpSupport.ts index 9cbb1f3ac37c..f47f2421a436 100644 --- a/apps/server/src/provider/acp/OhMyPiAcpSupport.ts +++ b/apps/server/src/provider/acp/OhMyPiAcpSupport.ts @@ -4,9 +4,7 @@ import { type ProviderApprovalDecision, type ProviderOptionSelection, type RuntimeMode, - type ServerProviderModel, } from "@t3tools/contracts"; -import { createModelCapabilities } from "@t3tools/shared/model"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner"; @@ -71,49 +69,6 @@ export function selectOhMyPiPermissionOption( ); } -export function ohMyPiModelsFromConfig( - options: ReadonlyArray, -): ReadonlyArray { - const model = options.find((option) => option.category === "model" || option.id === "model"); - if (model?.type !== "select") return []; - const thinking = options.find((option) => option.category === "thought_level"); - const capabilities = createModelCapabilities({ - optionDescriptors: - thinking?.type === "select" - ? [ - { - id: thinking.id, - label: thinking.name, - type: "select", - currentValue: thinking.currentValue, - options: thinking.options - .flatMap((entry) => ("value" in entry ? [entry] : entry.options)) - .map((entry) => ({ id: entry.value, label: entry.name })), - }, - ] - : [], - }); - const seen = new Set(); - return model.options - .flatMap((entry) => ("value" in entry ? [entry] : entry.options)) - .flatMap((entry): ServerProviderModel[] => { - if (!entry.value.trim() || seen.has(entry.value)) return []; - seen.add(entry.value); - return [ - { - slug: entry.value, - name: entry.name || entry.value, - subProvider: entry.value.split("/")[0], - isCustom: false, - ...(entry.value === model.currentValue - ? { isDefault: true, aliases: [OH_MY_PI_DEFAULT_MODEL] } - : {}), - capabilities, - }, - ]; - }); -} - export const applyOhMyPiAcpModelSelection = Effect.fn("applyOhMyPiAcpModelSelection")(function* < E, >(input: { @@ -137,6 +92,9 @@ export const applyOhMyPiAcpModelSelection = Effect.fn("applyOhMyPiAcpModelSelect (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 }))); From a68d58984b60fb2f8df867e0953753d1546d3e34 Mon Sep 17 00:00:00 2001 From: Luiz Ferraz Date: Sun, 13 Sep 2026 00:21:05 +0000 Subject: [PATCH 5/6] fix(providers): serialize OhMyPi prompt admission through dispatch --- .../src/provider/Drivers/OhMyPiDriver.test.ts | 159 +++++---- .../src/provider/Layers/OhMyPiAdapter.ts | 336 ++++++++++-------- 2 files changed, 263 insertions(+), 232 deletions(-) diff --git a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts index 5cef8495b368..84599d7dac01 100644 --- a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts @@ -122,81 +122,92 @@ it.layer(testLayer)("OhMyPi driver", (it) => { }).pipe(Effect.scoped), ); - it.effect("cancels an active ACP prompt before steering within the same turn", () => - 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), - ), + 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" }) - .pipe(Effect.forkChild); - const seen: ProviderRuntimeEvent[] = []; - // The mock emits a tool event only after receiving the first prompt, then - // waits for cancellation. No timer or premature approval releases it. - 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.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), - ); + ); + 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.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.skipIf(process.env.T3_OH_MY_PI_MODELS_PROBE !== "1")( "discovers models from the installed OhMyPi CLI", diff --git a/apps/server/src/provider/Layers/OhMyPiAdapter.ts b/apps/server/src/provider/Layers/OhMyPiAdapter.ts index 30d6b889d5fe..220d24fd01e2 100644 --- a/apps/server/src/provider/Layers/OhMyPiAdapter.ts +++ b/apps/server/src/provider/Layers/OhMyPiAdapter.ts @@ -728,178 +728,198 @@ export function makeOhMyPiAdapter( const sendTurn: OhMyPiAdapterShape["sendTurn"] = (input) => Effect.gen(function* () { - const ctx = yield* requireSession(input.threadId); - // 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.promptsInFlight += 1; - - return 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, - }; + 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); + // 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(); - 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, + 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), }); - if (!attachmentPath) { - return yield* new ProviderAdapterRequestError({ + 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, - method: "session/prompt", - detail: `Invalid attachment id '${attachment.id}'.`, + threadId: input.threadId, + turnId, + payload: { model: resolvedModel }, }); } - const bytes = yield* fileSystem.readFile(attachmentPath).pipe( - Effect.mapError( - (cause) => - new ProviderAdapterRequestError({ + + 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: cause.message, - cause, - }), - ), - ); - promptParts.push({ - type: "image", - data: Buffer.from(bytes).toString("base64"), - mimeType: attachment.mimeType, - }); - } - } + 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 (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), - ), - ); - } + 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 = yield* ctx.acp - .prompt({ - prompt: [ - ...promptParts, - { - type: "text", - text: buildRuntimeInstructions({ harness: "OhMyPi", model: resolvedModel }), - }, - ], - }) - .pipe( - Effect.mapError((error) => - mapAcpToAdapterError(PROVIDER, input.threadId, "session/prompt", error), - ), - ); + // ACP has no system-message field; keep runtime context separate from the user's text. + const result = 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, - }; + 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, - }, - }); - } + // 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); - }), - ), + 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* () { From c878f97857980c8fb00a0f85e5b4862da6085b2d Mon Sep 17 00:00:00 2001 From: Luiz Ferraz Date: Sun, 13 Sep 2026 00:35:21 +0000 Subject: [PATCH 6/6] fix(providers): honor OhMyPi stops during preparation and stream thoughts --- .../src/provider/Drivers/OhMyPiDriver.test.ts | 90 +++++++++++++++++++ .../src/provider/Layers/OhMyPiAdapter.ts | 53 +++++++---- 2 files changed, 124 insertions(+), 19 deletions(-) diff --git a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts index 84599d7dac01..b118ecb44de0 100644 --- a/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts +++ b/apps/server/src/provider/Drivers/OhMyPiDriver.test.ts @@ -8,6 +8,7 @@ import { 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"; @@ -188,6 +189,14 @@ it.layer(testLayer)("OhMyPi driver", (it) => { 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" } }, @@ -209,6 +218,87 @@ it.layer(testLayer)("OhMyPi driver", (it) => { ); } + 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", () => diff --git a/apps/server/src/provider/Layers/OhMyPiAdapter.ts b/apps/server/src/provider/Layers/OhMyPiAdapter.ts index 220d24fd01e2..7065ed074772 100644 --- a/apps/server/src/provider/Layers/OhMyPiAdapter.ts +++ b/apps/server/src/provider/Layers/OhMyPiAdapter.ts @@ -113,6 +113,7 @@ interface OhMyPiSessionContext { * >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; } @@ -570,6 +571,7 @@ export function makeOhMyPiAdapter( lastPlanFingerprint: undefined, activeTurnId: undefined, promptsInFlight: 0, + interruptionVersion: 0, stopped: false, }; @@ -661,6 +663,7 @@ export function makeOhMyPiAdapter( }), ); return; + case "ThoughtDelta": case "ContentDelta": yield* logNative( ctx.threadId, @@ -674,7 +677,10 @@ export function makeOhMyPiAdapter( provider: PROVIDER, threadId: ctx.threadId, turnId: ctx.activeTurnId, - ...(event.itemId ? { itemId: event.itemId } : {}), + ...(event._tag === "ContentDelta" && event.itemId + ? { itemId: event.itemId } + : {}), + ...(event._tag === "ThoughtDelta" ? { streamKind: "reasoning_text" } : {}), text: event.text, rawPayload: event.rawPayload, }), @@ -735,6 +741,7 @@ export function makeOhMyPiAdapter( 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); @@ -851,24 +858,30 @@ export function makeOhMyPiAdapter( } // ACP has no system-message field; keep runtime context separate from the user's text. - const result = 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), - ), - ); + 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); @@ -924,6 +937,8 @@ export function makeOhMyPiAdapter( 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(