mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-09 19:09:49 -04:00
Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e6bb88bc1d | |||
| fc0cf2a710 | |||
| 1025540fcc | |||
| eb9a683b40 | |||
| 48c26fa039 | |||
| 10d1e04e9b |
@@ -237,6 +237,10 @@ export const layer: Layer.Layer<Service, never, FSUtil.Service | Ripgrep.Service
|
||||
// and does not await the scan; the native background scan starts as soon as
|
||||
// the picker exists. The `wait` gate dedupes concurrent creation.
|
||||
const acquire = Effect.fn("Search.acquire")(function* (cwd: string) {
|
||||
// The opencode test runtime owns an isolated XDG tree that Windows must
|
||||
// remove before process exit, so use ripgrep instead of native FFF there.
|
||||
if (process.env.OPENCODE_TEST_HOME) return undefined
|
||||
|
||||
const available = yield* fffSync("check availability", () => Fff.available()).pipe(
|
||||
Effect.catch((error) => {
|
||||
log.warn("fff availability check failed", { error })
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
export * as Image from "./image"
|
||||
|
||||
import { Context, Effect, Layer, Schema } from "effect"
|
||||
import { Config } from "./config"
|
||||
import { FileSystem } from "./filesystem"
|
||||
|
||||
export class ResizerUnavailableError extends Schema.TaggedErrorClass<ResizerUnavailableError>()(
|
||||
"Image.ResizerUnavailableError",
|
||||
{},
|
||||
) {}
|
||||
|
||||
export class DecodeError extends Schema.TaggedErrorClass<DecodeError>()("Image.DecodeError", {
|
||||
resource: Schema.String,
|
||||
}) {
|
||||
override get message() {
|
||||
return `Image could not be decoded: ${this.resource}`
|
||||
}
|
||||
}
|
||||
|
||||
export class SizeError extends Schema.TaggedErrorClass<SizeError>()("Image.SizeError", {
|
||||
resource: Schema.String,
|
||||
width: Schema.Number,
|
||||
height: Schema.Number,
|
||||
bytes: Schema.Number,
|
||||
maxWidth: Schema.Number,
|
||||
maxHeight: Schema.Number,
|
||||
maxBytes: Schema.Number,
|
||||
}) {
|
||||
override get message() {
|
||||
return `Image ${this.resource} is ${this.width}x${this.height} with base64 size ${this.bytes}, exceeding configured limits ${this.maxWidth}x${this.maxHeight}/${this.maxBytes} bytes`
|
||||
}
|
||||
}
|
||||
|
||||
export interface Interface {
|
||||
readonly normalize: (
|
||||
resource: string,
|
||||
content: FileSystem.BinaryContent,
|
||||
) => Effect.Effect<FileSystem.BinaryContent, ResizerUnavailableError | DecodeError | SizeError>
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/Image") {}
|
||||
|
||||
export const layer = Layer.effect(
|
||||
Service,
|
||||
Effect.gen(function* () {
|
||||
const config = yield* Config.Service
|
||||
const loadAdapter = yield* Effect.cached(
|
||||
Effect.tryPromise({
|
||||
try: () => import("./image/photon"),
|
||||
catch: () => new ResizerUnavailableError(),
|
||||
}).pipe(Effect.flatMap((adapter) => adapter.make)),
|
||||
)
|
||||
const normalize = Effect.fn("Image.normalize")(function* (resource: string, content: FileSystem.BinaryContent) {
|
||||
const image = Object.assign(
|
||||
{},
|
||||
...(yield* config.entries()).flatMap((entry) =>
|
||||
entry.type === "document" && entry.info.attachments?.image ? [entry.info.attachments.image] : [],
|
||||
),
|
||||
)
|
||||
const normalize = yield* loadAdapter
|
||||
return yield* normalize(resource, content, {
|
||||
autoResize: image.auto_resize ?? true,
|
||||
maxWidth: image.max_width ?? 2_000,
|
||||
maxHeight: image.max_height ?? 2_000,
|
||||
maxBase64Bytes: image.max_base64_bytes ?? 5 * 1024 * 1024,
|
||||
})
|
||||
})
|
||||
return Service.of({ normalize })
|
||||
}),
|
||||
)
|
||||
|
||||
export const locationLayer = layer.pipe(Layer.provide(Config.locationLayer))
|
||||
@@ -0,0 +1,94 @@
|
||||
// @ts-ignore Bun's static file import is embedded by `bun build --compile`; some consumers also declare *.wasm.
|
||||
import photonWasm from "@silvia-odwyer/photon-node/photon_rs_bg.wasm" with { type: "file" }
|
||||
import { Effect } from "effect"
|
||||
import path from "node:path"
|
||||
import { fileURLToPath } from "node:url"
|
||||
import { FileSystem } from "../filesystem"
|
||||
import { DecodeError, ResizerUnavailableError, SizeError } from "../image"
|
||||
|
||||
const JPEG_QUALITIES = [80, 85, 70, 55, 40]
|
||||
|
||||
export const make = Effect.gen(function* () {
|
||||
;(globalThis as typeof globalThis & { __OPENCODE_PHOTON_WASM_PATH?: string }).__OPENCODE_PHOTON_WASM_PATH =
|
||||
path.isAbsolute(photonWasm) ? photonWasm : fileURLToPath(new URL(photonWasm, import.meta.url))
|
||||
const loadPhoton = yield* Effect.cached(
|
||||
Effect.tryPromise({
|
||||
try: () => import("@silvia-odwyer/photon-node"),
|
||||
catch: () => new ResizerUnavailableError(),
|
||||
}),
|
||||
)
|
||||
return Effect.fn("Image.Photon.normalize")(function* (
|
||||
resource: string,
|
||||
content: FileSystem.BinaryContent,
|
||||
limits: {
|
||||
readonly autoResize: boolean
|
||||
readonly maxWidth: number
|
||||
readonly maxHeight: number
|
||||
readonly maxBase64Bytes: number
|
||||
},
|
||||
) {
|
||||
const photon = yield* loadPhoton
|
||||
const decoded = yield* Effect.try({
|
||||
try: () => photon.PhotonImage.new_from_byteslice(Buffer.from(content.content, "base64")),
|
||||
catch: () => new DecodeError({ resource }),
|
||||
})
|
||||
try {
|
||||
const width = decoded.get_width()
|
||||
const height = decoded.get_height()
|
||||
const bytes = Buffer.byteLength(content.content, "utf-8")
|
||||
if (width <= limits.maxWidth && height <= limits.maxHeight && bytes <= limits.maxBase64Bytes) return content
|
||||
if (!limits.autoResize)
|
||||
return yield* new SizeError({
|
||||
resource,
|
||||
width,
|
||||
height,
|
||||
bytes,
|
||||
maxWidth: limits.maxWidth,
|
||||
maxHeight: limits.maxHeight,
|
||||
maxBytes: limits.maxBase64Bytes,
|
||||
})
|
||||
const scale = Math.min(1, limits.maxWidth / width, limits.maxHeight / height)
|
||||
const sizes = Array.from({ length: 32 }).reduce<Array<{ width: number; height: number }>>((acc) => {
|
||||
const previous = acc.at(-1) ?? {
|
||||
width: Math.max(1, Math.round(width * scale)),
|
||||
height: Math.max(1, Math.round(height * scale)),
|
||||
}
|
||||
const next =
|
||||
acc.length === 0
|
||||
? previous
|
||||
: {
|
||||
width: previous.width === 1 ? 1 : Math.max(1, Math.floor(previous.width * 0.75)),
|
||||
height: previous.height === 1 ? 1 : Math.max(1, Math.floor(previous.height * 0.75)),
|
||||
}
|
||||
return acc.some((item) => item.width === next.width && item.height === next.height) ? acc : [...acc, next]
|
||||
}, [])
|
||||
for (const size of sizes) {
|
||||
const resized = photon.resize(decoded, size.width, size.height, photon.SamplingFilter.Lanczos3)
|
||||
try {
|
||||
const encoders: Array<readonly [mime: string, encode: () => Uint8Array]> = [
|
||||
["image/png", () => resized.get_bytes()],
|
||||
...JPEG_QUALITIES.map((quality) => ["image/jpeg", () => resized.get_bytes_jpeg(quality)] as const),
|
||||
]
|
||||
for (const [mime, encode] of encoders) {
|
||||
const candidate = Buffer.from(encode()).toString("base64")
|
||||
if (Buffer.byteLength(candidate, "utf-8") <= limits.maxBase64Bytes)
|
||||
return new FileSystem.BinaryContent({ type: "binary", content: candidate, encoding: "base64", mime })
|
||||
}
|
||||
} finally {
|
||||
resized.free()
|
||||
}
|
||||
}
|
||||
return yield* new SizeError({
|
||||
resource,
|
||||
width,
|
||||
height,
|
||||
bytes,
|
||||
maxWidth: limits.maxWidth,
|
||||
maxHeight: limits.maxHeight,
|
||||
maxBytes: limits.maxBase64Bytes,
|
||||
})
|
||||
} finally {
|
||||
decoded.free()
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -28,6 +28,7 @@ import { Pty } from "./pty"
|
||||
import { SkillV2 } from "./skill"
|
||||
import { SkillGuidance } from "./skill/guidance"
|
||||
import { BuiltInTools } from "./tool/builtins"
|
||||
import { Image } from "./image"
|
||||
import { ToolRegistry } from "./tool/registry"
|
||||
import { ApplicationTools } from "./tool/application-tools"
|
||||
import { ToolOutputStore } from "./tool-output-store"
|
||||
@@ -71,6 +72,7 @@ export class LocationServiceMap extends LayerMap.Service<LocationServiceMap>()("
|
||||
Layer.provide(base),
|
||||
)
|
||||
const services = Layer.mergeAll(base, resources, permissionsAndTools)
|
||||
const image = Image.layer.pipe(Layer.provide(services))
|
||||
const mutation = FileMutation.locationLayer.pipe(Layer.provide(services))
|
||||
const searches = LocationSearch.layer.pipe(Layer.provide(Ripgrep.layer), Layer.provide(services))
|
||||
const skillGuidance = SkillGuidance.locationLayer.pipe(Layer.provide(services))
|
||||
@@ -83,6 +85,7 @@ export class LocationServiceMap extends LayerMap.Service<LocationServiceMap>()("
|
||||
Layer.provide(resources),
|
||||
Layer.provide(todos),
|
||||
Layer.provide(questions),
|
||||
Layer.provide(image),
|
||||
)
|
||||
const model = SessionRunnerModel.locationLayer.pipe(Layer.provide(services))
|
||||
const runner = SessionRunnerLLM.defaultLayer.pipe(
|
||||
@@ -90,9 +93,18 @@ export class LocationServiceMap extends LayerMap.Service<LocationServiceMap>()("
|
||||
Layer.provide(model),
|
||||
Layer.provide(skillGuidance),
|
||||
)
|
||||
return Layer.mergeAll(services, mutation, searches, resources, todos, questions, model, runner, builtInTools).pipe(
|
||||
Layer.fresh,
|
||||
)
|
||||
return Layer.mergeAll(
|
||||
services,
|
||||
image,
|
||||
mutation,
|
||||
searches,
|
||||
resources,
|
||||
todos,
|
||||
questions,
|
||||
model,
|
||||
runner,
|
||||
builtInTools,
|
||||
).pipe(Layer.fresh)
|
||||
},
|
||||
idleTimeToLive: "60 minutes",
|
||||
dependencies: [
|
||||
|
||||
@@ -3,7 +3,7 @@ export { Agent } from "./agent"
|
||||
export { Model } from "./model"
|
||||
export { OpenCode } from "./opencode"
|
||||
export { Session } from "./session"
|
||||
export * as Tool from "./tool"
|
||||
export { Tool } from "./tool"
|
||||
export { Location } from "./location"
|
||||
export { Prompt } from "../session/prompt"
|
||||
export { AbsolutePath } from "../schema"
|
||||
|
||||
@@ -4,7 +4,7 @@ import { Effect, Scope } from "effect"
|
||||
import type { AnyTool, RegistrationError } from "../tool/tool"
|
||||
|
||||
export { Failure, RegistrationError, make } from "../tool/tool"
|
||||
export type { AnyTool, Content, Context } from "../tool/tool"
|
||||
export type { AnyTool, Content, Context, Definition } from "../tool/tool"
|
||||
|
||||
export interface Interface {
|
||||
/**
|
||||
|
||||
@@ -1,268 +1,384 @@
|
||||
export * as SessionRunCoordinator from "./run-coordinator"
|
||||
|
||||
import { Cause, Context, Deferred, Effect, Exit, Fiber, FiberSet, Layer, Scope } from "effect"
|
||||
import {
|
||||
Cause,
|
||||
Context,
|
||||
Data,
|
||||
Deferred,
|
||||
Effect,
|
||||
Equal,
|
||||
Exit,
|
||||
Fiber,
|
||||
FiberSet,
|
||||
Layer,
|
||||
Scope,
|
||||
SynchronizedRef,
|
||||
} from "effect"
|
||||
import { SessionRunner } from "./runner"
|
||||
import { SessionSchema } from "./schema"
|
||||
|
||||
export type Mode = "run" | "wake"
|
||||
|
||||
/** Why one drain generation should run. Explicit runs dominate advisory wakes when demands coalesce. */
|
||||
type Demand = { readonly _tag: "run" } | { readonly _tag: "wake"; readonly seq?: number }
|
||||
|
||||
/**
|
||||
* Runs at most one drain chain per key while allowing different keys to drain concurrently.
|
||||
*
|
||||
* For each key:
|
||||
*
|
||||
* idle --run/wake--> draining --run/wake--> draining + one coalesced rerun --> idle
|
||||
*
|
||||
* `run` is an explicit drain request. It starts a chain or joins the current chain and
|
||||
* upgrades a pending follow-up so the caller receives explicit-run semantics.
|
||||
*
|
||||
* `wake` reports that durable work may now be available. It starts a chain while idle or
|
||||
* requests one coalesced follow-up while draining. Repeated wakes collapse together.
|
||||
*
|
||||
* `interrupt` stops the current ownership chain. Advisory wakes from before the interrupt
|
||||
* boundary are suppressed; advisory wakes after the boundary run after cleanup.
|
||||
*/
|
||||
export interface Coordinator<Key, A, E> {
|
||||
/** Starts or joins one explicit drain generation. */
|
||||
readonly run: (key: Key) => Effect.Effect<A, E>
|
||||
/** Coalesces one wake-up after durable work is recorded. */
|
||||
readonly wake: (key: Key, seq?: number) => Effect.Effect<void>
|
||||
/** Waits until the current ownership chain settles. */
|
||||
readonly awaitIdle: (key: Key) => Effect.Effect<void, E>
|
||||
/** Interrupts the active ownership chain without automatically draining pending wakes. */
|
||||
readonly interrupt: (key: Key, seq?: number) => Effect.Effect<void>
|
||||
}
|
||||
|
||||
/** One Session's process-local execution lane: one active demand and at most one coalesced follow-up. */
|
||||
type Entry<A, E> = {
|
||||
readonly done: Deferred.Deferred<A, E>
|
||||
readonly settled: Deferred.Deferred<Exit.Exit<A, E>>
|
||||
current: Demand
|
||||
pending?: Demand
|
||||
explicitWaiter?: Deferred.Deferred<A, E>
|
||||
interruptSeq?: number
|
||||
owner?: Fiber.Fiber<void, never>
|
||||
stopping: boolean
|
||||
/** @internal */
|
||||
export class Demand extends Data.Class<{
|
||||
readonly explicit: boolean
|
||||
readonly wakeSeq?: number
|
||||
readonly unsequencedWake: boolean
|
||||
}> {
|
||||
static readonly empty = new Demand({ explicit: false, wakeSeq: undefined, unsequencedWake: false })
|
||||
static readonly run = nonEmpty(new Demand({ explicit: true, wakeSeq: undefined, unsequencedWake: false }))
|
||||
|
||||
static wake(seq?: number) {
|
||||
return nonEmpty(new Demand({ explicit: false, wakeSeq: seq, unsequencedWake: seq === undefined }))
|
||||
}
|
||||
|
||||
combine(other: Demand) {
|
||||
return new Demand({
|
||||
explicit: this.explicit || other.explicit,
|
||||
wakeSeq:
|
||||
this.wakeSeq === undefined
|
||||
? other.wakeSeq
|
||||
: other.wakeSeq === undefined
|
||||
? this.wakeSeq
|
||||
: Math.max(this.wakeSeq, other.wakeSeq),
|
||||
unsequencedWake: this.unsequencedWake || other.unsequencedWake,
|
||||
})
|
||||
}
|
||||
|
||||
afterBoundary(boundary?: number) {
|
||||
return new Demand({
|
||||
explicit: false,
|
||||
wakeSeq:
|
||||
boundary !== undefined && this.wakeSeq !== undefined && this.wakeSeq > boundary ? this.wakeSeq : undefined,
|
||||
unsequencedWake: false,
|
||||
})
|
||||
}
|
||||
|
||||
isNonEmpty(): this is NonEmptyDemand {
|
||||
return this.explicit || this.wakeSeq !== undefined || this.unsequencedWake
|
||||
}
|
||||
|
||||
get mode(): Mode {
|
||||
return this.explicit ? "run" : "wake"
|
||||
}
|
||||
}
|
||||
|
||||
/** Combines follow-up demand: runs dominate, while wakes retain the newest durable admission sequence. */
|
||||
const coalesce = (left: Demand | undefined, right: Demand): Demand => {
|
||||
if (left?._tag === "run" || right._tag === "run") return { _tag: "run" }
|
||||
return { _tag: "wake", seq: maxSeq(left?.seq, right.seq) }
|
||||
type NonEmptyDemand = Demand &
|
||||
({ readonly explicit: true } | { readonly wakeSeq: number } | { readonly unsequencedWake: true })
|
||||
|
||||
function nonEmpty(demand: Demand): NonEmptyDemand {
|
||||
if (!demand.isNonEmpty()) throw new Error("Session run demand must not be empty")
|
||||
return demand
|
||||
}
|
||||
|
||||
const maxSeq = (left: number | undefined, right: number | undefined) => {
|
||||
if (left === undefined) return right
|
||||
if (right === undefined) return left
|
||||
return Math.max(left, right)
|
||||
type Lifecycle =
|
||||
| { readonly _tag: "Running"; readonly token: object; readonly owner: Deferred.Deferred<Fiber.Fiber<void>> }
|
||||
| {
|
||||
readonly _tag: "Stopping"
|
||||
readonly token: object
|
||||
readonly owner: Deferred.Deferred<Fiber.Fiber<void>>
|
||||
readonly boundary?: number
|
||||
}
|
||||
|
||||
type Lane<A, E> = {
|
||||
readonly current: NonEmptyDemand
|
||||
readonly pending: Demand
|
||||
readonly lifecycle: Lifecycle
|
||||
readonly terminal: Deferred.Deferred<Exit.Exit<A, E>>
|
||||
readonly waiter?: Deferred.Deferred<Exit.Exit<A, E>>
|
||||
}
|
||||
|
||||
type State<Key, A, E> = {
|
||||
readonly closed: boolean
|
||||
readonly lanes: ReadonlyMap<Key, Lane<A, E>>
|
||||
readonly interruptSeq: ReadonlyMap<Key, number>
|
||||
}
|
||||
|
||||
type Start<Key, A, E> = {
|
||||
readonly key: Key
|
||||
readonly demand: NonEmptyDemand
|
||||
readonly successor: boolean
|
||||
readonly token: object
|
||||
readonly owner: Deferred.Deferred<Fiber.Fiber<void>>
|
||||
readonly ready: Deferred.Deferred<void>
|
||||
readonly terminal: Deferred.Deferred<Exit.Exit<A, E>>
|
||||
}
|
||||
|
||||
type RunRequest<Key, A, E> =
|
||||
| { readonly _tag: "Closed" }
|
||||
| { readonly _tag: "Await"; readonly terminal: Deferred.Deferred<Exit.Exit<A, E>> }
|
||||
| { readonly _tag: "Retry"; readonly terminal: Deferred.Deferred<Exit.Exit<A, E>> }
|
||||
| { readonly _tag: "Start"; readonly start: Start<Key, A, E>; readonly terminal: Deferred.Deferred<Exit.Exit<A, E>> }
|
||||
|
||||
type Completion<Key, A, E> = {
|
||||
readonly start?: Start<Key, A, E>
|
||||
readonly terminal?: Deferred.Deferred<Exit.Exit<A, E>>
|
||||
readonly waiter?: Deferred.Deferred<Exit.Exit<A, E>>
|
||||
readonly report?: Cause.Cause<E>
|
||||
}
|
||||
|
||||
/** Constructs a scoped coordinator. Every in-memory transition is synchronous. */
|
||||
export const make = <Key, A, E>(options: {
|
||||
readonly drain: (key: Key, mode: Mode) => Effect.Effect<A, E>
|
||||
readonly onFailure?: (key: Key, cause: Cause.Cause<E>) => Effect.Effect<void>
|
||||
}): Effect.Effect<Coordinator<Key, A, E>, never, Scope.Scope> =>
|
||||
Effect.gen(function* () {
|
||||
const active = new Map<Key, Entry<A, E>>()
|
||||
const interruptSeq = new Map<Key, number>()
|
||||
const report = yield* FiberSet.makeRuntime<never, void, never>()
|
||||
const state = yield* SynchronizedRef.make<State<Key, A, E>>({
|
||||
closed: false,
|
||||
lanes: new Map(),
|
||||
interruptSeq: new Map(),
|
||||
})
|
||||
const fork = yield* FiberSet.makeRuntime<never, void, never>()
|
||||
const shutdown = Deferred.makeUnsafe<void>()
|
||||
let closed = false
|
||||
yield* Effect.addFinalizer(() =>
|
||||
Effect.sync(() => {
|
||||
closed = true
|
||||
Deferred.doneUnsafe(shutdown, Effect.void)
|
||||
active.clear()
|
||||
interruptSeq.clear()
|
||||
}),
|
||||
)
|
||||
|
||||
const makeEntry = (current: Demand, explicitWaiter?: Deferred.Deferred<A, E>): Entry<A, E> => ({
|
||||
done: Deferred.makeUnsafe<A, E>(),
|
||||
settled: Deferred.makeUnsafe<Exit.Exit<A, E>>(),
|
||||
current,
|
||||
explicitWaiter,
|
||||
stopping: false,
|
||||
})
|
||||
|
||||
const start = (key: Key, entry: Entry<A, E>, demand: Demand, successor = false) => {
|
||||
const ready = Deferred.makeUnsafe<void>()
|
||||
const drain = Effect.suspend(() => options.drain(key, demand._tag))
|
||||
// Initial work retains immediate-start behavior but cannot run before ownership is published.
|
||||
// Observer-started successors yield once so synchronous drains cannot recurse on the JS stack.
|
||||
const owner = fork(
|
||||
(successor
|
||||
? Effect.yieldNow.pipe(Effect.andThen(drain))
|
||||
: Deferred.await(ready).pipe(Effect.andThen(drain))
|
||||
).pipe(
|
||||
Effect.onExit((exit) => Effect.sync(() => settle(key, entry, demand, exit))),
|
||||
Effect.exit,
|
||||
Effect.asVoid,
|
||||
),
|
||||
)
|
||||
entry.owner = owner
|
||||
if (!successor) Deferred.doneUnsafe(ready, Effect.void)
|
||||
const updateLane = (current: State<Key, A, E>, key: Key, lane?: Lane<A, E>): State<Key, A, E> => {
|
||||
const lanes = new Map(current.lanes)
|
||||
if (lane === undefined) lanes.delete(key)
|
||||
else lanes.set(key, lane)
|
||||
return { ...current, lanes }
|
||||
}
|
||||
|
||||
const settle = (key: Key, entry: Entry<A, E>, demand: Demand, exit: Exit.Exit<A, E>) => {
|
||||
if (closed) {
|
||||
Deferred.doneUnsafe(entry.done, exit)
|
||||
Deferred.doneUnsafe(entry.settled, Effect.succeed(exit))
|
||||
return
|
||||
const start = (input: {
|
||||
readonly state: State<Key, A, E>
|
||||
readonly key: Key
|
||||
readonly demand: NonEmptyDemand
|
||||
readonly terminal?: Deferred.Deferred<Exit.Exit<A, E>>
|
||||
readonly waiter?: Deferred.Deferred<Exit.Exit<A, E>>
|
||||
readonly successor?: boolean
|
||||
}) => {
|
||||
const instruction: Start<Key, A, E> = {
|
||||
key: input.key,
|
||||
demand: input.demand,
|
||||
successor: input.successor ?? false,
|
||||
token: {},
|
||||
owner: Deferred.makeUnsafe<Fiber.Fiber<void>>(),
|
||||
ready: Deferred.makeUnsafe<void>(),
|
||||
terminal: input.terminal ?? Deferred.makeUnsafe<Exit.Exit<A, E>>(),
|
||||
}
|
||||
if (demand._tag === "run" && entry.explicitWaiter !== undefined) {
|
||||
Deferred.doneUnsafe(entry.explicitWaiter, exit)
|
||||
entry.explicitWaiter = undefined
|
||||
return {
|
||||
state: updateLane(input.state, input.key, {
|
||||
current: input.demand,
|
||||
pending: Demand.empty,
|
||||
lifecycle: { _tag: "Running", token: instruction.token, owner: instruction.owner },
|
||||
terminal: instruction.terminal,
|
||||
waiter: input.waiter,
|
||||
}),
|
||||
start: instruction,
|
||||
result: instruction.terminal,
|
||||
}
|
||||
if (entry.stopping && demand._tag === "wake" && entry.explicitWaiter !== undefined) {
|
||||
Deferred.doneUnsafe(entry.explicitWaiter, exit)
|
||||
entry.explicitWaiter = undefined
|
||||
}
|
||||
if (active.get(key) !== entry) {
|
||||
Deferred.doneUnsafe(entry.done, exit)
|
||||
Deferred.doneUnsafe(entry.settled, Effect.succeed(exit))
|
||||
return
|
||||
}
|
||||
if (exit._tag === "Success" && !entry.stopping) {
|
||||
if (entry.pending !== undefined) {
|
||||
const pending = entry.pending
|
||||
entry.pending = undefined
|
||||
entry.current = pending
|
||||
start(key, entry, pending, true)
|
||||
return
|
||||
}
|
||||
|
||||
const launch = (instruction: Start<Key, A, E>) =>
|
||||
Effect.gen(function* () {
|
||||
const fiber = fork(
|
||||
Deferred.await(instruction.ready).pipe(
|
||||
Effect.andThen(instruction.successor ? Effect.yieldNow : Effect.void),
|
||||
Effect.andThen(Effect.suspend(() => options.drain(instruction.key, instruction.demand.mode))),
|
||||
Effect.onExit((exit) => complete(instruction.key, instruction.token, exit)),
|
||||
Effect.exit,
|
||||
Effect.asVoid,
|
||||
),
|
||||
)
|
||||
yield* Deferred.succeed(instruction.owner, fiber)
|
||||
yield* Deferred.succeed(instruction.ready, undefined)
|
||||
})
|
||||
|
||||
const complete = (key: Key, token: object, exit: Exit.Exit<A, E>): Effect.Effect<void> => {
|
||||
return SynchronizedRef.modify(state, (current): readonly [Completion<Key, A, E>, State<Key, A, E>] => {
|
||||
const lane = current.lanes.get(key)
|
||||
if (lane === undefined || lane.lifecycle.token !== token) return [{}, current]
|
||||
|
||||
const deliberateInterrupt =
|
||||
lane.lifecycle._tag === "Stopping" && exit._tag === "Failure" && Cause.hasInterruptsOnly(exit.cause)
|
||||
const report =
|
||||
exit._tag === "Failure" && !deliberateInterrupt && !lane.current.explicit ? exit.cause : undefined
|
||||
const completesWaiter = lane.current.explicit || (lane.lifecycle._tag === "Stopping" && !lane.current.explicit)
|
||||
const waiter = completesWaiter ? undefined : lane.waiter
|
||||
|
||||
if (exit._tag === "Success" && lane.lifecycle._tag === "Running" && lane.pending.isNonEmpty()) {
|
||||
const next = start({
|
||||
state: current,
|
||||
key,
|
||||
demand: lane.pending,
|
||||
terminal: lane.terminal,
|
||||
waiter,
|
||||
successor: true,
|
||||
})
|
||||
return [{ start: next.start, waiter: completesWaiter ? lane.waiter : undefined, report }, next.state]
|
||||
}
|
||||
active.delete(key)
|
||||
Deferred.doneUnsafe(entry.done, exit)
|
||||
Deferred.doneUnsafe(entry.settled, Effect.succeed(exit))
|
||||
return
|
||||
}
|
||||
|
||||
const successor = entry.pending !== undefined ? makeEntry(entry.pending, entry.explicitWaiter) : undefined
|
||||
if (successor === undefined) active.delete(key)
|
||||
else active.set(key, successor)
|
||||
if (successor !== undefined) start(key, successor, successor.current, true)
|
||||
Deferred.doneUnsafe(entry.done, exit)
|
||||
Deferred.doneUnsafe(entry.settled, Effect.succeed(exit))
|
||||
if (
|
||||
exit._tag === "Failure" &&
|
||||
!(entry.stopping && Cause.hasInterruptsOnly(exit.cause)) &&
|
||||
demand._tag === "wake" &&
|
||||
options.onFailure !== undefined
|
||||
) {
|
||||
report(Effect.suspend(() => options.onFailure!(key, exit.cause)))
|
||||
}
|
||||
const next = lane.pending.isNonEmpty()
|
||||
? start({ state: current, key, demand: lane.pending, waiter, successor: true })
|
||||
: { state: updateLane(current, key) }
|
||||
return [
|
||||
{
|
||||
start: "start" in next ? next.start : undefined,
|
||||
terminal: lane.terminal,
|
||||
waiter: completesWaiter ? lane.waiter : undefined,
|
||||
report,
|
||||
},
|
||||
next.state,
|
||||
]
|
||||
}).pipe(Effect.flatMap((instruction) => executeCompletion(key, exit, instruction)))
|
||||
}
|
||||
|
||||
const executeCompletion = (key: Key, exit: Exit.Exit<A, E>, instruction: Completion<Key, A, E>) =>
|
||||
Effect.gen(function* () {
|
||||
if (instruction.start !== undefined) yield* launch(instruction.start)
|
||||
if (instruction.waiter !== undefined) yield* Deferred.succeed(instruction.waiter, exit)
|
||||
if (instruction.terminal !== undefined) yield* Deferred.succeed(instruction.terminal, exit)
|
||||
if (instruction.report !== undefined && options.onFailure !== undefined) {
|
||||
const onFailure = options.onFailure
|
||||
const cause = instruction.report
|
||||
fork(Effect.suspend(() => onFailure(key, cause)).pipe(Effect.exit, Effect.asVoid))
|
||||
}
|
||||
})
|
||||
|
||||
const awaitTerminal = (terminal: Deferred.Deferred<Exit.Exit<A, E>>) =>
|
||||
Effect.raceFirst(
|
||||
Deferred.await(terminal).pipe(
|
||||
Effect.flatMap(
|
||||
Exit.match({
|
||||
onSuccess: Effect.succeed,
|
||||
onFailure: Effect.failCause,
|
||||
}),
|
||||
),
|
||||
),
|
||||
Deferred.await(shutdown).pipe(Effect.andThen(Effect.interrupt)),
|
||||
)
|
||||
|
||||
const run = (key: Key): Effect.Effect<A, E> =>
|
||||
Effect.suspend(() =>
|
||||
Effect.uninterruptibleMask((restore) => {
|
||||
return SynchronizedRef.modify(state, (current): readonly [RunRequest<Key, A, E>, State<Key, A, E>] => {
|
||||
if (current.closed) return [{ _tag: "Closed" }, current]
|
||||
const lane = current.lanes.get(key)
|
||||
if (lane?.lifecycle._tag === "Stopping") return [{ _tag: "Retry", terminal: lane.terminal }, current]
|
||||
if (lane?.current.explicit) return [{ _tag: "Await", terminal: lane.terminal }, current]
|
||||
if (lane !== undefined) {
|
||||
const terminal = lane.waiter ?? Deferred.makeUnsafe<Exit.Exit<A, E>>()
|
||||
const pending = lane.pending.combine(Demand.run)
|
||||
if (Equal.equals(pending, lane.pending) && lane.waiter !== undefined)
|
||||
return [{ _tag: "Await", terminal }, current]
|
||||
return [{ _tag: "Await", terminal }, updateLane(current, key, { ...lane, pending, waiter: terminal })]
|
||||
}
|
||||
const next = start({ state: current, key, demand: Demand.run })
|
||||
return [{ _tag: "Start", start: next.start, terminal: next.result }, next.state]
|
||||
}).pipe(
|
||||
Effect.flatMap((request) => {
|
||||
if (request._tag === "Closed") return Effect.interrupt
|
||||
if (request._tag === "Start")
|
||||
return launch(request.start).pipe(Effect.andThen(awaitTerminal(request.terminal)))
|
||||
if (request._tag === "Await") return awaitTerminal(request.terminal)
|
||||
return Effect.raceFirst(
|
||||
Deferred.await(request.terminal).pipe(Effect.as(true)),
|
||||
Deferred.await(shutdown).pipe(Effect.as(false)),
|
||||
).pipe(Effect.flatMap((retry) => (retry ? run(key) : Effect.interrupt)))
|
||||
}),
|
||||
restore,
|
||||
)
|
||||
}),
|
||||
)
|
||||
|
||||
const wake = (key: Key, seq?: number) =>
|
||||
Effect.sync(() => {
|
||||
if (closed) return
|
||||
if (!isAfterInterrupt(key, seq)) return
|
||||
const entry = active.get(key)
|
||||
if (entry !== undefined) {
|
||||
if (!acceptsWake(entry, seq)) return
|
||||
entry.pending = coalesce(entry.pending, { _tag: "wake", seq })
|
||||
return
|
||||
}
|
||||
Effect.uninterruptible(
|
||||
Effect.suspend(() => {
|
||||
return SynchronizedRef.modify(state, (current): readonly [Start<Key, A, E> | undefined, State<Key, A, E>] => {
|
||||
if (current.closed) return [undefined, current]
|
||||
const boundary = current.interruptSeq.get(key)
|
||||
if (boundary !== undefined && (seq === undefined || seq <= boundary)) return [undefined, current]
|
||||
const lane = current.lanes.get(key)
|
||||
if (lane === undefined) {
|
||||
const next = start({ state: current, key, demand: Demand.wake(seq) })
|
||||
return [next.start, next.state]
|
||||
}
|
||||
if (
|
||||
lane.lifecycle._tag === "Stopping" &&
|
||||
(lane.lifecycle.boundary === undefined || seq === undefined || seq <= lane.lifecycle.boundary)
|
||||
)
|
||||
return [undefined, current]
|
||||
const pending = lane.pending.combine(Demand.wake(seq))
|
||||
if (Equal.equals(pending, lane.pending)) return [undefined, current]
|
||||
return [undefined, updateLane(current, key, { ...lane, pending })]
|
||||
}).pipe(Effect.flatMap((instruction) => (instruction === undefined ? Effect.void : launch(instruction))))
|
||||
}),
|
||||
)
|
||||
|
||||
const next = makeEntry({ _tag: "wake", seq })
|
||||
active.set(key, next)
|
||||
start(key, next, next.current)
|
||||
})
|
||||
const interrupt = (key: Key, seq?: number) =>
|
||||
Effect.uninterruptible(
|
||||
SynchronizedRef.modify(state, (current) => {
|
||||
if (current.closed) return [undefined, current] as const
|
||||
const latest = current.interruptSeq.get(key)
|
||||
const lane = current.lanes.get(key)
|
||||
if (seq !== undefined && latest !== undefined && seq <= latest)
|
||||
return [lane?.lifecycle._tag === "Stopping" ? lane.lifecycle.owner : undefined, current] as const
|
||||
|
||||
const bounded = (() => {
|
||||
if (seq === undefined) return current
|
||||
const interruptSeq = new Map(current.interruptSeq)
|
||||
interruptSeq.set(key, seq)
|
||||
return { ...current, interruptSeq }
|
||||
})()
|
||||
if (lane === undefined) return [undefined, bounded] as const
|
||||
if (
|
||||
!lane.current.explicit &&
|
||||
seq !== undefined &&
|
||||
lane.current.wakeSeq !== undefined &&
|
||||
lane.current.wakeSeq > seq
|
||||
)
|
||||
return [undefined, bounded] as const
|
||||
|
||||
const pending = lane.current.afterBoundary(seq).combine(lane.pending.afterBoundary(seq))
|
||||
const boundary =
|
||||
lane.lifecycle._tag === "Stopping" && lane.lifecycle.boundary !== undefined && seq !== undefined
|
||||
? Math.max(lane.lifecycle.boundary, seq)
|
||||
: lane.lifecycle._tag === "Stopping" && seq === undefined
|
||||
? lane.lifecycle.boundary
|
||||
: seq
|
||||
return [
|
||||
lane.lifecycle.owner,
|
||||
updateLane(bounded, key, {
|
||||
...lane,
|
||||
pending,
|
||||
lifecycle: { _tag: "Stopping", token: lane.lifecycle.token, owner: lane.lifecycle.owner, boundary },
|
||||
}),
|
||||
] as const
|
||||
}).pipe(
|
||||
Effect.flatMap((owner) =>
|
||||
owner === undefined ? Effect.void : Deferred.await(owner).pipe(Effect.flatMap(Fiber.interrupt)),
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
const awaitIdle = (key: Key): Effect.Effect<void, E> =>
|
||||
Effect.gen(function* () {
|
||||
let firstFailure: Cause.Cause<E> | undefined
|
||||
while (!closed) {
|
||||
const entry = active.get(key)
|
||||
if (entry === undefined) break
|
||||
let failure: Cause.Cause<E> | undefined
|
||||
while (true) {
|
||||
const terminal = (yield* SynchronizedRef.get(state)).lanes.get(key)?.terminal
|
||||
if (terminal === undefined) break
|
||||
const exit = yield* Effect.raceFirst(
|
||||
Deferred.await(entry.settled),
|
||||
Deferred.await(terminal),
|
||||
Deferred.await(shutdown).pipe(Effect.as(Exit.void)),
|
||||
)
|
||||
if (closed) break
|
||||
if (exit._tag === "Failure" && firstFailure === undefined) firstFailure = exit.cause
|
||||
if (exit._tag === "Failure" && failure === undefined) failure = exit.cause
|
||||
}
|
||||
if (firstFailure !== undefined) return yield* Effect.failCause(firstFailure)
|
||||
if (failure !== undefined) return yield* Effect.failCause(failure)
|
||||
})
|
||||
|
||||
const interrupt = (key: Key, seq?: number): Effect.Effect<void> =>
|
||||
Effect.suspend(() => {
|
||||
const entry = active.get(key)
|
||||
const latest = interruptSeq.get(key)
|
||||
if (seq !== undefined && latest !== undefined && seq <= latest)
|
||||
return entry?.stopping && entry.owner !== undefined ? Fiber.interrupt(entry.owner) : Effect.void
|
||||
if (seq !== undefined) interruptSeq.set(key, seq)
|
||||
if (entry?.owner === undefined) return Effect.void
|
||||
if (
|
||||
seq !== undefined &&
|
||||
entry.current._tag === "wake" &&
|
||||
entry.current.seq !== undefined &&
|
||||
entry.current.seq > seq
|
||||
)
|
||||
return Effect.void
|
||||
if (entry.stopping) {
|
||||
entry.interruptSeq = maxSeq(entry.interruptSeq, seq)
|
||||
suppressPendingAtOrBefore(entry, seq)
|
||||
return Fiber.interrupt(entry.owner)
|
||||
}
|
||||
entry.stopping = true
|
||||
entry.interruptSeq = seq
|
||||
suppressPendingAtOrBefore(entry, seq)
|
||||
return Fiber.interrupt(entry.owner)
|
||||
})
|
||||
yield* Effect.addFinalizer(() =>
|
||||
SynchronizedRef.modify(state, (_current) => [
|
||||
undefined,
|
||||
{ closed: true, lanes: new Map(), interruptSeq: new Map() } satisfies State<Key, A, E>,
|
||||
]).pipe(Effect.andThen(Deferred.succeed(shutdown, undefined))),
|
||||
)
|
||||
|
||||
return { run, wake, awaitIdle, interrupt }
|
||||
|
||||
function run(key: Key): Effect.Effect<A, E> {
|
||||
return Effect.uninterruptibleMask((restore) => {
|
||||
if (closed) return Effect.interrupt
|
||||
const entry = active.get(key)
|
||||
if (entry !== undefined) {
|
||||
if (entry.stopping) {
|
||||
return restore(Deferred.await(entry.settled).pipe(Effect.andThen(run(key))))
|
||||
}
|
||||
if (entry.current._tag === "wake") {
|
||||
entry.pending = coalesce(entry.pending, { _tag: "run" })
|
||||
entry.explicitWaiter ??= Deferred.makeUnsafe<A, E>()
|
||||
return restore(awaitRun(entry.explicitWaiter))
|
||||
}
|
||||
return restore(awaitRun(entry.done))
|
||||
}
|
||||
|
||||
const next = makeEntry({ _tag: "run" })
|
||||
active.set(key, next)
|
||||
start(key, next, next.current)
|
||||
return restore(awaitRun(next.done))
|
||||
})
|
||||
}
|
||||
|
||||
function awaitRun(done: Deferred.Deferred<A, E>): Effect.Effect<A, E> {
|
||||
return Effect.raceFirst(Deferred.await(done), Deferred.await(shutdown).pipe(Effect.andThen(Effect.interrupt)))
|
||||
}
|
||||
|
||||
function acceptsWake(entry: Entry<A, E>, seq: number | undefined) {
|
||||
return !entry.stopping || (entry.interruptSeq !== undefined && seq !== undefined && seq > entry.interruptSeq)
|
||||
}
|
||||
|
||||
function isAfterInterrupt(key: Key, seq: number | undefined) {
|
||||
const latest = interruptSeq.get(key)
|
||||
return latest === undefined || (seq !== undefined && seq > latest)
|
||||
}
|
||||
|
||||
function suppressPendingAtOrBefore(entry: Entry<A, E>, seq: number | undefined) {
|
||||
if (
|
||||
entry.pending?._tag === "wake" &&
|
||||
seq !== undefined &&
|
||||
entry.pending.seq !== undefined &&
|
||||
entry.pending.seq > seq
|
||||
)
|
||||
return
|
||||
entry.pending = undefined
|
||||
}
|
||||
return { run, wake, interrupt, awaitIdle }
|
||||
})
|
||||
|
||||
export interface Interface extends Coordinator<SessionSchema.ID, void, SessionRunner.RunError> {}
|
||||
|
||||
@@ -320,6 +320,11 @@ export const layer = Layer.effect(
|
||||
yield* FiberSet.clear(toolFibers)
|
||||
yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted"))
|
||||
}
|
||||
if (settled._tag === "Failure" && !Cause.hasInterrupts(settled.cause)) {
|
||||
const failure = Cause.squash(settled.cause)
|
||||
const message = failure instanceof Error ? failure.message : String(failure)
|
||||
yield* withPublication(publisher.failUnsettledTools(`Tool execution failed: ${message}`))
|
||||
}
|
||||
if (publisher.hasProviderError())
|
||||
yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted"))
|
||||
if (stream._tag === "Success" && !publisher.hasProviderError())
|
||||
|
||||
+23
-13
@@ -61,30 +61,40 @@ export function create<State extends Objectish, Editor>(options: Options<State,
|
||||
state = next
|
||||
})
|
||||
|
||||
const rebuild = Effect.fn("State.rebuild")(function* () {
|
||||
const rebuild = Effect.fnUntraced(function* () {
|
||||
const next = options.initial()
|
||||
const api = options.editor(next as Draft<State>)
|
||||
for (const transform of transforms)
|
||||
yield* Effect.sync(() => transform.update(api)).pipe(Effect.withSpan("State.rebuild.update", {}))
|
||||
yield* commit(next)
|
||||
}, semaphore.withPermit)
|
||||
})
|
||||
|
||||
return {
|
||||
get: () => state,
|
||||
transform: Effect.fn("State.transform")(function* () {
|
||||
const transform = { update: (_editor: Editor) => {} }
|
||||
transforms = [...transforms, transform]
|
||||
const scope = yield* Scope.Scope
|
||||
yield* Scope.addFinalizer(
|
||||
scope,
|
||||
Effect.sync(() => {
|
||||
transforms = transforms.filter((item) => item !== transform)
|
||||
}).pipe(Effect.andThen(rebuild())),
|
||||
return yield* Effect.uninterruptible(
|
||||
Effect.gen(function* () {
|
||||
const transform = { update: (_editor: Editor) => {} }
|
||||
transforms = [...transforms, transform]
|
||||
yield* Scope.addFinalizer(
|
||||
scope,
|
||||
semaphore.withPermit(
|
||||
Effect.sync(() => {
|
||||
transforms = transforms.filter((item) => item !== transform)
|
||||
}).pipe(Effect.andThen(rebuild())),
|
||||
),
|
||||
)
|
||||
return (update: Transform<Editor>) =>
|
||||
Effect.uninterruptible(
|
||||
semaphore.withPermit(
|
||||
Effect.sync(() => {
|
||||
transform.update = update
|
||||
}).pipe(Effect.andThen(rebuild())),
|
||||
),
|
||||
)
|
||||
}),
|
||||
)
|
||||
return Effect.fnUntraced(function* (update: Transform<Editor>) {
|
||||
transform.update = update
|
||||
yield* rebuild()
|
||||
})
|
||||
}),
|
||||
update: Effect.fn("State.update")(function* (update, reason) {
|
||||
const api = options.editor(state as Draft<State>)
|
||||
|
||||
@@ -11,7 +11,6 @@ import type { ToolOutput } from "@opencode-ai/llm"
|
||||
|
||||
export const MAX_LINES = 2_000
|
||||
export const MAX_BYTES = 50 * 1024
|
||||
export const MAX_INLINE_MEDIA_BYTES = 5 * 1024 * 1024
|
||||
export const RETENTION = Duration.days(7)
|
||||
|
||||
export const MANAGED_DIRECTORY = "tool-output"
|
||||
@@ -32,13 +31,7 @@ export class StorageError extends Schema.TaggedErrorClass<StorageError>()("ToolO
|
||||
cause: Schema.Defect,
|
||||
}) {}
|
||||
|
||||
export class MediaLimitError extends Schema.TaggedErrorClass<MediaLimitError>()("ToolOutputStore.MediaLimitError", {
|
||||
mime: Schema.String,
|
||||
bytes: Schema.Int,
|
||||
limit: Schema.Int,
|
||||
}) {}
|
||||
|
||||
export type Error = StorageError | MediaLimitError
|
||||
export type Error = StorageError
|
||||
|
||||
export interface Interface {
|
||||
readonly limits: () => Effect.Effect<{ readonly maxLines: number; readonly maxBytes: number }>
|
||||
@@ -139,37 +132,33 @@ export const layer = Layer.effect(
|
||||
const bound = Effect.fn("ToolOutputStore.bound")(function* (input: BoundInput) {
|
||||
const outputLimits = yield* limits()
|
||||
const media = input.output.content.filter((item) => item.type === "file")
|
||||
let mediaBytes = 0
|
||||
for (const item of media) {
|
||||
if (item.source.type !== "data") continue
|
||||
mediaBytes += Buffer.byteLength(item.source.data, "utf-8")
|
||||
if (mediaBytes > MAX_INLINE_MEDIA_BYTES)
|
||||
return yield* new MediaLimitError({ mime: item.mime, bytes: mediaBytes, limit: MAX_INLINE_MEDIA_BYTES })
|
||||
}
|
||||
const contextual = {
|
||||
structured: media.length > 0 ? {} : input.output.structured,
|
||||
content: input.output.content.filter((item) => item.type === "text"),
|
||||
}
|
||||
const encoded = yield* Effect.try({
|
||||
try: () => JSON.stringify(contextual, null, 2),
|
||||
catch: (cause) => new StorageError({ operation: "encode", cause }),
|
||||
})
|
||||
if (lineCount(encoded) <= outputLimits.maxLines && Buffer.byteLength(encoded, "utf-8") <= outputLimits.maxBytes)
|
||||
const text = input.output.content.filter((item) => item.type === "text")
|
||||
const contextual =
|
||||
input.output.content.length === 0
|
||||
? yield* Effect.try({
|
||||
try: () => JSON.stringify(input.output.structured, null, 2) ?? String(input.output.structured),
|
||||
catch: (cause) => new StorageError({ operation: "encode", cause }),
|
||||
})
|
||||
: text.map((item) => item.text).join("")
|
||||
if (
|
||||
lineCount(contextual) <= outputLimits.maxLines &&
|
||||
Buffer.byteLength(contextual, "utf-8") <= outputLimits.maxBytes
|
||||
)
|
||||
return {
|
||||
output: { structured: contextual.structured, content: input.output.content },
|
||||
output: input.output,
|
||||
outputPaths: [],
|
||||
}
|
||||
|
||||
const outputPath = yield* write(encoded)
|
||||
const outputPath = yield* write(contextual)
|
||||
const marker = `... output truncated; full content saved to ${outputPath} ...`
|
||||
|
||||
return {
|
||||
output: {
|
||||
structured: {},
|
||||
structured: input.output.structured,
|
||||
content: [
|
||||
{
|
||||
type: "text" as const,
|
||||
text: boundedPreview(encoded, marker, outputLimits.maxLines, outputLimits.maxBytes),
|
||||
text: boundedPreview(contextual, marker, outputLimits.maxLines, outputLimits.maxBytes),
|
||||
},
|
||||
...media,
|
||||
],
|
||||
|
||||
@@ -56,3 +56,4 @@ Producer capture limits are separate. For example, Bash keeps `AppProcess.maxOut
|
||||
|
||||
- Plugin boot has not been redesigned to register canonical tools through `Tools.Service`; do not redesign it as part of leaf migrations.
|
||||
- MCP and future Session-scoped registrations still need an explicit canonical registration design.
|
||||
- The public Session result shape currently exposes managed `outputPaths`; full storage encapsulation requires a future opaque managed-output reference design.
|
||||
|
||||
@@ -12,7 +12,7 @@ import { Tools } from "./tools"
|
||||
|
||||
export const name = "apply_patch"
|
||||
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
patchText: Schema.String.annotate({
|
||||
description: "The full patch text describing add, update, and delete operations",
|
||||
}),
|
||||
@@ -24,10 +24,10 @@ export const Applied = Schema.Struct({
|
||||
target: Schema.String,
|
||||
})
|
||||
|
||||
export const Success = Schema.Struct({ applied: Schema.Array(Applied) })
|
||||
export type Success = typeof Success.Type
|
||||
export const Output = Schema.Struct({ applied: Schema.Array(Applied) })
|
||||
export type Output = typeof Output.Type
|
||||
|
||||
export const toModelOutput = (output: Success) =>
|
||||
export const toModelOutput = (output: Output) =>
|
||||
[
|
||||
"Applied patch sequentially:",
|
||||
...output.applied.map(
|
||||
@@ -57,8 +57,8 @@ export const layer = Layer.effectDiscard(
|
||||
Tool.make({
|
||||
description:
|
||||
"Apply one patch containing add, update, and delete file operations. All targets are resolved and approved before target contents are read. Operations apply sequentially; if a later operation fails, earlier operations remain applied and the failure reports them explicitly. Moves and atomic rollback are not supported yet.",
|
||||
input: Parameters,
|
||||
output: Success,
|
||||
input: Input,
|
||||
output: Output,
|
||||
toModelOutput: ({ output }) => [toolText({ type: "text", text: toModelOutput(output) })],
|
||||
execute: (input, context) => {
|
||||
const applied: Array<typeof Applied.Type> = []
|
||||
|
||||
@@ -18,7 +18,7 @@ export const DEFAULT_TIMEOUT_MS = 2 * 60 * 1_000
|
||||
export const MAX_TIMEOUT_MS = 10 * 60 * 1_000
|
||||
export const MAX_CAPTURE_BYTES = 1024 * 1024
|
||||
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
command: Schema.String.annotate({ description: "Shell command string to execute" }),
|
||||
workdir: Schema.String.pipe(Schema.optional).annotate({
|
||||
description: "Working directory. Defaults to the active Location; relative paths resolve from that Location.",
|
||||
@@ -33,7 +33,7 @@ export const Parameters = Schema.Struct({
|
||||
}),
|
||||
})
|
||||
|
||||
const Success = Schema.Struct({
|
||||
const Output = Schema.Struct({
|
||||
command: Schema.String,
|
||||
cwd: Schema.String,
|
||||
exitCode: Schema.Number.pipe(Schema.optional),
|
||||
@@ -46,7 +46,7 @@ const Success = Schema.Struct({
|
||||
warnings: Schema.Array(Schema.String).pipe(Schema.optional),
|
||||
})
|
||||
|
||||
type Success = typeof Success.Type
|
||||
type Output = typeof Output.Type
|
||||
|
||||
const defaultShell = () => (process.platform === "win32" ? (process.env.COMSPEC ?? "cmd.exe") : "/bin/sh")
|
||||
|
||||
@@ -62,7 +62,7 @@ const captureNotice = (stdoutTruncated: boolean, stderrTruncated: boolean) => {
|
||||
return undefined
|
||||
}
|
||||
|
||||
const modelOutput = (output: Success) => {
|
||||
const modelOutput = (output: Output) => {
|
||||
const warnings = output.warnings?.length
|
||||
? `\n\nWarnings:\n${output.warnings.map((warning) => `- ${warning}`).join("\n")}`
|
||||
: ""
|
||||
@@ -117,8 +117,8 @@ export const layer = Layer.effectDiscard(
|
||||
.register({
|
||||
[name]: Tool.make({
|
||||
description: `Execute one shell command string with the host user's filesystem, process, and network authority. The active Location is the default working directory. Relative workdir values resolve from that Location. External workdir values require external_directory approval; best-effort command-argument path warnings are advisory only. Timeout values are milliseconds (default: ${DEFAULT_TIMEOUT_MS}; maximum: ${MAX_TIMEOUT_MS}). Uses the configured shell when set; otherwise uses /bin/sh on POSIX and COMSPEC or cmd.exe on Windows.`,
|
||||
input: Parameters,
|
||||
output: Success,
|
||||
input: Input,
|
||||
output: Output,
|
||||
toModelOutput: ({ output }) => [toolText({ type: "text", text: modelOutput(output) })],
|
||||
execute: (input, context) =>
|
||||
Effect.gen(function* () {
|
||||
|
||||
@@ -18,7 +18,7 @@ import { Tools } from "./tools"
|
||||
|
||||
export const name = "edit"
|
||||
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
path: Schema.String.annotate({
|
||||
description:
|
||||
"File path to edit. Relative paths resolve within the active Location. Absolute paths inside that Location are accepted; external absolute paths require external_directory approval. Named project references are read-oriented and are not accepted.",
|
||||
@@ -30,14 +30,14 @@ export const Parameters = Schema.Struct({
|
||||
}),
|
||||
})
|
||||
|
||||
export const Success = Schema.Struct({
|
||||
export const Output = Schema.Struct({
|
||||
operation: Schema.Literal("write"),
|
||||
target: Schema.String,
|
||||
resource: Schema.String,
|
||||
existed: Schema.Boolean,
|
||||
replacements: Schema.Number,
|
||||
})
|
||||
export type Success = typeof Success.Type
|
||||
export type Output = typeof Output.Type
|
||||
|
||||
const normalizeLineEndings = (text: string) => text.replaceAll("\r\n", "\n")
|
||||
const detectLineEnding = (text: string): "\n" | "\r\n" => (text.includes("\r\n") ? "\r\n" : "\n")
|
||||
@@ -70,7 +70,7 @@ const previewLines = (value: string, prefix: "+" | "-") => {
|
||||
return shown
|
||||
}
|
||||
|
||||
export const toModelOutput = (output: Success, oldString: string, newString: string) =>
|
||||
export const toModelOutput = (output: Output, oldString: string, newString: string) =>
|
||||
[
|
||||
`Edited file successfully: ${output.resource}`,
|
||||
`Replacements: ${output.replacements}`,
|
||||
@@ -101,8 +101,8 @@ export const layer = Layer.effectDiscard(
|
||||
Tool.make({
|
||||
description:
|
||||
"Replace exact text in one file. Relative paths resolve within the active Location. Absolute paths inside the Location are accepted. Explicit external absolute paths require external_directory approval before edit approval. Named project references are read-oriented and are not accepted.",
|
||||
input: Parameters,
|
||||
output: Success,
|
||||
input: Input,
|
||||
output: Output,
|
||||
toModelOutput: ({ input, output }) => [
|
||||
toolText({ type: "text", text: toModelOutput(output, input.oldString, input.newString) }),
|
||||
],
|
||||
@@ -188,7 +188,7 @@ export const layer = Layer.effectDiscard(
|
||||
content: joinBom(next.text, source.bom || next.bom),
|
||||
}),
|
||||
)
|
||||
return { ...result, replacements } satisfies Success
|
||||
return { ...result, replacements } satisfies Output
|
||||
})
|
||||
},
|
||||
}),
|
||||
|
||||
@@ -10,7 +10,7 @@ import { Tools } from "./tools"
|
||||
|
||||
export const name = "glob"
|
||||
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
pattern: LocationSearch.FilesInput.fields.pattern.annotate({ description: "Glob pattern to match files against" }),
|
||||
path: LocationSearch.FilesInput.fields.path.annotate({
|
||||
description: "Relative directory to search. Defaults to the active Location.",
|
||||
@@ -56,7 +56,7 @@ export const layer = Layer.effectDiscard(
|
||||
[name]: Tool.make({
|
||||
description:
|
||||
"Find files by glob pattern within the active Location or a named project reference. Returns concise relative file resources. Use a relative path to narrow the search and limit to bound the result count.",
|
||||
input: Parameters,
|
||||
input: Input,
|
||||
output: LocationSearch.FilesResult,
|
||||
toModelOutput: ({ output }) => [toolText({ type: "text", text: toModelOutput(output) })],
|
||||
execute: (input, context) =>
|
||||
|
||||
@@ -11,7 +11,7 @@ import { Tools } from "./tools"
|
||||
|
||||
export const name = "grep"
|
||||
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
pattern: LocationSearch.GrepInput.fields.pattern.annotate({
|
||||
description: "Regex pattern to search for in file contents",
|
||||
}),
|
||||
@@ -29,10 +29,10 @@ export const Parameters = Schema.Struct({
|
||||
}),
|
||||
})
|
||||
|
||||
type Success = typeof LocationSearch.GrepResult.Encoded
|
||||
type Output = typeof LocationSearch.GrepResult.Encoded
|
||||
|
||||
/** Format raw Location search matches into the familiar concise model output. */
|
||||
export const toModelOutput = (output: Success) => {
|
||||
export const toModelOutput = (output: Output) => {
|
||||
const lines = output.items.length === 0 ? ["No files found"] : [`Found ${output.items.length} matches`]
|
||||
let current = ""
|
||||
for (const match of output.items) {
|
||||
@@ -71,7 +71,7 @@ export const layer = Layer.effectDiscard(
|
||||
[name]: Tool.make({
|
||||
description:
|
||||
"Search file contents by regular expression within the active Location, a named project reference, or an absolute managed tool-output file. Use a path to narrow the search, include to filter files by glob, and limit to bound the match count. Returns concise file resources, line numbers, and bounded line previews.",
|
||||
input: Parameters,
|
||||
input: Input,
|
||||
output: LocationSearch.GrepResult,
|
||||
toModelOutput: ({ output }) => [toolText({ type: "text", text: toModelOutput(output) })],
|
||||
execute: (input, context) =>
|
||||
|
||||
@@ -20,14 +20,14 @@ Usage notes:
|
||||
- Answers are returned as arrays of labels; set \`multiple: true\` to allow selecting more than one
|
||||
- If you recommend a specific option, make that the first option in the list and add "(Recommended)" at the end of the label`
|
||||
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
questions: Schema.Array(QuestionV2.Prompt).annotate({ description: "Questions to ask" }),
|
||||
})
|
||||
|
||||
export const Success = Schema.Struct({
|
||||
export const Output = Schema.Struct({
|
||||
answers: Schema.Array(QuestionV2.Answer),
|
||||
})
|
||||
export type Success = typeof Success.Type
|
||||
export type Output = typeof Output.Type
|
||||
|
||||
export const toModelOutput = (
|
||||
questions: ReadonlyArray<QuestionV2.Prompt>,
|
||||
@@ -52,8 +52,8 @@ export const layer = Layer.effectDiscard(
|
||||
.register({
|
||||
[name]: Tool.make({
|
||||
description,
|
||||
input: Parameters,
|
||||
output: Success,
|
||||
input: Input,
|
||||
output: Output,
|
||||
toModelOutput: ({ input, output }) => [
|
||||
toolText({ type: "text", text: toModelOutput(input.questions, output.answers) }),
|
||||
],
|
||||
|
||||
@@ -1,47 +1,15 @@
|
||||
export * as ReadTool from "./read"
|
||||
|
||||
import { ToolFailure } from "@opencode-ai/llm"
|
||||
// @ts-ignore Bun's static file import is embedded by `bun build --compile`; some consumers also declare *.wasm.
|
||||
import photonWasm from "@silvia-odwyer/photon-node/photon_rs_bg.wasm" with { type: "file" }
|
||||
import { Effect, Layer, Schema } from "effect"
|
||||
import path from "node:path"
|
||||
import { fileURLToPath } from "node:url"
|
||||
import { Config } from "../config"
|
||||
import { FileSystem } from "../filesystem"
|
||||
import { Image } from "../image"
|
||||
import { PermissionV2 } from "../permission"
|
||||
import { Tool } from "./tool"
|
||||
import { Tools } from "./tools"
|
||||
|
||||
export const name = "read"
|
||||
const SUPPORTED_IMAGE_MIMES = new Set(["image/jpeg", "image/png", "image/gif", "image/webp"])
|
||||
const MAX_IMAGE_BASE64_BYTES = 5 * 1024 * 1024
|
||||
const MAX_IMAGE_WIDTH = 2_000
|
||||
const MAX_IMAGE_HEIGHT = 2_000
|
||||
const JPEG_QUALITIES = [80, 85, 70, 55, 40]
|
||||
|
||||
class ImageDecodeError extends Error {
|
||||
constructor(readonly resource: string) {
|
||||
super(`Image could not be decoded: ${resource}`)
|
||||
this.name = "ImageDecodeError"
|
||||
}
|
||||
}
|
||||
|
||||
class ImageSizeError extends Error {
|
||||
constructor(
|
||||
readonly resource: string,
|
||||
readonly width: number,
|
||||
readonly height: number,
|
||||
readonly bytes: number,
|
||||
readonly maxWidth: number,
|
||||
readonly maxHeight: number,
|
||||
readonly maxBytes: number,
|
||||
) {
|
||||
super(
|
||||
`Image ${resource} is ${width}x${height} with base64 size ${bytes}, exceeding configured limits ${maxWidth}x${maxHeight}/${maxBytes} bytes`,
|
||||
)
|
||||
this.name = "ImageSizeError"
|
||||
}
|
||||
}
|
||||
const LocationInput = Schema.Struct({
|
||||
...FileSystem.ReadInput.fields,
|
||||
offset: FileSystem.ListPageInput.fields.offset.annotate({
|
||||
@@ -52,20 +20,14 @@ const LocationInput = Schema.Struct({
|
||||
}),
|
||||
})
|
||||
const Input = LocationInput
|
||||
const Success = Schema.Union([FileSystem.Content, FileSystem.TextPage, FileSystem.ListPage])
|
||||
const Output = Schema.Union([FileSystem.Content, FileSystem.TextPage, FileSystem.ListPage])
|
||||
|
||||
export const layer = Layer.effectDiscard(
|
||||
Effect.gen(function* () {
|
||||
const tools = yield* Tools.Service
|
||||
const filesystem = yield* FileSystem.Service
|
||||
const config = yield* Config.Service
|
||||
const image = yield* Image.Service
|
||||
const permission = yield* PermissionV2.Service
|
||||
const loadPhoton = yield* Effect.cached(
|
||||
Effect.sync(() => {
|
||||
;(globalThis as typeof globalThis & { __OPENCODE_PHOTON_WASM_PATH?: string }).__OPENCODE_PHOTON_WASM_PATH =
|
||||
path.isAbsolute(photonWasm) ? photonWasm : fileURLToPath(new URL(photonWasm, import.meta.url))
|
||||
}).pipe(Effect.andThen(() => Effect.promise(() => import("@silvia-odwyer/photon-node")))),
|
||||
)
|
||||
|
||||
yield* tools
|
||||
.register({
|
||||
@@ -73,7 +35,7 @@ export const layer = Layer.effectDiscard(
|
||||
description:
|
||||
"Read a text file or supported image, page through a large UTF-8 text file by line offset, or list a directory page relative to the current location. Absolute paths are accepted only for managed tool-output files.",
|
||||
input: Input,
|
||||
output: Success,
|
||||
output: Output,
|
||||
toModelOutput: ({ input, output }) => {
|
||||
if (!("type" in output) || output.type !== "binary" || !SUPPORTED_IMAGE_MIMES.has(output.mime)) return []
|
||||
return [
|
||||
@@ -98,95 +60,9 @@ export const layer = Layer.effectDiscard(
|
||||
limit: input.limit,
|
||||
})
|
||||
if (content.type === "binary" && SUPPORTED_IMAGE_MIMES.has(content.mime)) {
|
||||
const mime = content.mime
|
||||
const base64 = content.content
|
||||
const image = Object.assign(
|
||||
{},
|
||||
...(yield* config.entries()).flatMap((entry) =>
|
||||
entry.type === "document" && entry.info.attachments?.image ? [entry.info.attachments.image] : [],
|
||||
),
|
||||
)
|
||||
const limits = {
|
||||
autoResize: image.auto_resize ?? true,
|
||||
maxWidth: image.max_width ?? MAX_IMAGE_WIDTH,
|
||||
maxHeight: image.max_height ?? MAX_IMAGE_HEIGHT,
|
||||
maxBase64Bytes: image.max_base64_bytes ?? MAX_IMAGE_BASE64_BYTES,
|
||||
}
|
||||
const photon = yield* loadPhoton
|
||||
const decoded = yield* Effect.try({
|
||||
try: () => photon.PhotonImage.new_from_byteslice(Buffer.from(base64, "base64")),
|
||||
catch: () => new ImageDecodeError(resolved.resource),
|
||||
})
|
||||
try {
|
||||
const width = decoded.get_width()
|
||||
const height = decoded.get_height()
|
||||
const bytes = Buffer.byteLength(base64, "utf-8")
|
||||
if (width <= limits.maxWidth && height <= limits.maxHeight && bytes <= limits.maxBase64Bytes)
|
||||
return new FileSystem.BinaryContent({ type: "binary", content: base64, encoding: "base64", mime })
|
||||
if (!limits.autoResize)
|
||||
return yield* Effect.fail(
|
||||
new ImageSizeError(
|
||||
resolved.resource,
|
||||
width,
|
||||
height,
|
||||
bytes,
|
||||
limits.maxWidth,
|
||||
limits.maxHeight,
|
||||
limits.maxBase64Bytes,
|
||||
),
|
||||
)
|
||||
const scale = Math.min(1, limits.maxWidth / width, limits.maxHeight / height)
|
||||
const sizes = Array.from({ length: 32 }).reduce<Array<{ width: number; height: number }>>((acc) => {
|
||||
const previous = acc.at(-1) ?? {
|
||||
width: Math.max(1, Math.round(width * scale)),
|
||||
height: Math.max(1, Math.round(height * scale)),
|
||||
}
|
||||
const next =
|
||||
acc.length === 0
|
||||
? previous
|
||||
: {
|
||||
width: previous.width === 1 ? 1 : Math.max(1, Math.floor(previous.width * 0.75)),
|
||||
height: previous.height === 1 ? 1 : Math.max(1, Math.floor(previous.height * 0.75)),
|
||||
}
|
||||
return acc.some((item) => item.width === next.width && item.height === next.height)
|
||||
? acc
|
||||
: [...acc, next]
|
||||
}, [])
|
||||
for (const size of sizes) {
|
||||
const resized = photon.resize(decoded, size.width, size.height, photon.SamplingFilter.Lanczos3)
|
||||
try {
|
||||
const candidate = [
|
||||
{ content: Buffer.from(resized.get_bytes()).toString("base64"), mime: "image/png" },
|
||||
...JPEG_QUALITIES.map((quality) => ({
|
||||
content: Buffer.from(resized.get_bytes_jpeg(quality)).toString("base64"),
|
||||
mime: "image/jpeg",
|
||||
})),
|
||||
].find((item) => Buffer.byteLength(item.content, "utf-8") <= limits.maxBase64Bytes)
|
||||
if (candidate)
|
||||
return new FileSystem.BinaryContent({
|
||||
type: "binary",
|
||||
content: candidate.content,
|
||||
encoding: "base64",
|
||||
mime: candidate.mime,
|
||||
})
|
||||
} finally {
|
||||
resized.free()
|
||||
}
|
||||
}
|
||||
return yield* Effect.fail(
|
||||
new ImageSizeError(
|
||||
resolved.resource,
|
||||
width,
|
||||
height,
|
||||
bytes,
|
||||
limits.maxWidth,
|
||||
limits.maxHeight,
|
||||
limits.maxBase64Bytes,
|
||||
),
|
||||
)
|
||||
} finally {
|
||||
decoded.free()
|
||||
}
|
||||
return yield* image
|
||||
.normalize(resolved.resource, content)
|
||||
.pipe(Effect.catchTag("Image.ResizerUnavailableError", () => Effect.succeed(content)))
|
||||
}
|
||||
if (content.type === "binary")
|
||||
return yield* Effect.fail(new FileSystem.BinaryFileError(resolved.resource))
|
||||
@@ -196,8 +72,8 @@ export const layer = Layer.effectDiscard(
|
||||
const message =
|
||||
error instanceof FileSystem.BinaryFileError ||
|
||||
error instanceof FileSystem.MediaIngestLimitError ||
|
||||
error instanceof ImageDecodeError ||
|
||||
error instanceof ImageSizeError
|
||||
error instanceof Image.DecodeError ||
|
||||
error instanceof Image.SizeError
|
||||
? error.message
|
||||
: `Unable to read ${input.path}`
|
||||
return new ToolFailure({ message })
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
export * as ToolRegistry from "./registry"
|
||||
|
||||
import { Tool as LlmTool, ToolOutput, type ToolCall, type ToolSettlement } from "@opencode-ai/llm"
|
||||
import { ToolOutput, type ToolCall, type ToolDefinition, type ToolSettlement } from "@opencode-ai/llm"
|
||||
import { Context, Effect, Layer, Scope } from "effect"
|
||||
import { AgentV2 } from "../agent"
|
||||
import { PermissionV2 } from "../permission"
|
||||
@@ -9,7 +9,7 @@ import { SessionSchema } from "../session/schema"
|
||||
import { ToolOutputStore } from "../tool-output-store"
|
||||
import { Wildcard } from "../util/wildcard"
|
||||
import { ApplicationTools } from "./application-tools"
|
||||
import { Tool } from "./tool"
|
||||
import { definition, permission, settle, validateName, type AnyTool, type RegistrationError } from "./tool"
|
||||
import { Tools } from "./tools"
|
||||
|
||||
export type ExecuteInput = {
|
||||
@@ -22,13 +22,11 @@ export type ExecuteInput = {
|
||||
export interface Interface {
|
||||
readonly materialize: (permissions?: PermissionV2.Ruleset) => Effect.Effect<Materialization>
|
||||
/** Internal registration capability exposed publicly only through Tools.Service. */
|
||||
readonly register: (
|
||||
tools: Readonly<Record<string, Tool.AnyTool>>,
|
||||
) => Effect.Effect<void, Tool.RegistrationError, Scope.Scope>
|
||||
readonly register: (tools: Readonly<Record<string, AnyTool>>) => Effect.Effect<void, RegistrationError, Scope.Scope>
|
||||
}
|
||||
|
||||
export interface Materialization {
|
||||
readonly definitions: ReadonlyArray<ReturnType<typeof LlmTool.toDefinitions>[number]>
|
||||
readonly definitions: ReadonlyArray<ToolDefinition>
|
||||
readonly settle: (input: ExecuteInput) => Effect.Effect<Settlement, ToolOutputStore.Error>
|
||||
}
|
||||
|
||||
@@ -43,7 +41,7 @@ const registryLayer = Layer.effect(
|
||||
Effect.gen(function* () {
|
||||
const applications = yield* ApplicationTools.Service
|
||||
const resources = yield* ToolOutputStore.Service
|
||||
type Registration = { readonly identity: object; readonly tool: Tool.AnyTool }
|
||||
type Registration = { readonly identity: object; readonly tool: AnyTool }
|
||||
const local = new Map<string, Array<{ readonly token: object; readonly registration: Registration }>>()
|
||||
|
||||
const settleWith = Effect.fn("ToolRegistry.settle")(function* (input: ExecuteInput, advertised?: object) {
|
||||
@@ -58,7 +56,7 @@ const registryLayer = Layer.effect(
|
||||
}
|
||||
if (advertised && registration.identity !== advertised)
|
||||
return { result: { type: "error" as const, value: `Stale tool call: ${input.call.name}` } }
|
||||
const pending = yield* Tool.settle(registration.tool, input.call, {
|
||||
const pending = yield* settle(registration.tool, input.call, {
|
||||
sessionID: input.sessionID,
|
||||
agent: input.agent,
|
||||
assistantMessageID: input.assistantMessageID,
|
||||
@@ -84,17 +82,21 @@ const registryLayer = Layer.effect(
|
||||
register: Effect.fn("ToolRegistry.register")(function* (tools) {
|
||||
const entries = Object.entries(tools)
|
||||
if (entries.length === 0) return
|
||||
yield* Effect.forEach(entries, ([name]) => Tool.validateName(name), { discard: true })
|
||||
const token = {}
|
||||
for (const [name, tool] of entries)
|
||||
local.set(name, [...(local.get(name) ?? []), { token, registration: { identity: {}, tool } }])
|
||||
yield* Effect.addFinalizer(() =>
|
||||
Effect.sync(() => {
|
||||
for (const [name] of entries) {
|
||||
const registrations = local.get(name)?.filter((registration) => registration.token !== token) ?? []
|
||||
if (registrations.length > 0) local.set(name, registrations)
|
||||
else local.delete(name)
|
||||
}
|
||||
yield* Effect.forEach(entries, ([name]) => validateName(name), { discard: true })
|
||||
yield* Effect.uninterruptible(
|
||||
Effect.gen(function* () {
|
||||
const token = {}
|
||||
for (const [name, tool] of entries)
|
||||
local.set(name, [...(local.get(name) ?? []), { token, registration: { identity: {}, tool } }])
|
||||
yield* Effect.addFinalizer(() =>
|
||||
Effect.sync(() => {
|
||||
for (const [name] of entries) {
|
||||
const registrations = local.get(name)?.filter((registration) => registration.token !== token) ?? []
|
||||
if (registrations.length > 0) local.set(name, registrations)
|
||||
else local.delete(name)
|
||||
}
|
||||
}),
|
||||
)
|
||||
}),
|
||||
)
|
||||
}),
|
||||
@@ -105,9 +107,9 @@ const registryLayer = Layer.effect(
|
||||
if (registration) registrations.set(name, registration)
|
||||
}
|
||||
for (const [name, registration] of registrations)
|
||||
if (whollyDisabled(Tool.permission(registration.tool, name), permissions)) registrations.delete(name)
|
||||
if (whollyDisabled(permission(registration.tool, name), permissions)) registrations.delete(name)
|
||||
return {
|
||||
definitions: Array.from(registrations, ([name, registration]) => Tool.definition(name, registration.tool)),
|
||||
definitions: Array.from(registrations, ([name, registration]) => definition(name, registration.tool)),
|
||||
settle: (input) => {
|
||||
const registration = registrations.get(input.call.name)
|
||||
if (registration) return settleWith(input, registration.identity)
|
||||
|
||||
@@ -14,11 +14,11 @@ import { Tools } from "./tools"
|
||||
export const name = "skill"
|
||||
const FILE_LIMIT = 10
|
||||
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
name: Schema.String.annotate({ description: "The name of the skill from the available skills list" }),
|
||||
})
|
||||
|
||||
export const Success = Schema.Struct({
|
||||
export const Output = Schema.Struct({
|
||||
name: Schema.String,
|
||||
directory: Schema.String,
|
||||
output: Schema.String,
|
||||
@@ -66,8 +66,8 @@ export const layer = Layer.effectDiscard(
|
||||
.register({
|
||||
[name]: Tool.make({
|
||||
description,
|
||||
input: Parameters,
|
||||
output: Success,
|
||||
input: Input,
|
||||
output: Output,
|
||||
toModelOutput: ({ output }) => [toolText({ type: "text", text: output.output })],
|
||||
execute: (input, context) =>
|
||||
Effect.gen(function* () {
|
||||
|
||||
@@ -9,16 +9,16 @@ import { Tools } from "./tools"
|
||||
|
||||
export const name = "todowrite"
|
||||
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
todos: Schema.Array(SessionTodo.Info).annotate({ description: "The updated todo list" }),
|
||||
})
|
||||
|
||||
export const Success = Schema.Struct({
|
||||
export const Output = Schema.Struct({
|
||||
todos: Schema.Array(SessionTodo.Info),
|
||||
})
|
||||
export type Success = typeof Success.Type
|
||||
export type Output = typeof Output.Type
|
||||
|
||||
export const toModelOutput = (output: Success) => JSON.stringify(output.todos, null, 2)
|
||||
export const toModelOutput = (output: Output) => JSON.stringify(output.todos, null, 2)
|
||||
|
||||
export const layer = Layer.effectDiscard(
|
||||
Effect.gen(function* () {
|
||||
@@ -31,8 +31,8 @@ export const layer = Layer.effectDiscard(
|
||||
[name]: Tool.make({
|
||||
description:
|
||||
"Create and maintain a structured task list for the current coding session. Use it to track progress during multi-step work and keep todo statuses current.",
|
||||
input: Parameters,
|
||||
output: Success,
|
||||
input: Input,
|
||||
output: Output,
|
||||
toModelOutput: ({ output }) => [toolText({ type: "text", text: toModelOutput(output) })],
|
||||
execute: (input, context) =>
|
||||
Effect.gen(function* () {
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
export * as Tool from "./tool"
|
||||
|
||||
import { Tool as LlmTool, ToolFailure, ToolOutput, type ToolCall } from "@opencode-ai/llm"
|
||||
import { Effect, Schema } from "effect"
|
||||
import { ToolDefinition, ToolFailure, ToolOutput, type ToolCall } from "@opencode-ai/llm"
|
||||
import { Effect, JsonSchema, Schema } from "effect"
|
||||
import type { AgentV2 } from "../agent"
|
||||
import type { SessionMessage } from "../session/message"
|
||||
import type { SessionSchema } from "../session/schema"
|
||||
@@ -17,14 +17,14 @@ export type SchemaType<A> = Schema.Codec<A, any, never, never>
|
||||
|
||||
declare const TypeId: unique symbol
|
||||
|
||||
export interface Tool<Input extends SchemaType<any>, Output extends SchemaType<any>> {
|
||||
export interface Definition<Input extends SchemaType<any>, Output extends SchemaType<any>> {
|
||||
readonly [TypeId]: {
|
||||
readonly _Input: Input
|
||||
readonly _Output: Output
|
||||
}
|
||||
}
|
||||
|
||||
export type AnyTool = Tool<any, any>
|
||||
export type AnyTool = Definition<any, any>
|
||||
export const Failure = ToolFailure
|
||||
export type Failure = ToolFailure
|
||||
|
||||
@@ -53,7 +53,7 @@ type Config<Input extends SchemaType<any>, Output extends SchemaType<any>> = {
|
||||
|
||||
type Runtime = {
|
||||
readonly permission?: string
|
||||
readonly definition: (name: string) => ReturnType<typeof LlmTool.toDefinitions>[number]
|
||||
readonly definition: (name: string) => ToolDefinition
|
||||
readonly settle: (call: ToolCall, context: Context) => Effect.Effect<ToolOutput, ToolFailure>
|
||||
}
|
||||
|
||||
@@ -61,16 +61,19 @@ const runtimes = new WeakMap<AnyTool, Runtime>()
|
||||
|
||||
export function make<Input extends SchemaType<any>, Output extends SchemaType<any>>(
|
||||
config: Config<Input, Output>,
|
||||
): Tool<Input, Output> {
|
||||
const tool = Object.freeze({}) as Tool<Input, Output>
|
||||
const definitions = new Map<string, ReturnType<typeof LlmTool.toDefinitions>[number]>()
|
||||
): Definition<Input, Output> {
|
||||
const tool = Object.freeze({}) as Definition<Input, Output>
|
||||
const definitions = new Map<string, ToolDefinition>()
|
||||
runtimes.set(tool, {
|
||||
definition: (name) => {
|
||||
const cached = definitions.get(name)
|
||||
if (cached) return cached
|
||||
const definition = LlmTool.toDefinitions({
|
||||
[name]: LlmTool.make({ description: config.description, parameters: config.input, success: config.output }),
|
||||
})[0]
|
||||
const definition = new ToolDefinition({
|
||||
name,
|
||||
description: config.description,
|
||||
inputSchema: toJsonSchema(config.input),
|
||||
outputSchema: toJsonSchema(config.output),
|
||||
})
|
||||
definitions.set(name, definition)
|
||||
return definition
|
||||
},
|
||||
@@ -117,10 +120,10 @@ export const validateName = (name: string) =>
|
||||
: Effect.fail(new RegistrationError({ name, message: `Invalid tool name: ${name}` }))
|
||||
|
||||
export const withPermission = <Input extends SchemaType<any>, Output extends SchemaType<any>>(
|
||||
tool: Tool<Input, Output>,
|
||||
tool: Definition<Input, Output>,
|
||||
permission: string,
|
||||
) => {
|
||||
const decorated = Object.freeze({}) as Tool<Input, Output>
|
||||
const decorated = Object.freeze({}) as Definition<Input, Output>
|
||||
runtimes.set(decorated, { ...runtimeOf(tool), permission })
|
||||
return decorated
|
||||
}
|
||||
@@ -134,3 +137,9 @@ function runtimeOf(tool: AnyTool) {
|
||||
if (!runtime) throw new TypeError("Invalid Core Tool value")
|
||||
return runtime
|
||||
}
|
||||
|
||||
function toJsonSchema(schema: Schema.Top): JsonSchema.JsonSchema {
|
||||
const document = Schema.toJsonSchemaDocument(schema)
|
||||
if (Object.keys(document.definitions).length === 0) return document.schema
|
||||
return { ...document.schema, $defs: document.definitions }
|
||||
}
|
||||
|
||||
@@ -16,11 +16,11 @@ export const MAX_TIMEOUT_SECONDS = 120
|
||||
|
||||
export const description = `Fetch content from an HTTP or HTTPS URL and return it as text, markdown, or HTML. Markdown is the default.
|
||||
|
||||
Use a more targeted tool when one is available. This tool is read-only. Large text results are truncated and saved to a managed file that ordinary Read, Grep, and Bash tools can inspect.`
|
||||
Use a more targeted tool when one is available. This tool is read-only. Large text results may be replaced with a preview while the complete output is retained in managed storage.`
|
||||
|
||||
const Timeout = Schema.Number.check(Schema.isGreaterThan(0), Schema.isLessThanOrEqualTo(MAX_TIMEOUT_SECONDS))
|
||||
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
url: Schema.String.annotate({ description: "The HTTP or HTTPS URL to fetch content from" }),
|
||||
format: Schema.Literals(["text", "markdown", "html"])
|
||||
.annotate({ description: "The format to return the content in. Defaults to markdown." })
|
||||
@@ -30,14 +30,14 @@ export const Parameters = Schema.Struct({
|
||||
}),
|
||||
})
|
||||
|
||||
const Success = Schema.Struct({
|
||||
const Output = Schema.Struct({
|
||||
url: Schema.String,
|
||||
contentType: Schema.String,
|
||||
format: Parameters.fields.format,
|
||||
format: Input.fields.format,
|
||||
output: Schema.String,
|
||||
})
|
||||
|
||||
type Format = (typeof Parameters.Type)["format"]
|
||||
type Format = (typeof Input.Type)["format"]
|
||||
|
||||
const acceptHeader = (format: Format) => {
|
||||
switch (format) {
|
||||
@@ -134,8 +134,8 @@ export const layer = Layer.effectDiscard(
|
||||
.register({
|
||||
[name]: Tool.make({
|
||||
description,
|
||||
input: Parameters,
|
||||
output: Success,
|
||||
input: Input,
|
||||
output: Output,
|
||||
toModelOutput: ({ output }) => [toolText({ type: "text", text: output.output })],
|
||||
execute: (input, context) =>
|
||||
Effect.gen(function* () {
|
||||
|
||||
@@ -33,7 +33,7 @@ Optional controls support result count, live crawling ('fallback' or 'preferred'
|
||||
|
||||
The current year is ${new Date().getFullYear()}. Use this year when searching for recent information or current events.`
|
||||
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
query: Schema.String.annotate({ description: "Websearch query" }),
|
||||
numResults: Schema.optional(PositiveInt.check(Schema.isLessThanOrEqualTo(MAX_NUM_RESULTS))).annotate({
|
||||
description: `Number of search results to return (default: 8, maximum: ${MAX_NUM_RESULTS})`,
|
||||
@@ -176,7 +176,7 @@ const callMcp = <F extends Schema.Struct.Fields>(
|
||||
)
|
||||
})
|
||||
|
||||
const Success = Schema.Struct({
|
||||
const Output = Schema.Struct({
|
||||
provider: Provider,
|
||||
text: Schema.String,
|
||||
})
|
||||
@@ -192,8 +192,8 @@ export const layer = Layer.effectDiscard(
|
||||
.register({
|
||||
[name]: Tool.make({
|
||||
description,
|
||||
input: Parameters,
|
||||
output: Success,
|
||||
input: Input,
|
||||
output: Output,
|
||||
toModelOutput: ({ output }) => [toolText({ type: "text", text: output.text })],
|
||||
execute: (input, context) => {
|
||||
const provider = selectProvider(context.sessionID, config, config.provider)
|
||||
|
||||
@@ -18,7 +18,7 @@ import { Tools } from "./tools"
|
||||
export const name = "write"
|
||||
|
||||
// TODO: Revisit whether model-facing mutation schemas should prefer absolute `filePath` naming for trained-in compatibility after evaluating model behavior.
|
||||
export const Parameters = Schema.Struct({
|
||||
export const Input = Schema.Struct({
|
||||
path: Schema.String.annotate({
|
||||
description:
|
||||
"File path to write. Relative paths resolve within the active Location. Absolute paths inside that Location are accepted; external absolute paths require external_directory approval. Named project references are read-oriented and are not accepted.",
|
||||
@@ -26,15 +26,15 @@ export const Parameters = Schema.Struct({
|
||||
content: Schema.String.annotate({ description: "Content to write to the file" }),
|
||||
})
|
||||
|
||||
export const Success = Schema.Struct({
|
||||
export const Output = Schema.Struct({
|
||||
operation: Schema.Literal("write"),
|
||||
target: Schema.String,
|
||||
resource: Schema.String,
|
||||
existed: Schema.Boolean,
|
||||
})
|
||||
export type Success = typeof Success.Type
|
||||
export type Output = typeof Output.Type
|
||||
|
||||
export const toModelOutput = (output: Success) =>
|
||||
export const toModelOutput = (output: Output) =>
|
||||
`${output.existed ? "Wrote" : "Created"} file successfully: ${output.resource}`
|
||||
|
||||
/** Deferred V2 write UX integrations remain visible at the model-facing seam. */
|
||||
@@ -56,8 +56,8 @@ export const layer = Layer.effectDiscard(
|
||||
Tool.make({
|
||||
description:
|
||||
"Write content to one file. Relative paths resolve within the active Location. Absolute paths inside the Location are accepted. Explicit external absolute paths require external_directory approval before edit approval. Named project references are read-oriented and are not accepted.",
|
||||
input: Parameters,
|
||||
output: Success,
|
||||
input: Input,
|
||||
output: Output,
|
||||
toModelOutput: ({ output }) => [toolText({ type: "text", text: toModelOutput(output) })],
|
||||
execute: (input, context) =>
|
||||
Effect.gen(function* () {
|
||||
|
||||
@@ -9,7 +9,7 @@ import { ToolRegistry } from "@opencode-ai/core/tool/registry"
|
||||
import { executeTool, settleTool, toolDefinitions } from "./lib/tool"
|
||||
import { ToolOutputStore } from "@opencode-ai/core/tool-output-store"
|
||||
import { Tools } from "@opencode-ai/core/tool/tools"
|
||||
import { Effect, Exit, Layer, Schema, Scope } from "effect"
|
||||
import { Deferred, Effect, Exit, Fiber, Layer, Schema, Scope } from "effect"
|
||||
import { testEffect } from "./lib/effect"
|
||||
|
||||
const permission = Layer.mock(PermissionV2.Service, {
|
||||
@@ -136,7 +136,7 @@ describe("ApplicationTools", () => {
|
||||
],
|
||||
},
|
||||
output: {
|
||||
structured: {},
|
||||
structured: { answer: "HELLO" },
|
||||
content: [
|
||||
{ type: "text", text: "HELLO" },
|
||||
{ type: "file", source: { type: "data", data: "aGVsbG8=" }, mime: "image/png", name: "result.png" },
|
||||
@@ -147,7 +147,7 @@ describe("ApplicationTools", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("removes an application tool when its attachment scope closes", () =>
|
||||
it.effect("removes an application tool when its registration scope closes", () =>
|
||||
Effect.gen(function* () {
|
||||
const applications = yield* ApplicationTools.Service
|
||||
const registry = yield* ToolRegistry.Service
|
||||
@@ -165,11 +165,11 @@ describe("ApplicationTools", () => {
|
||||
Effect.gen(function* () {
|
||||
const applications = yield* ApplicationTools.Service
|
||||
const registry = yield* ToolRegistry.Service
|
||||
const attachmentScope = yield* Scope.make()
|
||||
yield* applications.register({ contextual: contextual([]) }).pipe(Scope.provide(attachmentScope))
|
||||
const registrationScope = yield* Scope.make()
|
||||
yield* applications.register({ contextual: contextual([]) }).pipe(Scope.provide(registrationScope))
|
||||
expect((yield* toolDefinitions(registry)).map((tool) => tool.name)).toEqual(["contextual"])
|
||||
|
||||
yield* Scope.close(attachmentScope, Exit.void)
|
||||
yield* Scope.close(registrationScope, Exit.void)
|
||||
expect(
|
||||
yield* settleTool(registry, {
|
||||
sessionID,
|
||||
@@ -181,7 +181,7 @@ describe("ApplicationTools", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("does not leak an attachment into an already closed scope", () =>
|
||||
it.effect("does not leak a registration into an already closed scope", () =>
|
||||
Effect.gen(function* () {
|
||||
const applications = yield* ApplicationTools.Service
|
||||
const registry = yield* ToolRegistry.Service
|
||||
@@ -194,13 +194,36 @@ describe("ApplicationTools", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("captures the attached record before later State rebuilds", () =>
|
||||
it.effect("preserves an interrupted application registration until its scope closes", () =>
|
||||
Effect.gen(function* () {
|
||||
const applications = yield* ApplicationTools.Service
|
||||
const registry = yield* ToolRegistry.Service
|
||||
const attached = { stable: contextual([]) }
|
||||
yield* applications.register(attached)
|
||||
Object.assign(attached, { late: contextual([]) })
|
||||
const scope = yield* Scope.make()
|
||||
const registered = yield* Deferred.make<void>()
|
||||
const fiber = yield* applications
|
||||
.register({ interrupted: contextual([]) })
|
||||
.pipe(
|
||||
Effect.andThen(Deferred.succeed(registered, undefined)),
|
||||
Effect.andThen(Effect.never),
|
||||
Scope.provide(scope),
|
||||
Effect.forkChild,
|
||||
)
|
||||
yield* Deferred.await(registered)
|
||||
yield* Fiber.interrupt(fiber)
|
||||
|
||||
expect((yield* toolDefinitions(registry)).map((tool) => tool.name)).toEqual(["interrupted"])
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
expect(yield* toolDefinitions(registry)).toEqual([])
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("captures the registered record before later State rebuilds", () =>
|
||||
Effect.gen(function* () {
|
||||
const applications = yield* ApplicationTools.Service
|
||||
const registry = yield* ToolRegistry.Service
|
||||
const registered = { stable: contextual([]) }
|
||||
yield* applications.register(registered)
|
||||
Object.assign(registered, { late: contextual([]) })
|
||||
|
||||
yield* Effect.scoped(applications.register({ temporary: contextual([]) }))
|
||||
|
||||
@@ -208,7 +231,7 @@ describe("ApplicationTools", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("settles with the current same-name application tool and restores earlier attachments", () =>
|
||||
it.effect("settles with the current same-name application tool and restores earlier registrations", () =>
|
||||
Effect.gen(function* () {
|
||||
const applications = yield* ApplicationTools.Service
|
||||
const registry = yield* ToolRegistry.Service
|
||||
|
||||
@@ -161,6 +161,21 @@ describe("PermissionV2", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("allows managed output reads without granting external directory access", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup([
|
||||
{ action: "*", resource: "*", effect: "deny" },
|
||||
{ action: "read", resource: "*", effect: "allow" },
|
||||
])
|
||||
const service = yield* PermissionV2.Service
|
||||
|
||||
expect(yield* service.ask(assertion({ resources: ["tool_123"] }))).toMatchObject({ effect: "allow" })
|
||||
expect(
|
||||
yield* service.ask(assertion({ action: "external_directory", resources: ["/tmp/tool-output/*"] })),
|
||||
).toMatchObject({ effect: "deny" })
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("uses build permissions when the Session agent is omitted", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup()
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
import { describe, expect, it } from "bun:test"
|
||||
import { Tool } from "@opencode-ai/core/public"
|
||||
import { Effect } from "effect"
|
||||
|
||||
describe("public Tool API", () => {
|
||||
it("keeps the public registration capability narrow", () => {
|
||||
const tools = {
|
||||
register: () => Effect.void,
|
||||
} satisfies Tool.Interface
|
||||
|
||||
expect(Object.keys(tools)).toEqual(["register"])
|
||||
})
|
||||
})
|
||||
@@ -1,10 +1,34 @@
|
||||
import { describe, expect } from "bun:test"
|
||||
import { Cause, Deferred, Effect, Exit, Fiber, Layer, Scope } from "effect"
|
||||
import { describe, expect, test } from "bun:test"
|
||||
import { Cause, Deferred, Effect, Equal, Exit, Fiber, Layer, Scope } from "effect"
|
||||
import { SessionRunCoordinator } from "@opencode-ai/core/session/run-coordinator"
|
||||
import { testEffect } from "./lib/effect"
|
||||
|
||||
const it = testEffect(Layer.empty)
|
||||
|
||||
describe("SessionRunCoordinator.Demand", () => {
|
||||
const Demand = SessionRunCoordinator.Demand
|
||||
|
||||
test("combines associatively with an identity", () => {
|
||||
const left = Demand.run.combine(Demand.wake(1))
|
||||
const right = Demand.wake().combine(Demand.wake(3))
|
||||
|
||||
expect(Equal.equals(Demand.empty.combine(left), left)).toBeTrue()
|
||||
expect(Equal.equals(left.combine(Demand.empty), left)).toBeTrue()
|
||||
expect(Equal.equals(left.combine(right), right.combine(left))).toBeTrue()
|
||||
expect(
|
||||
Equal.equals(left.combine(right).combine(Demand.wake(2)), left.combine(right.combine(Demand.wake(2)))),
|
||||
).toBeTrue()
|
||||
})
|
||||
|
||||
test("keeps only sequenced wakes newer than an interrupt boundary", () => {
|
||||
const demand = Demand.run.combine(Demand.wake()).combine(Demand.wake(3))
|
||||
|
||||
expect(Equal.equals(demand.afterBoundary(2), Demand.wake(3))).toBeTrue()
|
||||
expect(Equal.equals(demand.afterBoundary(3), Demand.empty)).toBeTrue()
|
||||
expect(Equal.equals(demand.afterBoundary(), Demand.empty)).toBeTrue()
|
||||
})
|
||||
})
|
||||
|
||||
describe("SessionRunCoordinator", () => {
|
||||
it.effect("joins concurrent resumes for one key", () =>
|
||||
Effect.scoped(
|
||||
@@ -29,6 +53,51 @@ describe("SessionRunCoordinator", () => {
|
||||
),
|
||||
)
|
||||
|
||||
it.effect("allocates fresh ownership when one run effect is reused", () =>
|
||||
Effect.scoped(
|
||||
Effect.gen(function* () {
|
||||
let runs = 0
|
||||
const coordinator = yield* SessionRunCoordinator.make({ drain: () => Effect.sync(() => runs++) })
|
||||
const run = coordinator.run("session")
|
||||
|
||||
yield* run
|
||||
yield* run
|
||||
|
||||
expect(runs).toBe(2)
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
it.effect("captures awaitIdle chains safely while settlement races", () =>
|
||||
Effect.scoped(
|
||||
Effect.gen(function* () {
|
||||
const iterations = 500
|
||||
const gates = Array.from({ length: iterations }, () => Deferred.makeUnsafe<void>())
|
||||
let runs = 0
|
||||
const coordinator = yield* SessionRunCoordinator.make({
|
||||
drain: () =>
|
||||
Effect.suspend(() => {
|
||||
const gate = gates[runs++]
|
||||
return gate === undefined ? Effect.die("Missing test gate") : Deferred.await(gate)
|
||||
}),
|
||||
})
|
||||
|
||||
for (let index = 0; index < iterations; index++) {
|
||||
const run = yield* coordinator.run("session").pipe(Effect.forkChild)
|
||||
yield* Effect.yieldNow
|
||||
const idle = yield* coordinator.awaitIdle("session").pipe(Effect.forkChild({ startImmediately: true }))
|
||||
const gate = gates[index]
|
||||
if (gate === undefined) yield* Effect.die("Missing test gate")
|
||||
yield* Deferred.succeed(gate, undefined)
|
||||
yield* Effect.all([Fiber.join(run), Fiber.join(idle)])
|
||||
}
|
||||
|
||||
expect(runs).toBe(iterations)
|
||||
return undefined
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
it.effect("starts a drain when woken while idle", () =>
|
||||
Effect.scoped(
|
||||
Effect.gen(function* () {
|
||||
@@ -122,6 +191,106 @@ describe("SessionRunCoordinator", () => {
|
||||
),
|
||||
)
|
||||
|
||||
it.effect("preserves a newer wake coalesced behind a pending explicit run", () =>
|
||||
Effect.scoped(
|
||||
Effect.gen(function* () {
|
||||
const firstStarted = yield* Deferred.make<void>()
|
||||
const secondStarted = yield* Deferred.make<void>()
|
||||
const modes: SessionRunCoordinator.Mode[] = []
|
||||
const coordinator = yield* SessionRunCoordinator.make<string, void, never>({
|
||||
drain: (_key, mode) =>
|
||||
Effect.sync(() => modes.push(mode)).pipe(
|
||||
Effect.flatMap((run) =>
|
||||
run === 1
|
||||
? Deferred.succeed(firstStarted, undefined).pipe(Effect.andThen(Effect.never))
|
||||
: Deferred.succeed(secondStarted, undefined),
|
||||
),
|
||||
),
|
||||
})
|
||||
|
||||
yield* coordinator.wake("session", 1)
|
||||
yield* Deferred.await(firstStarted)
|
||||
const run = yield* coordinator.run("session").pipe(Effect.exit, Effect.forkChild)
|
||||
yield* Effect.yieldNow
|
||||
yield* coordinator.wake("session", 3)
|
||||
yield* coordinator.interrupt("session", 2)
|
||||
yield* Deferred.await(secondStarted)
|
||||
yield* coordinator.awaitIdle("session").pipe(Effect.exit)
|
||||
|
||||
const runExit = yield* Fiber.join(run)
|
||||
expect(Exit.isFailure(runExit) && Cause.hasInterruptsOnly(runExit.cause)).toBeTrue()
|
||||
expect(modes).toEqual(["wake", "wake"])
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
it.effect("preserves a newer wake from an interrupted active combined demand", () =>
|
||||
Effect.scoped(
|
||||
Effect.gen(function* () {
|
||||
const firstGate = yield* Deferred.make<void>()
|
||||
const secondStarted = yield* Deferred.make<void>()
|
||||
const thirdStarted = yield* Deferred.make<void>()
|
||||
const modes: SessionRunCoordinator.Mode[] = []
|
||||
const coordinator = yield* SessionRunCoordinator.make<string, void, never>({
|
||||
drain: (_key, mode) =>
|
||||
Effect.sync(() => modes.push(mode)).pipe(
|
||||
Effect.flatMap((run) => {
|
||||
if (run === 1) return Deferred.await(firstGate)
|
||||
if (run === 2) return Deferred.succeed(secondStarted, undefined).pipe(Effect.andThen(Effect.never))
|
||||
return Deferred.succeed(thirdStarted, undefined)
|
||||
}),
|
||||
),
|
||||
})
|
||||
|
||||
yield* coordinator.wake("session", 1)
|
||||
const run = yield* coordinator.run("session").pipe(Effect.exit, Effect.forkChild)
|
||||
yield* Effect.yieldNow
|
||||
yield* coordinator.wake("session", 3)
|
||||
yield* Deferred.succeed(firstGate, undefined)
|
||||
yield* Deferred.await(secondStarted)
|
||||
yield* coordinator.interrupt("session", 2)
|
||||
yield* Deferred.await(thirdStarted)
|
||||
yield* coordinator.awaitIdle("session").pipe(Effect.exit)
|
||||
|
||||
const runExit = yield* Fiber.join(run)
|
||||
expect(Exit.isFailure(runExit) && Cause.hasInterruptsOnly(runExit.cause)).toBeTrue()
|
||||
expect(modes).toEqual(["wake", "run", "wake"])
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
it.effect("suppresses an older wake from an interrupted active combined demand", () =>
|
||||
Effect.scoped(
|
||||
Effect.gen(function* () {
|
||||
const firstGate = yield* Deferred.make<void>()
|
||||
const secondStarted = yield* Deferred.make<void>()
|
||||
const modes: SessionRunCoordinator.Mode[] = []
|
||||
const coordinator = yield* SessionRunCoordinator.make<string, void, never>({
|
||||
drain: (_key, mode) =>
|
||||
Effect.sync(() => modes.push(mode)).pipe(
|
||||
Effect.flatMap((run) => {
|
||||
if (run === 1) return Deferred.await(firstGate)
|
||||
return Deferred.succeed(secondStarted, undefined).pipe(Effect.andThen(Effect.never))
|
||||
}),
|
||||
),
|
||||
})
|
||||
|
||||
yield* coordinator.wake("session", 1)
|
||||
const run = yield* coordinator.run("session").pipe(Effect.exit, Effect.forkChild)
|
||||
yield* Effect.yieldNow
|
||||
yield* coordinator.wake("session", 2)
|
||||
yield* Deferred.succeed(firstGate, undefined)
|
||||
yield* Deferred.await(secondStarted)
|
||||
yield* coordinator.interrupt("session", 2)
|
||||
yield* coordinator.awaitIdle("session").pipe(Effect.exit)
|
||||
|
||||
const runExit = yield* Fiber.join(run)
|
||||
expect(Exit.isFailure(runExit) && Cause.hasInterruptsOnly(runExit.cause)).toBeTrue()
|
||||
expect(modes).toEqual(["wake", "run"])
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
it.effect("interrupts only the requested key", () =>
|
||||
Effect.scoped(
|
||||
Effect.gen(function* () {
|
||||
@@ -316,6 +485,53 @@ describe("SessionRunCoordinator", () => {
|
||||
),
|
||||
)
|
||||
|
||||
it.effect("does not let an interrupted attempt completion settle its successor", () =>
|
||||
Effect.scoped(
|
||||
Effect.gen(function* () {
|
||||
const firstStarted = yield* Deferred.make<void>()
|
||||
const cleanupStarted = yield* Deferred.make<void>()
|
||||
const cleanupGate = yield* Deferred.make<void>()
|
||||
const secondStarted = yield* Deferred.make<void>()
|
||||
const secondGate = yield* Deferred.make<void>()
|
||||
const idleSettled = yield* Deferred.make<void>()
|
||||
let runs = 0
|
||||
const coordinator = yield* SessionRunCoordinator.make({
|
||||
drain: () =>
|
||||
Effect.sync(() => ++runs).pipe(
|
||||
Effect.flatMap((run) =>
|
||||
run === 1
|
||||
? Deferred.succeed(firstStarted, undefined).pipe(
|
||||
Effect.andThen(Effect.never),
|
||||
Effect.onInterrupt(() =>
|
||||
Deferred.succeed(cleanupStarted, undefined).pipe(Effect.andThen(Deferred.await(cleanupGate))),
|
||||
),
|
||||
)
|
||||
: Deferred.succeed(secondStarted, undefined).pipe(Effect.andThen(Deferred.await(secondGate))),
|
||||
),
|
||||
),
|
||||
})
|
||||
|
||||
yield* coordinator.wake("session", 1)
|
||||
yield* Deferred.await(firstStarted)
|
||||
const interrupt = yield* coordinator.interrupt("session", 2).pipe(Effect.forkChild)
|
||||
yield* Deferred.await(cleanupStarted)
|
||||
yield* coordinator.wake("session", 3)
|
||||
yield* Deferred.succeed(cleanupGate, undefined)
|
||||
yield* Fiber.join(interrupt)
|
||||
yield* Deferred.await(secondStarted)
|
||||
const idle = yield* coordinator
|
||||
.awaitIdle("session")
|
||||
.pipe(Effect.ensuring(Deferred.succeed(idleSettled, undefined)), Effect.forkChild)
|
||||
|
||||
yield* Effect.yieldNow
|
||||
expect(yield* Deferred.isDone(idleSettled)).toBeFalse()
|
||||
yield* Deferred.succeed(secondGate, undefined)
|
||||
yield* Fiber.join(idle)
|
||||
expect(runs).toBe(2)
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
it.effect("interrupts an explicit run queued before the interruption request", () =>
|
||||
Effect.scoped(
|
||||
Effect.gen(function* () {
|
||||
@@ -847,6 +1063,37 @@ describe("SessionRunCoordinator", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("settles a post-stop run waiter when its owning scope closes", () =>
|
||||
Effect.gen(function* () {
|
||||
const scope = yield* Scope.make()
|
||||
const started = yield* Deferred.make<void>()
|
||||
const cleanupStarted = yield* Deferred.make<void>()
|
||||
const cleanupGate = yield* Deferred.make<void>()
|
||||
const coordinator = yield* SessionRunCoordinator.make<string, void, never>({
|
||||
drain: () =>
|
||||
Deferred.succeed(started, undefined).pipe(
|
||||
Effect.andThen(Effect.never),
|
||||
Effect.onInterrupt(() =>
|
||||
Deferred.succeed(cleanupStarted, undefined).pipe(Effect.andThen(Deferred.await(cleanupGate))),
|
||||
),
|
||||
),
|
||||
}).pipe(Scope.provide(scope))
|
||||
|
||||
yield* coordinator.wake("session")
|
||||
yield* Deferred.await(started)
|
||||
const interrupt = yield* coordinator.interrupt("session").pipe(Effect.forkChild)
|
||||
yield* Deferred.await(cleanupStarted)
|
||||
const run = yield* coordinator.run("session").pipe(Effect.forkChild)
|
||||
const close = yield* Scope.close(scope, Exit.void).pipe(Effect.forkChild)
|
||||
|
||||
const runExit = yield* Fiber.await(run)
|
||||
expect(Exit.isFailure(runExit) && Cause.hasInterruptsOnly(runExit.cause)).toBeTrue()
|
||||
yield* Deferred.succeed(cleanupGate, undefined)
|
||||
yield* Fiber.join(interrupt)
|
||||
yield* Fiber.join(close)
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("does not start work after its owning scope closes", () =>
|
||||
Effect.gen(function* () {
|
||||
const scope = yield* Scope.make()
|
||||
|
||||
@@ -7,7 +7,7 @@ import { SessionMessage } from "@opencode-ai/core/session/message"
|
||||
import { ToolOutputStore } from "@opencode-ai/core/tool-output-store"
|
||||
import { ToolRegistry } from "@opencode-ai/core/tool/registry"
|
||||
import { executeTool, settleTool, toolDefinitions } from "./lib/tool"
|
||||
import { Cause, Deferred, Effect, Exit, Fiber, Layer, Option, Schema, Scope } from "effect"
|
||||
import { Cause, Deferred, Effect, Exit, Fiber, Layer, Option, Schema, SchemaGetter, SchemaIssue, Scope } from "effect"
|
||||
import { testEffect } from "./lib/effect"
|
||||
|
||||
const bounds: ToolOutputStore.BoundInput[] = []
|
||||
@@ -29,6 +29,7 @@ const outputStore = Layer.mock(ToolOutputStore.Service, {
|
||||
})
|
||||
const registry = ToolRegistry.layer.pipe(Layer.provide(ApplicationTools.layer), Layer.provide(outputStore))
|
||||
const it = testEffect(registry)
|
||||
const integrated = testEffect(Layer.mergeAll(ApplicationTools.layer, registry))
|
||||
const identity = {
|
||||
agent: AgentV2.ID.make("build"),
|
||||
assistantMessageID: SessionMessage.ID.make("msg_registry"),
|
||||
@@ -125,6 +126,28 @@ describe("ToolRegistry", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("preserves an interrupted registration until its scope closes", () =>
|
||||
Effect.gen(function* () {
|
||||
const service = yield* ToolRegistry.Service
|
||||
const scope = yield* Scope.make()
|
||||
const registered = yield* Deferred.make<void>()
|
||||
const fiber = yield* service
|
||||
.register({ echo: make() })
|
||||
.pipe(
|
||||
Effect.andThen(Deferred.succeed(registered, undefined)),
|
||||
Effect.andThen(Effect.never),
|
||||
Scope.provide(scope),
|
||||
Effect.forkChild,
|
||||
)
|
||||
yield* Deferred.await(registered)
|
||||
yield* Fiber.interrupt(fiber)
|
||||
|
||||
expect((yield* toolDefinitions(service)).map((tool) => tool.name)).toEqual(["echo"])
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
expect(yield* toolDefinitions(service)).toEqual([])
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("returns model errors without swallowing interruption or defects", () =>
|
||||
Effect.gen(function* () {
|
||||
const service = yield* ToolRegistry.Service
|
||||
@@ -237,6 +260,72 @@ describe("ToolRegistry", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("enforces transformed codecs at execution and projection boundaries", () =>
|
||||
Effect.gen(function* () {
|
||||
const service = yield* ToolRegistry.Service
|
||||
const executed: string[] = []
|
||||
const Transformed = Schema.Boolean.pipe(
|
||||
Schema.decodeTo(Schema.String, {
|
||||
decode: SchemaGetter.transform((value) => (value ? "yes" : "no")),
|
||||
encode: SchemaGetter.transform((value) => value === "yes"),
|
||||
}),
|
||||
)
|
||||
yield* service.register({
|
||||
transformed: Tool.make({
|
||||
description: "Transform values",
|
||||
input: Schema.Struct({ value: Transformed }),
|
||||
output: Schema.Struct({ value: Transformed }),
|
||||
execute: ({ value }) => Effect.sync(() => executed.push(value)).pipe(Effect.as({ value })),
|
||||
toModelOutput: ({ output }) => [{ type: "text", text: String(output.value) }],
|
||||
}),
|
||||
})
|
||||
|
||||
expect(
|
||||
yield* executeTool(service, {
|
||||
sessionID,
|
||||
...identity,
|
||||
call: { type: "tool-call", id: "transformed", name: "transformed", input: { value: true } },
|
||||
}),
|
||||
).toEqual({ type: "text", value: "true" })
|
||||
expect(executed).toEqual(["yes"])
|
||||
expect(
|
||||
yield* executeTool(service, {
|
||||
sessionID,
|
||||
...identity,
|
||||
call: { type: "tool-call", id: "invalid-input", name: "transformed", input: { value: "yes" } },
|
||||
}),
|
||||
).toMatchObject({ type: "error", value: expect.stringContaining("Invalid tool input") })
|
||||
expect(executed).toEqual(["yes"])
|
||||
|
||||
yield* service.register({
|
||||
invalid_output: Tool.make({
|
||||
description: "Return invalid output",
|
||||
input: Schema.Struct({}),
|
||||
output: Schema.Struct({
|
||||
value: Schema.Boolean.pipe(
|
||||
Schema.decodeTo(Schema.String, {
|
||||
decode: SchemaGetter.transform((value) => String(value)),
|
||||
encode: SchemaGetter.transformOrFail((value) =>
|
||||
value === "valid"
|
||||
? Effect.succeed(true)
|
||||
: Effect.fail(new SchemaIssue.InvalidValue(Option.some(value), { message: "invalid output" })),
|
||||
),
|
||||
}),
|
||||
),
|
||||
}),
|
||||
execute: () => Effect.succeed({ value: "invalid" }),
|
||||
}),
|
||||
})
|
||||
expect(
|
||||
yield* executeTool(service, {
|
||||
sessionID,
|
||||
...identity,
|
||||
call: { type: "tool-call", id: "invalid-output", name: "invalid_output", input: {} },
|
||||
}),
|
||||
).toMatchObject({ type: "error", value: expect.stringContaining("invalid value for its output schema") })
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("executes the unchanged registration advertised for a provider turn", () =>
|
||||
Effect.gen(function* () {
|
||||
const service = yield* ToolRegistry.Service
|
||||
@@ -293,6 +382,38 @@ describe("ToolRegistry", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
integrated.effect("rejects an application call after a Location override is registered", () =>
|
||||
Effect.gen(function* () {
|
||||
const applications = yield* ApplicationTools.Service
|
||||
const service = yield* ToolRegistry.Service
|
||||
yield* applications.register({ echo: make() })
|
||||
const materialized = yield* service.materialize()
|
||||
yield* service.register({ echo: make() })
|
||||
|
||||
expect((yield* materialized.settle(call("echo"))).result).toEqual({
|
||||
type: "error",
|
||||
value: "Stale tool call: echo",
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
integrated.effect("rejects a Location call after removal reveals an application registration", () =>
|
||||
Effect.gen(function* () {
|
||||
const applications = yield* ApplicationTools.Service
|
||||
const service = yield* ToolRegistry.Service
|
||||
yield* applications.register({ echo: make() })
|
||||
const scope = yield* Scope.make()
|
||||
yield* service.register({ echo: make() }).pipe(Scope.provide(scope))
|
||||
const materialized = yield* service.materialize()
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
|
||||
expect((yield* materialized.settle(call("echo"))).result).toEqual({
|
||||
type: "error",
|
||||
value: "Stale tool call: echo",
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("keeps captured execution running after registration mutation", () =>
|
||||
Effect.gen(function* () {
|
||||
const service = yield* ToolRegistry.Service
|
||||
|
||||
@@ -2955,6 +2955,22 @@ describe("SessionRunnerLLM", () => {
|
||||
expect(yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))).toBe("unexpected tool defect")
|
||||
|
||||
expect(requests).toHaveLength(1)
|
||||
expect(yield* session.context(sessionID)).toMatchObject([
|
||||
{ type: "user", text: "Call defect" },
|
||||
{
|
||||
type: "assistant",
|
||||
content: [
|
||||
{
|
||||
type: "tool",
|
||||
id: "call-defect",
|
||||
state: {
|
||||
status: "error",
|
||||
error: { type: "unknown", message: "Tool execution failed: unexpected tool defect" },
|
||||
},
|
||||
},
|
||||
],
|
||||
},
|
||||
])
|
||||
}),
|
||||
)
|
||||
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
import { describe, expect } from "bun:test"
|
||||
import { State } from "@opencode-ai/core/state"
|
||||
import { Deferred, Effect, Exit, Fiber, Layer, Scope } from "effect"
|
||||
import { testEffect } from "./lib/effect"
|
||||
|
||||
const it = testEffect(Layer.empty)
|
||||
|
||||
describe("State", () => {
|
||||
it.effect("commits a transform atomically when its updater is interrupted", () =>
|
||||
Effect.gen(function* () {
|
||||
const rebuilding = yield* Deferred.make<void>()
|
||||
const release = yield* Deferred.make<void>()
|
||||
let block = true
|
||||
const state = State.create({
|
||||
initial: () => ({ values: [] as string[] }),
|
||||
editor: (draft) => ({ add: (value: string) => draft.values.push(value) }),
|
||||
finalize: () =>
|
||||
block ? Deferred.succeed(rebuilding, undefined).pipe(Effect.andThen(Deferred.await(release))) : Effect.void,
|
||||
})
|
||||
const scope = yield* Scope.make()
|
||||
const update = yield* state.transform().pipe(Scope.provide(scope))
|
||||
const fiber = yield* update((editor) => editor.add("registered")).pipe(Effect.forkChild)
|
||||
yield* Deferred.await(rebuilding)
|
||||
const interruption = yield* Fiber.interrupt(fiber).pipe(Effect.forkChild)
|
||||
block = false
|
||||
yield* Deferred.succeed(release, undefined)
|
||||
yield* Fiber.join(interruption)
|
||||
|
||||
expect(state.get().values).toEqual(["registered"])
|
||||
yield* Scope.close(scope, Exit.void)
|
||||
expect(state.get().values).toEqual([])
|
||||
}),
|
||||
)
|
||||
})
|
||||
@@ -115,7 +115,7 @@ const withTool = <A, E, R>(
|
||||
}).pipe(Effect.provide(Layer.mergeAll(registry, bash)))
|
||||
}
|
||||
|
||||
const call = (input: typeof BashTool.Parameters.Type, id = "call-bash") => ({
|
||||
const call = (input: typeof BashTool.Input.Type, id = "call-bash") => ({
|
||||
sessionID,
|
||||
...toolIdentity,
|
||||
call: { type: "tool-call" as const, id, name: "bash", input },
|
||||
|
||||
@@ -93,7 +93,7 @@ const withTool = <A, E, R>(directory: string, body: (registry: ToolRegistry.Inte
|
||||
}).pipe(Effect.provide(Layer.mergeAll(registry, resolution, mutation, edit)))
|
||||
}
|
||||
|
||||
const call = (input: typeof EditTool.Parameters.Type, id = "call-edit") => ({
|
||||
const call = (input: typeof EditTool.Input.Type, id = "call-edit") => ({
|
||||
sessionID,
|
||||
...toolIdentity,
|
||||
call: { type: "tool-call" as const, id, name: "edit", input },
|
||||
|
||||
@@ -91,7 +91,7 @@ const reset = () => {
|
||||
result = new LocationSearch.FilesResult({ items: [], truncated: false, partial: false })
|
||||
}
|
||||
|
||||
const call = (input: typeof GlobTool.Parameters.Type, id = "call-glob") => ({
|
||||
const call = (input: typeof GlobTool.Input.Type, id = "call-glob") => ({
|
||||
sessionID,
|
||||
...toolIdentity,
|
||||
call: { type: "tool-call" as const, id, name: "glob", input },
|
||||
|
||||
@@ -151,7 +151,7 @@ function provideLive(directory: string, projectReferences = references({})) {
|
||||
}
|
||||
|
||||
describe("GrepTool", () => {
|
||||
it.effect("registers the grep contribution", () =>
|
||||
it.effect("registers grep", () =>
|
||||
Effect.gen(function* () {
|
||||
reset()
|
||||
expect(yield* toolDefinitions(yield* ToolRegistry.Service)).toMatchObject([{ name: "grep" }])
|
||||
|
||||
@@ -43,7 +43,7 @@ const withStore = <A, E, R>(
|
||||
const it = testEffect(Layer.empty)
|
||||
|
||||
describe("ToolOutputStore", () => {
|
||||
it.live("bounds aggregate text and structured output with one managed file", () =>
|
||||
it.live("bounds the provider-facing text channel with one managed file", () =>
|
||||
withStore(({ store, fs }) =>
|
||||
Effect.gen(function* () {
|
||||
const first = "HEAD-" + "x".repeat(30_000)
|
||||
@@ -59,15 +59,9 @@ describe("ToolOutputStore", () => {
|
||||
],
|
||||
},
|
||||
})
|
||||
expect(result.output.structured).toEqual({})
|
||||
expect(result.output.structured).toEqual({ kind: "report" })
|
||||
expect(result.outputPaths).toHaveLength(1)
|
||||
expect(JSON.parse(yield* fs.readFileString(result.outputPaths[0]))).toEqual({
|
||||
structured: { kind: "report" },
|
||||
content: [
|
||||
{ type: "text", text: first },
|
||||
{ type: "text", text: second },
|
||||
],
|
||||
})
|
||||
expect(yield* fs.readFileString(result.outputPaths[0])).toBe(first + second)
|
||||
if (result.output.content[0]?.type !== "text") throw new Error("expected text preview")
|
||||
expect(Buffer.byteLength(result.output.content[0].text)).toBeLessThanOrEqual(ToolOutputStore.MAX_BYTES)
|
||||
}),
|
||||
@@ -79,18 +73,18 @@ describe("ToolOutputStore", () => {
|
||||
Effect.gen(function* () {
|
||||
const structured = { text: "x".repeat(ToolOutputStore.MAX_BYTES) }
|
||||
const result = yield* store.bound({ sessionID, toolCallID: "call-json", output: { structured, content: [] } })
|
||||
expect(result.output.structured).toEqual({})
|
||||
expect(result.output.structured).toEqual(structured)
|
||||
expect(result.outputPaths).toHaveLength(1)
|
||||
expect(JSON.parse(yield* fs.readFileString(result.outputPaths[0]))).toEqual({ structured, content: [] })
|
||||
expect(JSON.parse(yield* fs.readFileString(result.outputPaths[0]))).toEqual(structured)
|
||||
expect(result.output.content).toHaveLength(1)
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
it.live("preserves oversized inline media without duplicating its data", () =>
|
||||
it.live("preserves native media and structured metadata without applying a settlement media limit", () =>
|
||||
withStore(({ store }) =>
|
||||
Effect.gen(function* () {
|
||||
const data = "a".repeat(ToolOutputStore.MAX_BYTES)
|
||||
const data = "a".repeat(6 * 1024 * 1024)
|
||||
const result = yield* store.bound({
|
||||
sessionID,
|
||||
toolCallID: "call-file",
|
||||
@@ -100,7 +94,7 @@ describe("ToolOutputStore", () => {
|
||||
},
|
||||
})
|
||||
expect(result.outputPaths).toEqual([])
|
||||
expect(result.output.structured).toEqual({})
|
||||
expect(result.output.structured).toEqual({ caption: "pixel" })
|
||||
expect(result.output.content).toHaveLength(1)
|
||||
expect(result.output.content[0]).toEqual({
|
||||
type: "file",
|
||||
@@ -112,51 +106,38 @@ describe("ToolOutputStore", () => {
|
||||
),
|
||||
)
|
||||
|
||||
it.live("rejects inline media beyond the settlement media limit", () =>
|
||||
withStore(({ store }) =>
|
||||
it.live("preserves structured metadata and native media when bounding text", () =>
|
||||
withStore(({ store, fs }) =>
|
||||
Effect.gen(function* () {
|
||||
const exit = yield* store
|
||||
.bound({
|
||||
sessionID,
|
||||
toolCallID: "call-file-too-large",
|
||||
output: {
|
||||
structured: {},
|
||||
content: [
|
||||
{
|
||||
type: "file",
|
||||
source: { type: "data", data: "a".repeat(ToolOutputStore.MAX_INLINE_MEDIA_BYTES + 1) },
|
||||
mime: "image/png",
|
||||
},
|
||||
],
|
||||
},
|
||||
})
|
||||
.pipe(Effect.exit)
|
||||
expect(Exit.isFailure(exit)).toBe(true)
|
||||
if (Exit.isFailure(exit))
|
||||
expect(Option.getOrUndefined(Cause.findErrorOption(exit.cause))?._tag).toBe("ToolOutputStore.MediaLimitError")
|
||||
const text = "x".repeat(ToolOutputStore.MAX_BYTES + 1)
|
||||
const media = {
|
||||
type: "file" as const,
|
||||
source: { type: "data" as const, data: "aGVsbG8=" },
|
||||
mime: "image/png",
|
||||
name: "pixel.png",
|
||||
}
|
||||
const result = yield* store.bound({
|
||||
sessionID,
|
||||
toolCallID: "call-text-and-media",
|
||||
output: { structured: { caption: "pixel" }, content: [{ type: "text", text }, media] },
|
||||
})
|
||||
|
||||
expect(result.output.structured).toEqual({ caption: "pixel" })
|
||||
expect(result.output.content[1]).toEqual(media)
|
||||
expect(yield* fs.readFileString(result.outputPaths[0])).toBe(text)
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
it.live("rejects inline media whose aggregate size exceeds the settlement limit", () =>
|
||||
it.live("does not double-count structured data duplicated in projected text", () =>
|
||||
withStore(({ store }) =>
|
||||
Effect.gen(function* () {
|
||||
const exit = yield* store
|
||||
.bound({
|
||||
sessionID,
|
||||
toolCallID: "call-files-too-large",
|
||||
output: {
|
||||
structured: {},
|
||||
content: [
|
||||
{ type: "file", source: { type: "data", data: "a".repeat(3 * 1024 * 1024) }, mime: "image/png" },
|
||||
{ type: "file", source: { type: "data", data: "b".repeat(3 * 1024 * 1024) }, mime: "image/png" },
|
||||
],
|
||||
},
|
||||
})
|
||||
.pipe(Effect.exit)
|
||||
expect(Exit.isFailure(exit)).toBe(true)
|
||||
if (Exit.isFailure(exit))
|
||||
expect(Option.getOrUndefined(Cause.findErrorOption(exit.cause))?._tag).toBe("ToolOutputStore.MediaLimitError")
|
||||
const text = "x".repeat(30_000)
|
||||
const output = { structured: { output: text }, content: [{ type: "text" as const, text }] }
|
||||
expect(yield* store.bound({ sessionID, toolCallID: "call-duplicated", output })).toEqual({
|
||||
output,
|
||||
outputPaths: [],
|
||||
})
|
||||
}),
|
||||
),
|
||||
)
|
||||
@@ -179,22 +160,14 @@ describe("ToolOutputStore", () => {
|
||||
),
|
||||
)
|
||||
|
||||
it.live("fails operationally when output cannot be encoded for bounding", () =>
|
||||
it.live("does not encode ignored structured metadata when projected content exists", () =>
|
||||
withStore(({ store }) =>
|
||||
Effect.gen(function* () {
|
||||
const exit = yield* store
|
||||
.bound({
|
||||
sessionID,
|
||||
toolCallID: "call-unencodable",
|
||||
output: {
|
||||
structured: { value: 1n },
|
||||
content: [{ type: "text", text: "readable text" }],
|
||||
},
|
||||
})
|
||||
.pipe(Effect.exit)
|
||||
expect(Exit.isFailure(exit)).toBe(true)
|
||||
if (Exit.isFailure(exit))
|
||||
expect(Option.getOrUndefined(Cause.findErrorOption(exit.cause))?._tag).toBe("ToolOutputStore.StorageError")
|
||||
const output = { structured: { value: 1n }, content: [{ type: "text" as const, text: "readable text" }] }
|
||||
expect(yield* store.bound({ sessionID, toolCallID: "call-unencodable", output })).toEqual({
|
||||
output,
|
||||
outputPaths: [],
|
||||
})
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
||||
@@ -3,6 +3,7 @@ import { Effect, Exit, Layer } from "effect"
|
||||
import { Config } from "@opencode-ai/core/config"
|
||||
import { ConfigAttachments } from "@opencode-ai/core/config/attachments"
|
||||
import { FileSystem } from "@opencode-ai/core/filesystem"
|
||||
import { Image } from "@opencode-ai/core/image"
|
||||
import { PermissionV2 } from "@opencode-ai/core/permission"
|
||||
import { SessionV2 } from "@opencode-ai/core/session"
|
||||
import { ToolRegistry } from "@opencode-ai/core/tool/registry"
|
||||
@@ -75,13 +76,29 @@ const permission = Layer.succeed(
|
||||
)
|
||||
const registry = ToolRegistry.defaultLayer.pipe(Layer.provide(permission))
|
||||
const config = Layer.succeed(Config.Service, Config.Service.of({ entries: () => Effect.succeed(configEntries) }))
|
||||
const image = Image.layer.pipe(Layer.provide(config))
|
||||
const unavailableImage = Layer.succeed(
|
||||
Image.Service,
|
||||
Image.Service.of({ normalize: () => Effect.fail(new Image.ResizerUnavailableError()) }),
|
||||
)
|
||||
const read = ReadTool.layer.pipe(
|
||||
Layer.provide(registry),
|
||||
Layer.provide(filesystem),
|
||||
Layer.provide(permission),
|
||||
Layer.provide(config),
|
||||
Layer.provide(image),
|
||||
)
|
||||
const it = testEffect(Layer.mergeAll(registry, filesystem, permission, config, image, read))
|
||||
const unavailableRead = ReadTool.layer.pipe(
|
||||
Layer.provide(registry),
|
||||
Layer.provide(filesystem),
|
||||
Layer.provide(permission),
|
||||
Layer.provide(config),
|
||||
Layer.provide(unavailableImage),
|
||||
)
|
||||
const itWithoutResizer = testEffect(
|
||||
Layer.mergeAll(registry, filesystem, permission, config, unavailableImage, unavailableRead),
|
||||
)
|
||||
const it = testEffect(Layer.mergeAll(registry, filesystem, permission, config, read))
|
||||
const sessionID = SessionV2.ID.make("ses_read_tool_test")
|
||||
|
||||
describe("ReadTool", () => {
|
||||
@@ -146,7 +163,7 @@ describe("ReadTool", () => {
|
||||
...toolIdentity,
|
||||
call: { type: "tool-call", id: "call-image-settle", name: "read", input: { path: "pixel.png" } },
|
||||
})
|
||||
expect(settled.output?.structured).toEqual({})
|
||||
expect(settled.output?.structured).toMatchObject({ type: "binary", mime: "image/png", encoding: "base64" })
|
||||
expect(settled.output?.content).toMatchObject([
|
||||
{ type: "text", text: "Image read successfully" },
|
||||
{ type: "file", mime: "image/png", source: { type: "data", data: png } },
|
||||
@@ -177,7 +194,7 @@ describe("ReadTool", () => {
|
||||
})
|
||||
|
||||
expect(settled.outputPaths).toBeUndefined()
|
||||
expect(settled.output?.structured).toEqual({})
|
||||
expect(settled.output?.structured).toMatchObject({ type: "binary", mime: "image/png", encoding: "base64" })
|
||||
expect(settled.result).toEqual({
|
||||
type: "content",
|
||||
value: [
|
||||
@@ -188,6 +205,30 @@ describe("ReadTool", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
itWithoutResizer.effect("returns the original image when the resizer is unavailable", () =>
|
||||
Effect.gen(function* () {
|
||||
const png = "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNk+A8AAQUBAScY42YAAAAASUVORK5CYII="
|
||||
readResult = new FileSystem.BinaryContent({
|
||||
type: "binary",
|
||||
content: png,
|
||||
encoding: "base64",
|
||||
mime: "image/png",
|
||||
})
|
||||
const registry = yield* ToolRegistry.Service
|
||||
|
||||
expect(
|
||||
yield* executeTool(registry, {
|
||||
sessionID,
|
||||
...toolIdentity,
|
||||
call: { type: "tool-call", id: "call-image-fallback", name: "read", input: { path: "pixel.png" } },
|
||||
}),
|
||||
).toMatchObject({
|
||||
type: "content",
|
||||
value: [{ type: "text" }, { type: "media", mediaType: "image/png", data: png }],
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("rejects invalid image data returned by the filesystem", () =>
|
||||
Effect.gen(function* () {
|
||||
readResult = new FileSystem.BinaryContent({
|
||||
|
||||
@@ -51,7 +51,7 @@ const reset = () => {
|
||||
respond = () => Effect.succeed(new Response("hello", { headers: { "content-type": "text/plain" } }))
|
||||
}
|
||||
|
||||
const call = (input: typeof WebFetchTool.Parameters.Type, id = "call-webfetch") => ({
|
||||
const call = (input: typeof WebFetchTool.Input.Type, id = "call-webfetch") => ({
|
||||
sessionID,
|
||||
...toolIdentity,
|
||||
call: { type: "tool-call" as const, id, name: "webfetch", input },
|
||||
@@ -59,7 +59,7 @@ const call = (input: typeof WebFetchTool.Parameters.Type, id = "call-webfetch")
|
||||
|
||||
describe("WebFetchTool helpers", () => {
|
||||
test("defaults format and rejects invalid timeout controls", () => {
|
||||
const decode = Schema.decodeUnknownSync(WebFetchTool.Parameters)
|
||||
const decode = Schema.decodeUnknownSync(WebFetchTool.Input)
|
||||
expect(decode({ url: "https://example.com" })).toEqual({ url: "https://example.com", format: "markdown" })
|
||||
expect(() => decode({ url: "https://example.com", timeout: 0 })).toThrow()
|
||||
expect(() => decode({ url: "https://example.com", timeout: WebFetchTool.MAX_TIMEOUT_SECONDS + 1 })).toThrow()
|
||||
@@ -72,7 +72,7 @@ describe("WebFetchTool helpers", () => {
|
||||
})
|
||||
})
|
||||
|
||||
describe("WebFetchTool contribution", () => {
|
||||
describe("WebFetchTool registration", () => {
|
||||
it.effect("registers and fetches an ordinary hostname HTTP URL without rewriting it", () =>
|
||||
Effect.gen(function* () {
|
||||
reset()
|
||||
|
||||
@@ -18,7 +18,7 @@ const payload = (text: string) =>
|
||||
|
||||
describe("WebSearchTool provider selection", () => {
|
||||
test("rejects out-of-range numeric controls", () => {
|
||||
const decode = Schema.decodeUnknownSync(WebSearchTool.Parameters)
|
||||
const decode = Schema.decodeUnknownSync(WebSearchTool.Input)
|
||||
expect(() => decode({ query: "x", numResults: 0 })).toThrow()
|
||||
expect(() => decode({ query: "x", numResults: WebSearchTool.MAX_NUM_RESULTS + 1 })).toThrow()
|
||||
expect(() => decode({ query: "x", contextMaxCharacters: WebSearchTool.MAX_CONTEXT_CHARACTERS + 1 })).toThrow()
|
||||
@@ -122,7 +122,7 @@ const websearch = WebSearchTool.layer.pipe(
|
||||
)
|
||||
const it = testEffect(Layer.mergeAll(registry, permission, http, websearchConfig, websearch))
|
||||
|
||||
describe("WebSearchTool contribution", () => {
|
||||
describe("WebSearchTool registration", () => {
|
||||
it.effect("registers websearch, asserts query permission, and calls Exa", () =>
|
||||
Effect.gen(function* () {
|
||||
requests.length = 0
|
||||
|
||||
@@ -76,7 +76,7 @@ const withTool = <A, E, R>(directory: string, body: (registry: ToolRegistry.Inte
|
||||
}).pipe(Effect.provide(Layer.mergeAll(registry, resolution, mutation, write)))
|
||||
}
|
||||
|
||||
const call = (input: typeof WriteTool.Parameters.Type, id = "call-write") => ({
|
||||
const call = (input: typeof WriteTool.Input.Type, id = "call-write") => ({
|
||||
sessionID,
|
||||
...toolIdentity,
|
||||
call: { type: "tool-call" as const, id, name: "write", input },
|
||||
|
||||
@@ -10,6 +10,9 @@ import { afterAll } from "bun:test"
|
||||
const dir = path.join(os.tmpdir(), "opencode-test-data-" + process.pid)
|
||||
await fs.mkdir(dir, { recursive: true })
|
||||
afterAll(async () => {
|
||||
const { AppRuntime } = await import("../src/effect/app-runtime")
|
||||
await AppRuntime.dispose()
|
||||
|
||||
const busy = (error: unknown) =>
|
||||
typeof error === "object" && error !== null && "code" in error && error.code === "EBUSY"
|
||||
const rm = async (left: number): Promise<void> => {
|
||||
|
||||
+12
-10
@@ -5,8 +5,8 @@
|
||||
V2 has one opaque type for locally executable tools:
|
||||
|
||||
```ts
|
||||
type Tool<Input, Output>
|
||||
type AnyTool = Tool<any, any>
|
||||
type Definition<Input, Output>
|
||||
type AnyTool = Definition<any, any>
|
||||
|
||||
const make: <
|
||||
Input extends Schema.Codec<any, any, never, never>,
|
||||
@@ -23,12 +23,12 @@ const make: <
|
||||
readonly input: Schema.Type<Input>
|
||||
readonly output: Output["Encoded"]
|
||||
}) => ReadonlyArray<Tool.Content>
|
||||
}) => Tool<Input, Output>
|
||||
}) => Definition<Input, Output>
|
||||
```
|
||||
|
||||
Application tools, built-ins, and statically authored plugin tools use this same constructor and execution contract.
|
||||
|
||||
`Tool` is opaque and has exactly one executor. Its schemas and executor are not public fields. The Tool module privately derives model definitions and interprets invocations for the registry; it never embeds another executable tool representation.
|
||||
`Tool.Definition` is opaque and has exactly one executor. Its schemas and executor are not public fields. The Tool module privately derives model definitions and interprets invocations for the registry; callers normally rely on `Tool.make` inference rather than naming the carrier type.
|
||||
|
||||
Input and output codecs are self-contained. Schema conversion cannot require services. Tool dependencies are acquired during construction and captured by `execute`.
|
||||
|
||||
@@ -83,9 +83,9 @@ A Location plugin receives only the narrow `Tools` registration capability, not
|
||||
Within one placement:
|
||||
|
||||
- The latest active registration for a name wins.
|
||||
- Closing a registration removes only that contribution.
|
||||
- Closing the winner reveals the next-latest active contribution.
|
||||
- Mutating the caller's registration record later does not change the captured contribution.
|
||||
- Closing a registration removes only that registration.
|
||||
- Closing the winner reveals the next-latest active registration.
|
||||
- Mutating the caller's registration record later does not change the captured registration.
|
||||
|
||||
Location registrations take precedence over process application registrations.
|
||||
|
||||
@@ -142,19 +142,19 @@ The Location-scoped registry owns effective lookup and settlement. For each loca
|
||||
4. Encodes the returned output with the output codec.
|
||||
5. Projects encoded output into model-facing content.
|
||||
6. Bounds the complete model-facing output.
|
||||
7. Persists the settlement and any internal managed-output references.
|
||||
7. Returns the settlement and managed-output references to the runner, which persists them durably.
|
||||
|
||||
Invalid input never invokes the tool. Invalid output never produces a successful settlement.
|
||||
|
||||
`toModelOutput` is pure and total. When omitted, the encoded output remains structured output; an encoded string is also projected as text. Projection does not receive invocation identity because presentation depends only on validated input and output.
|
||||
|
||||
Provider-turn materialization captures the effective registration identity for each advertised name without retaining its handler. Settlement rejects the call as stale if that registration was removed or replaced, including when closing an overlay reveals the previously effective registration. The current handler is captured only after this check; detaching or replacing it afterward does not affect the running invocation.
|
||||
Provider-turn materialization captures the effective registration identity for each advertised name without retaining its handler. Settlement rejects the call as stale if that registration was removed or replaced, including when closing an overlay reveals the previously effective registration. The current handler is captured only after this check; removing or replacing its registration afterward does not affect the running invocation.
|
||||
|
||||
## Output Bounding
|
||||
|
||||
Tools return complete validated domain output. They do not truncate model-facing output or manage retention files.
|
||||
|
||||
After projection, one generic settlement boundary bounds textual and structured provider context. Supported inline media remains native up to the producer's media limit and is never encoded into a text preview. Structured data duplicated by native media content is omitted from provider settlement accounting and storage. Oversized textual or structured values are materialized in managed storage and replaced with bounded previews or references; if complete retention fails, settlement fails operationally rather than publishing lossy success. Managed paths are internal settlement metadata and never appear in `Tool.make`, tool output schemas, or projection callbacks solely for retention bookkeeping.
|
||||
After projection, one generic settlement boundary bounds the channel actually sent to the provider. When content exists, only its textual parts are measured; structured metadata is retained unchanged without being double-counted, and native media remains unchanged under producer-owned limits. When content is empty, the structured output is measured. Oversized provider-facing text or structured output is retained in managed storage and replaced with a bounded text preview while structured metadata and media are preserved; if complete retention fails, settlement fails operationally rather than publishing lossy success. Managed paths never appear in `Tool.make`, tool output schemas, or projection callbacks solely for retention bookkeeping.
|
||||
|
||||
Model-output bounding is not producer memory management. Processes and streaming sources may need separate capture or spooling limits before a tool result exists. Those limits must be modeled at the producer boundary and must not masquerade as model-output truncation. A producer cannot claim a complete retained output after it has already discarded bytes.
|
||||
|
||||
@@ -182,3 +182,5 @@ Leaf tools translate only errors they deliberately classify as recoverable. Broa
|
||||
## Follow-Up
|
||||
|
||||
Location plugin installation should receive the same narrow `Tools` capability. That requires a separate Location-layer ordering change so built-ins register before plugins without introducing a `PluginBoot -> Tools -> PluginBoot` dependency cycle. The carrier, registrar, and plugin-owned Scope semantics are already suitable; no tool-specific plugin hook is needed.
|
||||
|
||||
Session's current public result shape still exposes managed `outputPaths`. Extending storage encapsulation across the public Session API requires a separate opaque managed-output reference design; paths are not entirely internal today.
|
||||
|
||||
Reference in New Issue
Block a user