fix(plugin): derive promise adapter from protocol schemas (#42669)

This commit is contained in:
Kit Langton 2026-08-14 20:52:06 -04:00 committed by GitHub
parent 014a364dfd
commit 8afcb3870e
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
7 changed files with 421 additions and 395 deletions

View file

@ -0,0 +1,5 @@
---
"@opencode-ai/plugin": patch
---
Derive Promise plugin API request and response conversion from the canonical protocol schemas.

View file

@ -568,6 +568,7 @@
"@ai-sdk/provider": "3.0.8",
"@opencode-ai/ai": "workspace:*",
"@opencode-ai/client": "workspace:*",
"@opencode-ai/protocol": "workspace:*",
"@opencode-ai/schema": "workspace:*",
"@opencode-ai/sdk": "1.18.5",
"@standard-schema/spec": "catalog:",

View file

@ -1,396 +1,3 @@
export * as PluginPromise from "./promise.js"
import { define } from "@opencode-ai/plugin/effect/plugin"
import type { Context, Plugin } from "@opencode-ai/plugin/promise/plugin"
import type { Info } from "@opencode-ai/plugin/promise/tool"
import { Agent } from "@opencode-ai/schema/agent"
import { Integration } from "@opencode-ai/schema/integration"
import { Location } from "@opencode-ai/schema/location"
import { Model } from "@opencode-ai/schema/model"
import { Provider } from "@opencode-ai/schema/provider"
import { AbsolutePath } from "@opencode-ai/schema/schema"
import { Session } from "@opencode-ai/schema/session"
import { SessionMessage } from "@opencode-ai/schema/session-message"
import { Skill } from "@opencode-ai/schema/skill"
import { Workspace } from "@opencode-ai/schema/workspace"
import { WebSearch } from "@opencode-ai/schema/websearch"
import { DateTime, Effect, Scope, Stream } from "effect"
import { Tool } from "../tool.js"
type HostRegistration = { readonly dispose: Effect.Effect<void> }
type Registration = { readonly dispose: () => Promise<void> }
type PromiseEvent = ReturnType<Context["event"]["subscribe"]> extends AsyncIterable<infer Event> ? Event : never
type JsonValue = null | boolean | number | string | Array<JsonValue> | { [key: string]: JsonValue }
/**
* Adapts a Promise plugin into an Effect plugin so the existing Effect-only
* loader (`Plugin` / `PluginSupervisor`) can run it unchanged.
*
* Hook registrations created during the async `setup` attach to the plugin's
* scope, so unloading the plugin disposes them. The captured fiber context
* preserves boot-time batching, so Promise-plugin transforms still coalesce
* into one reload per domain.
*/
export function fromPromise(plugin: Plugin) {
return define({
id: plugin.id,
effect: (host) =>
Effect.gen(function* () {
const scope = yield* Scope.Scope
const context = yield* Effect.context<Scope.Scope>()
// Run a hook registration on the plugin scope and resolve once it is registered.
const register = (effect: Effect.Effect<HostRegistration, never, Scope.Scope>): Promise<Registration> =>
Effect.runPromiseWith(context)(Scope.provide(scope)(effect)).then((registration) => ({
dispose: () => Effect.runPromiseWith(context)(registration.dispose),
}))
const run = <A, E>(effect: Effect.Effect<A, E>) => Effect.runPromiseWith(context)(effect).then(wire)
const transform =
<Draft>(domain: {
transform: (callback: (draft: Draft) => void) => Effect.Effect<HostRegistration, never, Scope.Scope>
}) =>
(callback: (draft: Draft) => void) =>
register(
domain.transform((draft) => {
callback(draft)
}),
)
const context2: Context = {
app: host.app,
options: host.options,
agent: {
get: (input) => run(host.agent.get({ ...input, agentID: Agent.ID.make(input.agentID) })),
list: (input) => run(host.agent.list(input)),
transform: transform(host.agent),
reload: () => run(host.agent.reload()),
},
aisdk: {
hook: (name, callback) =>
register(host.aisdk.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))),
},
catalog: {
provider: {
list: (input) => run(host.catalog.provider.list(input)),
get: (input) =>
run(host.catalog.provider.get({ ...input, providerID: Provider.ID.make(input.providerID) })),
},
model: {
list: (input) => run(host.catalog.model.list(input)),
default: (input) =>
run(host.catalog.model.default(input)).then((result) => ({ ...result, data: result.data ?? null })),
},
transform: transform(host.catalog),
reload: () => run(host.catalog.reload()),
},
command: {
list: (input) => run(host.command.list(input)),
transform: transform(host.command),
reload: () => run(host.command.reload()),
},
event: {
subscribe: () => Stream.toAsyncIterable(host.event.subscribe().pipe(Stream.map(wireEvent))),
},
integration: {
list: (input) => run(host.integration.list(input)),
get: (input) =>
run(host.integration.get({ ...input, integrationID: Integration.ID.make(input.integrationID) })).then(
(result) => ({ ...result, data: result.data ?? null }),
),
connect: {
key: (input) =>
run(
host.integration.connect.key({ ...input, integrationID: Integration.ID.make(input.integrationID) }),
),
},
oauth: {
connect: (input) =>
run(
host.integration.oauth.connect({
...input,
integrationID: Integration.ID.make(input.integrationID),
methodID: Integration.MethodID.make(input.methodID),
}),
),
status: (input) =>
run(
host.integration.oauth.status({
...input,
integrationID: Integration.ID.make(input.integrationID),
attemptID: Integration.AttemptID.make(input.attemptID),
}),
),
complete: (input) =>
run(
host.integration.oauth.complete({
...input,
integrationID: Integration.ID.make(input.integrationID),
attemptID: Integration.AttemptID.make(input.attemptID),
}),
),
cancel: (input) =>
run(
host.integration.oauth.cancel({
...input,
integrationID: Integration.ID.make(input.integrationID),
attemptID: Integration.AttemptID.make(input.attemptID),
}),
),
},
command: {
connect: (input) =>
run(
host.integration.command.connect({
...input,
integrationID: Integration.ID.make(input.integrationID),
methodID: Integration.MethodID.make(input.methodID),
}),
),
status: (input) =>
run(
host.integration.command.status({
...input,
integrationID: Integration.ID.make(input.integrationID),
attemptID: Integration.AttemptID.make(input.attemptID),
}),
),
cancel: (input) =>
run(
host.integration.command.cancel({
...input,
integrationID: Integration.ID.make(input.integrationID),
attemptID: Integration.AttemptID.make(input.attemptID),
}),
),
},
transform: (callback) =>
register(
host.integration.transform((draft) =>
callback({
list: draft.list,
get: draft.get,
update: draft.update,
remove: draft.remove,
method: {
list: draft.method.list,
update: (input) => {
if (!("authorize" in input)) return draft.method.update(input)
const refresh = input.refresh
draft.method.update({
...input,
authorize: (answer) =>
Effect.promise(() => input.authorize(answer)).pipe(
Effect.map((authorization) =>
authorization.mode === "auto"
? {
...authorization,
callback: Effect.promise(() => authorization.callback),
}
: {
...authorization,
callback: (code) => Effect.promise(() => authorization.callback(code)),
},
),
),
refresh:
refresh === undefined
? undefined
: (credential) => Effect.promise(() => refresh(credential)),
})
},
remove: draft.method.remove,
},
}),
),
),
reload: () => run(host.integration.reload()),
connection: {
active: (id) => Effect.runPromiseWith(context)(host.integration.connection.active(id)),
resolve: (connection) => Effect.runPromiseWith(context)(host.integration.connection.resolve(connection)),
},
},
plugin: {
list: (input) => run(host.plugin.list(input)),
},
reference: {
list: (input) => run(host.reference.list(input)),
transform: transform(host.reference),
reload: () => run(host.reference.reload()),
},
skill: {
list: (input) => run(host.skill.list(input)),
transform: transform(host.skill),
reload: () => run(host.skill.reload()),
},
tool: {
transform: (callback) =>
register(
host.tool.transform((draft) =>
callback({
add: (tool: Info) =>
draft.add({
...tool,
execute: (input, context) => executePromiseTool(tool, input, context),
}),
}),
),
),
hook: (name, callback) =>
register(host.tool.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))),
},
websearch: {
providers: (input) => run(host.websearch.providers(input)),
query: (input) =>
run(
host.websearch.query({
...input,
providerID: input.providerID === undefined ? undefined : WebSearch.ID.make(input.providerID),
}),
),
reload: () => run(host.websearch.reload()),
transform: (callback) =>
register(
host.websearch.transform((draft) => {
callback({
add: (definition) =>
draft.add({
id: definition.id,
name: definition.name,
execute: (input) => attempt((signal) => definition.execute(input, { signal })),
}),
default: draft.default,
})
}),
),
},
session: {
hook: (name, callback) =>
register(host.session.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))),
create: (input) =>
run(
host.session.create(
input === undefined
? undefined
: {
id: input.id == null ? undefined : Session.ID.make(input.id),
agent: input.agent == null ? undefined : Agent.ID.make(input.agent),
model: input.model == null ? undefined : model(input.model),
location:
input.location == null
? undefined
: Location.Ref.make({
directory: AbsolutePath.make(input.location.directory),
workspaceID:
input.location.workspaceID === undefined
? undefined
: Workspace.ID.make(input.location.workspaceID),
}),
},
),
),
get: (input) => run(host.session.get({ sessionID: Session.ID.make(input.sessionID) })),
prompt: (input) =>
run(
host.session.prompt({
...input,
sessionID: Session.ID.make(input.sessionID),
id: input.id == null ? undefined : SessionMessage.ID.make(input.id),
skills: input.skills?.map((skill) => ({ ...skill, id: Skill.ID.make(skill.id) })),
delivery: input.delivery ?? undefined,
resume: input.resume ?? undefined,
}),
),
generate: (input) =>
run(host.session.generate({ sessionID: Session.ID.make(input.sessionID), prompt: input.prompt })),
command: (input) =>
run(
host.session.command({
...input,
sessionID: Session.ID.make(input.sessionID),
id: input.id == null ? undefined : SessionMessage.ID.make(input.id),
agent: input.agent == null ? undefined : Agent.ID.make(input.agent),
model: input.model == null ? undefined : model(input.model),
skills: input.skills?.map((skill) => ({ ...skill, id: Skill.ID.make(skill.id) })),
arguments: input.arguments ?? undefined,
delivery: input.delivery ?? undefined,
resume: input.resume ?? undefined,
}),
),
synthetic: (input) =>
run(
host.session.synthetic({
...input,
sessionID: Session.ID.make(input.sessionID),
id: input.id == null ? undefined : SessionMessage.ID.make(input.id),
description: input.description ?? undefined,
delivery: input.delivery ?? undefined,
resume: input.resume ?? undefined,
}),
),
interrupt: (input) => run(host.session.interrupt({ sessionID: Session.ID.make(input.sessionID) })),
},
shell: {
hook: (name, callback) =>
register(host.shell.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))),
},
}
const cleanup = yield* Effect.promise(() => Promise.resolve(plugin.setup(context2)))
if (!cleanup) return
yield* Effect.addFinalizer(() => Effect.promise(() => Promise.resolve(cleanup())))
}),
})
}
function attempt<A>(evaluate: (signal: AbortSignal) => PromiseLike<A>) {
return Effect.tryPromise({ try: evaluate, catch: (cause) => cause })
}
function model(input: { readonly id: string; readonly providerID: string; readonly variant?: string }) {
return Model.Ref.make({
id: Model.ID.make(input.id),
providerID: Provider.ID.make(input.providerID),
variant: input.variant === undefined ? undefined : Model.VariantID.make(input.variant),
})
}
type Wire<Value> = unknown extends Value
? JsonValue
: Value extends string | number | boolean | bigint | symbol | null | undefined
? Value
: Value extends DateTime.DateTime
? number
: Value extends readonly [infer Head, ...infer Tail]
? [Wire<Head>, ...WireTuple<Tail>]
: Value extends ReadonlyArray<infer Item>
? Array<Wire<Item>>
: Value extends object
? { -readonly [Key in keyof Value]: Wire<Value[Key]> }
: Value
type WireTuple<Value extends ReadonlyArray<unknown>> = {
-readonly [Key in keyof Value]: Wire<Value[Key]>
}
function wire<Value>(value: Value): Wire<Value>
function wire(value: unknown): unknown {
if (DateTime.isDateTime(value)) return DateTime.toEpochMillis(value)
if (Array.isArray(value)) return value.map(wire)
if (typeof value !== "object" || value === null) return value
return Object.fromEntries(Object.entries(value).map(([key, item]) => [key, wire(item)]))
}
function wireEvent(value: unknown): PromiseEvent
function wireEvent(value: unknown): unknown {
return wire(value)
}
const executePromiseTool = (tool: Info, input: any, context: Tool.Context) =>
Effect.promise(() =>
tool.execute(input, {
...context,
progress: (update) => Effect.runPromise(context.progress(update)),
}),
)
export { fromPromise } from "@opencode-ai/plugin/promise/adapter"

View file

@ -4,6 +4,7 @@ import { DateTime, Effect, Schema } from "effect"
import { Agent } from "@opencode-ai/core/agent"
import { Catalog } from "@opencode-ai/core/catalog"
import { Model } from "@opencode-ai/core/model"
import { Location } from "@opencode-ai/core/location"
import { Plugin } from "@opencode-ai/core/plugin"
import { PluginHooks } from "@opencode-ai/core/plugin/hooks"
import { PluginHost } from "@opencode-ai/core/plugin/host"
@ -14,7 +15,10 @@ import { SessionMessage } from "@opencode-ai/core/session/message"
import { SessionInbox } from "@opencode-ai/core/session/inbox"
import { Tool } from "@opencode-ai/core/tool"
import { Provider } from "@opencode-ai/core/provider"
import { Project } from "@opencode-ai/core/project"
import { AbsolutePath } from "@opencode-ai/core/schema"
import { define } from "@opencode-ai/plugin/promise/plugin"
import { Money } from "@opencode-ai/schema/money"
import type { SessionHooks } from "@opencode-ai/plugin/effect/session"
import { testEffect } from "../lib/effect"
import { PluginTestLayer } from "./fixture"
@ -23,6 +27,53 @@ import { host as testHost } from "./host"
const it = testEffect(PluginTestLayer)
describe("fromPromise", () => {
it.effect("adapts session creation through the protocol schema", () =>
Effect.gen(function* () {
let seen: unknown
const host = testHost({
session: {
create: (input) => {
seen = input
return Effect.succeed(
Session.Info.make({
id: Session.ID.make("ses_protocol_adapter"),
projectID: Project.ID.make("project"),
cost: Money.USD.make(0),
tokens: { input: 1, output: 2, reasoning: 3, cache: { read: 4, write: 5 } },
time: { created: DateTime.makeUnsafe(10), updated: DateTime.makeUnsafe(20) },
title: input?.title,
location: Location.Ref.make({ directory: AbsolutePath.make("/workspace") }),
}),
)
},
},
})
yield* PluginPromise.fromPromise(
define({
id: "promise-session-create",
setup: async (ctx) => {
await expect(Reflect.apply(ctx.session.create, undefined, [{ title: 42 }])).rejects.toBeDefined()
const result = await ctx.session.create({
id: null,
title: "Promise title",
agent: null,
model: null,
location: null,
})
expect(result).toMatchObject({
id: "ses_protocol_adapter",
title: "Promise title",
time: { created: 10, updated: 20 },
})
},
}),
).effect(host)
expect(seen).toEqual({ title: "Promise title" })
}),
)
it.effect("forwards transient session generation", () =>
Effect.gen(function* () {
const host = testHost({
@ -44,6 +95,42 @@ describe("fromPromise", () => {
}),
)
it.effect("preserves no-content and rejected Promise behavior", () =>
Effect.gen(function* () {
const seen: unknown[] = []
const host = testHost({
session: {
interrupt: (input) => {
if (input.sessionID === Session.ID.make("ses_failure")) {
return Effect.fail(new Error("interrupt failed"))
}
expect(input.continue).toBe(true)
return Effect.void
},
rename: (input) => Effect.sync(() => seen.push(input)),
wait: (input) => Effect.sync(() => seen.push(input)),
},
})
yield* PluginPromise.fromPromise(
define({
id: "promise-session-interrupt",
setup: async (ctx) => {
expect(await ctx.session.interrupt({ sessionID: "ses_success", continue: true })).toBeUndefined()
await expect(ctx.session.interrupt({ sessionID: "ses_failure" })).rejects.toThrow("interrupt failed")
expect(await ctx.session.rename({ sessionID: "ses_success", title: "Renamed" })).toBeUndefined()
expect(await ctx.session.wait({ sessionID: "ses_success" })).toBeUndefined()
},
}),
).effect(host)
expect(seen).toEqual([
{ sessionID: Session.ID.make("ses_success"), title: "Renamed" },
{ sessionID: Session.ID.make("ses_success") },
])
}),
)
it.effect("forwards synthetic session input", () =>
Effect.gen(function* () {
const input = {
@ -114,6 +201,7 @@ describe("fromPromise", () => {
ctx.skill.list(),
])
seen.push(...results.map((result) => result.location.directory))
expect((await ctx.integration.get({ integrationID: "missing" })).data).toBeNull()
},
})

View file

@ -23,6 +23,7 @@
"@ai-sdk/provider": "3.0.8",
"@opencode-ai/ai": "workspace:*",
"@opencode-ai/client": "workspace:*",
"@opencode-ai/protocol": "workspace:*",
"@opencode-ai/schema": "workspace:*",
"@opencode-ai/sdk": "1.18.5",
"@standard-schema/spec": "catalog:",

View file

@ -0,0 +1,324 @@
import { Tool } from "@opencode-ai/schema/tool"
import { Effect, Schema, SchemaAST, Scope, Stream } from "effect"
import { HttpApiEndpoint, HttpApiSchema } from "effect/unstable/httpapi"
import { define } from "../effect/plugin.js"
import type { Context, Plugin } from "./plugin.js"
import type { Info } from "./tool.js"
type HostRegistration = { readonly dispose: Effect.Effect<void> }
type Registration = { readonly dispose: () => Promise<void> }
type PromiseEvent = ReturnType<Context["event"]["subscribe"]> extends AsyncIterable<infer Event> ? Event : never
interface CompiledEndpoint {
readonly decode: ReadonlyArray<(input: unknown) => Effect.Effect<unknown, Schema.SchemaError>>
readonly encode: (output: unknown) => Effect.Effect<unknown, Schema.SchemaError>
readonly noContent: boolean
}
const compiledEndpoints = new WeakMap<object, CompiledEndpoint>()
function compileEndpoint(endpoint: HttpApiEndpoint.Top) {
const cached = compiledEndpoints.get(endpoint)
if (cached) return cached
const payloadSchemas = Array.from(endpoint.payload.values()).flatMap(({ schemas }) => schemas)
const successSchemas = Array.from(endpoint.success)
if (payloadSchemas.length > 1 || successSchemas.length > 1) {
throw new Error(`Unsupported API schema cardinality: ${endpoint.identifier}`)
}
const inputs = [
endpoint.params,
endpoint.query === undefined ? undefined : Schema.toType(endpoint.query),
endpoint.headers,
...payloadSchemas,
].filter((schema): schema is Schema.Top => schema !== undefined) as Array<RuntimeSchema>
const success = (successSchemas[0] ?? HttpApiSchema.NoContent) as RuntimeSchema
const noContent = HttpApiSchema.isNoContent(success.ast)
const type = Schema.toType(success).ast
const data = SchemaAST.isObjects(success.ast)
? success.ast.propertySignatures.find((property) => property.name === "data")
: undefined
const output =
!noContent &&
SchemaAST.isObjects(type) &&
type.indexSignatures.length === 0 &&
type.propertySignatures.length === 1 &&
type.propertySignatures[0]?.name === "data" &&
data !== undefined
? (Schema.make<Schema.Top>(data.type) as RuntimeSchema)
: success
const compiled = {
decode: inputs.map((schema) => Schema.decodeUnknownEffect(schema)),
encode: Schema.encodeUnknownEffect(output),
noContent,
} satisfies CompiledEndpoint
compiledEndpoints.set(endpoint, compiled)
return compiled
}
/**
* Adapts a Promise plugin into an Effect plugin so the existing Effect-only
* loader (`Plugin` / `PluginSupervisor`) can run it unchanged.
*
* Hook registrations created during the async `setup` attach to the plugin's
* scope, so unloading the plugin disposes them. The captured fiber context
* preserves boot-time batching, so Promise-plugin transforms still coalesce
* into one reload per domain.
*/
export function fromPromise(plugin: Plugin) {
return define({
id: plugin.id,
effect: (host) =>
Effect.gen(function* () {
const [{ ClientApi }, { OpenCodeEvent }] = yield* Effect.promise(() =>
Promise.all([import("@opencode-ai/protocol/client"), import("@opencode-ai/protocol/groups/event")]),
)
const AgentEndpoints = ClientApi.groups["server.agent"].endpoints
const CommandEndpoints = ClientApi.groups["server.command"].endpoints
const IntegrationEndpoints = ClientApi.groups["server.integration"].endpoints
const ModelEndpoints = ClientApi.groups["server.model"].endpoints
const PluginEndpoints = ClientApi.groups["server.plugin"].endpoints
const ProviderEndpoints = ClientApi.groups["server.provider"].endpoints
const ReferenceEndpoints = ClientApi.groups["server.reference"].endpoints
const SessionEndpoints = ClientApi.groups["server.session"].endpoints
const SkillEndpoints = ClientApi.groups["server.skill"].endpoints
const WebSearchEndpoints = ClientApi.groups["server.websearch"].endpoints
const scope = yield* Scope.Scope
const context = yield* Effect.context<Scope.Scope>()
// Run a hook registration on the plugin scope and resolve once it is registered.
const register = (effect: Effect.Effect<HostRegistration, never, Scope.Scope>): Promise<Registration> =>
Effect.runPromiseWith(context)(Scope.provide(scope)(effect)).then((registration) => ({
dispose: () => Effect.runPromiseWith(context)(registration.dispose),
}))
const run = <A, E>(effect: Effect.Effect<A, E>) => Effect.runPromiseWith(context)(effect)
const adaptApiMethod = <PromiseMethod>(
endpoint: HttpApiEndpoint.Top,
method: (input: never) => Effect.Effect<unknown, unknown>,
) => {
const compiled = compileEndpoint(endpoint)
return ((input?: unknown) =>
Effect.gen(function* () {
const decoded = yield* Effect.forEach(compiled.decode, (decode) => decode(input ?? {}))
const result = yield* method(Object.assign({}, ...decoded) as never)
if (compiled.noContent) return undefined
return yield* compiled.encode(result)
}).pipe(Effect.runPromiseWith(context))) as PromiseMethod
}
const transform =
<Draft>(domain: {
transform: (callback: (draft: Draft) => void) => Effect.Effect<HostRegistration, never, Scope.Scope>
}) =>
(callback: (draft: Draft) => void) =>
register(
domain.transform((draft) => {
callback(draft)
}),
)
const context2: Context = {
app: host.app,
options: host.options,
agent: {
get: adaptApiMethod(AgentEndpoints["agent.get"], host.agent.get),
list: adaptApiMethod(AgentEndpoints["agent.list"], host.agent.list),
transform: transform(host.agent),
reload: () => run(host.agent.reload()),
},
aisdk: {
hook: (name, callback) =>
register(host.aisdk.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))),
},
catalog: {
provider: {
list: adaptApiMethod(ProviderEndpoints["provider.list"], host.catalog.provider.list),
get: adaptApiMethod(ProviderEndpoints["provider.get"], host.catalog.provider.get),
},
model: {
list: adaptApiMethod(ModelEndpoints["model.list"], host.catalog.model.list),
default: adaptApiMethod(ModelEndpoints["model.default"], host.catalog.model.default),
},
transform: transform(host.catalog),
reload: () => run(host.catalog.reload()),
},
command: {
list: adaptApiMethod(CommandEndpoints["command.list"], host.command.list),
transform: transform(host.command),
reload: () => run(host.command.reload()),
},
event: {
subscribe: () =>
Stream.toAsyncIterable(
host.event.subscribe().pipe(
Stream.mapEffect((event) => Schema.encodeUnknownEffect(OpenCodeEvent)(event)),
Stream.map((event) => event as unknown as PromiseEvent),
),
),
},
integration: {
list: adaptApiMethod(IntegrationEndpoints["integration.list"], host.integration.list),
get: adaptApiMethod(IntegrationEndpoints["integration.get"], host.integration.get),
connect: {
key: adaptApiMethod(IntegrationEndpoints["integration.connect.key"], host.integration.connect.key),
},
oauth: {
connect: adaptApiMethod(
IntegrationEndpoints["integration.oauth.connect"],
host.integration.oauth.connect,
),
status: adaptApiMethod(IntegrationEndpoints["integration.oauth.status"], host.integration.oauth.status),
complete: adaptApiMethod(
IntegrationEndpoints["integration.oauth.complete"],
host.integration.oauth.complete,
),
cancel: adaptApiMethod(IntegrationEndpoints["integration.oauth.cancel"], host.integration.oauth.cancel),
},
command: {
connect: adaptApiMethod(
IntegrationEndpoints["integration.command.connect"],
host.integration.command.connect,
),
status: adaptApiMethod(
IntegrationEndpoints["integration.command.status"],
host.integration.command.status,
),
cancel: adaptApiMethod(
IntegrationEndpoints["integration.command.cancel"],
host.integration.command.cancel,
),
},
transform: (callback) =>
register(
host.integration.transform((draft) =>
callback({
list: draft.list,
get: draft.get,
update: draft.update,
remove: draft.remove,
method: {
list: draft.method.list,
update: (input) => {
if (!("authorize" in input)) return draft.method.update(input)
const refresh = input.refresh
draft.method.update({
...input,
authorize: (answer) =>
Effect.promise(() => input.authorize(answer)).pipe(
Effect.map((authorization) =>
authorization.mode === "auto"
? {
...authorization,
callback: Effect.promise(() => authorization.callback),
}
: {
...authorization,
callback: (code) => Effect.promise(() => authorization.callback(code)),
},
),
),
refresh:
refresh === undefined
? undefined
: (credential) => Effect.promise(() => refresh(credential)),
})
},
remove: draft.method.remove,
},
}),
),
),
reload: () => run(host.integration.reload()),
connection: {
active: (id) => Effect.runPromiseWith(context)(host.integration.connection.active(id)),
resolve: (connection) => Effect.runPromiseWith(context)(host.integration.connection.resolve(connection)),
},
},
plugin: {
list: adaptApiMethod(PluginEndpoints["plugin.list"], host.plugin.list),
},
reference: {
list: adaptApiMethod(ReferenceEndpoints["reference.list"], host.reference.list),
transform: transform(host.reference),
reload: () => run(host.reference.reload()),
},
skill: {
list: adaptApiMethod(SkillEndpoints["skill.list"], host.skill.list),
transform: transform(host.skill),
reload: () => run(host.skill.reload()),
},
tool: {
transform: (callback) =>
register(
host.tool.transform((draft) =>
callback({
add: (tool: Info) =>
draft.add({
...tool,
execute: (input, context) => executePromiseTool(tool, input, context),
}),
}),
),
),
hook: (name, callback) =>
register(host.tool.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))),
},
websearch: {
providers: adaptApiMethod(WebSearchEndpoints["websearch.providers"], host.websearch.providers),
query: adaptApiMethod(WebSearchEndpoints["websearch.query"], host.websearch.query),
reload: () => run(host.websearch.reload()),
transform: (callback) =>
register(
host.websearch.transform((draft) => {
callback({
add: (definition) =>
draft.add({
id: definition.id,
name: definition.name,
execute: (input) => attempt((signal) => definition.execute(input, { signal })),
}),
default: draft.default,
})
}),
),
},
session: {
hook: (name, callback) =>
register(host.session.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))),
create: adaptApiMethod(SessionEndpoints["session.create"], host.session.create),
get: adaptApiMethod(SessionEndpoints["session.get"], host.session.get),
prompt: adaptApiMethod(SessionEndpoints["session.prompt"], host.session.prompt),
generate: adaptApiMethod(SessionEndpoints["session.generate"], host.session.generate),
command: adaptApiMethod(SessionEndpoints["session.command"], host.session.command),
synthetic: adaptApiMethod(SessionEndpoints["session.synthetic"], host.session.synthetic),
interrupt: adaptApiMethod(SessionEndpoints["session.interrupt"], host.session.interrupt),
rename: adaptApiMethod(SessionEndpoints["session.rename"], host.session.rename),
wait: adaptApiMethod(SessionEndpoints["session.wait"], host.session.wait),
},
shell: {
hook: (name, callback) =>
register(host.shell.hook(name, (event) => Effect.promise(() => Promise.resolve(callback(event))))),
},
}
const cleanup = yield* Effect.promise(() => Promise.resolve(plugin.setup(context2)))
if (!cleanup) return
yield* Effect.addFinalizer(() => Effect.promise(() => Promise.resolve(cleanup())))
}),
})
}
function attempt<A>(evaluate: (signal: AbortSignal) => PromiseLike<A>) {
return Effect.tryPromise({ try: evaluate, catch: (cause) => cause })
}
type RuntimeSchema = Schema.Codec<unknown, unknown>
const executePromiseTool = (tool: Info, input: any, context: Tool.Context) =>
Effect.promise(() =>
tool.execute(input, {
...context,
progress: (update) => Effect.runPromise(context.progress(update)),
}),
)

View file

@ -38,7 +38,7 @@ export interface SessionHooks {
export type SessionDomain = Pick<
SessionApi,
"create" | "get" | "prompt" | "generate" | "command" | "synthetic" | "interrupt"
"create" | "get" | "prompt" | "generate" | "command" | "synthetic" | "interrupt" | "rename" | "wait"
> & {
readonly hook: Hooks<SessionHooks>
}