Compare commits

...
10 changed files with 772 additions and 556 deletions

No files matched your search

+37 -33
View File
@@ -8,22 +8,34 @@ import {
type AgentRequestMethod,
type Stream,
} from "@agentclientprotocol/sdk"
import { ClientError, type OpenCodeClient } from "@opencode/client/promise"
import { Cause, Effect, type Scope } from "effect"
import type { OpenCodeClient } from "@opencode/client/promise"
import { Cause, Deferred, Effect, Ref, type Scope } from "effect"
import { ACPCatalog } from "./catalog"
import { ACPConnection } from "./connection"
import { ACPError } from "./error"
import { ACPPromise } from "./promise"
import { ACPService } from "./service"
import { ACPSessions } from "./sessions"
import { ACPTurn } from "./turn"
// Untraced so request spans parent to the caller's span instead of a setup span that has already ended.
export const connect = Effect.fnUntraced(function* (client: OpenCodeClient, stream: Stream) {
const run = Effect.runPromiseWith(yield* Effect.context<Scope.Scope>())
const catalog = yield* ACPCatalog.make(client)
// Requests can dispatch once the stream's read loop yields, which may be before the service below is built.
const ready = yield* Deferred.make<ACPService.Interface>()
const handle =
<Params, A>(call: (ctx: AgentHandlerContext<Params>) => Effect.Effect<A, ACPError.Error | RequestError>) =>
<Params, A>(
call: (service: ACPService.Interface, ctx: AgentHandlerContext<Params>) => Effect.Effect<A, ACPService.Failure>,
) =>
(name: string) => {
const handler = Effect.fn(name)(
call,
(ctx: AgentHandlerContext<Params>) =>
Deferred.await(ready).pipe(Effect.flatMap((service) => call(service, ctx))),
Effect.catchTags({
ACPCatalogLoadError: (error) => ACPPromise.classify(error.cause),
ACPCatalogNotReadyError: (error) => Effect.die(error),
}),
Effect.mapError((error) => (error instanceof RequestError ? error : ACPError.toRequestError(error))),
Effect.tapCauseIf(Cause.hasDies, (cause) => Effect.logError("ACP request failed", cause)),
Effect.catchDefect((defect) => Effect.fail(ACPError.toRequestError(ACPError.fromUnknown(defect)))),
@@ -42,78 +54,70 @@ export const connect = Effect.fnUntraced(function* (client: OpenCodeClient, stre
request(
"initialize",
handle((ctx) => promise(() => service.initialize(ctx.params))),
handle((service, ctx) => service.initialize(ctx.params)),
)
request(
"authenticate",
handle((ctx) => promise(() => service.authenticate(ctx.params))),
handle((service, ctx) => service.authenticate(ctx.params)),
)
request(
"session/new",
handle((ctx) => promise(() => service.newSession(ctx.params))),
handle((service, ctx) => service.newSession(ctx.params)),
)
request(
"session/load",
handle((ctx) => promise(() => service.loadSession(ctx.params))),
handle((service, ctx) => service.loadSession(ctx.params)),
)
request(
"session/list",
handle((ctx) => promise(() => service.listSessions(ctx.params))),
handle((service, ctx) => service.listSessions(ctx.params)),
)
request(
"session/delete",
handle((ctx) => promise(() => service.deleteSession(ctx.params))),
handle((service, ctx) => service.deleteSession(ctx.params)),
)
request(
"session/resume",
handle((ctx) => promise(() => service.resumeSession(ctx.params))),
handle((service, ctx) => service.resumeSession(ctx.params)),
)
request(
"session/close",
handle((ctx) => promise(() => service.closeSession(ctx.params))),
handle((service, ctx) => service.closeSession(ctx.params)),
)
request(
"session/fork",
handle((ctx) => promise(() => service.forkSession(ctx.params))),
handle((service, ctx) => service.forkSession(ctx.params)),
)
request(
"session/set_config_option",
handle((ctx) => promise(() => service.setSessionConfigOption(ctx.params))),
handle((service, ctx) => service.setSessionConfigOption(ctx.params)),
)
request(
"session/set_mode",
handle((ctx) => promise(() => service.setSessionMode(ctx.params))),
handle((service, ctx) => service.setSessionMode(ctx.params)),
)
// The SDK signal is passed through rather than interrupting the fiber: a cancelled turn still resolves with
// `stopReason: "cancelled"`.
request(
"session/prompt",
handle((ctx) => promise(() => service.prompt(ctx.params, ctx.signal))),
handle((service, ctx) => ACPPromise.promise(() => service.prompt(ctx.params, ctx.signal))),
)
notification(
"session/cancel",
handle((ctx) => promise(() => service.cancel(ctx.params))),
handle((service, ctx) => ACPPromise.promise(() => service.cancel(ctx.params))),
)
const connection = app.connect(stream)
// Inbound dispatch starts after the stream's async read loop yields, so handlers never observe this before assignment.
const service = ACPService.make({ client, connection: ACPConnection.make(connection), catalog, run })
const promiseConnection = ACPConnection.promise(connection)
const sessions = yield* ACPSessions.make({ client, connection: ACPConnection.make(connection), catalog })
const capabilities = yield* Ref.make({ childSessionUpdates: false })
const turn = ACPTurn.make({ client, connection: promiseConnection, sessions, catalog, capabilities, run })
yield* Deferred.succeed(
ready,
ACPService.make({ client, connection: promiseConnection, catalog, sessions, capabilities, turn }),
)
return connection
})
const spanName = (method: string) => `cli.acp.${method.replaceAll("/", ".")}`
const promise = <A>(evaluate: () => Promise<A>) =>
Effect.tryPromise({
try: evaluate,
// A catalog load failure is classified by the client error that caused it.
catch: (cause) => (cause instanceof ACPCatalog.LoadError ? cause.cause : cause),
}).pipe(
Effect.catch((cause) => {
if (cause instanceof RequestError || ACPError.is(cause)) return Effect.fail(cause)
if (cause instanceof ClientError && cause.reason === "Transport")
return Effect.fail(new ACPError.ServerUnavailableError())
return Effect.die(cause)
}),
)
export * as ACP from "./agent"
+8 -48
View File
@@ -1,14 +1,20 @@
import type { CommandInfo, ModelInfo, ModelRef, OpenCodeClient, OpenCodeEvent } from "@opencode/client/promise"
import { FSUtil } from "@opencode/util/fs-util"
import { Context, Deferred, Effect, Exit, Schedule, Schema, Scope, Semaphore, Stream, SubscriptionRef } from "effect"
import { Context, Deferred, Effect, Exit, Schedule, Schema, Semaphore, Stream, SubscriptionRef } from "effect"
import type { ConfigOptionProvider } from "./config-option"
// ACP runs these itself; they take precedence over server commands with the same name.
export const builtinCommands = new Map([
["compact", { description: "Compact the session", start: "compaction" as const }],
])
export type Catalog = {
readonly providers: ConfigOptionProvider[]
readonly models: ModelInfo[]
readonly defaultModel: ModelRef
readonly modes: Array<{ id: string; name: string; description?: string }>
readonly defaultModeID: string
/** Server commands, without those shadowed by a built-in. */
readonly commands: CommandInfo[]
}
@@ -133,52 +139,6 @@ export const make = Effect.fnUntraced(function* (client: OpenCodeClient) {
})
})
export type Live = {
readonly cwd: string
readonly current: Catalog
}
/** Temporary adapter for the promise-based `ACPService` until sessions consume the service directly. */
export function promise(input: {
readonly catalog: Interface
readonly run: <A, E>(effect: Effect.Effect<A, E, Scope.Scope>) => Promise<A>
readonly changed: (live: Live, previous: Catalog) => Promise<unknown>
}) {
const lives = new Map<string, Live>()
return {
get: (cwd: string) =>
input.run(
Effect.gen(function* () {
const current = yield* input.catalog.get(cwd)
const key = FSUtil.resolve(cwd)
const existing = lives.get(key)
if (existing) return existing
const live: Live = {
cwd,
// A loaded entry's read never suspends.
get current() {
return Effect.runSync(input.catalog.get(cwd))
},
}
lives.set(key, live)
yield* input.catalog.changes(cwd).pipe(
Stream.runFoldEffect(
() => current,
(previous, next) =>
next === previous
? Effect.succeed(previous)
: Effect.promise(() => input.changed(live, previous).catch(() => {})).pipe(Effect.as(next)),
),
Effect.ignore,
Effect.forkScoped({ startImmediately: true }),
)
return live
}),
),
reload: (live: Live) => input.run(input.catalog.reload(live.cwd)),
}
}
const load = (client: OpenCodeClient, cwd: string) =>
read(client, cwd).pipe(
// Some providers discover models in the background after plugin startup begins.
@@ -222,7 +182,7 @@ const read = Effect.fnUntraced(function* (client: OpenCodeClient, cwd: string) {
},
modes: agents.map((agent) => ({ id: agent.id, name: agent.name, description: agent.description })),
defaultModeID: defaultAgent.id,
commands: commandResult.data,
commands: commandResult.data.filter((command) => !builtinCommands.has(command.name)),
} satisfies Catalog
})
+30
View File
@@ -1,4 +1,6 @@
import type { SessionConfigOption } from "@agentclientprotocol/sdk"
import type { ModelRef } from "@opencode/client/promise"
import { builtinCommands, type Catalog } from "./catalog"
export const DEFAULT_VARIANT_VALUE = "default"
@@ -25,6 +27,34 @@ export type ModelSelection = {
variant?: string
}
/** A session's model and mode. Unset fields follow the catalog defaults. */
export type Selection = {
readonly model?: ModelRef
readonly modeID?: string
}
export function currentModel(catalog: Catalog, selection: Selection) {
return selection.model ?? catalog.defaultModel
}
export function configOptions(catalog: Catalog, selection: Selection) {
const model = currentModel(catalog, selection)
return buildConfigOptions({
providers: catalog.providers,
currentModel: { providerID: model.providerID, modelID: model.id },
currentVariant: model.variant,
modes: catalog.modes,
currentModeId: selection.modeID ?? catalog.defaultModeID,
})
}
export function availableCommands(catalog: Catalog) {
return [
...catalog.commands.map((command) => ({ name: command.name, description: command.description ?? "" })),
...Array.from(builtinCommands, ([name, command]) => ({ name, description: command.description })),
]
}
export function buildConfigOptions(input: {
providers: readonly ConfigOptionProvider[]
currentModel: ModelSelection["model"]
+19 -1
View File
@@ -1,11 +1,28 @@
import {
methods,
type AgentConnection,
type RequestError,
type RequestPermissionRequest,
type RequestPermissionResponse,
type SendRequestOptions,
type SessionNotification,
} from "@agentclientprotocol/sdk"
import { Context, type Effect } from "effect"
import type { ACPError } from "./error"
import { ACPPromise } from "./promise"
export interface Interface {
readonly sessionUpdate: (params: SessionNotification) => Effect.Effect<void, ACPError.Error | RequestError>
}
export class Service extends Context.Service<Service, Interface>()("@opencode/cli/acp/Connection") {}
export function make(connection: AgentConnection) {
return Service.of({
sessionUpdate: (params) =>
ACPPromise.promise(() => connection.client.notify(methods.client.session.update, params)),
})
}
export type Connection = {
readonly signal?: AbortSignal
@@ -14,7 +31,8 @@ export type Connection = {
extNotification?(method: string, params: Record<string, unknown>): Promise<void>
}
export function make(connection: AgentConnection): Connection {
/** Promise view for the turn, permission, and replay code until they run as effects. */
export function promise(connection: AgentConnection): Connection {
return {
signal: connection.signal,
sessionUpdate: (params) => connection.client.notify(methods.client.session.update, params),
+17
View File
@@ -0,0 +1,17 @@
import { RequestError } from "@agentclientprotocol/sdk"
import { ClientError } from "@opencode/client/promise"
import { Effect } from "effect"
import { ACPError } from "./error"
/** Runs a promise, keeping ACP failures typed. Any other rejection is a defect. */
export const promise = <A>(evaluate: (signal: AbortSignal) => Promise<A>) =>
Effect.tryPromise({ try: evaluate, catch: (cause) => cause }).pipe(Effect.catch(classify))
export function classify(cause: unknown): Effect.Effect<never, ACPError.Error | RequestError> {
if (cause instanceof RequestError || ACPError.is(cause)) return Effect.fail(cause)
if (cause instanceof ClientError && cause.reason === "Transport")
return Effect.fail(new ACPError.ServerUnavailableError())
return Effect.die(cause)
}
export * as ACPPromise from "./promise"
+165 -460
View File
@@ -1,14 +1,6 @@
import { isDeepStrictEqual } from "node:util"
import {
isSessionNotFoundError,
type CommandInfo,
type ModelRef,
type OpenCodeClient,
type SessionInfo,
type SessionMessageInfo,
} from "@opencode/client/promise"
import { isSessionNotFoundError, type ModelRef, type OpenCodeClient } from "@opencode/client/promise"
import { FSUtil } from "@opencode/util/fs-util"
import type { Effect, Scope } from "effect"
import { Effect, Option, Ref, Stream } from "effect"
import { withTimestampedFallback } from "@opencode/util/session-title-fallback"
import type {
AuthenticateRequest,
@@ -27,11 +19,11 @@ import type {
ListSessionsResponse,
LoadSessionRequest,
LoadSessionResponse,
McpServer,
NewSessionRequest,
NewSessionResponse,
PromptRequest,
PromptResponse,
RequestError,
ResumeSessionRequest,
ResumeSessionResponse,
SetSessionConfigOptionRequest,
@@ -40,177 +32,116 @@ import type {
SetSessionModeResponse,
} from "@agentclientprotocol/sdk"
import { OPENCODE_VERSION } from "../version"
import { SessionMessage } from "@opencode/schema/session-message"
import { ACPCatalog, type Catalog } from "./catalog"
import { buildConfigOptions, DEFAULT_VARIANT_VALUE, parseModelSelection } from "./config-option"
import type { ACPCatalog, Catalog } from "./catalog"
import { configOptions, currentModel, DEFAULT_VARIANT_VALUE, parseModelSelection } from "./config-option"
import type { ACPConnection } from "./connection"
import { promptContentToParts } from "./content"
import {
ChildSessionUpdateMethod,
ChildSessionUpdatesCapability,
replayMessages,
streamTurn,
type ChildSessionUpdate,
type TurnControl,
type TurnStart,
} from "./event"
import { ChildSessionUpdatesCapability, replayMessages } from "./event"
import { ACPError } from "./error"
import { ACPPromise } from "./promise"
import type { ACPSessions, Attached } from "./sessions"
import type { ACPTurn } from "./turn"
export const AuthMethodID = "opencode-login"
// ACP runs these itself; they take precedence over server commands with the same name.
const builtinCommands = new Map([["compact", { description: "Compact the session", start: "compaction" as const }]])
// Model and mode are unset while the session follows the server defaults.
type Attached = {
readonly id: string
readonly cwd: string
readonly abort: AbortController
readonly catalog: ACPCatalog.Live
model?: ModelRef
modeID?: string
}
type PreparedPrompt = {
readonly start: TurnStart
readonly text: string
readonly files: Array<{ readonly uri: string; readonly name?: string }>
readonly synthetic: ReadonlyArray<string>
readonly slash?: { readonly name: string; readonly args: string }
readonly command?: CommandInfo
}
export type Failure = ACPError.Error | RequestError | ACPCatalog.Error
export interface Interface {
initialize(input: InitializeRequest): Promise<InitializeResponse>
authenticate(input: AuthenticateRequest): Promise<AuthenticateResponse>
newSession(input: NewSessionRequest): Promise<NewSessionResponse>
loadSession(input: LoadSessionRequest): Promise<LoadSessionResponse>
listSessions(input: ListSessionsRequest): Promise<ListSessionsResponse>
deleteSession(input: DeleteSessionRequest): Promise<DeleteSessionResponse>
resumeSession(input: ResumeSessionRequest): Promise<ResumeSessionResponse>
closeSession(input: CloseSessionRequest): Promise<CloseSessionResponse>
forkSession(input: ForkSessionRequest): Promise<ForkSessionResponse>
setSessionConfigOption(input: SetSessionConfigOptionRequest): Promise<SetSessionConfigOptionResponse>
setSessionMode(input: SetSessionModeRequest): Promise<SetSessionModeResponse>
readonly initialize: (input: InitializeRequest) => Effect.Effect<InitializeResponse>
readonly authenticate: (input: AuthenticateRequest) => Effect.Effect<AuthenticateResponse, Failure>
readonly newSession: (input: NewSessionRequest) => Effect.Effect<NewSessionResponse, Failure>
readonly loadSession: (input: LoadSessionRequest) => Effect.Effect<LoadSessionResponse, Failure>
readonly listSessions: (input: ListSessionsRequest) => Effect.Effect<ListSessionsResponse, Failure>
readonly deleteSession: (input: DeleteSessionRequest) => Effect.Effect<DeleteSessionResponse, Failure>
readonly resumeSession: (input: ResumeSessionRequest) => Effect.Effect<ResumeSessionResponse, Failure>
readonly closeSession: (input: CloseSessionRequest) => Effect.Effect<CloseSessionResponse, Failure>
readonly forkSession: (input: ForkSessionRequest) => Effect.Effect<ForkSessionResponse, Failure>
readonly setSessionConfigOption: (
input: SetSessionConfigOptionRequest,
) => Effect.Effect<SetSessionConfigOptionResponse, Failure>
readonly setSessionMode: (input: SetSessionModeRequest) => Effect.Effect<SetSessionModeResponse, Failure>
prompt(input: PromptRequest, signal?: AbortSignal): Promise<PromptResponse>
cancel(input: CancelNotification): Promise<void>
}
export function make(input: {
readonly client: OpenCodeClient
/** Replay still writes through the promise view. */
readonly connection: ACPConnection.Connection
readonly catalog: ACPCatalog.Interface
readonly run: <A, E>(effect: Effect.Effect<A, E, Scope.Scope>) => Promise<A>
readonly sessions: ACPSessions.Interface
readonly capabilities: Ref.Ref<{ readonly childSessionUpdates: boolean }>
readonly turn: ACPTurn.Interface
}): Interface {
const sessions = new Map<string, Attached>()
const registeredMcp = new Map<string, Set<string>>()
const active = new Map<string, { readonly control: TurnControl; readonly turn: Promise<PromptResponse> }>()
const capabilities = { childSessionUpdates: false }
const catalogs = ACPCatalog.promise({
catalog: input.catalog,
run: input.run,
changed: (live, previous) =>
Promise.all(
Array.from(sessions.values())
.filter((state) => state.catalog === live)
.map(async (state) => {
const options = configOptions(state)
if (!isDeepStrictEqual(options, configOptions(state, previous))) {
await input.connection.sessionUpdate({
sessionId: state.id,
update: { sessionUpdate: "config_option_update", configOptions: options },
})
}
if (!isDeepStrictEqual(live.current.commands, previous.commands)) await sendCommands(state)
}),
),
const currentOptions = Effect.fnUntraced(function* (attached: Attached) {
return configOptions(yield* input.catalog.get(attached.cwd), yield* Ref.get(attached.selection))
})
const sendCommands = (state: Attached) =>
input.connection.sessionUpdate({
sessionId: state.id,
update: {
sessionUpdate: "available_commands_update",
availableCommands: [
...state.catalog.current.commands
.filter((command) => !builtinCommands.has(command.name))
.map((command) => ({ name: command.name, description: command.description ?? "" })),
...Array.from(builtinCommands, ([name, command]) => ({ name, description: command.description })),
],
},
})
// A selection the catalog has not seen may be new on the server, so reload once before rejecting it.
const withReload = <A>(attached: Attached, select: Effect.Effect<A, Failure>) => {
const retry = () => input.catalog.reload(attached.cwd).pipe(Effect.andThen(select))
return select.pipe(
Effect.catchTags({ ACPInvalidModelError: retry, ACPInvalidModeError: retry, ACPInvalidEffortError: retry }),
)
}
const withReload = <A>(state: Attached, select: () => Promise<A>) =>
select().catch(async (error: unknown) => {
if (
!(
error instanceof ACPError.InvalidModelError ||
error instanceof ACPError.InvalidModeError ||
error instanceof ACPError.InvalidEffortError
)
)
const selectOption = Effect.fnUntraced(function* (attached: Attached, configId: string, value: string) {
const catalog = yield* input.catalog.get(attached.cwd)
const current = currentModel(catalog, yield* Ref.get(attached.selection))
switch (configId) {
case "model":
return yield* selectModel(attached, yield* requireModel(catalog, value, current))
case "effort":
return yield* selectModel(attached, yield* requireEffort(catalog, value, current))
case "mode":
return yield* selectMode(attached, value)
default:
return yield* new ACPError.InvalidConfigOptionError({ configId })
}
})
const selectModel = Effect.fnUntraced(function* (attached: Attached, model: ModelRef) {
yield* Ref.update(attached.selection, (selection) => ({ ...selection, model }))
yield* ACPPromise.promise(() => input.client.session.switchModel({ sessionID: attached.id, model }))
})
const selectMode = Effect.fnUntraced(function* (attached: Attached, modeID: string) {
const catalog = yield* input.catalog.get(attached.cwd)
if (!catalog.modes.some((mode) => mode.id === modeID)) return yield* new ACPError.InvalidModeError({ mode: modeID })
yield* Ref.update(attached.selection, (selection) => ({ ...selection, modeID }))
yield* ACPPromise.promise(() => input.client.session.switchAgent({ sessionID: attached.id, agent: modeID }))
})
const getSession = Effect.fnUntraced(function* (sessionID: string, cwd: string) {
const session = yield* ACPPromise.promise(() =>
input.client.session.get({ sessionID }).catch((error) => {
if (isSessionNotFoundError(error)) throw new ACPError.SessionNotFoundError({ sessionId: sessionID })
throw error
await catalogs.reload(state.catalog)
return select()
})
}),
)
if (FSUtil.resolve(cwd) !== FSUtil.resolve(session.location.directory))
return yield* new ACPError.SessionDirectoryMismatchError({ sessionId: sessionID, cwd })
return session
})
const requireSession = async (sessionID: string) => {
const current = sessions.get(sessionID)
if (current) return current
throw new ACPError.SessionNotFoundError({ sessionId: sessionID })
}
const detach = (sessionID: string) => {
sessions.get(sessionID)?.abort.abort()
sessions.delete(sessionID)
registeredMcp.delete(sessionID)
}
const cancelTurn = (sessionID: string) => {
const turn = active.get(sessionID)
if (turn) {
turn.control.cancelled = true
turn.control.admission.abort()
}
return input.client.session.interrupt({ sessionID })
}
const attach = async (session: SessionInfo, cwd: string, mcpServers: readonly McpServer[]) => {
const catalog = await catalogs.get(cwd)
sessions.get(session.id)?.abort.abort()
const state: Attached = {
id: session.id,
cwd,
abort: new AbortController(),
catalog,
model: session.model,
modeID: session.agent,
}
sessions.set(session.id, state)
await registerMcpServers(input.client, registeredMcp, state, mcpServers)
await sendCommands(state)
return state
}
const replay = async (state: Attached) => {
await replayMessages(input.connection, state.id, state.cwd, await messages(input.client, state.id))
}
const configOptions = (state: Attached, catalog = state.catalog.current) => {
const model = currentModel(state, catalog)
return buildConfigOptions({
providers: catalog.providers,
currentModel: { providerID: model.providerID, modelID: model.id },
currentVariant: model.variant,
modes: catalog.modes,
currentModeId: state.modeID ?? catalog.defaultModeID,
})
}
const replay = (attached: Attached) =>
Stream.paginate(undefined, (cursor: string | undefined) =>
ACPPromise.promise(() =>
cursor
? input.client.message.list({ sessionID: attached.id, limit: 200, cursor })
: input.client.message.list({ sessionID: attached.id, limit: 200, order: "asc" }),
).pipe(Effect.map((page) => [page.data, Option.fromNullishOr(page.cursor.next)] as const)),
).pipe(
Stream.runCollect,
Effect.flatMap((messages) =>
ACPPromise.promise(() => replayMessages(input.connection, attached.id, attached.cwd, messages)),
),
)
return {
initialize: async (params) => {
capabilities.childSessionUpdates = params.clientCapabilities?._meta?.[ChildSessionUpdatesCapability] === true
initialize: Effect.fnUntraced(function* (params) {
yield* Ref.set(input.capabilities, {
childSessionUpdates: params.clientCapabilities?._meta?.[ChildSessionUpdatesCapability] === true,
})
const authMethod: AuthMethod = {
description: "Run `opencode auth login` in the terminal",
name: "Login with opencode",
@@ -233,32 +164,37 @@ export function make(input: {
authMethods: [authMethod],
agentInfo: { name: "OpenCode", version: OPENCODE_VERSION },
}
},
authenticate: async (params) => {
if (params.methodId !== AuthMethodID) throw new ACPError.UnknownAuthMethodError({ methodId: params.methodId })
}),
authenticate: Effect.fnUntraced(function* (params) {
if (params.methodId !== AuthMethodID)
return yield* new ACPError.UnknownAuthMethodError({ methodId: params.methodId })
return {}
},
newSession: async (params) => {
}),
newSession: Effect.fnUntraced(function* (params) {
// Load before creating so a catalog failure leaves no session behind. Agent and model stay unset
// so the server resolves its defaults after plugins activate.
await catalogs.get(params.cwd)
const created = await input.client.session.create({ location: { directory: params.cwd } })
const state = await attach(created, params.cwd, params.mcpServers)
return { sessionId: state.id, configOptions: configOptions(state) }
},
loadSession: async (params) => {
const session = await getSession(input.client, params.sessionId, params.cwd)
const state = await attach(session, session.location.directory, params.mcpServers)
await replay(state)
return { configOptions: configOptions(state) }
},
listSessions: async (params) => {
const page = await input.client.session.list({
...(params.cwd ? { directory: params.cwd } : {}),
order: "desc",
limit: 100,
...(params.cursor ? { cursor: params.cursor } : {}),
})
yield* input.catalog.get(params.cwd)
const created = yield* ACPPromise.promise(() =>
input.client.session.create({ location: { directory: params.cwd } }),
)
const attached = yield* input.sessions.attach(created, params.cwd, params.mcpServers)
return { sessionId: attached.id, configOptions: yield* currentOptions(attached) }
}),
loadSession: Effect.fnUntraced(function* (params) {
const session = yield* getSession(params.sessionId, params.cwd)
const attached = yield* input.sessions.attach(session, session.location.directory, params.mcpServers)
yield* replay(attached)
return { configOptions: yield* currentOptions(attached) }
}),
listSessions: Effect.fnUntraced(function* (params) {
const page = yield* ACPPromise.promise(() =>
input.client.session.list({
...(params.cwd ? { directory: params.cwd } : {}),
order: "desc",
limit: 100,
...(params.cursor ? { cursor: params.cursor } : {}),
}),
)
return {
sessions: page.data.map((session) => ({
sessionId: session.id,
@@ -268,178 +204,57 @@ export function make(input: {
})),
...(page.cursor.next ? { nextCursor: page.cursor.next } : {}),
}
},
deleteSession: async (params) => {
await input.client.session.remove({ sessionID: params.sessionId }).catch((error) => {
if (!isSessionNotFoundError(error)) throw error
})
detach(params.sessionId)
}),
deleteSession: Effect.fnUntraced(function* (params) {
yield* ACPPromise.promise(() =>
input.client.session.remove({ sessionID: params.sessionId }).catch((error) => {
if (!isSessionNotFoundError(error)) throw error
}),
)
yield* input.sessions.detach(params.sessionId)
return {}
},
resumeSession: async (params) => {
const session = await getSession(input.client, params.sessionId, params.cwd)
const state = await attach(session, session.location.directory, params.mcpServers ?? [])
return { configOptions: configOptions(state) }
},
closeSession: async (params) => {
const turn = active.get(params.sessionId)
await cancelTurn(params.sessionId).catch((error) => {
if (!isSessionNotFoundError(error)) throw error
})
await turn?.turn.catch(() => {})
detach(params.sessionId)
}),
resumeSession: Effect.fnUntraced(function* (params) {
const session = yield* getSession(params.sessionId, params.cwd)
const attached = yield* input.sessions.attach(session, session.location.directory, params.mcpServers ?? [])
return { configOptions: yield* currentOptions(attached) }
}),
closeSession: Effect.fnUntraced(function* (params) {
yield* input.turn.close(params.sessionId)
yield* input.sessions.detach(params.sessionId)
return {}
},
forkSession: async (params) => {
const forked = await input.client.session.fork({
sessionID: params.sessionId,
})
const state = await attach(forked, forked.location.directory, params.mcpServers ?? [])
await replay(state)
return { sessionId: state.id, configOptions: configOptions(state) }
},
setSessionConfigOption: async (params) => {
const state = await requireSession(params.sessionId)
}),
forkSession: Effect.fnUntraced(function* (params) {
const forked = yield* ACPPromise.promise(() => input.client.session.fork({ sessionID: params.sessionId }))
const attached = yield* input.sessions.attach(forked, forked.location.directory, params.mcpServers ?? [])
yield* replay(attached)
return { sessionId: attached.id, configOptions: yield* currentOptions(attached) }
}),
setSessionConfigOption: Effect.fnUntraced(function* (params) {
const attached = yield* input.sessions.require(params.sessionId)
const value = params.value
if (typeof value !== "string") throw new ACPError.InvalidConfigOptionError({ configId: params.configId })
await withReload(state, async () => {
switch (params.configId) {
case "model": {
const selected = requireModel(state.catalog.current, value, currentModel(state))
state.model = selected
await input.client.session.switchModel({ sessionID: state.id, model: selected })
return
}
case "effort": {
const current = currentModel(state)
const model = state.catalog.current.models.find(
(item) => item.providerID === current.providerID && item.id === current.id,
)
if (!model || (value !== DEFAULT_VARIANT_VALUE && !model.variants.some((variant) => variant.id === value)))
throw new ACPError.InvalidEffortError({ effort: value })
state.model = { ...current, variant: value }
await input.client.session.switchModel({ sessionID: state.id, model: state.model })
return
}
case "mode":
return selectMode(input.client, state, value)
default:
throw new ACPError.InvalidConfigOptionError({ configId: params.configId })
}
})
return { configOptions: configOptions(state) }
},
setSessionMode: async (params) => {
const state = await requireSession(params.sessionId)
await withReload(state, () => selectMode(input.client, state, params.modeId))
if (typeof value !== "string") return yield* new ACPError.InvalidConfigOptionError({ configId: params.configId })
yield* withReload(attached, selectOption(attached, params.configId, value))
return { configOptions: yield* currentOptions(attached) }
}),
setSessionMode: Effect.fnUntraced(function* (params) {
const attached = yield* input.sessions.require(params.sessionId)
yield* withReload(attached, selectMode(attached, params.modeId))
return {}
},
prompt: async (params, signal) => {
const state = await requireSession(params.sessionId)
if (active.has(state.id)) {
throw new ACPError.ServiceFailureError({
safeMessage: `Session already has an active ACP prompt: ${state.id}`,
service: "session",
})
}
const messageID = SessionMessage.ID.create()
const prepared = preparePrompt(state.catalog.current, params.prompt, messageID)
const control: TurnControl = { cancelled: false, admission: new AbortController() }
const extNotification = input.connection.extNotification
const childSessionUpdate =
capabilities.childSessionUpdates && extNotification
? (update: ChildSessionUpdate) => extNotification(ChildSessionUpdateMethod, update).then(() => {})
: undefined
// A `$/cancel_request` for this prompt behaves like `session/cancel` for its turn.
const cancel = () => void cancelTurn(state.id).catch(() => {})
const turn = streamTurn({
client: input.client,
connection: input.connection,
sessionID: state.id,
cwd: state.cwd,
start: prepared.start,
action: prepared.command !== undefined,
control,
connectionSignal: input.connection.signal,
sessionSignal: state.abort.signal,
submit: (signal) => submitPrompt(input.client, state, prepared, signal),
...(childSessionUpdate ? { childSessionUpdate } : {}),
})
.then(async (result) => {
await sendUsageUpdate(input.client, input.connection, state, result.contextTokens).catch(() => {})
return result.response
})
.finally(() => {
signal?.removeEventListener("abort", cancel)
if (active.get(state.id)?.control === control) active.delete(state.id)
})
active.set(state.id, { control, turn })
signal?.addEventListener("abort", cancel, { once: true })
// The cancel may already be buffered behind the awaits above.
if (signal?.aborted) cancel()
return turn
},
cancel: async (params) => {
await cancelTurn(params.sessionId).catch(() => {})
},
}),
prompt: input.turn.prompt,
cancel: input.turn.cancel,
}
}
function preparePrompt(catalog: Catalog, prompt: PromptRequest["prompt"], messageID: string): PreparedPrompt {
const parts = promptContentToParts(prompt)
const visible = parts.filter((part) => part.type !== "text" || (!part.synthetic && !part.ignored))
const synthetic = parts.flatMap((part) => (part.type === "text" && part.synthetic ? [part.text] : []))
const text = visible.flatMap((part) => (part.type === "text" ? [part.text] : [])).join("\n")
const files = visible.flatMap((part) => (part.type === "file" ? [{ uri: part.url, name: part.filename }] : []))
const slash = detectSlashCommand(text)
const command =
slash && !builtinCommands.has(slash.name) ? catalog.commands.find((item) => item.name === slash.name) : undefined
const start = turnStart(messageID, slash)
return { start, text, files, synthetic, slash, command }
}
async function submitPrompt(client: OpenCodeClient, session: Attached, prompt: PreparedPrompt, signal: AbortSignal) {
if (prompt.synthetic.length > 0) {
await client.session.synthetic({
sessionID: session.id,
text: prompt.synthetic.join("\n\n"),
description: "ACP embedded context",
delivery: "steer",
resume: false,
})
}
if (prompt.start.type === "compaction") return client.session.compact({ sessionID: session.id, id: prompt.start.id })
if (prompt.command) {
return client.session.command(
{
sessionID: session.id,
name: prompt.command.name,
text: prompt.slash?.args ?? "",
files: prompt.files,
delivery: "steer",
},
{ signal },
)
}
return client.session.prompt(
{ sessionID: session.id, id: prompt.start.id, text: prompt.text, files: prompt.files, delivery: "steer" },
{ signal },
)
}
function turnStart(messageID: string, slash: PreparedPrompt["slash"]): TurnStart {
if (slash && builtinCommands.get(slash.name)?.start === "compaction") return { type: "compaction", id: messageID }
return { type: "input", id: messageID }
}
function requireModel(catalog: Catalog, modelID: string, current: ModelRef): ModelRef {
const requireModel = Effect.fnUntraced(function* (catalog: Catalog, modelID: string, current: ModelRef) {
const selected = parseModelSelection(modelID, catalog.providers)
const model = catalog.models.find(
(item) => item.providerID === selected.model.providerID && item.id === selected.model.modelID,
)
if (!model) throw new ACPError.InvalidModelError({ providerId: selected.model.providerID, modelId: modelID })
if (!model) return yield* new ACPError.InvalidModelError({ providerId: selected.model.providerID, modelId: modelID })
if (selected.variant && !model.variants.some((variant) => variant.id === selected.variant))
throw new ACPError.InvalidEffortError({ effort: selected.variant })
return yield* new ACPError.InvalidEffortError({ effort: selected.variant })
const variant =
selected.variant ??
(current.providerID === model.providerID &&
@@ -447,124 +262,14 @@ function requireModel(catalog: Catalog, modelID: string, current: ModelRef): Mod
(current.variant === DEFAULT_VARIANT_VALUE || model.variants.some((variant) => variant.id === current.variant))
? current.variant
: undefined)
return { providerID: model.providerID, id: model.id, variant }
}
return { providerID: model.providerID, id: model.id, variant } satisfies ModelRef
})
function currentModel(state: Attached, catalog = state.catalog.current) {
return state.model ?? catalog.defaultModel
}
async function selectMode(client: OpenCodeClient, state: Attached, modeID: string) {
if (!state.catalog.current.modes.some((mode) => mode.id === modeID))
throw new ACPError.InvalidModeError({ mode: modeID })
state.modeID = modeID
await client.session.switchAgent({ sessionID: state.id, agent: modeID })
}
async function getSession(client: OpenCodeClient, sessionID: string, cwd: string) {
const session = await client.session.get({ sessionID }).catch((error) => {
if (isSessionNotFoundError(error)) throw new ACPError.SessionNotFoundError({ sessionId: sessionID })
throw error
})
if (FSUtil.resolve(cwd) !== FSUtil.resolve(session.location.directory)) {
throw new ACPError.SessionDirectoryMismatchError({ sessionId: sessionID, cwd })
}
return session
}
async function messages(client: OpenCodeClient, sessionID: string) {
const result: SessionMessageInfo[] = []
let cursor: string | undefined
do {
const page = cursor
? await client.message.list({ sessionID, limit: 200, cursor })
: await client.message.list({ sessionID, limit: 200, order: "asc" })
result.push(...page.data)
cursor = page.cursor.next ?? undefined
} while (cursor)
return result
}
async function registerMcpServers(
client: OpenCodeClient,
registered: Map<string, Set<string>>,
session: Attached,
servers: readonly McpServer[],
) {
const current = registered.get(session.id) ?? new Set<string>()
registered.set(session.id, current)
await Promise.all(
servers.flatMap((server) => {
const config = mcpConfig(server)
const key = `${server.name}:${stableStringify(config)}`
if (current.has(key)) return []
current.add(key)
return [
client.mcp.add({ server: server.name, location: { directory: session.cwd }, config }).catch((error) => {
current.delete(key)
throw error
}),
]
}),
)
}
function mcpConfig(server: McpServer) {
if ("type" in server) {
if (server.type === "acp") throw new Error("MCP-over-ACP is not supported")
return {
type: "remote" as const,
url: server.url,
headers: Object.fromEntries(server.headers.map((header) => [header.name, header.value])),
oauth: false as const,
}
}
return {
type: "local" as const,
command: [server.command, ...server.args],
environment: Object.fromEntries(server.env.map((entry) => [entry.name, entry.value])),
}
}
function stableStringify(value: unknown): string {
if (Array.isArray(value)) return `[${value.map(stableStringify).join(",")}]`
if (!value || typeof value !== "object") return JSON.stringify(value)
return `{${Object.entries(value)
.toSorted(([a], [b]) => a.localeCompare(b))
.map(([key, item]) => `${JSON.stringify(key)}:${stableStringify(item)}`)
.join(",")}}`
}
async function sendUsageUpdate(
client: OpenCodeClient,
connection: ACPConnection.Connection,
session: Attached,
used?: number,
) {
if (!used) return
const current = currentModel(session)
const model = session.catalog.current.models.find(
(item) => item.providerID === current.providerID && item.id === current.id,
)
if (!model?.limit.context) return
const info = await client.session.get({ sessionID: session.id })
await connection.sessionUpdate({
sessionId: session.id,
update: {
sessionUpdate: "usage_update",
used,
size: model.limit.context,
cost: { amount: info.cost, currency: "USD" },
},
})
}
function detectSlashCommand(text: string): { readonly name: string; readonly args: string } | undefined {
const value = text.trim()
if (!value.startsWith("/")) return undefined
const [name, ...rest] = value.slice(1).split(/\s+/)
if (!name) return undefined
return { name, args: rest.join(" ").trim() }
}
const requireEffort = Effect.fnUntraced(function* (catalog: Catalog, effort: string, current: ModelRef) {
const model = catalog.models.find((item) => item.providerID === current.providerID && item.id === current.id)
if (!model || (effort !== DEFAULT_VARIANT_VALUE && !model.variants.some((variant) => variant.id === effort)))
return yield* new ACPError.InvalidEffortError({ effort })
return { ...current, variant: effort } satisfies ModelRef
})
export * as ACPService from "./service"
+176
View File
@@ -0,0 +1,176 @@
import { isDeepStrictEqual } from "node:util"
import type { McpServer, RequestError } from "@agentclientprotocol/sdk"
import type { OpenCodeClient, SessionInfo } from "@opencode/client/promise"
import { Context, Effect, Exit, Ref, Scope, Stream } from "effect"
import type { ACPCatalog, Catalog } from "./catalog"
import { availableCommands, configOptions, type Selection } from "./config-option"
import type { ACPConnection } from "./connection"
import { ACPError } from "./error"
import { ACPPromise } from "./promise"
export type Attached = {
readonly id: string
readonly cwd: string
readonly selection: Ref.Ref<Selection>
/** Aborted when the session detaches, for the promise-based turn. */
readonly signal: AbortSignal
}
export interface Interface {
/**
* Attaches a session in its own scope, closing any previous attachment of the same ID. The scope follows the
* cwd's catalog and pushes config option and command updates while it is open. A failed attach leaves the
* session detached.
*/
readonly attach: (
session: SessionInfo,
cwd: string,
mcpServers: readonly McpServer[],
) => Effect.Effect<Attached, ACPError.Error | RequestError | ACPCatalog.Error>
/** Closes the session scope. No-op when the session is not attached. */
readonly detach: (sessionID: string) => Effect.Effect<void>
readonly require: (sessionID: string) => Effect.Effect<Attached, ACPError.SessionNotFoundError>
}
export class Service extends Context.Service<Service, Interface>()("@opencode/cli/acp/Sessions") {}
type Entry = { readonly attached: Attached; readonly scope: Scope.Closeable }
export const make = Effect.fnUntraced(function* (input: {
readonly client: OpenCodeClient
readonly connection: ACPConnection.Interface
readonly catalog: ACPCatalog.Interface
}) {
const scope = yield* Effect.scope
const sessions = new Map<string, Entry>()
// Kept across re-attachment so resuming with the same servers does not add them again.
const registeredMcp = new Map<string, Set<string>>()
const sendCommands = (sessionID: string, catalog: Catalog) =>
input.connection.sessionUpdate({
sessionId: sessionID,
update: { sessionUpdate: "available_commands_update", availableCommands: availableCommands(catalog) },
})
const changed = Effect.fnUntraced(function* (attached: Attached, previous: Catalog, next: Catalog) {
const selection = yield* Ref.get(attached.selection)
const options = configOptions(next, selection)
if (!isDeepStrictEqual(options, configOptions(previous, selection))) {
yield* input.connection.sessionUpdate({
sessionId: attached.id,
update: { sessionUpdate: "config_option_update", configOptions: options },
})
}
if (!isDeepStrictEqual(next.commands, previous.commands)) yield* sendCommands(attached.id, next)
})
const registerMcp = (attached: Attached, servers: readonly McpServer[]) =>
Effect.suspend(() => {
const registered = registeredMcp.get(attached.id) ?? new Set<string>()
registeredMcp.set(attached.id, registered)
return Effect.forEach(
servers,
(server) =>
Effect.suspend(() => {
const config = mcpConfig(server)
const key = `${server.name}:${stableStringify(config)}`
if (registered.has(key)) return Effect.void
registered.add(key)
return ACPPromise.promise(() =>
input.client.mcp.add({ server: server.name, location: { directory: attached.cwd }, config }),
).pipe(
Effect.onError(() => Effect.sync(() => registered.delete(key))),
Effect.uninterruptible,
)
}),
{ concurrency: "unbounded", discard: true },
)
})
const remove = (sessionID: string, entry: Entry) =>
Effect.suspend(() => {
if (sessions.get(sessionID) === entry) {
sessions.delete(sessionID)
registeredMcp.delete(sessionID)
}
return Scope.close(entry.scope, Exit.void)
})
return Service.of({
attach: Effect.fn("cli.acp.sessions.attach")(function* (session, cwd, mcpServers) {
const current = yield* input.catalog.get(cwd)
const abort = new AbortController()
const entry: Entry = {
attached: {
id: session.id,
cwd,
selection: yield* Ref.make<Selection>({ model: session.model, modeID: session.agent }),
signal: abort.signal,
},
scope: Scope.forkUnsafe(scope),
}
// Swap synchronously so concurrent attaches of one ID cannot both keep a scope.
const replaced = sessions.get(session.id)
sessions.set(session.id, entry)
if (replaced) yield* Scope.close(replaced.scope, Exit.void)
yield* Scope.addFinalizer(
entry.scope,
Effect.sync(() => abort.abort()),
)
yield* registerMcp(entry.attached, mcpServers).pipe(
Effect.andThen(sendCommands(session.id, current)),
Effect.onError(() => remove(session.id, entry)),
)
// `changes` emits the latest catalog first, so a reload since `current` is still pushed.
yield* input.catalog.changes(cwd).pipe(
Stream.runFoldEffect(
() => current,
(previous, next) =>
next === previous
? Effect.succeed(previous)
: changed(entry.attached, previous, next).pipe(Effect.ignore, Effect.as(next)),
),
Effect.ignore,
Effect.forkIn(entry.scope),
)
return entry.attached
}),
detach: Effect.fn("cli.acp.sessions.detach")(function* (sessionID) {
const entry = sessions.get(sessionID)
if (entry) yield* remove(sessionID, entry)
}),
require: Effect.fn("cli.acp.sessions.require")(function* (sessionID) {
const entry = sessions.get(sessionID)
if (!entry) return yield* new ACPError.SessionNotFoundError({ sessionId: sessionID })
return entry.attached
}),
})
})
function mcpConfig(server: McpServer) {
if ("type" in server) {
if (server.type === "acp") throw new Error("MCP-over-ACP is not supported")
return {
type: "remote" as const,
url: server.url,
headers: Object.fromEntries(server.headers.map((header) => [header.name, header.value])),
oauth: false as const,
}
}
return {
type: "local" as const,
command: [server.command, ...server.args],
environment: Object.fromEntries(server.env.map((entry) => [entry.name, entry.value])),
}
}
function stableStringify(value: unknown): string {
if (Array.isArray(value)) return `[${value.map(stableStringify).join(",")}]`
if (!value || typeof value !== "object") return JSON.stringify(value)
return `{${Object.entries(value)
.toSorted(([a], [b]) => a.localeCompare(b))
.map(([key, item]) => `${JSON.stringify(key)}:${stableStringify(item)}`)
.join(",")}}`
}
export * as ACPSessions from "./sessions"
+203
View File
@@ -0,0 +1,203 @@
import type { CancelNotification, PromptRequest, PromptResponse, RequestError } from "@agentclientprotocol/sdk"
import { isSessionNotFoundError, type CommandInfo, type OpenCodeClient } from "@opencode/client/promise"
import { SessionMessage } from "@opencode/schema/session-message"
import { Effect, Ref, type Scope } from "effect"
import { builtinCommands, type ACPCatalog, type Catalog } from "./catalog"
import { currentModel } from "./config-option"
import type { ACPConnection } from "./connection"
import { promptContentToParts } from "./content"
import { ACPError } from "./error"
import {
ChildSessionUpdateMethod,
streamTurn,
type ChildSessionUpdate,
type TurnControl,
type TurnStart,
} from "./event"
import { ACPPromise } from "./promise"
import type { ACPSessions, Attached } from "./sessions"
type PreparedPrompt = {
readonly start: TurnStart
readonly text: string
readonly files: Array<{ readonly uri: string; readonly name?: string }>
readonly synthetic: ReadonlyArray<string>
readonly slash?: { readonly name: string; readonly args: string }
readonly command?: CommandInfo
}
export interface Interface {
prompt(input: PromptRequest, signal?: AbortSignal): Promise<PromptResponse>
cancel(input: CancelNotification): Promise<void>
/** Cancels the session's active turn and waits for it to settle. */
readonly close: (sessionID: string) => Effect.Effect<void, ACPError.Error | RequestError>
}
export function make(input: {
readonly client: OpenCodeClient
readonly connection: ACPConnection.Connection
readonly sessions: ACPSessions.Interface
readonly catalog: ACPCatalog.Interface
readonly capabilities: Ref.Ref<{ readonly childSessionUpdates: boolean }>
readonly run: <A, E>(effect: Effect.Effect<A, E, Scope.Scope>) => Promise<A>
}): Interface {
const active = new Map<string, { readonly control: TurnControl; readonly turn: Promise<PromptResponse> }>()
const cancelTurn = (sessionID: string) => {
const turn = active.get(sessionID)
if (turn) {
turn.control.cancelled = true
turn.control.admission.abort()
}
return input.client.session.interrupt({ sessionID })
}
const sendUsageUpdate = async (state: Attached, used?: number) => {
if (!used) return
const model = await input.run(
Effect.gen(function* () {
const catalog = yield* input.catalog.get(state.cwd)
const current = currentModel(catalog, yield* Ref.get(state.selection))
return catalog.models.find((item) => item.providerID === current.providerID && item.id === current.id)
}),
)
if (!model?.limit.context) return
const info = await input.client.session.get({ sessionID: state.id })
await input.connection.sessionUpdate({
sessionId: state.id,
update: {
sessionUpdate: "usage_update",
used,
size: model.limit.context,
cost: { amount: info.cost, currency: "USD" },
},
})
}
return {
prompt: async (params, signal) => {
// Read everything first so the active check and registration below stay synchronous.
const resolved = await input.run(
Effect.gen(function* () {
const attached = yield* input.sessions.require(params.sessionId)
return {
attached,
catalog: yield* input.catalog.get(attached.cwd),
childSessionUpdates: (yield* Ref.get(input.capabilities)).childSessionUpdates,
}
}),
)
const state = resolved.attached
if (active.has(state.id)) {
throw new ACPError.ServiceFailureError({
safeMessage: `Session already has an active ACP prompt: ${state.id}`,
service: "session",
})
}
const messageID = SessionMessage.ID.create()
const prepared = preparePrompt(resolved.catalog, params.prompt, messageID)
const control: TurnControl = { cancelled: false, admission: new AbortController() }
const extNotification = input.connection.extNotification
const childSessionUpdate =
resolved.childSessionUpdates && extNotification
? (update: ChildSessionUpdate) => extNotification(ChildSessionUpdateMethod, update).then(() => {})
: undefined
// A `$/cancel_request` for this prompt behaves like `session/cancel` for its turn.
const cancel = () => void cancelTurn(state.id).catch(() => {})
const turn = streamTurn({
client: input.client,
connection: input.connection,
sessionID: state.id,
cwd: state.cwd,
start: prepared.start,
action: prepared.command !== undefined,
control,
connectionSignal: input.connection.signal,
sessionSignal: state.signal,
submit: (signal) => submitPrompt(input.client, state, prepared, signal),
...(childSessionUpdate ? { childSessionUpdate } : {}),
})
.then(async (result) => {
await sendUsageUpdate(state, result.contextTokens).catch(() => {})
return result.response
})
.finally(() => {
signal?.removeEventListener("abort", cancel)
if (active.get(state.id)?.control === control) active.delete(state.id)
})
active.set(state.id, { control, turn })
signal?.addEventListener("abort", cancel, { once: true })
// The cancel may already be buffered behind the awaits above.
if (signal?.aborted) cancel()
return turn
},
cancel: async (params) => {
await cancelTurn(params.sessionId).catch(() => {})
},
close: Effect.fnUntraced(function* (sessionID) {
const turn = active.get(sessionID)
yield* ACPPromise.promise(() =>
cancelTurn(sessionID).catch((error) => {
if (!isSessionNotFoundError(error)) throw error
}),
)
if (turn) yield* Effect.promise(() => turn.turn.catch(() => {}))
}),
}
}
function preparePrompt(catalog: Catalog, prompt: PromptRequest["prompt"], messageID: string): PreparedPrompt {
const parts = promptContentToParts(prompt)
const visible = parts.filter((part) => part.type !== "text" || (!part.synthetic && !part.ignored))
const synthetic = parts.flatMap((part) => (part.type === "text" && part.synthetic ? [part.text] : []))
const text = visible.flatMap((part) => (part.type === "text" ? [part.text] : [])).join("\n")
const files = visible.flatMap((part) => (part.type === "file" ? [{ uri: part.url, name: part.filename }] : []))
const slash = detectSlashCommand(text)
const command = slash ? catalog.commands.find((item) => item.name === slash.name) : undefined
const start = turnStart(messageID, slash)
return { start, text, files, synthetic, slash, command }
}
async function submitPrompt(client: OpenCodeClient, session: Attached, prompt: PreparedPrompt, signal: AbortSignal) {
if (prompt.synthetic.length > 0) {
await client.session.synthetic({
sessionID: session.id,
text: prompt.synthetic.join("\n\n"),
description: "ACP embedded context",
delivery: "steer",
resume: false,
})
}
if (prompt.start.type === "compaction") return client.session.compact({ sessionID: session.id, id: prompt.start.id })
if (prompt.command) {
return client.session.command(
{
sessionID: session.id,
name: prompt.command.name,
text: prompt.slash?.args ?? "",
files: prompt.files,
delivery: "steer",
},
{ signal },
)
}
return client.session.prompt(
{ sessionID: session.id, id: prompt.start.id, text: prompt.text, files: prompt.files, delivery: "steer" },
{ signal },
)
}
function turnStart(messageID: string, slash: PreparedPrompt["slash"]): TurnStart {
if (slash && builtinCommands.get(slash.name)?.start === "compaction") return { type: "compaction", id: messageID }
return { type: "input", id: messageID }
}
function detectSlashCommand(text: string): { readonly name: string; readonly args: string } | undefined {
const value = text.trim()
if (!value.startsWith("/")) return undefined
const [name, ...rest] = value.slice(1).split(/\s+/)
if (!name) return undefined
return { name, args: rest.join(" ").trim() }
}
export * as ACPTurn from "./turn"
+87 -13
View File
@@ -118,24 +118,71 @@ describe("acp catalog and config options over the wire", () => {
const session = await acp.newSession()
acp.server.catalog.models = [testModel, secondModel]
acp.server.catalog.agents = [buildAgent, planAgent, configured]
const set = (configId: string, value: string) =>
acp.request("session/set_config_option", { sessionId: session.sessionId, configId, value })
const initialModelReads = modelReads(acp)
const model = await acp.request("session/set_config_option", {
sessionId: session.sessionId,
configId: "model",
value: "test/second-model",
})
const model = await set("model", "test/second-model")
const reloadedModelReads = modelReads(acp)
const missingModel = await rpcError(set("model", "test/missing-model"))
const missingModelReads = modelReads(acp)
await acp.request("session/set_mode", { sessionId: session.sessionId, modeId: "copilot-build" })
const reads = agentReads(acp)
const missing = await rpcError(
acp.request("session/set_config_option", { sessionId: session.sessionId, configId: "mode", value: "missing" }),
)
const missing = await rpcError(set("mode", "missing"))
expect(currentValue(model, "model")).toBe("test/second-model")
expect([reloadedModelReads, missingModelReads]).toEqual([initialModelReads + 1, initialModelReads + 2])
expect(missingModel).toMatchObject({ code: -32602, data: { modelId: "test/missing-model" } })
expect(acp.server.selections).toContainEqual({ sessionID: session.sessionId, agent: "copilot-build" })
expect(missing).toMatchObject({ code: -32602, data: { mode: "missing" } })
expect(agentReads(acp)).toBeGreaterThan(reads)
})
test.each([
[
"a sibling session closes",
async (acp: Wire) => {
const closed = await acp.newSession()
const open = await acp.newSession()
await acp.request("session/close", { sessionId: closed.sessionId })
return open.sessionId
},
],
...(["session/load", "session/resume"] as const).map(
(method) =>
[
`${method} re-attaches the session`,
async (acp: Wire) => {
const session = await acp.newSession()
const params = { cwd: "/workspace", sessionId: session.sessionId, mcpServers: [] }
await acp.request(method, params)
await acp.request(method, params)
return session.sessionId
},
] as const,
),
])("pushes exactly one update per catalog change after %s", async (_, setup) => {
await using acp = await startWire()
acp.server.catalog.models = [testModel]
await acp.initialize()
const sessionId = await setup(acp)
const since = acp.updates.length
await change(acp, sessionId, "config_option_update", () => {
acp.server.catalog.models = [testModel, secondModel]
acp.server.send(ephemeralEvent("model.updated", {}))
})
await change(acp, sessionId, "available_commands_update", () => {
acp.server.catalog.commands = [reviewCommand, { name: "ship", description: "Ship it" }]
acp.server.send(ephemeralEvent("command.updated", {}, { directory: "/workspace" }))
})
expect(updateKinds(acp, since)).toEqual([
[sessionId, "config_option_update"],
[sessionId, "available_commands_update"],
])
})
test.each(["empty", "missing the default"])(
"retries when the model list is %s but the default is ready",
async (initial) => {
@@ -212,25 +259,52 @@ describe("acp catalog and config options over the wire", () => {
await using acp = await startSession()
const advertised = await acp.waitForUpdate((item) => commandNames(item) !== undefined)
acp.server.catalog.commands = [reviewCommand, { name: "compact", description: "Server compact" }]
acp.server.catalog.commands = [
reviewCommand,
{ name: "compact", description: "Server compact" },
{ name: "ship", description: "Ship it" },
]
acp.server.send(ephemeralEvent("command.updated", {}, { directory: "/workspace" }))
const replaced = await acp.waitForUpdate((item) => item !== advertised && commandNames(item) !== undefined)
const compacted = await acp.prompt(acp.sessionId, "/compact")
expect([advertised, replaced].map((item) => item.update)).toEqual(
Array.from({ length: 2 }, () => ({
expect([advertised, replaced].map((item) => item.update)).toEqual([
{
sessionUpdate: "available_commands_update",
availableCommands: [
{ name: "review", description: "Review changes" },
{ name: "compact", description: "Compact the session" },
],
})),
)
},
{
sessionUpdate: "available_commands_update",
availableCommands: [
{ name: "review", description: "Review changes" },
{ name: "ship", description: "Ship it" },
{ name: "compact", description: "Compact the session" },
],
},
])
expect(compacted.stopReason).toBe("end_turn")
expect(acp.server.submissions.map((item) => item.kind)).toEqual(["compact"])
})
})
// Each change waits on a catalog reload over HTTP, so a stray update for an earlier change lands before the next one.
async function change(acp: Wire, sessionId: string, kind: string, trigger: () => void) {
const seen = acp.updates.filter((item) => item.sessionId === sessionId && item.update.sessionUpdate === kind).length
trigger()
await acp.until(
() =>
acp.updates.filter((item) => item.sessionId === sessionId && item.update.sessionUpdate === kind).length > seen,
kind,
)
}
function updateKinds(acp: Wire, since: number) {
return acp.updates.slice(since).map((item) => [item.sessionId, item.update.sessionUpdate])
}
function commandNames(item: SessionNotification) {
if (item.update.sessionUpdate !== "available_commands_update") return undefined
return item.update.availableCommands.map((command) => command.name)
+30 -1
View File
@@ -1,7 +1,7 @@
import { describe, expect, test } from "bun:test"
import type { McpServer } from "@agentclientprotocol/sdk"
import { currentValue } from "./select-options"
import { makeSession, rpcError, secondModel, startSession, startWire } from "./wire-fixture"
import { ephemeralEvent, makeSession, rpcError, secondModel, startSession, startWire, testModel } from "./wire-fixture"
describe("acp session lifecycle over the wire", () => {
test("initialize advertises capabilities and terminal auth only when the client asks", async () => {
@@ -226,4 +226,33 @@ describe("acp session lifecycle over the wire", () => {
config: { type: "remote", url: "https://example.com/mcp", headers: { Authorization: "Bearer x" }, oauth: false },
})
})
test("leaves a session detached when re-attaching it fails", async () => {
const broken: McpServer = { name: "broken", command: "bun", args: [], env: [] }
await using acp = await startWire({
fetch: (request) =>
request.method === "PUT" && request.path === "/api/experimental/mcp/broken"
? new Response(null, { status: 500 })
: undefined,
})
await acp.initialize()
const failed = await acp.newSession()
const other = await acp.newSession()
expect(
await rpcError(
acp.request("session/resume", { cwd: "/workspace", sessionId: failed.sessionId, mcpServers: [broken] }),
),
).toMatchObject({ code: -32603 })
const since = acp.updates.length
acp.server.catalog.models = [testModel]
acp.server.send(ephemeralEvent("model.updated", {}))
await acp.until(() => acp.updates.length > since, "config options for the attached session")
expect(acp.updates.slice(since).map((item) => item.sessionId)).toEqual([other.sessionId])
expect(
await rpcError(
acp.request("session/set_config_option", { sessionId: failed.sessionId, configId: "mode", value: "plan" }),
),
).toMatchObject({ code: -32602, data: { sessionId: failed.sessionId } })
})
})