From 732bbde850cda6b2723ef0e866e4be64b92ced54 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Sat, 28 Mar 2026 19:58:14 -0400 Subject: [PATCH] effectify resolveCommand: use Command.Service, Agent, Plugin directly in layer --- packages/opencode/src/session/prompt.ts | 119 ++++- .../test/session/prompt-effect.test.ts | 435 ++++++++++++------ 2 files changed, 410 insertions(+), 144 deletions(-) diff --git a/packages/opencode/src/session/prompt.ts b/packages/opencode/src/session/prompt.ts index f8a358a573..e0e76bf3fa 100644 --- a/packages/opencode/src/session/prompt.ts +++ b/packages/opencode/src/session/prompt.ts @@ -96,6 +96,7 @@ export namespace SessionPrompt { const processor = yield* SessionProcessor.Service const compaction = yield* SessionCompaction.Service const plugin = yield* Plugin.Service + const commands = yield* Command.Service const fsys = yield* AppFileSystem.Service const scope = yield* Scope.Scope @@ -527,8 +528,121 @@ export namespace SessionPrompt { }) const command = Effect.fn("SessionPrompt.command")(function* (input: CommandInput) { - const resolved = yield* Effect.promise(() => resolveCommand(input)) - const result = yield* prompt(resolved.promptInput) + log.info("command", input) + const cmd = yield* commands.get(input.command) + if (!cmd) { + const available = (yield* commands.list()).map((c) => c.name) + const hint = available.length ? ` Available commands: ${available.join(", ")}` : "" + const error = new NamedError.Unknown({ message: `Command not found: "${input.command}".${hint}` }) + yield* bus.publish(Session.Event.Error, { sessionID: input.sessionID, error: error.toObject() }) + throw error + } + const agentName = cmd.agent ?? input.agent ?? (yield* agents.defaultAgent()) + + const raw = input.arguments.match(argsRegex) ?? [] + const args = raw.map((arg) => arg.replace(quoteTrimRegex, "")) + const templateCommand = yield* Effect.promise(() => Promise.resolve(cmd.template)) + + const placeholders = templateCommand.match(placeholderRegex) ?? [] + let last = 0 + for (const item of placeholders) { + const value = Number(item.slice(1)) + if (value > last) last = value + } + + const withArgs = templateCommand.replaceAll(placeholderRegex, (_, index) => { + const position = Number(index) + const argIndex = position - 1 + if (argIndex >= args.length) return "" + if (position === last) return args.slice(argIndex).join(" ") + return args[argIndex] + }) + const usesArgumentsPlaceholder = templateCommand.includes("$ARGUMENTS") + let template = withArgs.replaceAll("$ARGUMENTS", input.arguments) + + if (placeholders.length === 0 && !usesArgumentsPlaceholder && input.arguments.trim()) { + template = template + "\n\n" + input.arguments + } + + const shellMatches = ConfigMarkdown.shell(template) + if (shellMatches.length > 0) { + const sh = Shell.preferred() + const results = yield* Effect.promise(() => + Promise.all(shellMatches.map(async ([, cmd]) => (await Process.text([cmd], { shell: sh, nothrow: true })).text)), + ) + let index = 0 + template = template.replace(bashRegex, () => results[index++]) + } + template = template.trim() + + const taskModel = yield* Effect.promise(async () => { + if (cmd.model) return Provider.parseModel(cmd.model) + if (cmd.agent) { + const cmdAgent = await Agent.get(cmd.agent) + if (cmdAgent?.model) return cmdAgent.model + } + if (input.model) return Provider.parseModel(input.model) + return await lastModelImpl(input.sessionID) + }) + + yield* Effect.promise(() => + Provider.getModel(taskModel.providerID, taskModel.modelID).catch((e) => { + if (Provider.ModelNotFoundError.isInstance(e)) { + const hint = e.data.suggestions?.length ? ` Did you mean: ${e.data.suggestions.join(", ")}?` : "" + Bus.publish(Session.Event.Error, { + sessionID: input.sessionID, + error: new NamedError.Unknown({ message: `Model not found: ${e.data.providerID}/${e.data.modelID}.${hint}` }).toObject(), + }) + } + throw e + }), + ) + + const agent = yield* agents.get(agentName) + if (!agent) { + const available = (yield* agents.list()).filter((a) => !a.hidden).map((a) => a.name) + const hint = available.length ? ` Available agents: ${available.join(", ")}` : "" + const error = new NamedError.Unknown({ message: `Agent not found: "${agentName}".${hint}` }) + yield* bus.publish(Session.Event.Error, { sessionID: input.sessionID, error: error.toObject() }) + throw error + } + + const templateParts = yield* resolvePromptParts(template) + const isSubtask = (agent.mode === "subagent" && cmd.subtask !== false) || cmd.subtask === true + const parts = isSubtask + ? [ + { + type: "subtask" as const, + agent: agent.name, + description: cmd.description ?? "", + command: input.command, + model: { providerID: taskModel.providerID, modelID: taskModel.modelID }, + prompt: templateParts.find((y) => y.type === "text")?.text ?? "", + }, + ] + : [...templateParts, ...(input.parts ?? [])] + + const userAgent = isSubtask ? (input.agent ?? (yield* agents.defaultAgent())) : agentName + const userModel = isSubtask + ? input.model + ? Provider.parseModel(input.model) + : yield* Effect.promise(() => lastModelImpl(input.sessionID)) + : taskModel + + yield* plugin.trigger( + "command.execute.before", + { command: input.command, sessionID: input.sessionID, arguments: input.arguments }, + { parts }, + ) + + const result = yield* prompt({ + sessionID: input.sessionID, + messageID: input.messageID, + model: userModel, + agent: userAgent, + parts, + variant: input.variant, + }) yield* bus.publish(Command.Event.Executed, { name: input.command, sessionID: input.sessionID, @@ -556,6 +670,7 @@ export namespace SessionPrompt { Layer.provide(SessionStatus.layer), Layer.provide(SessionCompaction.defaultLayer), Layer.provide(SessionProcessor.defaultLayer), + Layer.provide(Command.defaultLayer), Layer.provide(AppFileSystem.defaultLayer), Layer.provide(Plugin.defaultLayer), Layer.provide(Session.defaultLayer), diff --git a/packages/opencode/test/session/prompt-effect.test.ts b/packages/opencode/test/session/prompt-effect.test.ts index 1b4bb2d19c..928089eac7 100644 --- a/packages/opencode/test/session/prompt-effect.test.ts +++ b/packages/opencode/test/session/prompt-effect.test.ts @@ -6,6 +6,7 @@ import path from "path" import type { Agent } from "../../src/agent/agent" import { Agent as AgentSvc } from "../../src/agent/agent" import { Bus } from "../../src/bus" +import { Command } from "../../src/command" import { Config } from "../../src/config/config" import { Permission } from "../../src/permission" import { Plugin } from "../../src/plugin" @@ -165,6 +166,7 @@ const deps = Layer.mergeAll( Session.defaultLayer, Snapshot.defaultLayer, AgentSvc.defaultLayer, + Command.defaultLayer, Permission.layer, Plugin.defaultLayer, Config.defaultLayer, @@ -375,108 +377,114 @@ it.effect("loop sets status to busy then idle", () => // Priority 2: Cancel safety -it.effect("cancel interrupts loop and returns last assistant", () => - provideTmpdirInstance( - (dir) => - Effect.gen(function* () { - const test = yield* TestLLM - const prompt = yield* SessionPrompt.Service - const sessions = yield* Session.Service +it.effect( + "cancel interrupts loop and returns last assistant", + () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const test = yield* TestLLM + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service - const chat = yield* sessions.create({}) - yield* seed(chat.id) + const chat = yield* sessions.create({}) + yield* seed(chat.id) - // Make LLM hang so the loop blocks - yield* test.push((input) => hang(input, start())) + // Make LLM hang so the loop blocks + yield* test.push((input) => hang(input, start())) - // Seed a new user message so the loop enters the LLM path - yield* user(chat.id, "more") + // Seed a new user message so the loop enters the LLM path + yield* user(chat.id, "more") - const fiber = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) - // Give the loop time to start - yield* Effect.promise(() => new Promise((r) => setTimeout(r, 200))) - yield* prompt.cancel(chat.id) + const fiber = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) + // Give the loop time to start + yield* Effect.promise(() => new Promise((r) => setTimeout(r, 200))) + yield* prompt.cancel(chat.id) - const exit = yield* Fiber.await(fiber) - expect(Exit.isSuccess(exit)).toBe(true) - if (Exit.isSuccess(exit)) { - expect(exit.value.info.role).toBe("assistant") - } - }), - { git: true, config: cfg }, - ), - 30_000, -) - -it.effect("cancel records MessageAbortedError on interrupted process", () => - provideTmpdirInstance( - (dir) => - Effect.gen(function* () { - const ready = defer() - const test = yield* TestLLM - const prompt = yield* SessionPrompt.Service - const sessions = yield* Session.Service - - yield* test.push((input) => - hang(input, start()).pipe( - Stream.tap((event) => (event.type === "start" ? Effect.sync(() => ready.resolve()) : Effect.void)), - ), - ) - - const chat = yield* sessions.create({}) - yield* user(chat.id, "hello") - - const fiber = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) - yield* Effect.promise(() => ready.promise) - yield* prompt.cancel(chat.id) - - const exit = yield* Fiber.await(fiber) - expect(Exit.isSuccess(exit)).toBe(true) - if (Exit.isSuccess(exit)) { - const info = exit.value.info - if (info.role === "assistant") { - expect(info.error?.name).toBe("MessageAbortedError") + const exit = yield* Fiber.await(fiber) + expect(Exit.isSuccess(exit)).toBe(true) + if (Exit.isSuccess(exit)) { + expect(exit.value.info.role).toBe("assistant") } - } - }), - { git: true, config: cfg }, - ), + }), + { git: true, config: cfg }, + ), 30_000, ) -it.effect("cancel with queued callers resolves all cleanly", () => - provideTmpdirInstance( - (dir) => - Effect.gen(function* () { - const ready = defer() - const test = yield* TestLLM - const prompt = yield* SessionPrompt.Service - const sessions = yield* Session.Service +it.effect( + "cancel records MessageAbortedError on interrupted process", + () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const ready = defer() + const test = yield* TestLLM + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service - yield* test.push((input) => - hang(input, start()).pipe( - Stream.tap((event) => (event.type === "start" ? Effect.sync(() => ready.resolve()) : Effect.void)), - ), - ) + yield* test.push((input) => + hang(input, start()).pipe( + Stream.tap((event) => (event.type === "start" ? Effect.sync(() => ready.resolve()) : Effect.void)), + ), + ) - const chat = yield* sessions.create({}) - yield* user(chat.id, "hello") + const chat = yield* sessions.create({}) + yield* user(chat.id, "hello") - const a = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) - yield* Effect.promise(() => ready.promise) - // Queue a second caller - const b = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) - yield* Effect.promise(() => new Promise((r) => setTimeout(r, 50))) + const fiber = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) + yield* Effect.promise(() => ready.promise) + yield* prompt.cancel(chat.id) - yield* prompt.cancel(chat.id) + const exit = yield* Fiber.await(fiber) + expect(Exit.isSuccess(exit)).toBe(true) + if (Exit.isSuccess(exit)) { + const info = exit.value.info + if (info.role === "assistant") { + expect(info.error?.name).toBe("MessageAbortedError") + } + } + }), + { git: true, config: cfg }, + ), + 30_000, +) - const [exitA, exitB] = yield* Effect.all([Fiber.await(a), Fiber.await(b)]) - // Both should resolve (success or interrupt, not error) - expect(Exit.isFailure(exitA) && !Cause.hasInterruptsOnly(exitA.cause)).toBe(false) - expect(Exit.isFailure(exitB) && !Cause.hasInterruptsOnly(exitB.cause)).toBe(false) - }), - { git: true, config: cfg }, - ), +it.effect( + "cancel with queued callers resolves all cleanly", + () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const ready = defer() + const test = yield* TestLLM + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service + + yield* test.push((input) => + hang(input, start()).pipe( + Stream.tap((event) => (event.type === "start" ? Effect.sync(() => ready.resolve()) : Effect.void)), + ), + ) + + const chat = yield* sessions.create({}) + yield* user(chat.id, "hello") + + const a = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) + yield* Effect.promise(() => ready.promise) + // Queue a second caller + const b = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) + yield* Effect.promise(() => new Promise((r) => setTimeout(r, 50))) + + yield* prompt.cancel(chat.id) + + const [exitA, exitB] = yield* Effect.all([Fiber.await(a), Fiber.await(b)]) + // Both should resolve (success or interrupt, not error) + expect(Exit.isFailure(exitA) && !Cause.hasInterruptsOnly(exitA.cause)).toBe(false) + expect(Exit.isFailure(exitB) && !Cause.hasInterruptsOnly(exitB.cause)).toBe(false) + }), + { git: true, config: cfg }, + ), 30_000, ) @@ -518,10 +526,9 @@ it.effect("concurrent loop callers all receive same error result", () => const chat = yield* sessions.create({}) yield* user(chat.id, "hello") - const [a, b] = yield* Effect.all( - [prompt.loop({ sessionID: chat.id }), prompt.loop({ sessionID: chat.id })], - { concurrency: "unbounded" }, - ) + const [a, b] = yield* Effect.all([prompt.loop({ sessionID: chat.id }), prompt.loop({ sessionID: chat.id })], { + concurrency: "unbounded", + }) // Both callers get the same assistant with an error recorded expect(a.info.id).toBe(b.info.id) @@ -534,35 +541,37 @@ it.effect("concurrent loop callers all receive same error result", () => ), ) -it.effect("assertNotBusy throws BusyError when loop running", () => - provideTmpdirInstance( - (dir) => - Effect.gen(function* () { - const ready = defer() - const test = yield* TestLLM - const prompt = yield* SessionPrompt.Service - const sessions = yield* Session.Service +it.effect( + "assertNotBusy throws BusyError when loop running", + () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const ready = defer() + const test = yield* TestLLM + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service - yield* test.push((input) => - hang(input, start()).pipe( - Stream.tap((event) => (event.type === "start" ? Effect.sync(() => ready.resolve()) : Effect.void)), - ), - ) + yield* test.push((input) => + hang(input, start()).pipe( + Stream.tap((event) => (event.type === "start" ? Effect.sync(() => ready.resolve()) : Effect.void)), + ), + ) - const chat = yield* sessions.create({}) - yield* user(chat.id, "hi") + const chat = yield* sessions.create({}) + yield* user(chat.id, "hi") - const fiber = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) - yield* Effect.promise(() => ready.promise) + const fiber = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) + yield* Effect.promise(() => ready.promise) - const exit = yield* prompt.assertNotBusy(chat.id).pipe(Effect.exit) - expect(Exit.isFailure(exit)).toBe(true) + const exit = yield* prompt.assertNotBusy(chat.id).pipe(Effect.exit) + expect(Exit.isFailure(exit)).toBe(true) - yield* prompt.cancel(chat.id) - yield* Fiber.await(fiber) - }), - { git: true, config: cfg }, - ), + yield* prompt.cancel(chat.id) + yield* Fiber.await(fiber) + }), + { git: true, config: cfg }, + ), 30_000, ) @@ -583,36 +592,178 @@ it.effect("assertNotBusy succeeds when idle", () => // Priority 4: Shell basics -it.effect("shell rejects with BusyError when loop running", () => - provideTmpdirInstance( - (dir) => - Effect.gen(function* () { - const ready = defer() - const test = yield* TestLLM - const prompt = yield* SessionPrompt.Service - const sessions = yield* Session.Service +it.effect( + "shell rejects with BusyError when loop running", + () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const ready = defer() + const test = yield* TestLLM + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service - yield* test.push((input) => - hang(input, start()).pipe( - Stream.tap((event) => (event.type === "start" ? Effect.sync(() => ready.resolve()) : Effect.void)), - ), - ) + yield* test.push((input) => + hang(input, start()).pipe( + Stream.tap((event) => (event.type === "start" ? Effect.sync(() => ready.resolve()) : Effect.void)), + ), + ) - const chat = yield* sessions.create({}) - yield* user(chat.id, "hi") + const chat = yield* sessions.create({}) + yield* user(chat.id, "hi") - const fiber = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) - yield* Effect.promise(() => ready.promise) + const fiber = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) + yield* Effect.promise(() => ready.promise) - const exit = yield* prompt - .shell({ sessionID: chat.id, agent: "build", command: "echo hi" }) - .pipe(Effect.exit) - expect(Exit.isFailure(exit)).toBe(true) + const exit = yield* prompt.shell({ sessionID: chat.id, agent: "build", command: "echo hi" }).pipe(Effect.exit) + expect(Exit.isFailure(exit)).toBe(true) - yield* prompt.cancel(chat.id) - yield* Fiber.await(fiber) - }), - { git: true, config: cfg }, - ), + yield* prompt.cancel(chat.id) + yield* Fiber.await(fiber) + }), + { git: true, config: cfg }, + ), + 30_000, +) + +it.effect( + "loop waits while shell runs and starts after shell exits", + () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const test = yield* TestLLM + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service + + yield* test.reply(start(), textStart(), textDelta("t", "after-shell"), textEnd(), finishStep(), finish()) + + const chat = yield* sessions.create({}) + + const sh = yield* prompt + .shell({ sessionID: chat.id, agent: "build", command: "sleep 0.2" }) + .pipe(Effect.forkChild) + yield* Effect.promise(() => new Promise((done) => setTimeout(done, 50))) + + const run = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) + yield* Effect.promise(() => new Promise((done) => setTimeout(done, 50))) + + expect(yield* test.calls).toBe(0) + + yield* Fiber.await(sh) + const exit = yield* Fiber.await(run) + + expect(Exit.isSuccess(exit)).toBe(true) + if (Exit.isSuccess(exit)) { + expect(exit.value.info.role).toBe("assistant") + expect(exit.value.parts.some((part) => part.type === "text" && part.text === "after-shell")).toBe(true) + } + expect(yield* test.calls).toBe(1) + }), + { git: true, config: cfg }, + ), + 30_000, +) + +it.effect( + "shell completion resumes queued loop callers", + () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const test = yield* TestLLM + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service + + yield* test.reply(start(), textStart(), textDelta("t", "done"), textEnd(), finishStep(), finish()) + + const chat = yield* sessions.create({}) + + const sh = yield* prompt + .shell({ sessionID: chat.id, agent: "build", command: "sleep 0.2" }) + .pipe(Effect.forkChild) + yield* Effect.promise(() => new Promise((done) => setTimeout(done, 50))) + + const a = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) + const b = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) + yield* Effect.promise(() => new Promise((done) => setTimeout(done, 50))) + + expect(yield* test.calls).toBe(0) + + yield* Fiber.await(sh) + const [ea, eb] = yield* Effect.all([Fiber.await(a), Fiber.await(b)]) + + expect(Exit.isSuccess(ea)).toBe(true) + expect(Exit.isSuccess(eb)).toBe(true) + if (Exit.isSuccess(ea) && Exit.isSuccess(eb)) { + expect(ea.value.info.id).toBe(eb.value.info.id) + expect(ea.value.info.role).toBe("assistant") + } + expect(yield* test.calls).toBe(1) + }), + { git: true, config: cfg }, + ), + 30_000, +) + +it.effect( + "cancel interrupts shell and resolves cleanly", + () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service + + const chat = yield* sessions.create({}) + + const sh = yield* prompt + .shell({ sessionID: chat.id, agent: "build", command: "sleep 30" }) + .pipe(Effect.forkChild) + yield* Effect.promise(() => new Promise((done) => setTimeout(done, 50))) + + yield* prompt.cancel(chat.id) + + const exit = yield* Fiber.await(sh) + expect(Exit.isSuccess(exit)).toBe(true) + if (Exit.isSuccess(exit)) { + expect(exit.value.info.role).toBe("assistant") + expect(exit.value.parts.some((part) => part.type === "tool")).toBe(true) + } + + const status = yield* SessionStatus.Service + expect((yield* status.get(chat.id)).type).toBe("idle") + const busy = yield* prompt.assertNotBusy(chat.id).pipe(Effect.exit) + expect(Exit.isSuccess(busy)).toBe(true) + }), + { git: true, config: cfg }, + ), + 30_000, +) + +it.effect( + "shell rejects when another shell is already running", + () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service + + const chat = yield* sessions.create({}) + + const a = yield* prompt + .shell({ sessionID: chat.id, agent: "build", command: "sleep 30" }) + .pipe(Effect.forkChild) + yield* Effect.promise(() => new Promise((done) => setTimeout(done, 50))) + + const exit = yield* prompt.shell({ sessionID: chat.id, agent: "build", command: "echo hi" }).pipe(Effect.exit) + expect(Exit.isFailure(exit)).toBe(true) + + yield* prompt.cancel(chat.id) + yield* Fiber.await(a) + }), + { git: true, config: cfg }, + ), 30_000, )