diff --git a/apps/mobile/src/features/threads/use-composer-command-menu.ts b/apps/mobile/src/features/threads/use-composer-command-menu.ts index b7a9741bf3f5..f2ff30fcd1f5 100644 --- a/apps/mobile/src/features/threads/use-composer-command-menu.ts +++ b/apps/mobile/src/features/threads/use-composer-command-menu.ts @@ -29,8 +29,10 @@ import { import { dedupeProviderSkillsByName, getProviderSkillsForSlashMenu, + getProviderSlashCommandsForSlashMenu, isProviderSkillUserInvocable, resolveProviderSkillsForCwd, + resolveProviderSlashCommandsForCwd, } from "@t3tools/client-runtime/providerSkills"; import { useCallback, useEffect, useMemo, useRef, useState } from "react"; @@ -329,6 +331,7 @@ export function useComposerCommandMenu({ if (trigger.kind === "slash-command") { const q = trigger.query.toLowerCase(); + const visibleSkills = getProviderSkillsForSlashMenu(skills, true); const commandItems = buildComposerSlashCommandItems({ query: q, atMessageStart: trigger.rangeStart === 0, @@ -336,10 +339,18 @@ export function useComposerCommandMenu({ hasCompactableConversation, offersUsageLimits, allowInteractionMode: onUpdateInteractionMode !== undefined, - selectedProviderStatus, + selectedProviderStatus: selectedProviderStatus + ? { + ...selectedProviderStatus, + slashCommands: getProviderSlashCommandsForSlashMenu( + resolveProviderSlashCommandsForCwd(selectedProviderStatus, projectCwd), + visibleSkills, + ), + } + : null, }); - const skillItems = getProviderSkillsForSlashMenu(skills, true) + const skillItems = visibleSkills .filter((skill) => matchesSlashSkillQuery(skill, q)) .map((skill) => ({ id: `skill:${skill.name}`, @@ -456,6 +467,7 @@ export function useComposerCommandMenu({ onUpdateInteractionMode, pathSearch.entries, pullRequestSearch.entries, + projectCwd, selectedProviderStatus, skills, trigger, diff --git a/apps/server/scripts/acp-mock-agent.ts b/apps/server/scripts/acp-mock-agent.ts index e4fd848ab6e1..9fbdeee03c86 100644 --- a/apps/server/scripts/acp-mock-agent.ts +++ b/apps/server/scripts/acp-mock-agent.ts @@ -784,7 +784,7 @@ const program = Effect.gen(function* () { sessionId: requestedSessionId, update: { sessionUpdate: "agent_message_chunk", - content: { type: "text", text: "hello from " }, + content: { type: "text", text: "hello from" }, }, }); @@ -826,13 +826,15 @@ const program = Effect.gen(function* () { }); } - writeJsonRpcNotification("session/update", { - sessionId: requestedSessionId, - update: { - sessionUpdate: "agent_message_chunk", - content: { type: "text", text: "mock" }, - }, - }); + for (const text of [" ", "mo", "ck"]) { + writeJsonRpcNotification("session/update", { + sessionId: requestedSessionId, + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text }, + }, + }); + } return yield* Effect.never; } diff --git a/apps/server/src/provider/Drivers/CursorDriver.ts b/apps/server/src/provider/Drivers/CursorDriver.ts index 5466af802e50..26d5dc4742f3 100644 --- a/apps/server/src/provider/Drivers/CursorDriver.ts +++ b/apps/server/src/provider/Drivers/CursorDriver.ts @@ -31,6 +31,7 @@ import { checkCursorProviderStatus, makeCursorModelDiscovery, enrichCursorSnapshot, + makeCursorCommandCatalog, } from "../Layers/CursorProvider.ts"; import { ProviderEventLoggers } from "../Layers/ProviderEventLoggers.ts"; import { makeManagedServerProvider } from "../makeManagedServerProvider.ts"; @@ -53,7 +54,7 @@ import { makeProviderSnapshotSettingsSource, type ProviderSnapshotSettings, } from "../providerUpdateSettings.ts"; -import { probeCursorSkills } from "./CursorSkills.ts"; +import { discoverCursorSkills, probeCursorSkills } from "./CursorSkills.ts"; const decodeCursorSettings = Schema.decodeSync(CursorSettings); const DRIVER_KIND = ProviderDriverKind.make("cursor"); @@ -130,11 +131,6 @@ export const CursorDriver: ProviderDriver = { ), ); - const adapter = yield* makeCursorAdapter(effectiveConfig, { - environment: processEnv, - ...(eventLoggers.native ? { nativeEventLogger: eventLoggers.native } : {}), - instanceId, - }); const textGeneration = yield* makeCursorTextGeneration(effectiveConfig, processEnv); const discoverModels = yield* makeCursorModelDiscovery(effectiveConfig, processEnv); @@ -151,7 +147,9 @@ export const CursorDriver: ProviderDriver = { ); const snapshotSettings = makeProviderSnapshotSettingsSource(effectiveConfig, serverSettings); - const snapshot = yield* makeManagedServerProvider>({ + const managedSnapshot = yield* makeManagedServerProvider< + ProviderSnapshotSettings + >({ resolveMaintenance, getSettings: snapshotSettings.getSettings, streamSettings: snapshotSettings.streamSettings, @@ -188,6 +186,20 @@ export const CursorDriver: ProviderDriver = { ), ); + const { snapshot, onAvailableCommands, snapshotForCwd } = + yield* makeCursorCommandCatalog(managedSnapshot); + const adapter = yield* makeCursorAdapter(effectiveConfig, { + environment: processEnv, + ...(eventLoggers.native ? { nativeEventLogger: eventLoggers.native } : {}), + instanceId, + onAvailableCommands: (commands, cwd) => + discoverCursorSkills(cwd, processEnv).pipe( + Effect.provideService(FileSystem.FileSystem, fileSystem), + Effect.provideService(Path.Path, path), + Effect.flatMap((skills) => onAvailableCommands(commands, cwd, skills)), + ), + }); + return { instanceId, driverKind: DRIVER_KIND, @@ -199,22 +211,20 @@ export const CursorDriver: ProviderDriver = { snapshotForCwd: (cwd) => !effectiveConfig.enabled ? snapshot.getSnapshot - : Effect.all([ - snapshot.getSnapshot, - probeCursorSkills(cwd, processEnv).pipe( - Effect.provideService(FileSystem.FileSystem, fileSystem), - Effect.provideService(Path.Path, path), - Effect.mapError( - (cause) => - new ProviderDriverError({ - driver: DRIVER_KIND, - instanceId, - detail: `Failed to discover Cursor skills for '${cwd}'`, - cause, - }), - ), + : probeCursorSkills(cwd, processEnv).pipe( + Effect.provideService(FileSystem.FileSystem, fileSystem), + Effect.provideService(Path.Path, path), + Effect.mapError( + (cause) => + new ProviderDriverError({ + driver: DRIVER_KIND, + instanceId, + detail: `Failed to discover Cursor skills for '${cwd}'`, + cause, + }), ), - ]).pipe(Effect.map(([machineSnapshot, skills]) => ({ ...machineSnapshot, skills }))), + Effect.flatMap((skills) => snapshotForCwd(cwd, skills)), + ), adapter, textGeneration, } satisfies ProviderInstance; diff --git a/apps/server/src/provider/Drivers/OpenCodeDriver.ts b/apps/server/src/provider/Drivers/OpenCodeDriver.ts index 72c1c0683de5..7cb956aefb44 100644 --- a/apps/server/src/provider/Drivers/OpenCodeDriver.ts +++ b/apps/server/src/provider/Drivers/OpenCodeDriver.ts @@ -31,10 +31,11 @@ import { checkOpenCodeProviderStatus, makePendingOpenCodeProvider, openCodeSkillsToServerProviderSkills, + openCodeCommandsToServerProviderSlashCommands, } from "../Layers/OpenCodeProvider.ts"; import { ProviderEventLoggers } from "../Layers/ProviderEventLoggers.ts"; import { makeManagedServerProvider } from "../makeManagedServerProvider.ts"; -import { OpenCodeRuntime } from "../opencodeRuntime.ts"; +import { OpenCodeRuntime, loadOpenCodeCommands } from "../opencodeRuntime.ts"; import * as OpenCodeServerOwner from "../OpenCodeServerOwner.ts"; import { defaultProviderContinuationIdentity, @@ -164,7 +165,18 @@ export const OpenCodeDriver: ProviderDriver // empty skill list and poisons the workspace snapshot the `$` picker // reads. The SDK `app.skills` endpoint honors the per-request directory // and returns complete results regardless of size. - const loadSkillsForCwd = (cwd: string) => + const loadWorkspaceInventory = (client: Parameters[0]) => + Effect.all( + { + skills: openCodeRuntime.loadOpenCodeSkills(client), + commands: loadOpenCodeCommands(client).pipe( + Effect.timeout("10 seconds"), + Effect.orElseSucceed(() => []), + ), + }, + { concurrency: "unbounded" }, + ); + const loadWorkspaceForCwd = (cwd: string) => effectiveConfig.serverUrl.trim().length > 0 ? Effect.scoped( Effect.gen(function* () { @@ -184,11 +196,11 @@ export const OpenCodeDriver: ProviderDriver ? { serverPassword: effectiveConfig.serverPassword } : {}), }); - return yield* openCodeRuntime.loadOpenCodeSkills(client); + return yield* loadWorkspaceInventory(client); }), ) : serverOwner.withServer((server) => - openCodeRuntime.loadOpenCodeSkills( + loadWorkspaceInventory( openCodeRuntime.createOpenCodeSdkClient({ baseUrl: server.url, directory: cwd, @@ -247,18 +259,19 @@ export const OpenCodeDriver: ProviderDriver ? snapshot.getSnapshot : Effect.all([ snapshot.getSnapshot, - loadSkillsForCwd(cwd).pipe(Effect.timeout("20 seconds")), + loadWorkspaceForCwd(cwd).pipe(Effect.timeout("20 seconds")), ]).pipe( - Effect.map(([machineSnapshot, skills]) => ({ + Effect.map(([machineSnapshot, { skills, commands }]) => ({ ...machineSnapshot, skills: openCodeSkillsToServerProviderSkills(skills), + slashCommands: openCodeCommandsToServerProviderSlashCommands(commands), })), Effect.mapError( (cause) => new ProviderDriverError({ driver: DRIVER_KIND, instanceId, - detail: `Failed to probe OpenCode skills for '${cwd}'`, + detail: `Failed to probe OpenCode commands and skills for '${cwd}'`, cause, }), ), diff --git a/apps/server/src/provider/Layers/CursorAdapter.test.ts b/apps/server/src/provider/Layers/CursorAdapter.test.ts index bdc818994a9a..4c3ad4b04551 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.test.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.test.ts @@ -302,7 +302,7 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { }), ); - it.effect("sends selected project skills in Cursor's native slash form", () => + it.effect("sends skills in Cursor's native form and preserves exact slash command input", () => Effect.gen(function* () { const adapter = yield* CursorAdapter; const settings = yield* ServerSettingsService; @@ -347,6 +347,7 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { ], ], ); + yield* adapter.sendTurn({ threadId, input: "/copy-request-id" }); yield* adapter.stopSession(threadId); const requests = yield* Effect.promise(() => readJsonLines(requestLogPath)); @@ -360,6 +361,7 @@ cursorAdapterTestLayer("CursorAdapterLive", (it) => { { type: "text", text: "please /review this" }, { type: "text", text: buildRuntimeInstructions({ harness: "Cursor" }) }, ], + [{ type: "text", text: "/copy-request-id" }], ], ); }), diff --git a/apps/server/src/provider/Layers/CursorAdapter.ts b/apps/server/src/provider/Layers/CursorAdapter.ts index 925d585e5838..38e6009737f4 100644 --- a/apps/server/src/provider/Layers/CursorAdapter.ts +++ b/apps/server/src/provider/Layers/CursorAdapter.ts @@ -118,6 +118,10 @@ export interface CursorAdapterLiveOptions { * the latest snapshot so the closure isn't stale. */ readonly resolveSettings?: Effect.Effect; + readonly onAvailableCommands?: ( + commands: ReadonlyArray, + cwd: string, + ) => Effect.Effect; } interface PendingApproval { @@ -810,6 +814,11 @@ export function makeCursorAdapter( return; case "ModeChanged": return; + case "AvailableCommandsUpdated": + yield* ( + options?.onAvailableCommands?.(event.availableCommands, cwd) ?? Effect.void + ); + return; case "AssistantItemStarted": ctx.assistantReply = new CursorTransportFailure(); yield* offerRuntimeEvent( @@ -1057,16 +1066,19 @@ export function makeCursorAdapter( }); } - // ACP has no system-message field; keep runtime context separate from the user's text. + // ACP commands parse the complete text. Extra context can turn an exact + // command into an ordinary model prompt or change its arguments. const result = yield* ctx.acp .prompt({ - prompt: [ - ...promptParts, - { - type: "text", - text: buildRuntimeInstructions({ harness: "Cursor", model: resolvedModel }), - }, - ], + prompt: /^\/[^\s/]+(?:\s|$)/.test(rawPrompt) + ? promptParts + : [ + ...promptParts, + { + type: "text", + text: buildRuntimeInstructions({ harness: "Cursor", model: resolvedModel }), + }, + ], }) .pipe( Effect.mapError((error) => diff --git a/apps/server/src/provider/Layers/CursorProvider.test.ts b/apps/server/src/provider/Layers/CursorProvider.test.ts index 5b71c64244b8..937831cbd9e7 100644 --- a/apps/server/src/provider/Layers/CursorProvider.test.ts +++ b/apps/server/src/provider/Layers/CursorProvider.test.ts @@ -1,14 +1,16 @@ import * as NodeOS from "node:os"; import * as NodeServices from "@effect/platform-node/NodeServices"; +import { it as effectIt } from "@effect/vitest"; import type * as Crypto from "effect/Crypto"; import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; import * as Path from "effect/Path"; +import * as Stream from "effect/Stream"; import type * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner"; import { describe, expect, it } from "vite-plus/test"; import type * as EffectAcpSchema from "effect-acp/schema"; -import type { CursorSettings } from "@t3tools/contracts"; +import { ProviderDriverKind, ProviderInstanceId, type CursorSettings } from "@t3tools/contracts"; import { createModelCapabilities } from "@t3tools/shared/model"; import { @@ -17,6 +19,7 @@ import { checkCursorProviderStatus, discoverCursorModelsViaAcp, makeCursorModelDiscovery, + makeCursorCommandCatalog, getCursorParameterizedModelPickerUnsupportedMessage, parseCursorAboutOutput, parseCursorCliConfigChannel, @@ -475,6 +478,87 @@ describe("Cursor skills", () => { }); }); +describe("Cursor command catalog", () => { + effectIt.effect( + "publishes workspace commands without leaking them globally and retains them across refreshes", + () => + Effect.gen(function* () { + const base = { + ...buildCursorProviderSnapshot({ + checkedAt: "2026-01-01T00:00:00.000Z", + cursorSettings: baseCursorSettings, + parsed: { version: null, status: "ready", auth: { status: "authenticated" } }, + }), + instanceId: ProviderInstanceId.make("cursor-catalog"), + driver: ProviderDriverKind.make("cursor"), + }; + const catalog = yield* makeCursorCommandCatalog({ + getSnapshot: Effect.succeed(base), + refresh: Effect.succeed(base), + streamChanges: Stream.empty, + resolveMaintenance: () => Effect.die("Not used"), + applyUsageLimits: () => Effect.void, + }); + const skills = [ + { name: "review", path: "/one/.cursor/skills/review/SKILL.md", enabled: true }, + ]; + const probedSkills = [ + { name: "explain", path: "/probed/.cursor/skills/explain/SKILL.md", enabled: true }, + ]; + yield* catalog.snapshotForCwd("/probed", probedSkills); + yield* catalog.onAvailableCommands( + [ + { name: "review", description: "Review changes", input: { hint: "target" } }, + { name: "compact", description: "Native duplicate" }, + { name: "review", description: "Duplicate" }, + ], + "/one", + skills, + ); + yield* catalog.onAvailableCommands([{ name: "deploy", description: "Deploy" }], "/two", []); + const reprobed = yield* catalog.snapshotForCwd("/one", skills); + expect(reprobed.slashCommands.map((command) => command.name)).toEqual([ + "compact", + "review", + ]); + const published = yield* catalog.snapshot.streamChanges.pipe( + Stream.take(1), + Stream.runCollect, + ); + expect(published[0]?.slashCommands.map((command) => command.name)).toEqual(["compact"]); + const refreshed = yield* catalog.snapshot.refresh; + expect( + refreshed.workspaceSnapshots?.find((entry) => entry.cwd === "/probed"), + ).toMatchObject({ + slashCommands: [{ name: "compact" }], + skills: probedSkills, + }); + expect(refreshed.workspaceSnapshots?.find((entry) => entry.cwd === "/one")).toMatchObject({ + slashCommands: [ + { name: "compact" }, + { name: "review", description: "Review changes", input: { hint: "target" } }, + ], + skills, + }); + yield* catalog.onAvailableCommands([], "/one", skills); + const updated = yield* catalog.snapshot.getSnapshot; + expect( + updated.workspaceSnapshots?.find((entry) => entry.cwd === "/probed")?.skills, + ).toEqual(probedSkills); + expect( + updated.workspaceSnapshots + ?.find((entry) => entry.cwd === "/one") + ?.slashCommands.map((command) => command.name), + ).toEqual(["compact"]); + expect( + updated.workspaceSnapshots + ?.find((entry) => entry.cwd === "/two") + ?.slashCommands.map((command) => command.name), + ).toEqual(["compact", "deploy"]); + }), + ); +}); + describe("buildCursorProviderSnapshot", () => { it("downgrades ready status to warning when ACP model discovery times out", () => { expect( diff --git a/apps/server/src/provider/Layers/CursorProvider.ts b/apps/server/src/provider/Layers/CursorProvider.ts index 7a18eef55e6b..6a8b93c88e90 100644 --- a/apps/server/src/provider/Layers/CursorProvider.ts +++ b/apps/server/src/provider/Layers/CursorProvider.ts @@ -22,6 +22,8 @@ import * as Option from "effect/Option"; import * as Path from "effect/Path"; import * as Result from "effect/Result"; import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; +import * as SubscriptionRef from "effect/SubscriptionRef"; import { HttpClient } from "effect/unstable/http"; import * as ChildProcess from "effect/unstable/process/ChildProcess"; import * as ChildProcessSpawner from "effect/unstable/process/ChildProcessSpawner"; @@ -49,6 +51,94 @@ import { } from "../providerMaintenance.ts"; import * as AcpSessionRuntime from "../acp/AcpSessionRuntime.ts"; import { CursorListAvailableModelsResponse } from "../acp/CursorAcpExtension.ts"; +import type { ServerProviderShape } from "../Services/ServerProvider.ts"; + +/** Session command catalogs stay scoped to their workspace across health refreshes. */ +export const makeCursorCommandCatalog = Effect.fn("makeCursorCommandCatalog")(function* ( + provider: ServerProviderShape, +) { + const workspaces = yield* SubscriptionRef.make>( + [], + ); + const getSnapshot = Effect.all([provider.getSnapshot, SubscriptionRef.get(workspaces)]).pipe( + Effect.map(([snapshot, workspaceSnapshots]) => + workspaceSnapshots.length > 0 ? { ...snapshot, workspaceSnapshots } : snapshot, + ), + ); + const snapshotForCwd = Effect.fn("CursorCommandCatalog.snapshotForCwd")(function* ( + cwd: string, + skills: ServerProvider["skills"], + ) { + const machineSnapshot = yield* provider.getSnapshot; + const checkedAt = DateTime.formatIso(yield* DateTime.now); + yield* SubscriptionRef.update(workspaces, (entries) => + [ + ...entries.filter((entry) => entry.cwd !== cwd), + { + cwd, + checkedAt, + slashCommands: + entries.find((entry) => entry.cwd === cwd)?.slashCommands ?? + machineSnapshot.slashCommands, + skills, + }, + ].slice(-16), + ); + const snapshot = yield* getSnapshot; + return { + ...snapshot, + checkedAt, + slashCommands: + snapshot.workspaceSnapshots?.find((entry) => entry.cwd === cwd)?.slashCommands ?? + snapshot.slashCommands, + skills, + }; + }); + const onAvailableCommands = Effect.fn("CursorCommandCatalog.onAvailableCommands")(function* ( + commands: ReadonlyArray, + cwd: string, + skills: ServerProvider["skills"], + ) { + const seen = new Set([COMPACT_SLASH_COMMAND.name]); + const slashCommands = [ + COMPACT_SLASH_COMMAND, + ...commands.flatMap((command) => { + const name = command.name.trim(); + if (!name || seen.has(name)) return []; + seen.add(name); + const description = command.description.trim(); + const hint = command.input?.hint.trim(); + return [ + { + name, + ...(description ? { description } : {}), + ...(hint ? { input: { hint } } : {}), + }, + ]; + }), + ]; + const checkedAt = DateTime.formatIso(yield* DateTime.now); + yield* SubscriptionRef.update(workspaces, (entries) => + [ + ...entries.filter((entry) => entry.cwd !== cwd), + { cwd, checkedAt, slashCommands, skills }, + ].slice(-16), + ); + }); + return { + onAvailableCommands, + snapshotForCwd, + snapshot: { + ...provider, + getSnapshot, + refresh: provider.refresh.pipe(Effect.andThen(getSnapshot)), + streamChanges: Stream.merge( + provider.streamChanges.pipe(Stream.map(() => undefined)), + SubscriptionRef.changes(workspaces).pipe(Stream.map(() => undefined)), + ).pipe(Stream.mapEffect(() => getSnapshot)), + } satisfies ServerProviderShape, + }; +}); const decodeCursorListAvailableModelsResponse = Schema.decodeUnknownEffect( CursorListAvailableModelsResponse, diff --git a/apps/server/src/provider/Layers/GrokAdapter.test.ts b/apps/server/src/provider/Layers/GrokAdapter.test.ts index 545a6a094e4e..ecd73af72dbe 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.test.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.test.ts @@ -292,7 +292,7 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { ); } - it.effect("sends runtime context with the current model without changing saved prompts", () => + it.effect("keeps runtime context out of native command arguments", () => Effect.gen(function* () { const threadId = ThreadId.make("grok-runtime-context"); const tempDir = yield* Effect.promise(() => @@ -337,6 +337,14 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { ], ], ); + const permissionError = yield* adapter + .sendTurn({ threadId, input: "/always-approve on" }) + .pipe(Effect.flip); + if (permissionError._tag !== "ProviderAdapterRequestError") { + assert.fail(`Unexpected error: ${permissionError._tag}`); + } + assert.include(permissionError.detail, "permission selector"); + yield* adapter.sendTurn({ threadId, input: "/goal status" }); yield* adapter.stopSession(threadId); const requests = yield* Effect.promise(() => readJsonLines(requestLogPath)); const prompts = requests @@ -344,7 +352,8 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { .map( (request) => (request.params as { prompt: Array<{ type: string; text: string }> }).prompt, ); - assert.equal(prompts.length, 2); + assert.equal(prompts.length, 3); + assert.deepEqual(prompts[2], [{ type: "text", text: "/goal status" }]); assert.deepEqual(prompts[0]?.[0], { type: "text", text: "First prompt" }); assert.include(prompts[0]?.[1]?.text, "Grok harness, as grok-mock-alt"); assert.deepEqual(prompts[1]?.[0], { type: "text", text: "Second prompt" }); @@ -571,13 +580,19 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { const runtimeEvents: ProviderRuntimeEvent[] = []; const turnCompleted = yield* Deferred.make(); + const secondTurnCompleted = yield* Deferred.make(); const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => Effect.sync(() => { runtimeEvents.push(event); }).pipe( - Effect.andThen( + Effect.andThen(() => event.type === "turn.completed" - ? Deferred.succeed(turnCompleted, undefined) + ? Deferred.succeed( + runtimeEvents.filter((entry) => entry.type === "turn.completed").length === 2 + ? secondTurnCompleted + : turnCompleted, + undefined, + ) : Effect.void, ), ), @@ -644,6 +659,44 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { assert.equal(readySession?.status, "ready"); assert.isUndefined(readySession?.activeTurnId); + const firstItemCompletion = runtimeEvents.findIndex( + (event) => + event.type === "item.completed" && event.payload.itemType === "assistant_message", + ); + assert.isAtLeast(firstItemCompletion, 0); + assert.isBelow(firstItemCompletion, terminalIndex); + + const secondTurn = yield* adapter.sendTurn({ threadId, input: "/goal status" }); + yield* Deferred.await(secondTurnCompleted); + const assistantItems = runtimeEvents.filter( + (event) => + (event.type === "item.started" || event.type === "item.completed") && + event.payload.itemType === "assistant_message", + ); + const startedItems = assistantItems.filter((event) => event.type === "item.started"); + const completedItems = assistantItems.filter((event) => event.type === "item.completed"); + assert.equal(completedItems.length, startedItems.length); + assert.equal(new Set(startedItems.map((event) => event.itemId)).size, startedItems.length); + for (const turnId of [sendTurnResult.turnId, secondTurn.turnId]) { + assert.lengthOf( + startedItems.filter((event) => event.turnId === turnId), + 1, + ); + } + for (const started of startedItems) { + const completions = completedItems.filter((event) => event.itemId === started.itemId); + assert.lengthOf(completions, 1); + const completed = completions[0]!; + assert.equal(completed.turnId, started.turnId); + assert.isAbove(runtimeEvents.indexOf(completed), runtimeEvents.indexOf(started)); + assert.isBelow( + runtimeEvents.indexOf(completed), + runtimeEvents.findIndex( + (event) => event.type === "turn.completed" && event.turnId === started.turnId, + ), + ); + } + yield* Fiber.interrupt(runtimeEventsFiber); yield* adapter.stopSession(threadId); }), @@ -1669,6 +1722,7 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { const runtimeEvents: ProviderRuntimeEvent[] = []; const activeTurnIdRef = yield* Ref.make(undefined); const trailingChunkTurnId = yield* Deferred.make(); + let receivedText = ""; const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) => Effect.gen(function* () { runtimeEvents.push(event); @@ -1678,7 +1732,11 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => { if (event.type === "turn.started") { yield* Ref.set(activeTurnIdRef, event.turnId); } - if (event.type !== "content.delta" || event.payload.delta !== "mock") { + if (event.type !== "content.delta") { + return; + } + receivedText += event.payload.delta; + if (receivedText !== "hello from mock") { return; } const turnId = event.turnId ?? (yield* Ref.get(activeTurnIdRef)); diff --git a/apps/server/src/provider/Layers/GrokAdapter.ts b/apps/server/src/provider/Layers/GrokAdapter.ts index 01ec3d118a3e..83b915285792 100644 --- a/apps/server/src/provider/Layers/GrokAdapter.ts +++ b/apps/server/src/provider/Layers/GrokAdapter.ts @@ -1505,6 +1505,13 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const sendTurn: GrokAdapterShape["sendTurn"] = (input) => Effect.gen(function* () { + if (/^\/always-approve(?:\s|$)/i.test(input.input?.trim() ?? "")) { + return yield* new ProviderAdapterRequestError({ + provider: PROVIDER, + method: "session/prompt", + detail: "Change permissions with T3's permission selector instead of /always-approve.", + }); + } const prepared = yield* withThreadLock( input.threadId, Effect.gen(function* () { @@ -1616,11 +1623,15 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte const displayModel = currentModelId ? resolveGrokAcpBaseModelId(currentModelId) : undefined; - const runtimeInstructions = buildRuntimeInstructions({ - harness: "Grok", - model: displayModel, - reasoningEffort: normalizeGrokReasoningEffort(requestedTurnReasoningEffort), - }); + // ACP slash commands must receive only their own arguments. + const runtimeInstructions = + text && /^\/[^\s/]+(?:\s|$)/.test(text) + ? undefined + : buildRuntimeInstructions({ + harness: "Grok", + model: displayModel, + reasoningEffort: normalizeGrokReasoningEffort(requestedTurnReasoningEffort), + }); for (let yieldAttempt = 0; yieldAttempt < 8; yieldAttempt += 1) { yield* Effect.yieldNow; } @@ -1737,7 +1748,9 @@ export function makeGrokAdapter(grokSettings: GrokSettings, options?: GrokAdapte { prompt: [ ...prepared.promptParts, - { type: "text", text: prepared.runtimeInstructions }, + ...(prepared.runtimeInstructions + ? [{ type: "text" as const, text: prepared.runtimeInstructions }] + : []), ], }, { dispatched }, diff --git a/apps/server/src/provider/Layers/GrokProvider.test.ts b/apps/server/src/provider/Layers/GrokProvider.test.ts index d0e5010bc2bf..127295af5e41 100644 --- a/apps/server/src/provider/Layers/GrokProvider.test.ts +++ b/apps/server/src/provider/Layers/GrokProvider.test.ts @@ -14,6 +14,7 @@ import { buildGrokModelsFromSessionModelState, buildInitialGrokProviderSnapshot, checkGrokProviderStatus, + grokSlashCommandsFromInitialize, parseGrokModelsCliOutput, } from "./GrokProvider.ts"; import { execScriptSource, writeFakeCli } from "../../testUtils/fakeCli.ts"; @@ -56,6 +57,69 @@ describe("parseGrokModelsCliOutput", () => { }); }); +describe("grokSlashCommandsFromInitialize", () => { + it("publishes native ACP commands and input hints without permission overrides", () => { + const commands = grokSlashCommandsFromInitialize({ + protocolVersion: 1, + _meta: { + availableCommands: [ + { name: "compact", description: "Compress history", input: { hint: "what to preserve" } }, + { + name: "always-approve", + description: "Skip permission prompts", + input: { hint: "on|off" }, + }, + { name: "context", description: "Show context usage", input: null }, + { name: "session-info", description: "Show session details" }, + { name: "deep-research", description: "Research a topic", input: { hint: "" } }, + { name: "workflow", description: "Manage workflows", input: { hint: "runs" } }, + { name: "goal", description: "Manage an autonomous goal", input: { hint: "status" } }, + ], + }, + }); + expect(commands.map((command) => command.name)).toEqual([ + "compact", + "session-info", + "deep-research", + "workflow", + "goal", + ]); + expect(commands[0]?.input).toEqual({ hint: "what to preserve" }); + expect(commands[2]).toEqual({ + name: "deep-research", + description: "Research a topic", + input: { hint: "" }, + }); + }); + + it("keeps valid commands when other metadata entries are malformed", () => { + const commands = grokSlashCommandsFromInitialize({ + protocolVersion: 1, + _meta: { + availableCommands: [ + null, + { name: "broken", description: 42 }, + { name: " ", description: "Empty name" }, + { name: " session-info ", description: " Session details " }, + { name: "session-info", description: "Updated session details" }, + ], + }, + }); + expect(commands.map((command) => command.name)).toEqual(["compact", "session-info"]); + expect(commands[1]?.description).toBe("Updated session details"); + }); + + it("keeps compact available for older agents without command metadata", () => { + for (const _meta of [undefined, {}, { availableCommands: "invalid" }]) { + expect( + grokSlashCommandsFromInitialize({ protocolVersion: 1, ...(_meta ? { _meta } : {}) }).map( + (command) => command.name, + ), + ).toEqual(["compact"]); + } + }); +}); + describe("buildGrokModelsFromSessionModelState", () => { it("marks the agent's current model as default and keeps reasoning options", () => { const models = buildGrokModelsFromSessionModelState({ @@ -438,6 +502,7 @@ it.layer(NodeServices.layer)("checkGrokProviderStatus", (it) => { ["grok-4.5", false], ]); expect(snapshot.message).toContain("ACP initialize failed"); + expect(snapshot.slashCommands.map((command) => command.name)).toEqual(["compact"]); }), ); diff --git a/apps/server/src/provider/Layers/GrokProvider.ts b/apps/server/src/provider/Layers/GrokProvider.ts index 18a77334b9ac..61cd8a9a3822 100644 --- a/apps/server/src/provider/Layers/GrokProvider.ts +++ b/apps/server/src/provider/Layers/GrokProvider.ts @@ -5,8 +5,9 @@ import { type ServerProvider, type ServerProviderAuth, type ServerProviderModel, + type ServerProviderSlashCommand, } from "@t3tools/contracts"; -import type * as EffectAcpSchema from "effect-acp/schema"; +import * as EffectAcpSchema from "effect-acp/schema"; import { causeErrorTag } from "@t3tools/shared/observability"; import * as Crypto from "effect/Crypto"; import * as DateTime from "effect/DateTime"; @@ -14,6 +15,7 @@ import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Option from "effect/Option"; import * as Result from "effect/Result"; +import * as Schema from "effect/Schema"; import { HttpClient } from "effect/unstable/http"; import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; import { createModelCapabilities } from "@t3tools/shared/model"; @@ -307,11 +309,41 @@ const runGrokCliCommand = ( ); }); +const decodeAvailableCommands = Schema.decodeUnknownOption(Schema.Array(Schema.Unknown)); +const decodeAvailableCommand = Schema.decodeUnknownOption(EffectAcpSchema.AvailableCommand); + +export function grokSlashCommandsFromInitialize( + initialized: EffectAcpSchema.InitializeResponse, +): ReadonlyArray { + const commands = decodeAvailableCommands(initialized._meta?.availableCommands); + const byName = new Map([ + [COMPACT_SLASH_COMMAND.name, COMPACT_SLASH_COMMAND], + ]); + for (const entry of Option.getOrElse(commands, () => [])) { + const decoded = decodeAvailableCommand(entry); + if (Option.isNone(decoded)) continue; + const command = decoded.value; + const name = command.name.trim(); + // Permission changes must go through T3 so the client and provider agree. + if (!name || name.toLowerCase() === "always-approve") continue; + // Grok advertises /context, but its ACP handler completes without emitting output. + if (name.toLowerCase() === "context") continue; + const description = command.description.trim(); + const hint = command.input?.hint.trim(); + byName.set(name, { + name, + ...(description ? { description } : {}), + ...(hint ? { input: { hint } } : {}), + }); + } + return [...byName.values()]; +} + /** - * Reads model metadata from `initialize._meta.modelState`. This never calls `authenticate` + * Reads model and command metadata from `initialize._meta`. This never calls `authenticate` * or `session/new`, so it cannot open a browser login or boot the workspace's MCP servers. */ -const discoverGrokModelsViaAcpInitialize = ( +const discoverGrokMetadataViaAcpInitialize = ( grokSettings: GrokSettings, environment: NodeJS.ProcessEnv, ) => @@ -325,7 +357,10 @@ const discoverGrokModelsViaAcpInitialize = ( clientInfo: { name: "t3-code-provider-probe", version: "0.0.0" }, }); const initialized = yield* acp.initialize(); - return buildGrokModelsFromSessionModelState(sessionModelStateFromInitialize(initialized)); + return { + models: buildGrokModelsFromSessionModelState(sessionModelStateFromInitialize(initialized)), + slashCommands: grokSlashCommandsFromInitialize(initialized), + }; }).pipe(Effect.scoped); export const checkGrokProviderStatus = Effect.fn("checkGrokProviderStatus")(function* ( @@ -461,11 +496,12 @@ export const checkGrokProviderStatus = Effect.fn("checkGrokProviderStatus")(func Effect.orElseSucceed(() => []), ); - const acpExit = yield* discoverGrokModelsViaAcpInitialize(grokSettings, environment).pipe( + const acpExit = yield* discoverGrokMetadataViaAcpInitialize(grokSettings, environment).pipe( Effect.timeoutOption(GROK_ACP_INITIALIZE_TIMEOUT_MS), Effect.exit, ); - const acpModels = Exit.isSuccess(acpExit) ? Option.getOrElse(acpExit.value, () => []) : []; + const acpMetadata = Exit.isSuccess(acpExit) ? Option.getOrUndefined(acpExit.value) : undefined; + const acpModels = acpMetadata?.models ?? []; const acpFailed = Exit.isFailure(acpExit) || Option.isNone(acpExit.value); if (acpFailed) { yield* Effect.logWarning("Grok ACP initialize probe failed or timed out.", { @@ -502,7 +538,7 @@ export const checkGrokProviderStatus = Effect.fn("checkGrokProviderStatus")(func checkedAt, models, skills, - slashCommands: [COMPACT_SLASH_COMMAND], + slashCommands: acpMetadata?.slashCommands ?? [COMPACT_SLASH_COMMAND], probe: { installed: true, version, diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index a6636fcd69d0..96d6b10d3839 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -88,6 +88,10 @@ const runtimeMock = { messageCalls: [] as Array<{ sessionID: string; messageID: string }>, messageFailures: 0, promptCalls: [] as Array, + commandCalls: [] as Array>, + commandImplementation: null as + | ((input: Record, signal?: AbortSignal) => Promise) + | null, summarizeCalls: [] as Array, promptAsyncError: null as Error | null, promptAsyncImplementation: null as (() => Promise) | null, @@ -150,6 +154,8 @@ const runtimeMock = { this.state.messageCalls.length = 0; this.state.messageFailures = 0; this.state.promptCalls.length = 0; + this.state.commandCalls.length = 0; + this.state.commandImplementation = null; this.state.summarizeCalls.length = 0; this.state.promptAsyncError = null; this.state.promptAsyncImplementation = null; @@ -236,6 +242,11 @@ const OpenCodeRuntimeTestDouble: OpenCodeRuntimeShape = { runOpenCodeCommand: () => Effect.succeed({ stdout: "", stderr: "", code: 0 }), createOpenCodeSdkClient: ({ baseUrl, serverPassword }) => ({ + command: { + list: async () => ({ + data: [{ name: "review", source: "command", hints: ["$ARGUMENTS"] }], + }), + }, session: { create: async (input: Record) => { runtimeMock.state.sessionCreateUrls.push(baseUrl); @@ -354,6 +365,10 @@ const OpenCodeRuntimeTestDouble: OpenCodeRuntimeShape = { : { "http://127.0.0.1:9999/session": { type: "busy" as const } }, }; }, + command: async (input: Record, options?: { signal?: AbortSignal }) => { + runtimeMock.state.commandCalls.push(input); + await runtimeMock.state.commandImplementation?.(input, options?.signal); + }, promptAsync: async (input: unknown) => { runtimeMock.state.promptCalls.push(input); await runtimeMock.state.promptAsyncImplementation?.(); @@ -1443,6 +1458,263 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it.effect("admits native commands before generation completes and keeps them interruptible", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-native-command"); + const publish = makeOpenCodeEventQueue(); + const completion = promiseWithResolvers(); + let responseSettled = false; + runtimeMock.state.commandImplementation = async (input) => { + publish({ + type: "message.updated", + properties: { sessionID: input.sessionID, info: { id: input.messageID, role: "user" } }, + }); + await completion.promise; + responseSettled = true; + }; + runtimeMock.state.abortImplementation = async () => { + completion.resolve(undefined); + }; + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const result = yield* adapter.sendTurn({ + threadId, + input: "/review main\nfocus on authentication", + modelSelection: createModelSelection(ProviderInstanceId.make("opencode"), "openai/gpt-5", [ + { id: "agent", value: "build" }, + { id: "variant", value: "high" }, + ]), + }); + NodeAssert.equal(responseSettled, false); + const { messageID, ...command } = runtimeMock.state.commandCalls[0]!; + NodeAssert.equal(typeof messageID, "string"); + NodeAssert.deepEqual(command, { + sessionID: "http://127.0.0.1:9999/session", + command: "review", + arguments: "main\nfocus on authentication", + model: "openai/gpt-5", + agent: "build", + variant: "high", + parts: [], + }); + NodeAssert.equal(runtimeMock.state.promptCalls.length, 0); + yield* advanceTestClock(11_000); + NodeAssert.equal(runtimeMock.state.abortCalls.length, 0); + yield* adapter.interruptTurn(threadId, result.turnId); + NodeAssert.equal(runtimeMock.state.abortCalls.length, 1); + NodeAssert.equal((yield* adapter.listSessions())[0]?.activeTurnId, undefined); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("recovers a native command receipt when the user-message event is lost", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-native-command-recovered"); + const started = promiseWithResolvers(); + const completion = promiseWithResolvers(); + runtimeMock.state.sessionStatus = "busy"; + runtimeMock.state.commandImplementation = async (input) => { + NodeAssert.ok(typeof input.messageID === "string"); + runtimeMock.state.messages.push({ info: { id: input.messageID, role: "user" }, parts: [] }); + started.resolve(undefined); + await completion.promise; + }; + runtimeMock.state.abortImplementation = async () => { + completion.resolve(undefined); + }; + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const sendFiber = yield* adapter + .sendTurn({ + threadId, + input: "/review", + modelSelection: createModelSelection(ProviderInstanceId.make("opencode"), "openai/gpt-5"), + }) + .pipe(Effect.forkChild); + yield* Effect.promise(() => started.promise); + yield* advanceTestClock(250); + const result = yield* Fiber.join(sendFiber); + NodeAssert.ok( + runtimeMock.state.messageCalls.some( + ({ messageID }) => messageID === runtimeMock.state.commandCalls[0]?.messageID, + ), + ); + yield* advanceTestClock(11_000); + NodeAssert.equal(runtimeMock.state.abortCalls.length, 0); + NodeAssert.equal( + (yield* adapter.listSessions()).find((session) => session.threadId === threadId) + ?.activeTurnId, + result.turnId, + ); + yield* adapter.interruptTurn(threadId, result.turnId); + NodeAssert.equal(runtimeMock.state.abortCalls.length, 1); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("bounds native admission recovery when timeout cleanup cannot abort the session", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-native-command-timeout-abort-failure"); + const started = promiseWithResolvers(); + runtimeMock.state.commandImplementation = async () => { + started.resolve(undefined); + await new Promise(() => {}); + }; + runtimeMock.state.abortImplementation = async () => { + throw new Error("abort failed"); + }; + runtimeMock.state.sessionStatus = "idle"; + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const exitedFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "session.exited"), + Stream.runHead, + Effect.forkChild, + ); + const sendFiber = yield* adapter + .sendTurn({ + threadId, + input: "/review", + modelSelection: createModelSelection(ProviderInstanceId.make("opencode"), "openai/gpt-5"), + }) + .pipe(Effect.exit, Effect.forkChild); + yield* Effect.promise(() => started.promise); + yield* advanceTestClock(10_000); + NodeAssert.equal(Exit.isFailure(yield* Fiber.join(sendFiber)), true); + yield* advanceTestClock(6_000); + const exited = Option.getOrThrow(yield* Fiber.join(exitedFiber)); + NodeAssert.ok(exited.type === "session.exited"); + NodeAssert.equal(exited.payload.exitKind, "error"); + NodeAssert.ok(runtimeMock.state.sessionStatusCalls > 0); + NodeAssert.equal( + (yield* adapter.listSessions()).some((session) => session.threadId === threadId), + false, + ); + const messageCalls = runtimeMock.state.messageCalls.length; + yield* advanceTestClock(5_000); + NodeAssert.equal(runtimeMock.state.messageCalls.length, messageCalls); + }), + ); + + for (const nativeStartsTurn of [true, false]) { + it.effect( + `reports a late native command failure after another steer (starts turn: ${nativeStartsTurn})`, + () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId(`thread-command-late-error-${nativeStartsTurn}`); + const publish = makeOpenCodeEventQueue(); + const completion = promiseWithResolvers(); + const modelSelection = createModelSelection( + ProviderInstanceId.make("opencode"), + "openai/gpt-5", + ); + runtimeMock.state.commandImplementation = async (input) => { + publish({ + type: "message.updated", + properties: { + sessionID: input.sessionID, + info: { id: input.messageID, role: "user" }, + }, + }); + await completion.promise; + }; + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + if (!nativeStartsTurn) + yield* adapter.sendTurn({ threadId, input: "Start work", modelSelection }); + const command = yield* adapter.sendTurn({ threadId, input: "/review", modelSelection }); + yield* adapter.sendTurn({ threadId, input: "Focus on authentication", modelSelection }); + const warningFiber = yield* adapter.streamEvents.pipe( + Stream.filter( + (event) => event.threadId === threadId && event.type === "runtime.warning", + ), + Stream.runHead, + Effect.forkChild, + ); + completion.reject(new Error("command failed after admission")); + const warning = yield* Fiber.join(warningFiber); + NodeAssert.equal(warning._tag, "Some"); + if (warning._tag === "Some" && warning.value.type === "runtime.warning") { + NodeAssert.equal(warning.value.payload.detail, "command failed after admission"); + } + const session = (yield* adapter.listSessions()).find( + (entry) => entry.threadId === threadId, + ); + NodeAssert.equal(session?.activeTurnId, command.turnId); + yield* adapter.stopSession(threadId); + }), + ); + } + + it.effect("surfaces native command rejection and leaves the session ready", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-native-command-error"); + runtimeMock.state.commandImplementation = async () => { + throw new Error("command unavailable"); + }; + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + const error = yield* adapter + .sendTurn({ + threadId, + input: "/review", + modelSelection: createModelSelection(ProviderInstanceId.make("opencode"), "openai/gpt-5"), + }) + .pipe(Effect.flip); + NodeAssert.equal(error._tag, "ProviderAdapterRequestError"); + if (error._tag !== "ProviderAdapterRequestError") throw new Error("Unexpected error type"); + NodeAssert.equal(error.method, "session.command"); + NodeAssert.equal(error.detail, "command unavailable"); + NodeAssert.equal((yield* adapter.listSessions())[0]?.status, "ready"); + yield* adapter.stopSession(threadId); + }), + ); + + it.effect("keeps unknown slash text on the ordinary prompt path", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-unknown-command"); + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId, + input: "/unknown explain this", + modelSelection: createModelSelection(ProviderInstanceId.make("opencode"), "openai/gpt-5"), + }); + NodeAssert.equal(runtimeMock.state.commandCalls.length, 0); + const prompt = runtimeMock.state.promptCalls[0] as { parts: unknown; system: string }; + NodeAssert.deepEqual(prompt.parts, [{ type: "text", text: "/unknown explain this" }]); + NodeAssert.equal( + prompt.system, + buildRuntimeInstructions({ harness: "OpenCode", model: "openai/gpt-5" }), + ); + yield* adapter.stopSession(threadId); + }), + ); + it.effect("rolls back session state when sendTurn fails before OpenCode accepts the prompt", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 3d216bb1167b..41bf634c0d3b 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -51,6 +51,7 @@ import { OpenCodeRuntimeError, openCodeQuestionId, openCodeRuntimeErrorDetail, + loadOpenCodeCommands, parseOpenCodeModelSlug, runOpenCodeSdk, toOpenCodeFileParts, @@ -215,6 +216,7 @@ interface OpenCodePromptAdmission { readonly generation: number; readonly turnId: TurnId; readonly messageId: string; + requiresMessageReceipt: boolean; readonly priorAwaitingBusy: boolean; readonly priorIdle: { readonly turnId: TurnId; readonly raw: unknown } | undefined; idleDuringAdmission: { readonly turnId: TurnId; readonly raw: unknown } | undefined; @@ -225,6 +227,7 @@ interface OpenCodePromptAdmission { accepted: boolean; cancelled: boolean; readonly acceptance: Deferred.Deferred; + readonly messageReceipt: Deferred.Deferred; readonly submissionSettled: Deferred.Deferred; promptFiber?: Fiber.Fiber; recoveryFiber?: Fiber.Fiber; @@ -362,6 +365,7 @@ interface OpenCodeSessionContext { pendingRequestRecovery: OpenCodePendingRequestRecovery | undefined; promptGeneration: number; promptAdmission: OpenCodePromptAdmission | undefined; + readonly commandFibers: Set>; readonly promptSemaphore: Semaphore.Semaphore; readonly firstConnection: Deferred.Deferred; /** @@ -1350,8 +1354,14 @@ export function makeOpenCodeAdapter( return; } const recover = Effect.gen(function* () { - yield* Deferred.await(promptAdmission.acceptance); - for (let retryCount = 0; retryCount < 5; retryCount += 1) { + if (!promptAdmission.requiresMessageReceipt) { + yield* Deferred.await(promptAdmission.acceptance); + } + for ( + let retryCount = 0; + retryCount < 5 || (promptAdmission.requiresMessageReceipt && !promptAdmission.accepted); + retryCount += 1 + ) { if ( context.promptAdmission !== promptAdmission || context.activeTurnId !== promptAdmission.turnId || @@ -1386,11 +1396,23 @@ export function makeOpenCodeAdapter( const message = Option.isSome(response) ? response.value.data : undefined; if (message?.info.id === promptAdmission.messageId && message.info.role === "user") { promptAdmission.messageObserved = true; + yield* Deferred.succeed(promptAdmission.messageReceipt, undefined); context.messageRoleById.set(promptAdmission.messageId, "user"); context.textPartsByMessageId.delete(promptAdmission.messageId); } } + // Native command responses wait for generation. Recover their receipt + // first, then let sendTurn acknowledge admission before reconciling idle. + if (promptAdmission.requiresMessageReceipt && !promptAdmission.accepted) { + if (!promptAdmission.messageObserved) { + yield* Effect.sleep(`${Math.min(250 * 2 ** retryCount, 2_000)} millis`); + continue; + } + retryCount = 0; + } + yield* Deferred.await(promptAdmission.acceptance); + const statusResponse = yield* runOpenCodeSdk("session.status", (signal) => context.client.session.status(undefined, { signal }), ).pipe(Effect.timeout("1 second"), Effect.option); @@ -2312,6 +2334,7 @@ export function makeOpenCodeAdapter( promptAdmission?.messageId === event.properties.info.id ) { promptAdmission.messageObserved = true; + yield* Deferred.succeed(promptAdmission.messageReceipt, undefined); if (promptAdmission.accepted) { const idle = promptAdmission.idleDuringAdmission; context.awaitingBusyAfterInterruption = false; @@ -3006,6 +3029,7 @@ export function makeOpenCodeAdapter( pendingRequestRecovery: undefined, promptGeneration: 0, promptAdmission: undefined, + commandFibers: new Set(), promptSemaphore: Semaphore.makeUnsafe(1), firstConnection: Deferred.makeUnsafe(), stopped: yield* Ref.make(false), @@ -3093,6 +3117,13 @@ export function makeOpenCodeAdapter( } const text = input.input?.trim(); + const commandMatch = text?.match(/^\/([^\s/]+)(?:\s+([\s\S]*))?$/); + const nativeCommand = commandMatch + ? (yield* loadOpenCodeCommands(context.client).pipe( + Effect.timeout("10 seconds"), + Effect.orElseSucceed(() => []), + )).find((command) => command.name === commandMatch[1]) + : undefined; // OpenCode ingests images, text, and PDFs natively; formats its model // paths reject ride only as the prompt's file path line. const fileParts = toOpenCodeFileParts({ @@ -3150,6 +3181,7 @@ export function makeOpenCodeAdapter( generation: promptGeneration, turnId, messageId, + requiresMessageReceipt: nativeCommand !== undefined, priorAwaitingBusy, priorIdle: priorIdleCandidate, idleDuringAdmission: undefined, @@ -3160,6 +3192,7 @@ export function makeOpenCodeAdapter( accepted: false, cancelled: false, acceptance: Deferred.makeUnsafe(), + messageReceipt: Deferred.makeUnsafe(), submissionSettled: Deferred.makeUnsafe(), recoveryRaw: undefined, }; @@ -3211,25 +3244,52 @@ export function makeOpenCodeAdapter( } let promptTimedOut = false; - const promptEffect = runOpenCodeSdk("session.promptAsync", (signal) => - context.client.session.promptAsync( - { - sessionID: context.openCodeSessionId, - messageID: messageId, - model: parsedModel, - ...(context.activeAgent ? { agent: context.activeAgent } : {}), - ...(context.activeVariant ? { variant: context.activeVariant } : {}), - // OpenCode appends this after its own agent/provider prompts. - system: buildRuntimeInstructions({ - harness: "OpenCode", - model: `${parsedModel.providerID}/${parsedModel.modelID}`, - }), - parts: [...(text ? [{ type: "text" as const, text }] : []), ...fileParts], - }, - { signal }, - ), - ).pipe( - Effect.timeout("10 seconds"), + const submissionMethod = nativeCommand ? "session.command" : "session.promptAsync"; + // Native commands expand provider-owned templates. Their API does not + // accept the per-turn system addendum supported by ordinary prompts. + const submission = nativeCommand + ? Effect.raceFirst( + runOpenCodeSdk("session.command", (signal) => + context.client.session.command( + { + sessionID: context.openCodeSessionId, + messageID: messageId, + command: nativeCommand.name, + arguments: commandMatch?.[2] ?? "", + model: `${parsedModel.providerID}/${parsedModel.modelID}`, + ...(context.activeAgent ? { agent: context.activeAgent } : {}), + ...(context.activeVariant ? { variant: context.activeVariant } : {}), + parts: fileParts, + }, + { signal }, + ), + ).pipe(Effect.asVoid), + // A command response waits for generation. Only bound admission; + // the user-message receipt proves OpenCode accepted the command. + Deferred.await(promptAdmission.messageReceipt).pipe( + Effect.timeout("10 seconds"), + Effect.andThen(Effect.never), + ), + ) + : runOpenCodeSdk("session.promptAsync", (signal) => + context.client.session.promptAsync( + { + sessionID: context.openCodeSessionId, + messageID: messageId, + model: parsedModel, + ...(context.activeAgent ? { agent: context.activeAgent } : {}), + ...(context.activeVariant ? { variant: context.activeVariant } : {}), + // OpenCode appends this after its own agent/provider prompts. + system: buildRuntimeInstructions({ + harness: "OpenCode", + model: `${parsedModel.providerID}/${parsedModel.modelID}`, + }), + parts: [...(text ? [{ type: "text" as const, text }] : []), ...fileParts], + }, + { signal }, + ), + ).pipe(Effect.timeout("10 seconds"), Effect.asVoid); + const promptEffect = submission.pipe( Effect.catchTags({ OpenCodeRuntimeError: (cause) => Effect.fail(toRequestError(cause)), TimeoutError: (cause) => { @@ -3237,15 +3297,47 @@ export function makeOpenCodeAdapter( return Effect.fail( new ProviderAdapterRequestError({ provider: PROVIDER, - method: "session.promptAsync", + method: submissionMethod, detail: "OpenCode prompt submission did not complete within 10 seconds.", cause, }), ); }, }), - Effect.tapError((requestError) => - context.promptAdmission !== promptAdmission || context.activeTurnId !== turnId + Effect.tapError(() => { + promptAdmission.requiresMessageReceipt = false; + return nativeCommand && promptAdmission.recoveryFiber + ? Fiber.interrupt(promptAdmission.recoveryFiber) + : Effect.void; + }), + Effect.tapError((requestError) => { + if ( + nativeCommand && + (promptAdmission.cancelled || context.cancellation?.turnId === turnId) + ) { + return Effect.void; + } + if ( + nativeCommand && + promptAdmission.accepted && + context.activeTurnId === turnId && + (steeringTurnId !== undefined || + context.promptGeneration !== promptAdmission.generation) + ) { + return Effect.gen(function* () { + yield* emit({ + ...(yield* buildEventBase({ threadId: input.threadId, turnId })), + type: "runtime.warning", + payload: { + message: `OpenCode /${nativeCommand.name} failed after it was accepted.`, + detail: requestError.detail, + }, + }); + }); + } + return (nativeCommand + ? context.promptGeneration !== promptAdmission.generation + : context.promptAdmission !== promptAdmission) || context.activeTurnId !== turnId ? Effect.void : Effect.gen(function* () { if (!promptTimedOut) { @@ -3334,8 +3426,8 @@ export function makeOpenCodeAdapter( tokenUsage, }, }); - }), - ), + }); + }), Effect.onExit((exit) => Effect.gen(function* () { yield* Deferred.succeed(promptAdmission.submissionSettled, undefined).pipe( @@ -3352,8 +3444,20 @@ export function makeOpenCodeAdapter( ); const promptFiber = yield* promptEffect.pipe(Effect.forkIn(context.sessionScope)); promptAdmission.promptFiber = promptFiber; - const promptExit = yield* Effect.exit(Fiber.join(promptFiber)); - delete promptAdmission.promptFiber; + if (nativeCommand) { + context.commandFibers.add(promptFiber); + promptFiber.addObserver(() => context.commandFibers.delete(promptFiber)); + yield* schedulePromptAdmissionRecovery(context, undefined); + } + const promptExit = yield* Effect.exit( + nativeCommand + ? Effect.raceFirst( + Fiber.join(promptFiber), + Deferred.await(promptAdmission.messageReceipt), + ) + : Fiber.join(promptFiber), + ); + if (!nativeCommand) delete promptAdmission.promptFiber; const intentionallyCancelled = promptAdmission.cancelled || @@ -3527,6 +3631,8 @@ export function makeOpenCodeAdapter( yield* Deferred.await(promptAdmission.submissionSettled); } + yield* Effect.forEach([...context.commandFibers], Fiber.interrupt, { discard: true }); + const parentAbortOutcome = yield* Effect.raceFirst( runOpenCodeSdk("session.abort", (signal) => context.client.session.abort({ sessionID: context.openCodeSessionId }, { signal }), diff --git a/apps/server/src/provider/Layers/OpenCodeProvider.test.ts b/apps/server/src/provider/Layers/OpenCodeProvider.test.ts index 0c0bf0c28801..ec00d0399d39 100644 --- a/apps/server/src/provider/Layers/OpenCodeProvider.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeProvider.test.ts @@ -18,7 +18,10 @@ import { type OpenCodeRuntimeShape, } from "../opencodeRuntime.ts"; import * as OpenCodeServerOwner from "../OpenCodeServerOwner.ts"; -import { checkOpenCodeProviderStatus } from "./OpenCodeProvider.ts"; +import { + checkOpenCodeProviderStatus, + openCodeCommandsToServerProviderSlashCommands, +} from "./OpenCodeProvider.ts"; import type { OpenCodeInventory } from "../opencodeRuntime.ts"; const decodeOpenCodeSettings = Schema.decodeSync(OpenCodeSettings); @@ -163,6 +166,22 @@ beforeEach(() => { runtimeMock.reset(); }); +it("keeps native and MCP commands while preserving compaction and separate skills", () => { + NodeAssert.deepEqual( + openCodeCommandsToServerProviderSlashCommands([ + { name: "review", description: "Review changes", source: "command", hints: ["$ARGUMENTS"] }, + { name: "review", source: "command", hints: [] }, + { name: "compact", source: "command", hints: [] }, + { name: "skill", source: "skill", hints: [] }, + { name: "mcp:search", source: "mcp", hints: ["query"] }, + ]).slice(1), + [ + { name: "review", description: "Review changes", input: { hint: "$ARGUMENTS" } }, + { name: "mcp:search", input: { hint: "query" } }, + ], + ); +}); + const testLayer = Layer.succeed(OpenCodeRuntime, OpenCodeRuntimeTestDouble).pipe( Layer.provideMerge(ServerConfig.layerTest(process.cwd(), process.cwd())), Layer.provideMerge(NodeServices.layer), diff --git a/apps/server/src/provider/Layers/OpenCodeProvider.ts b/apps/server/src/provider/Layers/OpenCodeProvider.ts index 7fc33d2bb9f2..0495f8b30737 100644 --- a/apps/server/src/provider/Layers/OpenCodeProvider.ts +++ b/apps/server/src/provider/Layers/OpenCodeProvider.ts @@ -3,6 +3,7 @@ import { type OpenCodeSettings, type ServerProviderModel, type ServerProviderSkill, + type ServerProviderSlashCommand, } from "@t3tools/contracts"; import * as Cause from "effect/Cause"; import * as Data from "effect/Data"; @@ -315,6 +316,26 @@ export function openCodeSkillsToServerProviderSkills( return skills.toSorted((left, right) => left.name.localeCompare(right.name)); } +export function openCodeCommandsToServerProviderSlashCommands( + input: OpenCodeInventory["commands"], +): ReadonlyArray { + const commands: ServerProviderSlashCommand[] = [COMPACT_SLASH_COMMAND]; + const names = new Set([COMPACT_SLASH_COMMAND.name]); + for (const command of input ?? []) { + const name = trimOptional(command.name); + if (!name || names.has(name) || command.source === "skill") continue; + names.add(name); + const description = trimOptional(command.description); + const hint = trimOptional(command.hints.join(" ")); + commands.push({ + name, + ...(description ? { description } : {}), + ...(hint ? { input: { hint } } : {}), + }); + } + return commands; +} + export const makePendingOpenCodeProvider = ( openCodeSettings: OpenCodeSettings, ): Effect.Effect => @@ -526,7 +547,9 @@ export const checkOpenCodeProviderStatus = Effect.fn("checkOpenCodeProviderStatu checkedAt, models, skills, - slashCommands: [COMPACT_SLASH_COMMAND], + slashCommands: openCodeCommandsToServerProviderSlashCommands( + inventoryExit.value.inventory.commands, + ), probe: { installed: true, version, diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index b5894192eed9..70106bd16783 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -358,6 +358,7 @@ export const make = ( const promptSerializationSemaphore = yield* Semaphore.make(1); const promptDispatchSemaphore = yield* Semaphore.make(1); const activePromptRef = yield* Ref.make>(Option.none()); + const assistantUpdatesOpenRef = yield* Ref.make(true); const sessionLoadGateRef = yield* Ref.make>(Option.none()); const ensureConnected = Effect.gen(function* () { @@ -547,6 +548,13 @@ export const make = ( ) { return; } + if ( + !(yield* Ref.get(assistantUpdatesOpenRef)) && + (notification.update.sessionUpdate === "agent_message_chunk" || + notification.update.sessionUpdate === "agent_thought_chunk") + ) { + return; + } yield* processSessionUpdate(notification); }), ), @@ -905,7 +913,16 @@ export const make = ( return; } const acknowledge = yield* Deferred.make(); - yield* Queue.offer(eventQueue, { _tag: "EventStreamBarrier", acknowledge }); + yield* notificationSemaphore.withPermit( + Effect.gen(function* () { + // Keep a provider's final flushed chunks together until the adapter settles the turn. + if (Option.isNone(yield* Ref.get(activePromptRef))) { + yield* Ref.set(assistantUpdatesOpenRef, false); + yield* closeActiveAssistantSegment({ queue: eventQueue, assistantSegmentRef }); + } + yield* Queue.offer(eventQueue, { _tag: "EventStreamBarrier", acknowledge }); + }), + ); yield* Effect.raceFirst(Deferred.await(acknowledge), Deferred.await(runtimeClosed)); }); @@ -984,6 +1001,7 @@ export const make = ( Effect.gen(function* () { const started = yield* getStartedState; yield* closeActiveAssistantSegment({ queue: eventQueue, assistantSegmentRef }); + yield* Ref.set(assistantUpdatesOpenRef, true); const requestPayload = { sessionId: started.sessionId, ...payload, @@ -1000,7 +1018,7 @@ export const make = ( yield* Deferred.succeed(promptOptions.dispatched, undefined); } return active; - }), + }).pipe(notificationSemaphore.withPermit), ), (activePrompt) => Fiber.join(activePrompt.fiber).pipe( diff --git a/apps/server/src/provider/opencodeRuntime.inventory.test.ts b/apps/server/src/provider/opencodeRuntime.inventory.test.ts index 39ffec7436d4..4613a8bd6082 100644 --- a/apps/server/src/provider/opencodeRuntime.inventory.test.ts +++ b/apps/server/src/provider/opencodeRuntime.inventory.test.ts @@ -47,17 +47,70 @@ it.layer(testLayer)("OpenCodeRuntime inventory", (it) => { }); const inventoryFiber = yield* runtime.loadOpenCodeInventory(client).pipe(Effect.forkChild); - yield* Queue.takeN(started, 3); + yield* Queue.takeN(started, 4); yield* Fiber.interrupt(inventoryFiber); NodeAssert.deepEqual((yield* Queue.takeAll(aborted)).toSorted(), [ "/agent", + "/command", "/provider", "/skill", ]); }), ); + it.effect("discovers directory-scoped commands without retaining prompt templates", () => + Effect.gen(function* () { + const runtime = yield* OpenCodeRuntime; + const requests: Request[] = []; + const client = createOpencodeClient({ + baseUrl: "http://opencode.test", + directory: "/workspace/project", + fetch: Object.assign( + async (input: string | Request | URL) => { + const request = input instanceof Request ? input : new Request(input.toString()); + requests.push(request); + const route = new URL(request.url).pathname; + return Response.json( + route === "/provider" + ? { connected: ["openai"], all: [], default: {} } + : route === "/command" + ? [ + { + name: "review", + description: "Review changes", + source: "command", + hints: ["$ARGUMENTS"], + template: "private native template", + }, + ] + : [], + ); + }, + { preconnect: () => undefined }, + ), + }); + const inventory = yield* runtime.loadOpenCodeInventory(client); + NodeAssert.deepEqual(inventory.commands, [ + { + name: "review", + description: "Review changes", + source: "command", + hints: ["$ARGUMENTS"], + }, + ]); + const commandRequest = requests.find( + (request) => new URL(request.url).pathname === "/command", + ); + NodeAssert.ok(commandRequest); + NodeAssert.equal( + new URL(commandRequest.url).searchParams.get("directory") ?? + decodeURIComponent(commandRequest.headers.get("x-opencode-directory") ?? ""), + "/workspace/project", + ); + }), + ); + it.effect("keeps provider inventory when agent discovery fails", () => Effect.gen(function* () { const runtime = yield* OpenCodeRuntime; diff --git a/apps/server/src/provider/opencodeRuntime.ts b/apps/server/src/provider/opencodeRuntime.ts index a79eff843cc7..e87f758e0e3d 100644 --- a/apps/server/src/provider/opencodeRuntime.ts +++ b/apps/server/src/provider/opencodeRuntime.ts @@ -4,6 +4,7 @@ import type { ChatAttachment, ProviderApprovalDecision, RuntimeMode } from "@t3t import { createOpencodeClient, type Agent, + type Command, type FilePartInput, type Model, type OpencodeClient, @@ -187,8 +188,24 @@ export interface OpenCodeInventory { readonly providerList: ProviderListResponse; readonly agents: ReadonlyArray; readonly skills: ReadonlyArray; + readonly commands?: ReadonlyArray; } +export type OpenCodeSlashCommand = Pick; + +/** Command templates stay in OpenCode, which expands arguments and runs MCP prompts. */ +export const loadOpenCodeCommands = (client: OpencodeClient) => + runOpenCodeSdk("command.list", (signal) => client.command.list(undefined, { signal })).pipe( + Effect.map((result): ReadonlyArray => + (result.data ?? []).map(({ name, description, source, hints }) => ({ + name, + ...(description === undefined ? {} : { description }), + ...(source === undefined ? {} : { source }), + hints, + })), + ), + ); + export interface ParsedOpenCodeModelSlug { readonly providerID: string; readonly modelID: string; @@ -923,9 +940,24 @@ const makeOpenCodeRuntime = Effect.gen(function* () { loadOpenCodeSkills(client).pipe(Effect.orElseSucceed((): ReadonlyArray => [])); const loadOpenCodeInventory: OpenCodeRuntimeShape["loadOpenCodeInventory"] = (client) => - Effect.all([loadProviders(client), loadAgents(client), loadSkills(client)], { - concurrency: "unbounded", - }).pipe(Effect.map(([providerList, agents, skills]) => ({ providerList, agents, skills }))); + Effect.all( + [ + loadProviders(client), + loadAgents(client), + loadSkills(client), + loadOpenCodeCommands(client).pipe(Effect.orElseSucceed(() => [])), + ], + { + concurrency: "unbounded", + }, + ).pipe( + Effect.map(([providerList, agents, skills, commands]) => ({ + providerList, + agents, + skills, + commands, + })), + ); const loadInventoryFromCli: OpenCodeRuntimeShape["loadInventoryFromCli"] = (input) => Effect.gen(function* () {