Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion packages/core/src/config/plugin/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,7 @@ export const Plugin = define({
.pipe(
Stream.filterEffect((update) => Effect.map(config.entries(), (entries) => isAgentSource(entries, update.path))),
)
const configUpdates = ctx.event.subscribe().pipe(Stream.filter((event) => event.type === "config.updated"))
const configUpdates = ctx.event.subscribe("config.updated")
yield* Stream.merge(sourceChanges, configUpdates).pipe(
Stream.debounce("100 millis"),
Stream.runForEach(() => reload),
Expand Down
2 changes: 1 addition & 1 deletion packages/core/src/config/plugin/command.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ export const Plugin = define({
Effect.map(config.entries(), (entries) => isCommandSource(entries, update.path)),
),
)
const configUpdates = ctx.event.subscribe().pipe(Stream.filter((event) => event.type === "config.updated"))
const configUpdates = ctx.event.subscribe("config.updated")
yield* Stream.merge(sourceChanges, configUpdates).pipe(
Stream.debounce("100 millis"),
Stream.runForEach(() => reload),
Expand Down
3 changes: 1 addition & 2 deletions packages/core/src/config/plugin/policy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,7 @@ export const Plugin = define({
if (policy?.effect === "deny") catalog.provider.remove(record.provider.id)
}
})
yield* ctx.event.subscribe().pipe(
Stream.filter((event) => event.type === "config.updated"),
yield* ctx.event.subscribe("config.updated").pipe(
Stream.runForEach(() =>
config.entries().pipe(
Effect.tap((entries) => Effect.sync(() => (loaded.entries = entries))),
Expand Down
3 changes: 1 addition & 2 deletions packages/core/src/config/plugin/provider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -97,8 +97,7 @@ export const Plugin = define({
}
}
})
yield* ctx.event.subscribe().pipe(
Stream.filter((event) => event.type === "config.updated"),
yield* ctx.event.subscribe("config.updated").pipe(
Stream.runForEach(() =>
config.entries().pipe(
Effect.tap((entries) => Effect.sync(() => (loaded.entries = entries))),
Expand Down
3 changes: 1 addition & 2 deletions packages/core/src/config/plugin/reference.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,8 +49,7 @@ export const Plugin = define({
}
for (const [name, source] of entries) draft.add(name, source)
})
yield* ctx.event.subscribe().pipe(
Stream.filter((event) => event.type === "config.updated"),
yield* ctx.event.subscribe("config.updated").pipe(
Stream.runForEach(() =>
config.entries().pipe(
Effect.tap((entries) => Effect.sync(() => (loaded.entries = entries))),
Expand Down
3 changes: 1 addition & 2 deletions packages/core/src/config/plugin/skill.ts
Original file line number Diff line number Diff line change
Expand Up @@ -180,8 +180,7 @@ export const Plugin = define({
yield* ctx.skill.transform((draft) => {
for (const skill of loaded.skills) draft.add(skill)
})
yield* ctx.event.subscribe().pipe(
Stream.filter((event) => event.type === "config.updated"),
yield* ctx.event.subscribe("config.updated").pipe(
Stream.runForEach(() =>
config.entries().pipe(
Effect.tap((entries) => Effect.sync(() => (loaded.entries = entries))),
Expand Down
3 changes: 1 addition & 2 deletions packages/core/src/config/plugin/websearch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,7 @@ export const Plugin = define({
if (selection === false) websearch.default.set(false)
if (selection) websearch.default.set(selection.provider)
})
yield* ctx.event.subscribe().pipe(
Stream.filter((event) => event.type === "config.updated"),
yield* ctx.event.subscribe("config.updated").pipe(
Stream.runForEach(() =>
config.entries().pipe(
Effect.tap((entries) => Effect.sync(() => (loaded.entries = entries))),
Expand Down
8 changes: 7 additions & 1 deletion packages/core/src/plugin/host.ts
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,12 @@ export const make = Effect.fn("PluginHost.make")(function* (plugin: import("../p
ref.directory === location.directory && ref.workspaceID === location.workspaceID
const response = <A, E, R>(effect: Effect.Effect<A, E, R>) =>
effect.pipe(Effect.map((data) => ({ location: locationInfo(), data })))
const subscribe: Plugin.Context["event"]["subscribe"] = (type?: EventManifest.ServerEvent["type"]) => {
if (type === undefined) return bus.subscribe().pipe(Stream.filter(EventManifest.isServer))
const definition = EventManifest.Server.get(type)
if (!definition) return Stream.fail(new Error(`Unknown plugin event type: ${type}`))
return bus.subscribe(definition).pipe(Stream.filter(EventManifest.isServer))
}

return {
app,
Expand Down Expand Up @@ -180,7 +186,7 @@ export const make = Effect.fn("PluginHost.make")(function* (plugin: import("../p
}),
},
event: {
subscribe: () => bus.subscribe().pipe(Stream.filter(EventManifest.isServer)),
subscribe,
},
integration: {
list: () => response(integration.list()),
Expand Down
43 changes: 37 additions & 6 deletions packages/core/test/plugin.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,24 +20,55 @@ class Secret extends Context.Service<Secret, string>()("@opencode/test/PluginSec
const versioned = <R>(plugin: EffectPlugin.Plugin<R>, version = "1") => ({ ...plugin, version })

describe("Plugin", () => {
it.live("exposes public events through the plugin context", () =>
it.live("selects one public event type through the plugin context", () =>
Effect.gen(function* () {
const plugins = yield* Plugin.Service
const bus = yield* Bus.Service
const host = yield* PluginHost.make(plugins)
const received = yield* host.event.subscribe().pipe(
Stream.filter((event) => event.type === "config.updated"),
Stream.runHead,
Effect.forkScoped({ startImmediately: true }),
)
const received = yield* host.event
.subscribe("config.updated")
.pipe(Stream.runHead, Effect.forkScoped({ startImmediately: true }))
yield* Effect.sleep("10 millis")

yield* bus.publish(Plugin.Event.Updated, {})
yield* bus.publish(ConfigSchema.Event.Updated, {})

expect((yield* Fiber.join(received)).valueOrUndefined?.type).toBe("config.updated")
}),
)

it.live("exposes all public events through a wildcard plugin subscription", () =>
Effect.gen(function* () {
const plugins = yield* Plugin.Service
const bus = yield* Bus.Service
const host = yield* PluginHost.make(plugins)
const received = yield* host.event
.subscribe()
.pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped({ startImmediately: true }))
yield* Effect.sleep("10 millis")

yield* bus.publish(Plugin.Event.Updated, {})
yield* bus.publish(ConfigSchema.Event.Updated, {})

expect(Array.from(yield* Fiber.join(received), (event) => event.type)).toEqual([
"plugin.updated",
"config.updated",
])
}),
)

it.effect("rejects unknown runtime plugin event types", () =>
Effect.gen(function* () {
const plugins = yield* Plugin.Service
const host = yield* PluginHost.make(plugins)
const subscribe = host.event.subscribe as unknown as (type: string) => Stream.Stream<never, Error>

const failure = yield* subscribe("unknown.event").pipe(Stream.runDrain, Effect.flip)

expect(failure.message).toBe("Unknown plugin event type: unknown.event")
}),
)

it.effect("replaces plugins by ID and version", () =>
Effect.gen(function* () {
const plugins = yield* Plugin.Service
Expand Down
26 changes: 25 additions & 1 deletion packages/core/test/plugin/promise.test.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { describe, expect } from "bun:test"
import { Message, SystemPart } from "@opencode-ai/ai"
import { DateTime, Effect, Schema } from "effect"
import { DateTime, Effect, Schema, Stream } from "effect"
import { Agent } from "@opencode-ai/core/agent"
import { Catalog } from "@opencode-ai/core/catalog"
import { Model } from "@opencode-ai/core/model"
Expand All @@ -18,6 +18,8 @@ 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 { Plugin as EffectPlugin } from "@opencode-ai/plugin/effect"
import type { PluginEventType } from "@opencode-ai/plugin/effect/event"
import { Money } from "@opencode-ai/schema/money"
import type { SessionHooks } from "@opencode-ai/plugin/effect/session"
import { testEffect } from "../lib/effect"
Expand All @@ -27,6 +29,28 @@ import { host as testHost } from "./host"
const it = testEffect(PluginTestLayer)

describe("fromPromise", () => {
it.effect("forwards a selected event type", () =>
Effect.gen(function* () {
let selected: string | undefined
const subscribe: EffectPlugin.Context["event"]["subscribe"] = (type?: PluginEventType) => {
selected = type
return Stream.empty
}
const host = testHost({ event: { subscribe } })

yield* PluginPromise.fromPromise(
define({
id: "promise-event-subscribe",
setup: (ctx) => {
ctx.event.subscribe("config.updated")
},
}),
).effect(host)

expect(selected).toBe("config.updated")
}),
)

it.effect("adapts session creation through the protocol schema", () =>
Effect.gen(function* () {
let seen: unknown
Expand Down
14 changes: 13 additions & 1 deletion packages/plugin/src/effect/event.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,15 @@
import type { EventApi } from "@opencode-ai/client/effect/api"
import type { OpenCodeEvent } from "@opencode-ai/client/effect"
import type { Stream } from "effect"

export interface EventDomain extends Pick<EventApi<unknown>, "subscribe"> {}
export type PluginEvent = Exclude<OpenCodeEvent, { readonly type: "server.connected" }>
export type PluginEventType = PluginEvent["type"]

export interface EventSubscribe {
(): Stream.Stream<PluginEvent, unknown>
(type: PluginEventType): Stream.Stream<PluginEvent, unknown>
}

export interface EventDomain extends Omit<EventApi<unknown>, "subscribe"> {
readonly subscribe: EventSubscribe
}
11 changes: 7 additions & 4 deletions packages/plugin/src/promise/adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ 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 { PluginEventType } from "./event.js"
import type { Context, Plugin } from "./plugin.js"
import type { Info } from "./tool.js"

Expand Down Expand Up @@ -149,13 +150,15 @@ export function fromPromise(plugin: Plugin) {
reload: () => run(host.command.reload()),
},
event: {
subscribe: () =>
Stream.toAsyncIterable(
host.event.subscribe().pipe(
subscribe: (type?: PluginEventType) => {
const events = type === undefined ? host.event.subscribe() : host.event.subscribe(type)
return Stream.toAsyncIterable(
events.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),
Expand Down
13 changes: 12 additions & 1 deletion packages/plugin/src/promise/event.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,14 @@
import type { OpenCodeEvent } from "@opencode-ai/client"
import type { EventApi } from "@opencode-ai/client/promise/api"

export interface EventDomain extends Pick<EventApi, "subscribe"> {}
export type PluginEvent = Exclude<OpenCodeEvent, { readonly type: "server.connected" }>
export type PluginEventType = PluginEvent["type"]

export interface EventSubscribe {
(): AsyncIterable<PluginEvent>
(type: PluginEventType): AsyncIterable<PluginEvent>
}

export interface EventDomain extends Omit<EventApi, "subscribe"> {
readonly subscribe: EventSubscribe
}
26 changes: 26 additions & 0 deletions packages/plugin/test/event-types.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
import { expect, test } from "bun:test"
import type { Context as EffectContext } from "../src/effect/plugin.js"
import type { Context as PromiseContext } from "../src/promise/plugin.js"

function effectSubscriptions(ctx: EffectContext) {
ctx.event.subscribe()
ctx.event.subscribe("config.updated")
// @ts-expect-error server.connected is a network-only marker
ctx.event.subscribe("server.connected")
// @ts-expect-error plugin subscriptions select at most one event type
ctx.event.subscribe(["config.updated"])
}

function promiseSubscriptions(ctx: PromiseContext) {
ctx.event.subscribe()
ctx.event.subscribe("config.updated")
// @ts-expect-error server.connected is a network-only marker
ctx.event.subscribe("server.connected")
// @ts-expect-error plugin subscriptions select at most one event type
ctx.event.subscribe(["config.updated"])
}

test("event subscription types support wildcard and one public event", () => {
expect(effectSubscriptions).toBeFunction()
expect(promiseSubscriptions).toBeFunction()
})
8 changes: 8 additions & 0 deletions packages/www/content/docs/build/plugins.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -184,6 +184,14 @@ and plugin options.
| `ctx.event` | `subscribe` to the current public server event stream |
| `ctx.options` | Readonly options from the matching config object |

Event subscriptions can receive every plugin-visible public event, or select
one event type:

```ts
ctx.event.subscribe()
ctx.event.subscribe("config.updated")
```

### Transform hooks

Transform hooks let a plugin modify how OpenCode is configured. Use them to add
Expand Down
Loading