effectify resolveCommand: use Command.Service, Agent, Plugin directly in layer

This commit is contained in:
Kit Langton
2026-03-28 19:58:14 -04:00
parent 4f784276b4
commit 732bbde850
2 changed files with 410 additions and 144 deletions
+117 -2
View File
@@ -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),
@@ -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<void>((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<void>((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<void>()
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<void>()
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<void>()
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<void>((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<void>()
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<void>((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<void>()
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<void>()
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<void>()
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<void>()
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<void>((done) => setTimeout(done, 50)))
const run = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild)
yield* Effect.promise(() => new Promise<void>((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<void>((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<void>((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<void>((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<void>((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,
)