mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-16 01:19:19 -04:00
Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 29015efb95 | |||
| 1cab383d3b | |||
| 8f183dc5a9 | |||
| bbd09e7da0 | |||
| 8251934007 | |||
| 42e345e1bc |
@@ -27,6 +27,7 @@ const requestedTarget = process.argv.find((arg) => arg.startsWith("--target="))?
|
||||
const skipInstall = process.argv.includes("--skip-install")
|
||||
const skipWebUi = process.argv.includes("--skip-web-ui")
|
||||
const solidPlugin = createSolidTransformPlugin()
|
||||
const releaseAssets = new Map<string, Promise<Map<string, string>>>()
|
||||
|
||||
const allTargets: {
|
||||
os: string
|
||||
@@ -161,7 +162,13 @@ async function compileExecutable(item: (typeof allTargets)[number]) {
|
||||
if (!release) return
|
||||
|
||||
const platform = item.os === "win32" ? "windows" : item.os
|
||||
const name = ["bun", platform, item.arch, item.abi, item.avx2 === false ? "baseline" : undefined]
|
||||
const name = [
|
||||
"bun",
|
||||
platform,
|
||||
item.arch === "arm64" ? "aarch64" : item.arch,
|
||||
item.abi,
|
||||
item.avx2 === false ? "baseline" : undefined,
|
||||
]
|
||||
.filter(Boolean)
|
||||
.join("-")
|
||||
const cache = path.join(outdir, ".bun", release)
|
||||
@@ -170,7 +177,13 @@ async function compileExecutable(item: (typeof allTargets)[number]) {
|
||||
|
||||
await mkdir(cache, { recursive: true })
|
||||
const archive = path.join(cache, `${name}.zip`)
|
||||
const response = await fetch(`https://github.com/oven-sh/bun/releases/download/${release}/${name}.zip`)
|
||||
const assets = await compileReleaseAssets(release)
|
||||
const url = assets.get(`${name}.zip`)
|
||||
if (!url) throw new Error(`Bun release ${release} does not include ${name}.zip`)
|
||||
const token = process.env.GH_TOKEN ?? process.env.GITHUB_TOKEN
|
||||
const response = await fetch(url, {
|
||||
headers: { Accept: "application/octet-stream", ...(token ? { Authorization: `Bearer ${token}` } : {}) },
|
||||
})
|
||||
if (!response.ok) throw new Error(`Failed to download ${name} from Bun release ${release}: ${response.status}`)
|
||||
await Bun.write(archive, response)
|
||||
await $`unzip -oq ${archive} -d ${cache}`
|
||||
@@ -178,6 +191,38 @@ async function compileExecutable(item: (typeof allTargets)[number]) {
|
||||
return executable
|
||||
}
|
||||
|
||||
function compileReleaseAssets(release: string) {
|
||||
const existing = releaseAssets.get(release)
|
||||
if (existing) return existing
|
||||
const pending = fetch(`https://api.github.com/repos/oven-sh/bun/releases/tags/${release}?cache=${Date.now()}`)
|
||||
.then(async (response) => {
|
||||
if (!response.ok) throw new Error(`Failed to resolve Bun release ${release}: ${response.status}`)
|
||||
const data: unknown = await response.json()
|
||||
if (typeof data !== "object" || data === null || !("assets" in data) || !Array.isArray(data.assets)) {
|
||||
throw new Error(`Bun release ${release} returned invalid metadata`)
|
||||
}
|
||||
return new Map(
|
||||
data.assets
|
||||
.filter(
|
||||
(asset): asset is { name: string; url: string } =>
|
||||
typeof asset === "object" &&
|
||||
asset !== null &&
|
||||
"name" in asset &&
|
||||
typeof asset.name === "string" &&
|
||||
"url" in asset &&
|
||||
typeof asset.url === "string",
|
||||
)
|
||||
.map((asset) => [asset.name, asset.url]),
|
||||
)
|
||||
})
|
||||
.catch((error) => {
|
||||
releaseAssets.delete(release)
|
||||
throw error
|
||||
})
|
||||
releaseAssets.set(release, pending)
|
||||
return pending
|
||||
}
|
||||
|
||||
function targetName(item: (typeof allTargets)[number]) {
|
||||
return [
|
||||
binary,
|
||||
|
||||
@@ -435,7 +435,22 @@ function prompt(request: LLMRequest): LanguageModelV3Prompt {
|
||||
.map((part) => part.text)
|
||||
.filter(Boolean)
|
||||
.join("\n\n")
|
||||
const messages = request.messages.flatMap(message)
|
||||
const pending: UserContent = []
|
||||
const messages = request.messages.flatMap((input, index) => {
|
||||
if (input.role !== "tool") return message(input)
|
||||
const lowered = toolMessage(input)
|
||||
pending.push(...lowered.media)
|
||||
if (request.messages[index + 1]?.role === "tool" || pending.length === 0) return lowered.messages
|
||||
const media = [...pending]
|
||||
pending.length = 0
|
||||
return [
|
||||
...lowered.messages,
|
||||
{
|
||||
role: "user" as const,
|
||||
content: [{ type: "text" as const, text: "Attached media from tool result:" }, ...media],
|
||||
},
|
||||
]
|
||||
})
|
||||
if (!system.length) return messages
|
||||
return [{ role: "system", content: system }, ...messages]
|
||||
}
|
||||
@@ -448,10 +463,33 @@ function message(input: LLMRequest["messages"][number]): LanguageModelV3Message[
|
||||
return [{ role: "user", content: input.content.flatMap(userPart) }]
|
||||
case "assistant":
|
||||
return [{ role: "assistant", content: input.content.flatMap(assistantPart) }]
|
||||
case "tool": {
|
||||
const content = input.content.flatMap(toolResultPart)
|
||||
return content.length ? [{ role: "tool", content }] : []
|
||||
}
|
||||
case "tool":
|
||||
return toolMessage(input).messages
|
||||
}
|
||||
}
|
||||
|
||||
function toolMessage(input: LLMRequest["messages"][number]) {
|
||||
const media: UserContent = []
|
||||
const content = input.content.flatMap((part) => {
|
||||
if (part.type !== "tool-result" || part.result.type !== "content") return toolResultPart(part)
|
||||
const value = part.result.value.filter((item) => {
|
||||
if (item.type !== "file") return true
|
||||
if (!item.mime.startsWith("image/") && item.mime !== "application/pdf") return true
|
||||
const data = /^data:[^;,]+(?:;[^,]*)*;base64,(.*)$/s.exec(item.uri)?.[1] ?? item.uri
|
||||
media.push({ type: "file", mediaType: item.mime, data, filename: item.name })
|
||||
return false
|
||||
})
|
||||
return toolResultPart({
|
||||
...part,
|
||||
result:
|
||||
value.length === 0
|
||||
? { type: "text", value: "Media attached in the following user message." }
|
||||
: { ...part.result, value },
|
||||
})
|
||||
})
|
||||
return {
|
||||
messages: content.length ? ([{ role: "tool", content }] satisfies LanguageModelV3Message[]) : [],
|
||||
media,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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("config.updated")
|
||||
const configUpdates = ctx.event.subscribe().pipe(Stream.filter((event) => event.type === "config.updated"))
|
||||
yield* Stream.merge(sourceChanges, configUpdates).pipe(
|
||||
Stream.debounce("100 millis"),
|
||||
Stream.runForEach(() => reload),
|
||||
|
||||
@@ -43,7 +43,7 @@ export const Plugin = define({
|
||||
Effect.map(config.entries(), (entries) => isCommandSource(entries, update.path)),
|
||||
),
|
||||
)
|
||||
const configUpdates = ctx.event.subscribe("config.updated")
|
||||
const configUpdates = ctx.event.subscribe().pipe(Stream.filter((event) => event.type === "config.updated"))
|
||||
yield* Stream.merge(sourceChanges, configUpdates).pipe(
|
||||
Stream.debounce("100 millis"),
|
||||
Stream.runForEach(() => reload),
|
||||
|
||||
@@ -22,7 +22,8 @@ export const Plugin = define({
|
||||
if (policy?.effect === "deny") catalog.provider.remove(record.provider.id)
|
||||
}
|
||||
})
|
||||
yield* ctx.event.subscribe("config.updated").pipe(
|
||||
yield* ctx.event.subscribe().pipe(
|
||||
Stream.filter((event) => event.type === "config.updated"),
|
||||
Stream.runForEach(() =>
|
||||
config.entries().pipe(
|
||||
Effect.tap((entries) => Effect.sync(() => (loaded.entries = entries))),
|
||||
|
||||
@@ -97,7 +97,8 @@ export const Plugin = define({
|
||||
}
|
||||
}
|
||||
})
|
||||
yield* ctx.event.subscribe("config.updated").pipe(
|
||||
yield* ctx.event.subscribe().pipe(
|
||||
Stream.filter((event) => event.type === "config.updated"),
|
||||
Stream.runForEach(() =>
|
||||
config.entries().pipe(
|
||||
Effect.tap((entries) => Effect.sync(() => (loaded.entries = entries))),
|
||||
|
||||
@@ -49,7 +49,8 @@ export const Plugin = define({
|
||||
}
|
||||
for (const [name, source] of entries) draft.add(name, source)
|
||||
})
|
||||
yield* ctx.event.subscribe("config.updated").pipe(
|
||||
yield* ctx.event.subscribe().pipe(
|
||||
Stream.filter((event) => event.type === "config.updated"),
|
||||
Stream.runForEach(() =>
|
||||
config.entries().pipe(
|
||||
Effect.tap((entries) => Effect.sync(() => (loaded.entries = entries))),
|
||||
|
||||
@@ -180,7 +180,8 @@ export const Plugin = define({
|
||||
yield* ctx.skill.transform((draft) => {
|
||||
for (const skill of loaded.skills) draft.add(skill)
|
||||
})
|
||||
yield* ctx.event.subscribe("config.updated").pipe(
|
||||
yield* ctx.event.subscribe().pipe(
|
||||
Stream.filter((event) => event.type === "config.updated"),
|
||||
Stream.runForEach(() =>
|
||||
config.entries().pipe(
|
||||
Effect.tap((entries) => Effect.sync(() => (loaded.entries = entries))),
|
||||
|
||||
@@ -14,7 +14,8 @@ export const Plugin = define({
|
||||
if (selection === false) websearch.default.set(false)
|
||||
if (selection) websearch.default.set(selection.provider)
|
||||
})
|
||||
yield* ctx.event.subscribe("config.updated").pipe(
|
||||
yield* ctx.event.subscribe().pipe(
|
||||
Stream.filter((event) => event.type === "config.updated"),
|
||||
Stream.runForEach(() =>
|
||||
config.entries().pipe(
|
||||
Effect.tap((entries) => Effect.sync(() => (loaded.entries = entries))),
|
||||
|
||||
@@ -59,12 +59,6 @@ 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,
|
||||
@@ -186,7 +180,7 @@ export const make = Effect.fn("PluginHost.make")(function* (plugin: import("../p
|
||||
}),
|
||||
},
|
||||
event: {
|
||||
subscribe,
|
||||
subscribe: () => bus.subscribe().pipe(Stream.filter(EventManifest.isServer)),
|
||||
},
|
||||
integration: {
|
||||
list: () => response(integration.list()),
|
||||
|
||||
@@ -150,6 +150,10 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
|
||||
if (!current) return yield* Effect.die(new Error(`${name} delta before start: ${id}`))
|
||||
if (!current.pending) return undefined
|
||||
const now = yield* Clock.currentTimeMillis
|
||||
if (!force && current.publishedAt === undefined) {
|
||||
current.publishedAt = now
|
||||
return undefined
|
||||
}
|
||||
if (!force && current.publishedAt !== undefined && now - current.publishedAt < deltaBatchInterval)
|
||||
return undefined
|
||||
yield* delta(id, current.pending, current.ordinal)
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { APICallError } from "@ai-sdk/provider"
|
||||
import type { LanguageModelV3, LanguageModelV3StreamPart } from "@ai-sdk/provider"
|
||||
import { createMistral } from "@ai-sdk/mistral"
|
||||
import { AISDK } from "@opencode-ai/core/aisdk"
|
||||
import { SessionRunnerRetry } from "@opencode-ai/core/session/runner/retry"
|
||||
import { toSessionError } from "@opencode-ai/core/session/to-session-error"
|
||||
@@ -277,67 +278,88 @@ it.effect("projects replay metadata onto AI SDK prompt parts", () =>
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("preserves tool result content in AI SDK prompts", () =>
|
||||
it.effect("moves a tool image through the real Mistral provider as a user message", () =>
|
||||
Effect.gen(function* () {
|
||||
const aisdk = yield* AISDK.Service
|
||||
let body: { messages?: unknown[] } | undefined
|
||||
const mockFetch = Object.assign(
|
||||
async (_input: Parameters<typeof fetch>[0], init?: RequestInit) => {
|
||||
body = JSON.parse(String(init?.body))
|
||||
const chunks = [
|
||||
{
|
||||
id: "response-1",
|
||||
created: 0,
|
||||
model: "pixtral-large-latest",
|
||||
choices: [{ index: 0, delta: { content: [{ type: "text", text: "I see it." }] } }],
|
||||
},
|
||||
{
|
||||
id: "response-1",
|
||||
created: 0,
|
||||
model: "pixtral-large-latest",
|
||||
choices: [{ index: 0, delta: {}, finish_reason: "stop" }],
|
||||
usage: { prompt_tokens: 1, completion_tokens: 1, total_tokens: 2 },
|
||||
},
|
||||
]
|
||||
return new Response(chunks.map((chunk) => `data: ${JSON.stringify(chunk)}\n\n`).join(""), {
|
||||
headers: { "Content-Type": "text/event-stream" },
|
||||
})
|
||||
},
|
||||
{ preconnect: fetch.preconnect },
|
||||
)
|
||||
yield* aisdk.hook.sdk((event) => {
|
||||
event.sdk = { languageModel: () => ({ provider: event.model.providerID }) }
|
||||
event.sdk = createMistral({ apiKey: "test", fetch: mockFetch })
|
||||
})
|
||||
|
||||
const resolved = yield* aisdk.model(model("test-ai-sdk"))
|
||||
const prepared = yield* compileRequest(
|
||||
const resolved = yield* aisdk.model({
|
||||
...model("@ai-sdk/mistral"),
|
||||
modelID: Model.ID.make("pixtral-large-latest"),
|
||||
})
|
||||
yield* LLMClient.generate(
|
||||
LLM.request({
|
||||
model: resolved,
|
||||
messages: [
|
||||
Message.user("Inspect the screenshot."),
|
||||
Message.assistant({ type: "tool-call", id: "call_1", name: "screenshot", input: {} }),
|
||||
Message.tool({
|
||||
type: "tool-result",
|
||||
id: "call_1",
|
||||
name: "read",
|
||||
name: "screenshot",
|
||||
result: {
|
||||
type: "content",
|
||||
value: [
|
||||
{ type: "text", text: "attachments" },
|
||||
{ type: "file", uri: "data:image/png;base64,AAAA", mime: "image/png", name: "pixel.png" },
|
||||
{
|
||||
type: "file",
|
||||
uri: "data:application/pdf;charset=utf-8;base64,JVBERg==",
|
||||
mime: "application/pdf",
|
||||
name: "document.pdf",
|
||||
},
|
||||
{ type: "file", uri: "data:audio/mpeg;base64,SUQz", mime: "audio/mpeg", name: "clip.mp3" },
|
||||
{ type: "file", uri: "https://example.com/pixel.png", mime: "image/png" },
|
||||
{ type: "file", uri: "https://example.com/document.pdf", mime: "application/pdf" },
|
||||
{ type: "text", text: "Screenshot captured" },
|
||||
{ type: "file", uri: "data:image/png;base64,AAAA", mime: "image/png", name: "screen.png" },
|
||||
],
|
||||
},
|
||||
}),
|
||||
],
|
||||
}),
|
||||
)
|
||||
).pipe(Effect.provide(client))
|
||||
|
||||
expect(prepared.body.prompt).toEqual([
|
||||
expect(body?.messages).toEqual([
|
||||
{ role: "user", content: [{ type: "text", text: "Inspect the screenshot." }] },
|
||||
{
|
||||
role: "assistant",
|
||||
content: "",
|
||||
tool_calls: [
|
||||
{
|
||||
id: "call_1",
|
||||
type: "function",
|
||||
function: { name: "screenshot", arguments: "{}" },
|
||||
},
|
||||
],
|
||||
},
|
||||
{
|
||||
role: "tool",
|
||||
name: "screenshot",
|
||||
tool_call_id: "call_1",
|
||||
content: '[{"type":"text","text":"Screenshot captured"}]',
|
||||
},
|
||||
{
|
||||
role: "user",
|
||||
content: [
|
||||
{
|
||||
type: "tool-result",
|
||||
toolCallId: "call_1",
|
||||
toolName: "read",
|
||||
output: {
|
||||
type: "content",
|
||||
value: [
|
||||
{ type: "text", text: "attachments" },
|
||||
{ type: "image-data", data: "AAAA", mediaType: "image/png" },
|
||||
{
|
||||
type: "file-data",
|
||||
data: "JVBERg==",
|
||||
mediaType: "application/pdf",
|
||||
filename: "document.pdf",
|
||||
},
|
||||
{ type: "file-data", data: "SUQz", mediaType: "audio/mpeg", filename: "clip.mp3" },
|
||||
{ type: "image-url", url: "https://example.com/pixel.png" },
|
||||
{ type: "file-url", url: "https://example.com/document.pdf" },
|
||||
],
|
||||
},
|
||||
},
|
||||
{ type: "text", text: "Attached media from tool result:" },
|
||||
{ type: "image_url", image_url: "data:image/png;base64,AAAA" },
|
||||
],
|
||||
},
|
||||
])
|
||||
|
||||
@@ -20,55 +20,24 @@ 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("selects one public event type through the plugin context", () =>
|
||||
it.live("exposes public events 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("config.updated")
|
||||
.pipe(Stream.runHead, Effect.forkScoped({ startImmediately: true }))
|
||||
const received = yield* host.event.subscribe().pipe(
|
||||
Stream.filter((event) => event.type === "config.updated"),
|
||||
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
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { describe, expect } from "bun:test"
|
||||
import { Message, SystemPart } from "@opencode-ai/ai"
|
||||
import { DateTime, Effect, Schema, Stream } from "effect"
|
||||
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"
|
||||
@@ -18,8 +18,6 @@ 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"
|
||||
@@ -29,28 +27,6 @@ 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
|
||||
|
||||
@@ -217,16 +217,13 @@ it.effect("batches text deltas and flushes pending text before the terminal even
|
||||
{ discard: true },
|
||||
)
|
||||
|
||||
expect(published.filter((event) => event.type === "session.text.delta").map((event) => event.data)).toMatchObject([
|
||||
{ delta: "one" },
|
||||
])
|
||||
expect(published.filter((event) => event.type === "session.text.delta")).toHaveLength(0)
|
||||
yield* TestClock.adjust("99 millis")
|
||||
expect(published.filter((event) => event.type === "session.text.delta")).toHaveLength(1)
|
||||
expect(published.filter((event) => event.type === "session.text.delta")).toHaveLength(0)
|
||||
yield* TestClock.adjust("1 millis")
|
||||
yield* publisher.publish(LLMEvent.textDelta({ id: "text", text: " four" }))
|
||||
expect(published.filter((event) => event.type === "session.text.delta").map((event) => event.data)).toMatchObject([
|
||||
{ delta: "one" },
|
||||
{ delta: " two three four" },
|
||||
{ delta: "one two three four" },
|
||||
])
|
||||
|
||||
yield* publisher.publish(LLMEvent.textDelta({ id: "text", text: " five" }))
|
||||
@@ -253,7 +250,7 @@ it.effect("batches reasoning deltas and flushes pending reasoning before the ter
|
||||
|
||||
expect(
|
||||
published.filter((event) => event.type === "session.reasoning.delta").map((event) => event.data),
|
||||
).toMatchObject([{ delta: "one" }, { delta: " two three" }])
|
||||
).toMatchObject([{ delta: "one two three" }])
|
||||
expect(published.slice(-2).map((event) => event.type)).toEqual([
|
||||
"session.reasoning.delta",
|
||||
"session.reasoning.ended.1",
|
||||
|
||||
@@ -767,7 +767,7 @@ const verifyEphemeralDeltas = (kind: FragmentKind) =>
|
||||
yield* admit(session, prompt)
|
||||
const bus = yield* Bus.Service
|
||||
const live = fixture.delta
|
||||
? yield* bus.subscribe(fixture.delta).pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
|
||||
? yield* bus.subscribe(fixture.delta).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
|
||||
: undefined
|
||||
yield* Effect.yieldNow
|
||||
yield* TestLLM.push(fixture.completeEvents)
|
||||
@@ -785,7 +785,7 @@ const verifyEphemeralDeltas = (kind: FragmentKind) =>
|
||||
: []
|
||||
if (live) {
|
||||
const streamed = Array.from(yield* Fiber.join(live))
|
||||
expect(streamed).toHaveLength(2)
|
||||
expect(streamed).toHaveLength(1)
|
||||
expect(
|
||||
streamed
|
||||
.map((event) => {
|
||||
|
||||
@@ -1,15 +1,3 @@
|
||||
import type { EventApi } from "@opencode-ai/client/effect/api"
|
||||
import type { OpenCodeEvent } from "@opencode-ai/client/effect"
|
||||
import type { Stream } from "effect"
|
||||
|
||||
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
|
||||
}
|
||||
export interface EventDomain extends Pick<EventApi<unknown>, "subscribe"> {}
|
||||
|
||||
@@ -2,7 +2,6 @@ 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"
|
||||
|
||||
@@ -150,15 +149,13 @@ export function fromPromise(plugin: Plugin) {
|
||||
reload: () => run(host.command.reload()),
|
||||
},
|
||||
event: {
|
||||
subscribe: (type?: PluginEventType) => {
|
||||
const events = type === undefined ? host.event.subscribe() : host.event.subscribe(type)
|
||||
return Stream.toAsyncIterable(
|
||||
events.pipe(
|
||||
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),
|
||||
|
||||
@@ -1,14 +1,3 @@
|
||||
import type { OpenCodeEvent } from "@opencode-ai/client"
|
||||
import type { EventApi } from "@opencode-ai/client/promise/api"
|
||||
|
||||
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
|
||||
}
|
||||
export interface EventDomain extends Pick<EventApi, "subscribe"> {}
|
||||
|
||||
@@ -1,26 +0,0 @@
|
||||
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
@@ -184,14 +184,6 @@ 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
|
||||
|
||||
Reference in New Issue
Block a user