mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-26 03:26:12 +00:00
Compare commits
7
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e342333eea | ||
|
|
40f11b9956 | ||
|
|
50a8539e4b | ||
|
|
7b47589225 | ||
|
|
84275c6e9d | ||
|
|
793ea52fa7 | ||
|
|
40380ad9b5 |
@@ -89,6 +89,12 @@ jobs:
|
||||
working-directory: packages/codemode
|
||||
run: bun run script/publish.ts --dry-run
|
||||
|
||||
- name: Verify packed workerd SDK
|
||||
if: runner.os == 'Linux'
|
||||
timeout-minutes: 15
|
||||
working-directory: packages/sdk
|
||||
run: bun run verify:package
|
||||
|
||||
- name: Verify compiled service lifecycle
|
||||
if: always()
|
||||
timeout-minutes: 10
|
||||
|
||||
@@ -595,8 +595,8 @@
|
||||
},
|
||||
"peerDependencies": {
|
||||
"@opencode-ai/theme": "workspace:*",
|
||||
"@opentui/core": ">=0.5.7",
|
||||
"@opentui/solid": ">=0.5.7",
|
||||
"@opentui/core": ">=0.5.8",
|
||||
"@opentui/solid": ">=0.5.8",
|
||||
"solid-js": ">=1.9.0",
|
||||
},
|
||||
"optionalPeers": [
|
||||
@@ -1090,9 +1090,9 @@
|
||||
"@npmcli/arborist": "9.4.0",
|
||||
"@octokit/rest": "22.0.0",
|
||||
"@openauthjs/openauth": "0.0.0-20250322224806",
|
||||
"@opentui/core": "0.5.7",
|
||||
"@opentui/keymap": "0.5.7",
|
||||
"@opentui/solid": "0.5.7",
|
||||
"@opentui/core": "0.5.8",
|
||||
"@opentui/keymap": "0.5.8",
|
||||
"@opentui/solid": "0.5.8",
|
||||
"@pierre/diffs": "1.2.10",
|
||||
"@playwright/test": "1.59.1",
|
||||
"@sentry/solid": "10.36.0",
|
||||
@@ -2218,27 +2218,27 @@
|
||||
|
||||
"@opentelemetry/semantic-conventions": ["@opentelemetry/semantic-conventions@1.43.0", "", {}, "sha512-eSYWTm620tTk45EKSedaUL8MFYI8hW164hIXsgIHyxu3VobUB3fFCu5t0hQby6OoWRPsG1KkKUG2M5UadiLiVg=="],
|
||||
|
||||
"@opentui/core": ["@opentui/core@0.5.7", "", { "dependencies": { "bun-ffi-structs": "0.3.1", "diff": "9.0.0", "marked": "17.0.1", "string-width": "7.2.0", "strip-ansi": "7.1.2" }, "optionalDependencies": { "@opentui/core-darwin-arm64": "0.5.7", "@opentui/core-darwin-x64": "0.5.7", "@opentui/core-linux-arm64": "0.5.7", "@opentui/core-linux-arm64-musl": "0.5.7", "@opentui/core-linux-x64": "0.5.7", "@opentui/core-linux-x64-musl": "0.5.7", "@opentui/core-win32-arm64": "0.5.7", "@opentui/core-win32-x64": "0.5.7" }, "peerDependencies": { "web-tree-sitter": "0.25.10" } }, "sha512-/XDabTkfBs2Wy2FhlC4jvzbpphYAs4SlnCRKYEvi+metlXNTSAshSs43wqZ9O4IPd0E4EcQqdfeQ0sdc3yPpYw=="],
|
||||
"@opentui/core": ["@opentui/core@0.5.8", "", { "dependencies": { "bun-ffi-structs": "0.3.1", "diff": "9.0.0", "marked": "17.0.1", "string-width": "7.2.0", "strip-ansi": "7.1.2" }, "optionalDependencies": { "@opentui/core-darwin-arm64": "0.5.8", "@opentui/core-darwin-x64": "0.5.8", "@opentui/core-linux-arm64": "0.5.8", "@opentui/core-linux-arm64-musl": "0.5.8", "@opentui/core-linux-x64": "0.5.8", "@opentui/core-linux-x64-musl": "0.5.8", "@opentui/core-win32-arm64": "0.5.8", "@opentui/core-win32-x64": "0.5.8" }, "peerDependencies": { "web-tree-sitter": "0.25.10" } }, "sha512-GbZ+nSLYZqxj2Z5TU19Mx6IsAvVsn2+7WEXz+6OlMGoootvt3TxP3vWwfGI0jWk9qp1ftRlML1JPNEkzy+9I8g=="],
|
||||
|
||||
"@opentui/core-darwin-arm64": ["@opentui/core-darwin-arm64@0.5.7", "", { "os": "darwin", "cpu": "arm64" }, "sha512-75TDJgFD6hDoCElIX35Yg3TIRF/jhrtVQO/9snhjGuTsZ7bd0W88jUlvSXSNhqB3CX431rwgi47B+asWPDI1lQ=="],
|
||||
"@opentui/core-darwin-arm64": ["@opentui/core-darwin-arm64@0.5.8", "", { "os": "darwin", "cpu": "arm64" }, "sha512-c9Y1FBrSnA4sKUCMETsrLYOmsMTyJae8mU9cE6M4o9rXcr3ZLPA7o9AkPA3S0+kH+kZ0o1Fn0wLejgNxvEp5mg=="],
|
||||
|
||||
"@opentui/core-darwin-x64": ["@opentui/core-darwin-x64@0.5.7", "", { "os": "darwin", "cpu": "x64" }, "sha512-j+Dwu2yV8zahBFjnVIssYdPrHz/+pkd7xWxI+nHjV41V3c+scvH6vJY1U85VCVaoQs/zCH5xAJbmAAn7Sgeh8g=="],
|
||||
"@opentui/core-darwin-x64": ["@opentui/core-darwin-x64@0.5.8", "", { "os": "darwin", "cpu": "x64" }, "sha512-oZ/6Iz1KN+4volMFKmmvziYJhMgyyJ99LfK1S+uPRyIRqzT3CESoJ58D4h04M1g0dHeHE4PvkU1p3uQKrXiP3g=="],
|
||||
|
||||
"@opentui/core-linux-arm64": ["@opentui/core-linux-arm64@0.5.7", "", { "os": "linux", "cpu": "arm64" }, "sha512-gWO9NWaivXRPc0XhEbxehj3ApyC1SW8Bf1cnHnME1zhw/EZNCNE3TQHgKwOnnNKtvlbhzLQjOI4OqyPYV8UcEQ=="],
|
||||
"@opentui/core-linux-arm64": ["@opentui/core-linux-arm64@0.5.8", "", { "os": "linux", "cpu": "arm64" }, "sha512-N6i/ocrsTjIq9aUQyrfJkqUo+tc4P5ZS6xz38Cm1MhtDzwBE+CrNIJdHfoJfmRLhJkYHlw6Z0ToICn8ZNj4Bpw=="],
|
||||
|
||||
"@opentui/core-linux-arm64-musl": ["@opentui/core-linux-arm64-musl@0.5.7", "", { "os": "linux", "cpu": "arm64" }, "sha512-RumSHTasIAWU7jPBKiTuEZG/m+prfQKCHhL3JAQnBLgp7eAB20XZGfSLVWDBgxkyKq1rHCUgLtZktC5ZAyKfTA=="],
|
||||
"@opentui/core-linux-arm64-musl": ["@opentui/core-linux-arm64-musl@0.5.8", "", { "os": "linux", "cpu": "arm64" }, "sha512-eFMB41AWODaYf8PsCx3vtTMX33tFgtSpR9tKNnUErCXlcnZ+WCRHVaJ9F6DrVNX3ir6HmgbwAmaWXvaVF+kLtg=="],
|
||||
|
||||
"@opentui/core-linux-x64": ["@opentui/core-linux-x64@0.5.7", "", { "os": "linux", "cpu": "x64" }, "sha512-fXDmrIlfp9xaoVkk3BNn3yUO/b7plBSOfUh2oWOezRA6E/g5z+hTO1alGFQCy+VV+1kfA6iuYadAww0xNtfACg=="],
|
||||
"@opentui/core-linux-x64": ["@opentui/core-linux-x64@0.5.8", "", { "os": "linux", "cpu": "x64" }, "sha512-/2QM7/wMnML/sxchzbwgoU5tUu/7k836/kSOKMti8opjuecv1K+WWNKGXufhTNeRcXFZaba5rsCdYrf3VqnVsQ=="],
|
||||
|
||||
"@opentui/core-linux-x64-musl": ["@opentui/core-linux-x64-musl@0.5.7", "", { "os": "linux", "cpu": "x64" }, "sha512-1opDV+W7C1F4iWivCNKuYl5Hen+IeFXAuUPm3hVyN0UtdcP8LESpFJTr8QRwHpGSE9Jz1K2g2S2ipzI3u5mhgw=="],
|
||||
"@opentui/core-linux-x64-musl": ["@opentui/core-linux-x64-musl@0.5.8", "", { "os": "linux", "cpu": "x64" }, "sha512-YXo+qUHYmep2uvv3ECvTeqr10aD7+lBsavYmsTLBzS5hHabbzlQ10oX/99nIRPC3Au1BMWD6d6zQP5czPx23eg=="],
|
||||
|
||||
"@opentui/core-win32-arm64": ["@opentui/core-win32-arm64@0.5.7", "", { "os": "win32", "cpu": "arm64" }, "sha512-gbiZUyttjq8s+bsnsaopygWN7LTS/RyzD6GOmpgOB2TVitvo0P6mxM+AebBGuu/7xri17Rqaxf3uqrgdKiJAeg=="],
|
||||
"@opentui/core-win32-arm64": ["@opentui/core-win32-arm64@0.5.8", "", { "os": "win32", "cpu": "arm64" }, "sha512-7qBdhEAlh4tLFzW7nWLPlREtNiF6NZMDMi+4uDpUlAMhRavLr6wjcIcgfhNAF/puq06DjuZlPTwaMCgB3qOuwA=="],
|
||||
|
||||
"@opentui/core-win32-x64": ["@opentui/core-win32-x64@0.5.7", "", { "os": "win32", "cpu": "x64" }, "sha512-bDwon45lUxbV3rMcekTCg4mXcBu9uhdG6HLeLK5TVZf05/h+Ibkf2UeZ8A2jjQt+MK6jpYmxWwDCVCTzZ7EU7Q=="],
|
||||
"@opentui/core-win32-x64": ["@opentui/core-win32-x64@0.5.8", "", { "os": "win32", "cpu": "x64" }, "sha512-Z76YaTKnmRDSHKdKa7iTBCXBQdInQLq7UG3qIE84nvyHCDtkhYtPWOaCzDisSgXC8Y6hJt6NiFbaNjQyxPkfpQ=="],
|
||||
|
||||
"@opentui/keymap": ["@opentui/keymap@0.5.7", "", { "dependencies": { "@opentui/core": "0.5.7" }, "peerDependencies": { "@opentui/react": "0.5.7", "@opentui/solid": "0.5.7", "react": ">=19.2.0", "solid-js": "1.9.12" }, "optionalPeers": ["@opentui/react", "@opentui/solid", "react", "solid-js"] }, "sha512-BqfSbjlLuctnew2aoPPdK3wpzJfEPp7ttLBkTu1wJTkps9AcfUgXqwloqvY4DUKjUBNFYYEpU/0CPxI4Blj4MA=="],
|
||||
"@opentui/keymap": ["@opentui/keymap@0.5.8", "", { "dependencies": { "@opentui/core": "0.5.8" }, "peerDependencies": { "@opentui/react": "0.5.8", "@opentui/solid": "0.5.8", "react": ">=19.2.0", "solid-js": "1.9.12" }, "optionalPeers": ["@opentui/react", "@opentui/solid", "react", "solid-js"] }, "sha512-KQHKRnLZroSZIbHmGmSeDPsXi7Yykdu9x9ACmDzBFkgPGx0EIuXHqRcmrqmUaX485MQEZuCD9t6FdDsOMGyIFg=="],
|
||||
|
||||
"@opentui/solid": ["@opentui/solid@0.5.7", "", { "dependencies": { "@babel/core": "7.28.0", "@babel/preset-typescript": "7.27.1", "@opentui/core": "0.5.7", "babel-plugin-module-resolver": "5.0.2", "babel-preset-solid": "1.9.12", "entities": "7.0.1", "s-js": "^0.4.9" }, "peerDependencies": { "solid-js": "1.9.12" } }, "sha512-qrKAZd9xt4D67LXZUARg1Aw0dwEWi29GyOxpM6abqmyExGPMS7jih5Z9oMQJ9EXsouDKpW7Mejjq7WbX2AaecQ=="],
|
||||
"@opentui/solid": ["@opentui/solid@0.5.8", "", { "dependencies": { "@babel/core": "7.28.0", "@babel/preset-typescript": "7.27.1", "@opentui/core": "0.5.8", "babel-plugin-module-resolver": "5.0.2", "babel-preset-solid": "1.9.12", "entities": "7.0.1", "s-js": "^0.4.9" }, "peerDependencies": { "solid-js": "1.9.12" } }, "sha512-L0NxuAU8XT+jlE5G90oA3kspqkof48b0hmzi5XLw+1gxnkxrkTb+YfKys+GzVK4UqhgwY9aW+TeDfrKbnfwCMw=="],
|
||||
|
||||
"@oslojs/asn1": ["@oslojs/asn1@1.0.0", "", { "dependencies": { "@oslojs/binary": "1.0.0" } }, "sha512-zw/wn0sj0j0QKbIXfIlnEcTviaCzYOY3V5rAyjR6YtOByFtJiT574+8p9Wlach0lZH9fddD4yb9laEAIl4vXQA=="],
|
||||
|
||||
|
||||
+4
-4
@@ -1,8 +1,8 @@
|
||||
{
|
||||
"nodeModules": {
|
||||
"x86_64-linux": "sha256-LvDHCOm8OAZfvb0I0L6AbdOevRoQmEJnnqrSgAiNHv8=",
|
||||
"aarch64-linux": "sha256-O0L0iHjb4cwl9xWHIna8VFHyQzoAKzfY8oMpVNayMOg=",
|
||||
"aarch64-darwin": "sha256-ETP8FE71NqufYDUbR7tBdsMEOVQ44wLmsZBeZiSRBRY=",
|
||||
"x86_64-darwin": "sha256-WUcoLldDriT3QxcdlnBQhuPrxDNub0EDvvZXk/pDMpY="
|
||||
"x86_64-linux": "sha256-phyTF0/jQZ3L0B66PSLdpH//kyPc1M6j5a40wCSx7TA=",
|
||||
"aarch64-linux": "sha256-1Zb/Is0ujIslCbPPusAVhcuzAPyIauQyeIIRRGtzpAk=",
|
||||
"aarch64-darwin": "sha256-DDsVm7z+PSDry6QqrwVDFSmEnq6jIKb709Y4ymAv9f8=",
|
||||
"x86_64-darwin": "sha256-S+5LI2J+WRhRP7jp2PAv6AesXk238wEYoyIO1oKdF3w="
|
||||
}
|
||||
}
|
||||
|
||||
+3
-3
@@ -49,9 +49,9 @@
|
||||
"@octokit/rest": "22.0.0",
|
||||
"@hono/standard-validator": "0.2.0",
|
||||
"@hono/zod-validator": "0.4.2",
|
||||
"@opentui/core": "0.5.7",
|
||||
"@opentui/keymap": "0.5.7",
|
||||
"@opentui/solid": "0.5.7",
|
||||
"@opentui/core": "0.5.8",
|
||||
"@opentui/keymap": "0.5.8",
|
||||
"@opentui/solid": "0.5.8",
|
||||
"@tanstack/solid-virtual": "3.13.37",
|
||||
"@shikijs/stream": "4.2.0",
|
||||
"@standard-schema/spec": "1.1.0",
|
||||
|
||||
@@ -1,7 +1,14 @@
|
||||
import { Argument, Flag } from "effect/unstable/cli"
|
||||
import { Argument, Flag, GlobalFlag } from "effect/unstable/cli"
|
||||
import { Schema } from "effect"
|
||||
import { Spec } from "../framework/spec"
|
||||
|
||||
export const PrintLogs = GlobalFlag.setting("print-logs")({
|
||||
flag: Flag.boolean("print-logs").pipe(
|
||||
Flag.withDescription("Print logs to stderr (server logs require --standalone)"),
|
||||
Flag.withDefault(false),
|
||||
),
|
||||
})
|
||||
|
||||
declare const OPENCODE_CLI_NAME: string | undefined
|
||||
|
||||
const ServerParams = {
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import { Effect, FileSystem, Scope } from "effect"
|
||||
import { Command } from "effect/unstable/cli"
|
||||
import { PrintLogs } from "../commands/commands"
|
||||
import { Spec } from "./spec"
|
||||
import { Global } from "@opencode-ai/util/global"
|
||||
import { Updater } from "../services/updater"
|
||||
@@ -77,7 +78,11 @@ export function handlers<const Root extends Spec.Any>(root: Root, handlers: Hand
|
||||
}
|
||||
|
||||
export function run(commands: Spec.Any, handlers: ReadonlyArray<LazyHandler>, options: { readonly version: string }) {
|
||||
return Command.run(provide(commands, handlers), options) as Effect.Effect<void, unknown, Command.Environment>
|
||||
return Command.run(provide(commands, handlers).pipe(Command.withGlobalFlags([PrintLogs])), options) as Effect.Effect<
|
||||
void,
|
||||
unknown,
|
||||
Command.Environment
|
||||
>
|
||||
}
|
||||
|
||||
function provide(node: Spec.Any, handlers: ReadonlyArray<LazyHandler>): ProvidedCommand {
|
||||
@@ -86,6 +91,7 @@ function provide(node: Spec.Any, handlers: ReadonlyArray<LazyHandler>): Provided
|
||||
? node.spec.pipe(
|
||||
Command.withHandler((input) =>
|
||||
Effect.gen(function* () {
|
||||
if (yield* PrintLogs) process.env.OPENCODE_PRINT_LOGS = "1"
|
||||
const module = yield* Effect.promise(handler.load)
|
||||
return yield* module.default(input)
|
||||
}),
|
||||
|
||||
@@ -27,7 +27,7 @@ function command(password: string, options: Options) {
|
||||
// The server treats EOF on this pipe as the end of its ownership lease.
|
||||
// The OS closes it even when the TUI is killed before Effect finalizers run.
|
||||
stdin: "pipe",
|
||||
stderr: "ignore",
|
||||
stderr: process.env.OPENCODE_PRINT_LOGS === "1" ? "inherit" : "ignore",
|
||||
killSignal: "SIGTERM",
|
||||
forceKillAfter: "3 seconds",
|
||||
})
|
||||
|
||||
@@ -32,3 +32,38 @@ const result = await Bun.build({
|
||||
},
|
||||
})
|
||||
if (!result.success) throw new AggregateError(result.logs, "Failed to build Core")
|
||||
|
||||
// Bun's Node target eagerly creates its shared require helper, so every split
|
||||
// entry evaluates import.meta.url even when it never requires a module. Keep
|
||||
// the helper lazy until Bun stops hoisting it into workerd-reachable chunks.
|
||||
// https://github.com/oven-sh/bun/issues/12615
|
||||
const eagerRequire = "var __require = /* @__PURE__ */ createRequire(import.meta.url);"
|
||||
const lazyRequire = `var __require = (specifier) => createRequire(import.meta.url ?? "file:///worker.js")(specifier);
|
||||
__require.resolve = (specifier, options) => createRequire(import.meta.url ?? "file:///worker.js").resolve(specifier, options);`
|
||||
const rewritten = await Promise.all(
|
||||
result.outputs.map(async (output) => {
|
||||
if (!output.path.endsWith(".js")) return false
|
||||
const source = await output.text()
|
||||
|
||||
const generatedUses = source
|
||||
.replace(/import\s*\{[^}]*\b__require\b[^}]*\}\s*from\s*["'][^"']+["'];/g, "")
|
||||
.replace(/export\s*\{[^}]*\b__require\b[^}]*\};/g, "")
|
||||
.replace(eagerRequire, "")
|
||||
if (/\bnew\s+__require\s*\(/.test(generatedUses))
|
||||
throw new Error(`Unsupported generated require constructor in ${output.path}`)
|
||||
const unsupported = generatedUses
|
||||
.replace(/\b__require\.resolve\s*\(/g, "")
|
||||
.replace(/\b__require\s*\(/g, "")
|
||||
if (/\b__require\b/.test(unsupported)) throw new Error(`Unsupported generated require usage in ${output.path}`)
|
||||
|
||||
if (!source.includes(eagerRequire)) return false
|
||||
if (source.indexOf(eagerRequire) !== source.lastIndexOf(eagerRequire))
|
||||
throw new Error(`Multiple eager require helpers in ${output.path}`)
|
||||
const rewrittenSource = source.replace(eagerRequire, lazyRequire)
|
||||
if (rewrittenSource.includes(eagerRequire)) throw new Error(`Failed to rewrite eager require helper in ${output.path}`)
|
||||
await Bun.write(output.path, rewrittenSource)
|
||||
return true
|
||||
}),
|
||||
)
|
||||
if (rewritten.filter(Boolean).length !== 1)
|
||||
throw new Error("Expected exactly one eager require helper; Bun may have fixed #12615 and made this shim removable")
|
||||
|
||||
+194
-24
@@ -19,10 +19,17 @@ export type Info = {
|
||||
metadata?: Record<string, unknown>
|
||||
}
|
||||
|
||||
export type SettledInfo = Info & {
|
||||
status: Exclude<Status, "running">
|
||||
completed_at: number
|
||||
}
|
||||
|
||||
type Active = {
|
||||
info: Info
|
||||
done: Deferred.Deferred<Info>
|
||||
backgrounded: Deferred.Deferred<Info>
|
||||
onBackgroundSettled?: (info: SettledInfo) => Effect.Effect<void>
|
||||
backgroundNotification?: Deferred.Deferred<void>
|
||||
scope: Scope.Closeable
|
||||
token: object
|
||||
blockingSessions: Map<SessionSchema.ID, number>
|
||||
@@ -31,18 +38,29 @@ type Active = {
|
||||
|
||||
type State = {
|
||||
jobs: SynchronizedRef.SynchronizedRef<Map<string, Active>>
|
||||
notifications: Set<Deferred.Deferred<void>>
|
||||
scope: Scope.Scope
|
||||
shuttingDown: boolean
|
||||
}
|
||||
|
||||
type Notification = {
|
||||
jobID: string
|
||||
effect: Effect.Effect<void>
|
||||
done: Deferred.Deferred<void>
|
||||
}
|
||||
|
||||
type FinishResult = {
|
||||
info?: Info
|
||||
done?: Deferred.Deferred<Info>
|
||||
notify?: Notification
|
||||
scope?: Scope.Closeable
|
||||
}
|
||||
|
||||
type BackgroundResult = {
|
||||
info?: Info
|
||||
backgrounded?: Deferred.Deferred<Info>
|
||||
notify?: Notification
|
||||
cancel?: { id: string; token: object }
|
||||
}
|
||||
|
||||
type StartResult = { info: Info } | { info: Info; scope: Scope.Closeable; token: object }
|
||||
@@ -64,6 +82,7 @@ export type StartInput = {
|
||||
title?: string
|
||||
metadata?: Record<string, unknown>
|
||||
run: Effect.Effect<string, unknown>
|
||||
onBackgroundSettled?: (info: SettledInfo) => Effect.Effect<void>
|
||||
}
|
||||
|
||||
export type WaitInput = {
|
||||
@@ -96,6 +115,8 @@ export interface Interface {
|
||||
readonly background: (id: string) => Effect.Effect<Info | undefined>
|
||||
readonly backgroundAll: (input: BackgroundAllInput) => Effect.Effect<Info[]>
|
||||
readonly cancel: (id: string) => Effect.Effect<Info | undefined>
|
||||
/** Cancels detached work and awaits its terminal callbacks before application teardown. */
|
||||
readonly shutdown: Effect.Effect<void>
|
||||
}
|
||||
|
||||
export class Service extends Context.Service<Service, Interface>()("@opencode/Job") {}
|
||||
@@ -125,6 +146,23 @@ function decrementSession(input: Map<SessionSchema.ID, number>, sessionID: Sessi
|
||||
return next
|
||||
}
|
||||
|
||||
function clearNotification(job: Active) {
|
||||
if (!job.onBackgroundSettled && !job.backgroundNotification) return job
|
||||
return { ...job, onBackgroundSettled: undefined, backgroundNotification: undefined }
|
||||
}
|
||||
|
||||
function claimNotification(job: Active, info: SettledInfo) {
|
||||
if (!job.isBackgrounded || !job.onBackgroundSettled || !job.backgroundNotification) return { job }
|
||||
return {
|
||||
job: clearNotification(job),
|
||||
notify: {
|
||||
jobID: info.id,
|
||||
effect: job.onBackgroundSettled(info),
|
||||
done: job.backgroundNotification,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Makes one scoped, process-local registry. Entries are intentionally not
|
||||
* durable: process restart or owner-scope closure loses status and interrupts
|
||||
@@ -135,9 +173,28 @@ function decrementSession(input: Map<SessionSchema.ID, number>, sessionID: Sessi
|
||||
export const make = Effect.gen(function* () {
|
||||
const state: State = {
|
||||
jobs: yield* SynchronizedRef.make(new Map()),
|
||||
notifications: new Set(),
|
||||
scope: yield* Scope.Scope,
|
||||
shuttingDown: false,
|
||||
}
|
||||
|
||||
const notify = Effect.fnUntraced(function* (notification: Notification) {
|
||||
yield* notification.effect.pipe(
|
||||
Effect.catchCause((cause) =>
|
||||
Effect.logError("Failed to notify background Job settlement", { jobID: notification.jobID, cause }),
|
||||
),
|
||||
Effect.ensuring(
|
||||
Effect.sync(() => state.notifications.delete(notification.done)).pipe(
|
||||
Effect.andThen(Deferred.succeed(notification.done, undefined)),
|
||||
),
|
||||
),
|
||||
)
|
||||
})
|
||||
|
||||
const launchNotification = Effect.fnUntraced(function* (notification: Notification) {
|
||||
yield* notify(notification).pipe(Effect.forkIn(state.scope, { startImmediately: true }))
|
||||
})
|
||||
|
||||
const settle = Effect.fnUntraced(function* (id: string, token: object, exit: Exit.Exit<string, unknown>) {
|
||||
const completed_at = yield* Clock.currentTimeMillis
|
||||
const result = yield* SynchronizedRef.modify(state.jobs, (jobs): readonly [FinishResult, Map<string, Active>] => {
|
||||
@@ -150,22 +207,37 @@ export const make = Effect.gen(function* () {
|
||||
: Cause.hasInterruptsOnly(exit.cause)
|
||||
? "cancelled"
|
||||
: "error"
|
||||
const info = {
|
||||
...job.info,
|
||||
status,
|
||||
completed_at,
|
||||
...(Exit.isSuccess(exit) ? { output: exit.value } : {}),
|
||||
...(Exit.isFailure(exit) ? { error: errorText(Cause.squash(exit.cause)) } : {}),
|
||||
...(job.info.metadata ? { metadata: { ...job.info.metadata } } : {}),
|
||||
} satisfies SettledInfo
|
||||
const next = {
|
||||
...job,
|
||||
blockingSessions: new Map<SessionSchema.ID, number>(),
|
||||
info: {
|
||||
...job.info,
|
||||
status,
|
||||
completed_at,
|
||||
...(Exit.isSuccess(exit) ? { output: exit.value } : {}),
|
||||
...(Exit.isFailure(exit) ? { error: errorText(Cause.squash(exit.cause)) } : {}),
|
||||
},
|
||||
info,
|
||||
}
|
||||
return [{ info: snapshot(next), done: job.done, scope: job.scope }, new Map(jobs).set(id, next)]
|
||||
const notification = claimNotification(next, info)
|
||||
return [
|
||||
{
|
||||
info,
|
||||
done: job.done,
|
||||
notify: notification.notify,
|
||||
scope: job.scope,
|
||||
},
|
||||
new Map(jobs).set(id, notification.job),
|
||||
]
|
||||
})
|
||||
if (result.info && result.done) yield* Deferred.succeed(result.done, result.info).pipe(Effect.ignore)
|
||||
if (result.notify) yield* launchNotification(result.notify)
|
||||
if (result.scope) {
|
||||
yield* Scope.close(result.scope, Exit.void).pipe(Effect.forkIn(state.scope, { startImmediately: true }))
|
||||
yield* Scope.close(result.scope, Exit.void).pipe(
|
||||
Effect.catchCause((cause) => Effect.logError("Failed to close settled Job scope", { id, cause })),
|
||||
Effect.forkIn(state.scope, { startImmediately: true }),
|
||||
)
|
||||
}
|
||||
return result.info
|
||||
})
|
||||
@@ -197,15 +269,43 @@ export const make = Effect.gen(function* () {
|
||||
Effect.gen(function* () {
|
||||
const id = input.id ?? Identifier.ascending("job")
|
||||
const started_at = yield* Clock.currentTimeMillis
|
||||
const done = yield* Deferred.make<Info>()
|
||||
const backgrounded = yield* Deferred.make<Info>()
|
||||
const result = yield* SynchronizedRef.modifyEffect(
|
||||
state.jobs,
|
||||
Effect.fnUntraced(function* (jobs) {
|
||||
if (state.shuttingDown)
|
||||
return [
|
||||
{
|
||||
info: {
|
||||
id,
|
||||
type: input.type,
|
||||
title: input.title,
|
||||
status: "cancelled",
|
||||
started_at,
|
||||
completed_at: started_at,
|
||||
metadata: input.metadata,
|
||||
},
|
||||
},
|
||||
jobs,
|
||||
] as readonly [StartResult, Map<string, Active>]
|
||||
const existing = jobs.get(id)
|
||||
if (existing?.info.status === "running") {
|
||||
return [{ info: snapshot(existing) }, jobs] as readonly [StartResult, Map<string, Active>]
|
||||
if (existing.onBackgroundSettled || !input.onBackgroundSettled)
|
||||
return [{ info: snapshot(existing) }, jobs] as readonly [StartResult, Map<string, Active>]
|
||||
const backgroundNotification = yield* Deferred.make<void>()
|
||||
const adopted = {
|
||||
...existing,
|
||||
onBackgroundSettled: input.onBackgroundSettled,
|
||||
backgroundNotification,
|
||||
}
|
||||
if (adopted.isBackgrounded) state.notifications.add(backgroundNotification)
|
||||
return [{ info: snapshot(adopted) }, new Map(jobs).set(id, adopted)] as readonly [
|
||||
StartResult,
|
||||
Map<string, Active>,
|
||||
]
|
||||
}
|
||||
const done = yield* Deferred.make<Info>()
|
||||
const backgrounded = yield* Deferred.make<Info>()
|
||||
const backgroundNotification = input.onBackgroundSettled ? yield* Deferred.make<void>() : undefined
|
||||
const scope = yield* Scope.fork(state.scope, "parallel")
|
||||
const token = {}
|
||||
const job = {
|
||||
@@ -219,6 +319,8 @@ export const make = Effect.gen(function* () {
|
||||
},
|
||||
done,
|
||||
backgrounded,
|
||||
onBackgroundSettled: input.onBackgroundSettled,
|
||||
backgroundNotification,
|
||||
scope,
|
||||
token,
|
||||
blockingSessions: new Map<SessionSchema.ID, number>(),
|
||||
@@ -250,7 +352,11 @@ export const make = Effect.gen(function* () {
|
||||
const removeBlock = Effect.fnUntraced(function* (input: BlockInput) {
|
||||
yield* SynchronizedRef.update(state.jobs, (jobs) => {
|
||||
const job = jobs.get(input.id)
|
||||
if (!job || job.info.status !== "running" || job.isBackgrounded) return jobs
|
||||
if (!job || job.isBackgrounded) return jobs
|
||||
if (job.info.status !== "running") {
|
||||
if (!job.onBackgroundSettled) return jobs
|
||||
return new Map(jobs).set(input.id, clearNotification(job))
|
||||
}
|
||||
return new Map(jobs).set(input.id, {
|
||||
...job,
|
||||
blockingSessions: decrementSession(job.blockingSessions, input.sessionID),
|
||||
@@ -262,7 +368,11 @@ export const make = Effect.gen(function* () {
|
||||
const result = yield* SynchronizedRef.modify(state.jobs, (jobs): readonly [BlockStart, Map<string, Active>] => {
|
||||
const job = jobs.get(input.id)
|
||||
if (!job) return [{ type: "missing" }, jobs]
|
||||
if (job.info.status !== "running") return [{ type: "finished", info: snapshot(job) }, jobs]
|
||||
if (job.info.status !== "running")
|
||||
return [
|
||||
{ type: "finished", info: snapshot(job) },
|
||||
job.onBackgroundSettled ? new Map(jobs).set(input.id, clearNotification(job)) : jobs,
|
||||
]
|
||||
if (job.isBackgrounded) return [{ type: "backgrounded", info: snapshot(job) }, jobs]
|
||||
return [
|
||||
{ type: "wait", wait: { done: job.done, backgrounded: job.backgrounded } },
|
||||
@@ -286,16 +396,38 @@ export const make = Effect.gen(function* () {
|
||||
state.jobs,
|
||||
(jobs): readonly [BackgroundResult, Map<string, Active>] => {
|
||||
const job = jobs.get(id)
|
||||
if (!job || job.info.status !== "running") return [{}, jobs]
|
||||
if (!job) return [{}, jobs]
|
||||
if (state.shuttingDown) {
|
||||
if (job.info.status === "running") return [{ info: snapshot(job), cancel: { id, token: job.token } }, jobs]
|
||||
return [
|
||||
{ info: snapshot(job) },
|
||||
job.onBackgroundSettled ? new Map(jobs).set(id, clearNotification(job)) : jobs,
|
||||
]
|
||||
}
|
||||
if (job.info.status !== "running") {
|
||||
if (!job.onBackgroundSettled || !job.backgroundNotification) return [{}, jobs]
|
||||
const info = {
|
||||
...snapshot(job),
|
||||
status: job.info.status,
|
||||
completed_at: job.info.completed_at ?? job.info.started_at,
|
||||
} satisfies SettledInfo
|
||||
const next = { ...job, info, isBackgrounded: true }
|
||||
state.notifications.add(job.backgroundNotification)
|
||||
const notification = claimNotification(next, info)
|
||||
return [{ info, notify: notification.notify }, new Map(jobs).set(id, notification.job)]
|
||||
}
|
||||
if (job.isBackgrounded) return [{ info: snapshot(job) }, jobs]
|
||||
const next = {
|
||||
...job,
|
||||
isBackgrounded: true,
|
||||
blockingSessions: new Map<SessionSchema.ID, number>(),
|
||||
}
|
||||
if (job.backgroundNotification) state.notifications.add(job.backgroundNotification)
|
||||
return [{ info: snapshot(next), backgrounded: job.backgrounded }, new Map(jobs).set(id, next)]
|
||||
},
|
||||
)
|
||||
if (result.cancel) return yield* cancelGeneration(result.cancel.id, result.cancel.token)
|
||||
if (result.notify) yield* launchNotification(result.notify)
|
||||
if (result.info && result.backgrounded)
|
||||
yield* Deferred.succeed(result.backgrounded, result.info).pipe(Effect.ignore)
|
||||
return result.info
|
||||
@@ -305,6 +437,7 @@ export const make = Effect.gen(function* () {
|
||||
const result = yield* SynchronizedRef.modify(
|
||||
state.jobs,
|
||||
(jobs): readonly [BackgroundResult[], Map<string, Active>] => {
|
||||
if (state.shuttingDown) return [[], jobs]
|
||||
const results: BackgroundResult[] = []
|
||||
const next = new Map(jobs)
|
||||
for (const [id, job] of jobs) {
|
||||
@@ -317,6 +450,7 @@ export const make = Effect.gen(function* () {
|
||||
isBackgrounded: true,
|
||||
blockingSessions: new Map<SessionSchema.ID, number>(),
|
||||
}
|
||||
if (job.backgroundNotification) state.notifications.add(job.backgroundNotification)
|
||||
results.push({ info: snapshot(updated), backgrounded: job.backgrounded })
|
||||
next.set(id, updated)
|
||||
}
|
||||
@@ -331,29 +465,65 @@ export const make = Effect.gen(function* () {
|
||||
return result.flatMap((item) => (item.info ? [item.info] : []))
|
||||
})
|
||||
|
||||
const cancel: Interface["cancel"] = Effect.fn("Job.cancel")(function* (id) {
|
||||
const cancelGeneration = Effect.fnUntraced(function* (id: string, token?: object) {
|
||||
const completed_at = yield* Clock.currentTimeMillis
|
||||
const result = yield* SynchronizedRef.modify(state.jobs, (jobs): readonly [FinishResult, Map<string, Active>] => {
|
||||
const job = jobs.get(id)
|
||||
if (!job) return [{}, jobs]
|
||||
if (token && job.token !== token) return [{}, jobs]
|
||||
if (job.info.status !== "running") return [{ info: snapshot(job) }, jobs]
|
||||
const info = {
|
||||
...job.info,
|
||||
status: "cancelled" as const,
|
||||
completed_at,
|
||||
...(job.info.metadata ? { metadata: { ...job.info.metadata } } : {}),
|
||||
} satisfies SettledInfo
|
||||
const next = {
|
||||
...job,
|
||||
blockingSessions: new Map<SessionSchema.ID, number>(),
|
||||
info: {
|
||||
...job.info,
|
||||
status: "cancelled" as const,
|
||||
completed_at,
|
||||
},
|
||||
info,
|
||||
}
|
||||
return [{ info: snapshot(next), done: job.done, scope: job.scope }, new Map(jobs).set(id, next)]
|
||||
const notification = claimNotification(next, info)
|
||||
return [
|
||||
{
|
||||
info,
|
||||
done: job.done,
|
||||
notify: notification.notify,
|
||||
scope: job.scope,
|
||||
},
|
||||
new Map(jobs).set(id, notification.notify ? notification.job : clearNotification(notification.job)),
|
||||
]
|
||||
})
|
||||
if (result.scope)
|
||||
yield* Scope.close(result.scope, Exit.void).pipe(
|
||||
Effect.catchCause((cause) => Effect.logError("Failed to close cancelled Job scope", { id, cause })),
|
||||
)
|
||||
if (result.info && result.done) yield* Deferred.succeed(result.done, result.info).pipe(Effect.ignore)
|
||||
if (result.scope) yield* Scope.close(result.scope, Exit.void)
|
||||
if (result.notify) yield* launchNotification(result.notify)
|
||||
return result.info
|
||||
})
|
||||
|
||||
return Service.of({ get, start, wait, block, background, backgroundAll, cancel })
|
||||
const cancel: Interface["cancel"] = Effect.fn("Job.cancel")((id) => cancelGeneration(id))
|
||||
|
||||
const shutdown: Interface["shutdown"] = Effect.gen(function* () {
|
||||
const drain = yield* SynchronizedRef.modify(state.jobs, (jobs) => {
|
||||
state.shuttingDown = true
|
||||
return [
|
||||
{
|
||||
running: Array.from(jobs.values()).filter((job) => job.info.status === "running" && job.isBackgrounded),
|
||||
notifications: Array.from(state.notifications),
|
||||
},
|
||||
jobs,
|
||||
] as const
|
||||
})
|
||||
yield* Effect.forEach(drain.running, (job) => cancelGeneration(job.info.id, job.token), {
|
||||
concurrency: "unbounded",
|
||||
discard: true,
|
||||
})
|
||||
yield* Effect.forEach(drain.notifications, Deferred.await, { concurrency: "unbounded", discard: true })
|
||||
}).pipe(Effect.withSpan("Job.shutdown"))
|
||||
|
||||
return Service.of({ get, start, wait, block, background, backgroundAll, cancel, shutdown })
|
||||
})
|
||||
|
||||
const layer = Layer.effect(Service, make)
|
||||
|
||||
@@ -572,7 +572,7 @@ function cacheKey(source: string) {
|
||||
}
|
||||
|
||||
export function bodyDigest(text: string) {
|
||||
return new Bun.CryptoHasher("sha256").update(text).digest("hex")
|
||||
return Hash.sha256(text)
|
||||
}
|
||||
|
||||
export const layer = (options?: Options) =>
|
||||
|
||||
@@ -156,9 +156,13 @@ export const providerLayerWithCell = (cell: Cell) =>
|
||||
}
|
||||
cell.runtime = runtime
|
||||
yield* Effect.addFinalizer(() =>
|
||||
Effect.sync(() => {
|
||||
if (cell.runtime === runtime) cell.runtime = undefined
|
||||
}),
|
||||
jobs.shutdown.pipe(
|
||||
Effect.ensuring(
|
||||
Effect.sync(() => {
|
||||
if (cell.runtime === runtime) cell.runtime = undefined
|
||||
}),
|
||||
),
|
||||
),
|
||||
)
|
||||
}),
|
||||
)
|
||||
|
||||
@@ -1,2 +1,2 @@
|
||||
export * as SessionMessage from "./message.js"
|
||||
export * as SessionMessage from "@opencode-ai/schema/session-message"
|
||||
export * from "@opencode-ai/schema/session-message"
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
export * as SessionSchema from "./schema.js"
|
||||
export * as SessionSchema from "@opencode-ai/schema/session"
|
||||
|
||||
import { Session } from "@opencode-ai/schema/session"
|
||||
|
||||
|
||||
@@ -49,7 +49,12 @@ type ToolRow = {
|
||||
duration: number | null
|
||||
}
|
||||
|
||||
type ToolSummaryRow = { status: string | null; count: number }
|
||||
type ToolSummaryRow = {
|
||||
calls: number
|
||||
succeeded: number
|
||||
failed: number
|
||||
unfinished: number
|
||||
}
|
||||
|
||||
type ModelAggregate = {
|
||||
model: Model.Ref
|
||||
@@ -188,25 +193,38 @@ export const get = Effect.fn("SessionStats.get")(function* (input: Input = {}) {
|
||||
(range) => {
|
||||
if (toolMode === "summary")
|
||||
return db
|
||||
.all<ToolSummaryRow>(
|
||||
.get<ToolSummaryRow>(
|
||||
sql`
|
||||
SELECT json_extract(content.value, '$.state.status') AS status, count(*) AS count
|
||||
FROM ${SessionMessageTable} AS message
|
||||
JOIN ${SessionTable} AS session ON session.id = message.session_id,
|
||||
json_each(message.data, '$.content') AS content
|
||||
WHERE message.type = 'assistant'
|
||||
AND message.time_created >= ${range.from}
|
||||
AND message.time_created < ${range.to}
|
||||
AND (session.fork_session_id IS NULL OR message.time_created >= session.time_created)
|
||||
AND json_extract(content.value, '$.type') = 'tool'
|
||||
${project}
|
||||
GROUP BY status
|
||||
WITH calls AS MATERIALIZED (
|
||||
SELECT json_extract(content.value, '$.state.status') AS status
|
||||
FROM ${SessionMessageTable} AS message
|
||||
JOIN ${SessionTable} AS session ON session.id = message.session_id,
|
||||
json_each(message.data, '$.content') AS content
|
||||
WHERE message.type = 'assistant'
|
||||
AND message.time_created >= ${range.from}
|
||||
AND message.time_created < ${range.to}
|
||||
AND (session.fork_session_id IS NULL OR message.time_created >= session.time_created)
|
||||
AND json_extract(content.value, '$.type') = 'tool'
|
||||
${project}
|
||||
)
|
||||
SELECT
|
||||
count(*) AS calls,
|
||||
count(*) FILTER (WHERE status = 'completed') AS succeeded,
|
||||
count(*) FILTER (WHERE status = 'error') AS failed,
|
||||
count(*) FILTER (WHERE status IS NULL OR status NOT IN ('completed', 'error')) AS unfinished
|
||||
FROM calls
|
||||
`,
|
||||
)
|
||||
.pipe(
|
||||
Effect.orDie,
|
||||
Effect.tap((rows) =>
|
||||
Effect.sync(() => rows.forEach((row) => addToolStatus(toolTotals, row.status, row.count))),
|
||||
Effect.tap((row) =>
|
||||
Effect.sync(() => {
|
||||
if (!row) return
|
||||
toolTotals.calls += row.calls
|
||||
toolTotals.succeeded += row.succeeded
|
||||
toolTotals.failed += row.failed
|
||||
toolTotals.unfinished += row.unfinished
|
||||
}),
|
||||
),
|
||||
Effect.asVoid,
|
||||
)
|
||||
|
||||
@@ -3,9 +3,10 @@ export * as ShellTool from "./shell.js"
|
||||
import path from "path"
|
||||
import { ToolFailure } from "@opencode-ai/ai"
|
||||
import type { Context as PluginContext } from "@opencode-ai/plugin/effect/plugin"
|
||||
import { Deferred, Effect, Schema, Scope } from "effect"
|
||||
import { Deferred, Effect, Schema } from "effect"
|
||||
import { Config } from "../../config.js"
|
||||
import { Environment } from "../../environment/index.js"
|
||||
import { Job } from "../../job.js"
|
||||
import { LocationMutation } from "../../location-mutation.js"
|
||||
import { Permission } from "../../permission.js"
|
||||
import { PluginRuntime } from "../../plugin/runtime.js"
|
||||
@@ -105,62 +106,49 @@ export const Plugin = {
|
||||
id: "opencode.tool.shell",
|
||||
effect: Effect.fn("ShellTool.Plugin")(function* (ctx: PluginContext) {
|
||||
const runtime = yield* PluginRuntime.Service
|
||||
const scope = yield* Scope.Scope
|
||||
const environment = yield* Environment.Service
|
||||
const mutation = yield* LocationMutation.Service
|
||||
const shell = yield* Shell.Service
|
||||
const permission = yield* Permission.Service
|
||||
const config = yield* Config.Service
|
||||
|
||||
const notifyWhenDone = Effect.fn("ShellTool.notifyWhenDone")(function* (
|
||||
sessionID: SessionSchema.ID,
|
||||
id: string,
|
||||
shellID: string,
|
||||
command: string,
|
||||
settled: Deferred.Deferred<Output>,
|
||||
const notifyWhenSettled = Effect.fn("ShellTool.notifyWhenSettled")(function* (
|
||||
input: {
|
||||
sessionID: SessionSchema.ID
|
||||
id: string
|
||||
shellID: string
|
||||
command: string
|
||||
settled: Deferred.Deferred<Output>
|
||||
},
|
||||
info: Job.SettledInfo,
|
||||
) {
|
||||
yield* runtime.job.wait({ id: id }).pipe(
|
||||
Effect.flatMap((result) =>
|
||||
Effect.gen(function* () {
|
||||
const info = result.info
|
||||
if (!info) return
|
||||
const state =
|
||||
info.status === "completed"
|
||||
? "completed"
|
||||
: info.status === "error"
|
||||
? "error"
|
||||
: info.status === "cancelled"
|
||||
? "cancelled"
|
||||
: undefined
|
||||
if (state === undefined) return
|
||||
const output = state === "completed" ? yield* Deferred.await(settled) : undefined
|
||||
const text = output
|
||||
? resultMessages(output).join("\n\n")
|
||||
: state === "error"
|
||||
? (info.error ?? "Command failed")
|
||||
: "Command cancelled"
|
||||
yield* runtime.session.synthetic({
|
||||
sessionID,
|
||||
text: `<shell id="${id}" state="${state}" command="${command}">\n${text}\n</shell>`,
|
||||
description: command,
|
||||
metadata: {
|
||||
source: "shell",
|
||||
jobID: id,
|
||||
shellID,
|
||||
state,
|
||||
...(output
|
||||
? {
|
||||
truncated: output.truncated,
|
||||
...(output.exit !== undefined ? { exit: output.exit } : {}),
|
||||
...(output.timeout !== undefined ? { timeout: output.timeout } : {}),
|
||||
}
|
||||
: {}),
|
||||
},
|
||||
})
|
||||
}),
|
||||
),
|
||||
Effect.forkIn(scope, { startImmediately: true }),
|
||||
)
|
||||
const output = info.status === "completed" ? yield* Deferred.await(input.settled) : undefined
|
||||
const text = output
|
||||
? resultMessages(output).join("\n\n")
|
||||
: info.status === "error"
|
||||
? (info.error ?? "Command failed")
|
||||
: "Command cancelled"
|
||||
yield* runtime.session
|
||||
.synthetic({
|
||||
sessionID: input.sessionID,
|
||||
text: `<shell id="${input.id}" state="${info.status}" command="${input.command}">\n${text}\n</shell>`,
|
||||
description: input.command,
|
||||
metadata: {
|
||||
source: "shell",
|
||||
jobID: input.id,
|
||||
shellID: input.shellID,
|
||||
state: info.status,
|
||||
...(output
|
||||
? {
|
||||
truncated: output.truncated,
|
||||
...(output.exit !== undefined ? { exit: output.exit } : {}),
|
||||
...(output.timeout !== undefined ? { timeout: output.timeout } : {}),
|
||||
}
|
||||
: {}),
|
||||
},
|
||||
...(info.status === "cancelled" ? { resume: false } : {}),
|
||||
})
|
||||
.pipe(Effect.ignore)
|
||||
})
|
||||
|
||||
yield* ctx.tool
|
||||
@@ -294,6 +282,13 @@ export const Plugin = {
|
||||
})
|
||||
|
||||
const settled = yield* Deferred.make<Output>()
|
||||
const notification = {
|
||||
sessionID: context.sessionID,
|
||||
id: context.id,
|
||||
shellID: info.id,
|
||||
command: info.command,
|
||||
settled,
|
||||
}
|
||||
const run = settleShell().pipe(
|
||||
Effect.tap((output) => Deferred.succeed(settled, output)),
|
||||
Effect.map((output) => output.output),
|
||||
@@ -305,11 +300,16 @@ export const Plugin = {
|
||||
title: info.command,
|
||||
metadata: { sessionID: context.sessionID, shellID: info.id },
|
||||
run,
|
||||
onBackgroundSettled: (result) => notifyWhenSettled(notification, result),
|
||||
})
|
||||
if (job.status === "cancelled") {
|
||||
yield* shell.remove(info.id).pipe(Effect.ignore)
|
||||
return yield* Effect.fail(new Error("Command cancelled"))
|
||||
}
|
||||
|
||||
if (input.background === true) {
|
||||
yield* runtime.job.background(job.id)
|
||||
yield* notifyWhenDone(context.sessionID, context.id, info.id, info.command, settled)
|
||||
const background = yield* runtime.job.background(job.id)
|
||||
if (background?.status === "cancelled") return yield* Effect.fail(new Error("Command cancelled"))
|
||||
return backgroundResult(info.id)
|
||||
}
|
||||
|
||||
@@ -318,7 +318,6 @@ export const Plugin = {
|
||||
.pipe(Effect.onInterrupt(() => runtime.job.cancel(job.id).pipe(Effect.ignore)))
|
||||
if (result?.type === "backgrounded") {
|
||||
yield* shell.timeout(info.id, 0)
|
||||
yield* notifyWhenDone(context.sessionID, context.id, info.id, info.command, settled)
|
||||
return backgroundResult(info.id)
|
||||
}
|
||||
if (result?.info.status === "error")
|
||||
|
||||
@@ -2,7 +2,7 @@ export * as SubagentTool from "./subagent.js"
|
||||
|
||||
import { ToolFailure } from "@opencode-ai/ai"
|
||||
import type { Context as PluginContext } from "@opencode-ai/plugin/effect/plugin"
|
||||
import { Effect, Schema, Scope } from "effect"
|
||||
import { Effect, Schema } from "effect"
|
||||
import { Agent } from "../../agent.js"
|
||||
import { Config } from "../../config.js"
|
||||
import { PluginRuntime } from "../../plugin/runtime.js"
|
||||
@@ -57,11 +57,6 @@ export const Plugin = {
|
||||
const agents = yield* Agent.Service
|
||||
const config = yield* Config.Service
|
||||
const permission = yield* Permission.Service
|
||||
const scope = yield* Scope.Scope
|
||||
// One completion observer per job generation. Keyed by child plus start time so a fresh
|
||||
// continuation job is observable even while a settled generation's observer is finalizing.
|
||||
const notifications = new Set<string>()
|
||||
|
||||
// Concatenate the child's final completed assistant text. Distinguishes "completed with no
|
||||
// text" (generic string) from "failed" (the run effect fails, surfaced as a job error).
|
||||
const latestAssistantText = Effect.fn("SubagentTool.latestAssistantText")(function* (sessionID: SessionSchema.ID) {
|
||||
@@ -86,44 +81,15 @@ export const Plugin = {
|
||||
state: "completed" | "error" | "cancelled",
|
||||
text: string,
|
||||
) {
|
||||
yield* runtime.session.synthetic({
|
||||
sessionID: parentID,
|
||||
text: `<subagent sessionID="${childID}" state="${state}" description="${description}">\n${text}\n</subagent>`,
|
||||
description,
|
||||
metadata: { source: "subagent", childID, agent, state },
|
||||
})
|
||||
})
|
||||
|
||||
const notifyWhenDone = Effect.fn("SubagentTool.notifyWhenDone")(function* (
|
||||
parentID: SessionSchema.ID,
|
||||
childID: SessionSchema.ID,
|
||||
agent: string,
|
||||
description: string,
|
||||
startedAt: number,
|
||||
) {
|
||||
const key = `${childID}:${startedAt}`
|
||||
if (notifications.has(key)) return
|
||||
notifications.add(key)
|
||||
yield* runtime.job.wait({ id: childID }).pipe(
|
||||
Effect.flatMap((result) => {
|
||||
if (result.info?.status === "completed")
|
||||
return injectCompletion(parentID, childID, agent, description, "completed", result.info.output ?? NO_TEXT)
|
||||
if (result.info?.status === "error")
|
||||
return injectCompletion(
|
||||
parentID,
|
||||
childID,
|
||||
agent,
|
||||
description,
|
||||
"error",
|
||||
result.info.error ?? "Subagent failed",
|
||||
)
|
||||
if (result.info?.status === "cancelled")
|
||||
return injectCompletion(parentID, childID, agent, description, "cancelled", "Subagent cancelled")
|
||||
return Effect.void
|
||||
}),
|
||||
Effect.ensuring(Effect.sync(() => notifications.delete(key))),
|
||||
Effect.forkIn(scope, { startImmediately: true }),
|
||||
)
|
||||
yield* runtime.session
|
||||
.synthetic({
|
||||
sessionID: parentID,
|
||||
text: `<subagent sessionID="${childID}" state="${state}" description="${description}">\n${text}\n</subagent>`,
|
||||
description,
|
||||
metadata: { source: "subagent", childID, agent, state },
|
||||
...(state === "cancelled" ? { resume: false } : {}),
|
||||
})
|
||||
.pipe(Effect.ignore)
|
||||
})
|
||||
|
||||
yield* ctx.tool
|
||||
@@ -257,11 +223,32 @@ export const Plugin = {
|
||||
title: input.description,
|
||||
metadata: {},
|
||||
run,
|
||||
onBackgroundSettled: (result) => {
|
||||
const text =
|
||||
result.status === "completed"
|
||||
? (result.output ?? NO_TEXT)
|
||||
: result.status === "error"
|
||||
? (result.error ?? "Subagent failed")
|
||||
: "Subagent cancelled"
|
||||
return injectCompletion(
|
||||
context.sessionID,
|
||||
child.id,
|
||||
agent.name,
|
||||
input.description,
|
||||
result.status,
|
||||
text,
|
||||
)
|
||||
},
|
||||
})
|
||||
if (info.status === "cancelled") {
|
||||
yield* runtime.session.interrupt(child.id)
|
||||
return yield* new ToolFailure({ message: `Subagent cancelled (sessionID: ${child.id})` })
|
||||
}
|
||||
|
||||
if (background) {
|
||||
yield* runtime.job.background(info.id)
|
||||
yield* notifyWhenDone(context.sessionID, child.id, agent.name, input.description, info.started_at)
|
||||
const result = yield* runtime.job.background(info.id)
|
||||
if (result?.status === "cancelled")
|
||||
return yield* new ToolFailure({ message: `Subagent cancelled (sessionID: ${child.id})` })
|
||||
return backgroundResult(child.id)
|
||||
}
|
||||
|
||||
@@ -273,13 +260,6 @@ export const Plugin = {
|
||||
),
|
||||
)
|
||||
if (result?.type === "backgrounded") {
|
||||
yield* notifyWhenDone(
|
||||
context.sessionID,
|
||||
child.id,
|
||||
agent.name,
|
||||
input.description,
|
||||
result.info.started_at,
|
||||
)
|
||||
return backgroundResult(child.id)
|
||||
}
|
||||
// Failure surfaces keep the sessionID visible so the model can continue the child.
|
||||
|
||||
@@ -3,25 +3,25 @@ import { Global } from "@opencode-ai/util/global"
|
||||
import { Effect, Layer } from "effect"
|
||||
import { tmpdir } from "./tmpdir"
|
||||
|
||||
export function globalLayer(root: string) {
|
||||
const data = path.join(root, "data")
|
||||
const cache = path.join(root, "cache")
|
||||
return Global.layerWith({
|
||||
home: path.join(root, "home"),
|
||||
data,
|
||||
cache,
|
||||
config: path.join(root, "config"),
|
||||
state: path.join(root, "state"),
|
||||
tmp: path.join(root, "tmp"),
|
||||
bin: path.join(cache, "bin"),
|
||||
log: path.join(data, "log"),
|
||||
repos: path.join(data, "repos"),
|
||||
})
|
||||
}
|
||||
|
||||
export const tempGlobalLayer = Layer.unwrap(
|
||||
Effect.acquireRelease(
|
||||
Effect.promise(() => tmpdir()),
|
||||
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
|
||||
).pipe(
|
||||
Effect.map((tmp) => {
|
||||
const data = path.join(tmp.path, "data")
|
||||
const cache = path.join(tmp.path, "cache")
|
||||
return Global.layerWith({
|
||||
home: path.join(tmp.path, "home"),
|
||||
data,
|
||||
cache,
|
||||
config: path.join(tmp.path, "config"),
|
||||
state: path.join(tmp.path, "state"),
|
||||
tmp: path.join(tmp.path, "tmp"),
|
||||
bin: path.join(cache, "bin"),
|
||||
log: path.join(data, "log"),
|
||||
repos: path.join(data, "repos"),
|
||||
})
|
||||
}),
|
||||
),
|
||||
).pipe(Effect.map((tmp) => globalLayer(tmp.path))),
|
||||
)
|
||||
|
||||
@@ -145,6 +145,111 @@ describe("Job", () => {
|
||||
}),
|
||||
)
|
||||
|
||||
it.live("shutdown cancels only background jobs and awaits their terminal callbacks", () =>
|
||||
Effect.gen(function* () {
|
||||
const jobs = yield* Job.Service
|
||||
const foregroundWork = yield* Deferred.make<void>()
|
||||
const interrupted = yield* Deferred.make<void>()
|
||||
const callbackStarted = yield* Deferred.make<Job.SettledInfo>()
|
||||
const releaseCallback = yield* Deferred.make<void>()
|
||||
const background = yield* jobs.start({
|
||||
id: "job_background_shutdown",
|
||||
type: "test",
|
||||
run: Effect.never.pipe(Effect.ensuring(Deferred.succeed(interrupted, undefined))),
|
||||
onBackgroundSettled: (info) =>
|
||||
Deferred.await(interrupted).pipe(
|
||||
Effect.andThen(Deferred.succeed(callbackStarted, info)),
|
||||
Effect.andThen(Deferred.await(releaseCallback)),
|
||||
),
|
||||
})
|
||||
const foreground = yield* jobs.start({
|
||||
id: "job_foreground_shutdown",
|
||||
type: "test",
|
||||
run: Deferred.await(foregroundWork).pipe(Effect.as("foreground")),
|
||||
})
|
||||
yield* jobs.background(background.id)
|
||||
|
||||
const shutdown = yield* jobs.shutdown.pipe(Effect.forkIn(yield* Scope.Scope, { startImmediately: true }))
|
||||
expect(yield* Deferred.await(callbackStarted)).toMatchObject({
|
||||
id: background.id,
|
||||
status: "cancelled",
|
||||
})
|
||||
expect(shutdown.pollUnsafe()).toBeUndefined()
|
||||
expect(yield* jobs.get(foreground.id)).toMatchObject({ status: "running" })
|
||||
|
||||
yield* Deferred.succeed(releaseCallback, undefined)
|
||||
yield* Fiber.join(shutdown)
|
||||
expect(yield* jobs.get(background.id)).toMatchObject({ status: "cancelled" })
|
||||
expect(yield* jobs.get(foreground.id)).toMatchObject({ status: "running" })
|
||||
expect(yield* jobs.background(foreground.id)).toMatchObject({ status: "cancelled" })
|
||||
expect(yield* jobs.start({ id: "job_started_during_shutdown", type: "test", run: Effect.never })).toMatchObject({
|
||||
status: "cancelled",
|
||||
})
|
||||
expect(yield* jobs.get("job_started_during_shutdown")).toBeUndefined()
|
||||
}),
|
||||
)
|
||||
|
||||
it.live("shutdown awaits a background callback already running after completion", () =>
|
||||
Effect.gen(function* () {
|
||||
const jobs = yield* Job.Service
|
||||
const work = yield* Deferred.make<void>()
|
||||
const callbackStarted = yield* Deferred.make<void>()
|
||||
const releaseCallback = yield* Deferred.make<void>()
|
||||
const job = yield* jobs.start({
|
||||
id: "job_completed_before_shutdown",
|
||||
type: "test",
|
||||
run: Deferred.await(work).pipe(Effect.as("done")),
|
||||
onBackgroundSettled: () =>
|
||||
Deferred.succeed(callbackStarted, undefined).pipe(Effect.andThen(Deferred.await(releaseCallback))),
|
||||
})
|
||||
yield* jobs.background(job.id)
|
||||
yield* Deferred.succeed(work, undefined)
|
||||
yield* Deferred.await(callbackStarted)
|
||||
expect(yield* jobs.get(job.id)).toMatchObject({ status: "completed" })
|
||||
|
||||
const shutdown = yield* jobs.shutdown.pipe(Effect.forkIn(yield* Scope.Scope, { startImmediately: true }))
|
||||
yield* Effect.yieldNow
|
||||
expect(shutdown.pollUnsafe()).toBeUndefined()
|
||||
|
||||
yield* Deferred.succeed(releaseCallback, undefined)
|
||||
yield* Fiber.join(shutdown)
|
||||
}),
|
||||
)
|
||||
|
||||
it.live("shutdown rejects late background notification registration", () =>
|
||||
Effect.gen(function* () {
|
||||
const jobs = yield* Job.Service
|
||||
const callback = yield* Deferred.make<void>()
|
||||
const job = yield* jobs.start({
|
||||
id: "job_completed_before_background_shutdown",
|
||||
type: "test",
|
||||
run: Effect.succeed("done"),
|
||||
onBackgroundSettled: () => Deferred.succeed(callback, undefined),
|
||||
})
|
||||
expect(yield* jobs.wait({ id: job.id })).toMatchObject({ info: { status: "completed" } })
|
||||
|
||||
yield* jobs.shutdown
|
||||
expect(yield* jobs.background(job.id)).toMatchObject({ status: "completed" })
|
||||
expect((yield* Deferred.poll(callback))._tag).toBe("None")
|
||||
}),
|
||||
)
|
||||
|
||||
it.live("shutdown completes when a background callback defects", () =>
|
||||
Effect.gen(function* () {
|
||||
const jobs = yield* Job.Service
|
||||
const job = yield* jobs.start({
|
||||
id: "job_defective_shutdown_callback",
|
||||
type: "test",
|
||||
run: Effect.never,
|
||||
onBackgroundSettled: () => Effect.die("callback defect"),
|
||||
})
|
||||
yield* jobs.background(job.id)
|
||||
|
||||
yield* jobs.shutdown.pipe(Effect.timeout("1 second"))
|
||||
expect(yield* jobs.get(job.id)).toMatchObject({ status: "cancelled" })
|
||||
}),
|
||||
)
|
||||
|
||||
it.live("interrupts live work without promising settlement after the owning process-local scope closes", () =>
|
||||
Effect.gen(function* () {
|
||||
const scope = yield* Scope.make()
|
||||
|
||||
@@ -113,6 +113,13 @@ describe("SessionStats", () => {
|
||||
completed: DateTime.makeUnsafe(Date.UTC(2026, 0, 2, 10, 0, 2)),
|
||||
},
|
||||
}),
|
||||
SessionMessage.AssistantTool.make({
|
||||
type: "tool",
|
||||
id: "call_pending",
|
||||
name: "pending",
|
||||
state: SessionMessage.ToolStateRunning.make({ status: "running", input: {}, metadata: {} }),
|
||||
time: { created: DateTime.makeUnsafe(Date.UTC(2026, 0, 2, 10)) },
|
||||
}),
|
||||
]),
|
||||
),
|
||||
messageRow(childID, 1, assistant("msg_stats_child", Date.UTC(2026, 0, 3, 10), [], "large", 2)),
|
||||
@@ -210,7 +217,7 @@ describe("SessionStats", () => {
|
||||
expect(stats.cost).toBe(Money.USD.make(6.25))
|
||||
expect(stats.tools).toMatchObject({
|
||||
mode: "detail",
|
||||
totals: { calls: 2, succeeded: 1, failed: 1, unfinished: 0 },
|
||||
totals: { calls: 3, succeeded: 1, failed: 1, unfinished: 1 },
|
||||
})
|
||||
expect(stats.activity).toEqual([
|
||||
{ date: "2026-01-02", steps: 1 },
|
||||
@@ -224,6 +231,7 @@ describe("SessionStats", () => {
|
||||
expect(stats.tools.usage).toMatchObject([
|
||||
{ name: "read", calls: 1, succeeded: 1, failed: 0, durationP50: 250 },
|
||||
{ name: "edit", calls: 1, succeeded: 0, failed: 1, durationP50: 2_000 },
|
||||
{ name: "pending", calls: 1, succeeded: 0, failed: 0, unfinished: 1 },
|
||||
])
|
||||
|
||||
const summary = yield* SessionStats.get({
|
||||
@@ -234,7 +242,7 @@ describe("SessionStats", () => {
|
||||
expect(summary.models.map((model) => String(model.model.id))).toEqual(["large", "sonnet", "fork-new"])
|
||||
expect(summary.tools).toEqual({
|
||||
mode: "summary",
|
||||
totals: { calls: 2, succeeded: 1, failed: 1, unfinished: 0 },
|
||||
totals: { calls: 3, succeeded: 1, failed: 1, unfinished: 1 },
|
||||
})
|
||||
|
||||
const withoutTools = yield* SessionStats.get({
|
||||
|
||||
@@ -3,7 +3,7 @@ import { realpathSync } from "node:fs"
|
||||
import os from "os"
|
||||
import path from "path"
|
||||
import { describe, expect } from "bun:test"
|
||||
import { Deferred, Duration, Effect, Fiber, Layer, Scope, Stream } from "effect"
|
||||
import { Context, Deferred, Duration, Effect, Exit, Fiber, Layer, Scope, Stream } from "effect"
|
||||
import { Money } from "@opencode-ai/schema/money"
|
||||
import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
|
||||
import { LayerNode } from "@opencode-ai/util/effect/layer-node"
|
||||
@@ -26,6 +26,7 @@ import { Job } from "@opencode-ai/core/job"
|
||||
import { Session } from "@opencode-ai/core/session"
|
||||
import { SessionEvent } from "@opencode-ai/core/session/event"
|
||||
import { SessionExecution } from "@opencode-ai/core/session/execution"
|
||||
import { SessionInbox } from "@opencode-ai/core/session/inbox"
|
||||
import { SessionMessage } from "@opencode-ai/core/session/message"
|
||||
import { SessionStore } from "@opencode-ai/core/session/store"
|
||||
import { Permission } from "@opencode-ai/core/permission"
|
||||
@@ -37,7 +38,7 @@ import { ShellTool } from "@opencode-ai/core/tool/plugin/shell"
|
||||
import { ToolOutput } from "@opencode-ai/core/tool-output"
|
||||
import { Tool } from "@opencode-ai/core/tool"
|
||||
import { tmpdir } from "./fixture/tmpdir"
|
||||
import { tempGlobalLayer } from "./fixture/global"
|
||||
import { globalLayer, tempGlobalLayer } from "./fixture/global"
|
||||
import { testEffect } from "./lib/effect"
|
||||
import { permissionLayer } from "./lib/permission"
|
||||
import { toolIdentity, executeTool, registerToolPlugin, toolDefinitions } from "./lib/tool"
|
||||
@@ -159,6 +160,7 @@ const replacements = [
|
||||
] satisfies LayerNode.Replacements
|
||||
const productionIt = testEffect(AppNodeBuilder.build(nodes, replacements))
|
||||
const it = testEffect(AppNodeBuilder.build(nodes, [...replacements, [PluginSupervisor.node, shellPluginSupervisor]]))
|
||||
const lifecycleIt = testEffect(Layer.empty)
|
||||
|
||||
const call = (input: typeof ShellTool.Input.Type, id = "call-shell") => ({
|
||||
sessionID,
|
||||
@@ -777,6 +779,83 @@ describe("ShellTool", () => {
|
||||
),
|
||||
)
|
||||
|
||||
lifecycleIt.live("notifies the session when application shutdown cancels a background command", () =>
|
||||
Effect.acquireUseRelease(
|
||||
Effect.all([Effect.promise(() => tmpdir()), Scope.make()]),
|
||||
([tmp, applicationScope]) =>
|
||||
Effect.gen(function* () {
|
||||
reset()
|
||||
const databasePath = path.join(tmp.path, "opencode.sqlite")
|
||||
const testGlobalLayer = globalLayer(tmp.path)
|
||||
const context = yield* Layer.buildWithScope(
|
||||
AppNodeBuilder.build(nodes, [
|
||||
[SessionExecution.node, executionNode],
|
||||
[Permission.node, permission],
|
||||
[PluginSupervisor.node, shellPluginSupervisor],
|
||||
[Database.node, Database.configured({ path: databasePath })],
|
||||
[Bus.node, Bus.configured({ persist: true })],
|
||||
[Global.node, testGlobalLayer],
|
||||
]),
|
||||
applicationScope,
|
||||
)
|
||||
const sessions = Context.get(context, Session.Service)
|
||||
const location = Location.Ref.make({ directory: AbsolutePath.make(tmp.path) })
|
||||
yield* sessions.create({
|
||||
id: sessionID,
|
||||
title: "shell shutdown test",
|
||||
location,
|
||||
model: sessionModel,
|
||||
})
|
||||
const locations = Context.get(context, LocationServiceMap.Service)
|
||||
const settled = yield* Effect.gen(function* () {
|
||||
yield* (yield* PluginSupervisor.Service).flush
|
||||
return yield* executeTool(
|
||||
yield* Tool.Service,
|
||||
call({ command: idleCommand, background: true }, "call-shutdown-background"),
|
||||
)
|
||||
}).pipe(Effect.provide(locations.get(location)))
|
||||
|
||||
expect(settled.metadata).toMatchObject({ status: "running" })
|
||||
yield* Scope.close(applicationScope, Exit.void)
|
||||
|
||||
const pending = yield* Layer.build(
|
||||
AppNodeBuilder.build(Database.node, [
|
||||
[Database.node, Database.configured({ path: databasePath })],
|
||||
[Global.node, testGlobalLayer],
|
||||
]),
|
||||
).pipe(
|
||||
Effect.flatMap((verification) =>
|
||||
SessionInbox.list(Context.get(verification, Database.Service).db, sessionID),
|
||||
),
|
||||
Effect.scoped,
|
||||
)
|
||||
const cancellation = pending.find(
|
||||
(item) =>
|
||||
item.type === "synthetic" &&
|
||||
item.payload.metadata?.source === "shell" &&
|
||||
item.payload.metadata.jobID === "call-shutdown-background",
|
||||
)
|
||||
expect(cancellation).toMatchObject({
|
||||
type: "synthetic",
|
||||
payload: {
|
||||
text: expect.stringContaining("Command cancelled"),
|
||||
description: idleCommand,
|
||||
metadata: {
|
||||
source: "shell",
|
||||
jobID: "call-shutdown-background",
|
||||
shellID: settled.metadata?.shellID,
|
||||
state: "cancelled",
|
||||
},
|
||||
},
|
||||
})
|
||||
}),
|
||||
([tmp, applicationScope]) =>
|
||||
Scope.close(applicationScope, Exit.void).pipe(
|
||||
Effect.andThen(Effect.promise(() => tmp[Symbol.asyncDispose]().then(() => undefined))),
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
it.live("preserves a background command's non-zero exit", () =>
|
||||
Effect.acquireUseRelease(
|
||||
Effect.promise(() => tmpdir()),
|
||||
|
||||
@@ -30,8 +30,8 @@
|
||||
},
|
||||
"peerDependencies": {
|
||||
"@opencode-ai/theme": "workspace:*",
|
||||
"@opentui/core": ">=0.5.7",
|
||||
"@opentui/solid": ">=0.5.7",
|
||||
"@opentui/core": ">=0.5.8",
|
||||
"@opentui/solid": ">=0.5.8",
|
||||
"solid-js": ">=1.9.0"
|
||||
},
|
||||
"peerDependenciesMeta": {
|
||||
|
||||
@@ -22,7 +22,8 @@
|
||||
"scripts": {
|
||||
"build": "bun run script/build.ts",
|
||||
"test": "bun test --timeout 5000",
|
||||
"typecheck": "tsgo -b"
|
||||
"typecheck": "tsgo -b",
|
||||
"verify:package": "bun run script/verify-package.ts"
|
||||
},
|
||||
"dependencies": {
|
||||
"@opencode-ai/client": "workspace:*",
|
||||
|
||||
@@ -0,0 +1,179 @@
|
||||
#!/usr/bin/env bun
|
||||
|
||||
import { $ } from "bun"
|
||||
import { mkdtemp, rm } from "node:fs/promises"
|
||||
import { tmpdir } from "node:os"
|
||||
import { join } from "node:path"
|
||||
import { fileURLToPath } from "node:url"
|
||||
|
||||
const root = fileURLToPath(new URL("../../..", import.meta.url))
|
||||
const names = ["schema", "codemode", "ai", "util", "protocol", "client", "plugin", "core", "simulation", "server", "sdk"]
|
||||
const temporary = await mkdtemp(join(tmpdir(), "opencode-sdk-package-"))
|
||||
const archives = new Map<string, string>()
|
||||
|
||||
try {
|
||||
for (const name of names) {
|
||||
const directory = join(root, "packages", name)
|
||||
await $`bun run build`.cwd(directory)
|
||||
const original = await Bun.file(join(directory, "package.json")).text()
|
||||
// oxlint-disable-next-line typescript-eslint/no-unsafe-type-assertion -- package manifests are validated by their package builds.
|
||||
const pkg = JSON.parse(original) as {
|
||||
name: string
|
||||
dependencies?: Record<string, string>
|
||||
exports?: Record<string, string | { import: string; types: string }>
|
||||
imports?: Record<string, Record<string, string>>
|
||||
}
|
||||
const archive = join(temporary, `${name}.tgz`)
|
||||
|
||||
if (pkg.dependencies) {
|
||||
const unpacked = Object.keys(pkg.dependencies).filter(
|
||||
(dependency) => dependency.startsWith("@opencode-ai/") && !archives.has(dependency),
|
||||
)
|
||||
if (unpacked.length > 0) throw new Error(`${pkg.name} has unpacked workspace dependencies: ${unpacked.join(", ")}`)
|
||||
pkg.dependencies = Object.fromEntries(
|
||||
Object.entries(pkg.dependencies).map(([dependency, version]) => {
|
||||
const local = archives.get(dependency)
|
||||
return [dependency, local ? `file:${local}` : version]
|
||||
}),
|
||||
)
|
||||
}
|
||||
if (pkg.exports) {
|
||||
pkg.exports = Object.fromEntries(
|
||||
Object.entries(pkg.exports).map(([key, value]) => {
|
||||
if (typeof value !== "string") return [key, value]
|
||||
return [key, { import: output(name, value), types: output(name, value, true) }]
|
||||
}),
|
||||
)
|
||||
}
|
||||
if (pkg.imports) {
|
||||
pkg.imports = Object.fromEntries(
|
||||
Object.entries(pkg.imports).map(([key, conditions]) => [
|
||||
key,
|
||||
Object.fromEntries(
|
||||
Object.entries(conditions).map(([condition, value]) => [condition, output(name, value, condition === "types")]),
|
||||
),
|
||||
]),
|
||||
)
|
||||
}
|
||||
|
||||
await Bun.write(join(directory, "package.json"), JSON.stringify(pkg, null, 2) + "\n")
|
||||
try {
|
||||
await $`bun pm pack --filename ${archive} --ignore-scripts --quiet`.cwd(directory)
|
||||
} finally {
|
||||
await Bun.write(join(directory, "package.json"), original)
|
||||
}
|
||||
archives.set(pkg.name, archive)
|
||||
}
|
||||
|
||||
const consumer = join(temporary, "consumer")
|
||||
await Bun.write(
|
||||
join(consumer, "package.json"),
|
||||
JSON.stringify({ name: "opencode-sdk-consumer", private: true, type: "module" }),
|
||||
)
|
||||
await Promise.all([
|
||||
Bun.write(
|
||||
join(consumer, "wrangler.jsonc"),
|
||||
JSON.stringify({
|
||||
name: "opencode-sdk-packed-consumer",
|
||||
main: "worker.js",
|
||||
compatibility_date: "2026-07-15",
|
||||
compatibility_flags: ["nodejs_compat"],
|
||||
durable_objects: { bindings: [{ name: "OPENCODE", class_name: "OpenCodeDO" }] },
|
||||
migrations: [{ tag: "v1", new_sqlite_classes: ["OpenCodeDO"] }],
|
||||
}),
|
||||
),
|
||||
Bun.write(
|
||||
join(consumer, "worker.js"),
|
||||
`import { bodyDigest } from "@opencode-ai/core/models-dev"
|
||||
import { OpenCodeWorkerd } from "@opencode-ai/sdk/workerd"
|
||||
import { Effect } from "effect"
|
||||
|
||||
export class OpenCodeDO {
|
||||
constructor(state) {
|
||||
this.state = state
|
||||
}
|
||||
|
||||
fetch() {
|
||||
if (bodyDigest("packed-workerd") !== "5fc174bf63e8dd108ebb6c53d85e7bbc4525b2f4c1c43280364cdbfd9b37aaf5") {
|
||||
throw new Error("Packed workerd SHA-256 mismatch")
|
||||
}
|
||||
const storage = this.state.storage
|
||||
return Effect.runPromise(
|
||||
Effect.gen(function* () {
|
||||
const sdk = yield* OpenCodeWorkerd.create({
|
||||
storage,
|
||||
app: { version: "packed-workerd" },
|
||||
config: { content: "{}" },
|
||||
})
|
||||
return Response.json(yield* sdk.health.get())
|
||||
}).pipe(Effect.scoped),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
export default {
|
||||
fetch(request, env) {
|
||||
return env.OPENCODE.get(env.OPENCODE.idFromName("packed-consumer")).fetch(request)
|
||||
},
|
||||
}
|
||||
`,
|
||||
),
|
||||
Bun.write(
|
||||
join(consumer, "boot.mjs"),
|
||||
`import { Miniflare } from "miniflare"
|
||||
|
||||
const miniflare = new Miniflare({
|
||||
compatibilityDate: "2026-07-15",
|
||||
compatibilityFlags: ["nodejs_compat"],
|
||||
modules: true,
|
||||
scriptPath: new URL("./dist/worker.js", import.meta.url).pathname,
|
||||
durableObjects: { OPENCODE: { className: "OpenCodeDO", useSQLite: true } },
|
||||
})
|
||||
|
||||
try {
|
||||
const response = await miniflare.dispatchFetch("http://opencode.local/health")
|
||||
if (response.status !== 200) throw new Error(
|
||||
"Packed workerd health returned " + response.status + ": " + await response.text(),
|
||||
)
|
||||
const body = await response.json()
|
||||
if (body.healthy !== true || body.version !== "packed-workerd") {
|
||||
throw new Error("Unexpected packed workerd health: " + JSON.stringify(body))
|
||||
}
|
||||
} finally {
|
||||
await miniflare.dispose()
|
||||
}
|
||||
`,
|
||||
),
|
||||
])
|
||||
|
||||
const sdk = archives.get("@opencode-ai/sdk")
|
||||
if (!sdk) throw new Error("Packed SDK archive was not created")
|
||||
await $`npm install --ignore-scripts --no-audit --no-fund --package-lock=false ${sdk} wrangler@4.110.0`.cwd(consumer)
|
||||
await $`node_modules/.bin/wrangler deploy --dry-run --config wrangler.jsonc --outdir dist`.cwd(consumer)
|
||||
|
||||
const transpiler = new Bun.Transpiler({ loader: "js" })
|
||||
const bundled = await Bun.file(join(consumer, "dist/worker.js")).text()
|
||||
if (/createRequire\s*\(\s*import\.meta\.url\s*\)/.test(bundled)) {
|
||||
throw new Error("Packed workerd bundle contains Bun's eager Node require initializer")
|
||||
}
|
||||
const bunGlobals = Array.from(new Set(bundled.match(/\bBun\.[A-Za-z_$][\w$]*/g) ?? []))
|
||||
if (bunGlobals.length > 0) throw new Error(`Packed workerd bundle references Bun globals: ${bunGlobals.join(", ")}`)
|
||||
const leaked = [
|
||||
...transpiler.scanImports(bundled)
|
||||
.filter((imported) => imported.kind !== "dynamic-import")
|
||||
.map((imported) => imported.path),
|
||||
...Array.from(bundled.matchAll(/\brequire\(\s*["']([^"']+)["']\s*\)/g), (match) => match[1]),
|
||||
]
|
||||
.filter((specifier) => specifier === "bun" || specifier.startsWith("bun:"))
|
||||
if (leaked.length > 0) throw new Error(`Packed workerd bundle statically imports Bun builtins: ${leaked.join(", ")}`)
|
||||
|
||||
await $`node boot.mjs`.cwd(consumer)
|
||||
console.log("packed SDK consumer OK")
|
||||
} finally {
|
||||
await rm(temporary, { recursive: true, force: true })
|
||||
}
|
||||
|
||||
function output(name: string, value: string, types = false) {
|
||||
const root = name === "core" && types ? "./dist/types/" : "./dist/"
|
||||
return value.replace("./src/", root).replace(/\.ts$/, types ? ".d.ts" : ".js")
|
||||
}
|
||||
@@ -25,7 +25,7 @@
|
||||
"default": "./src/global-roots.ts"
|
||||
},
|
||||
"#runtime-import": {
|
||||
"workerd": "./src/runtime/import.bun.ts",
|
||||
"workerd": "./src/runtime/import.workerd.ts",
|
||||
"bun": "./src/runtime/import.bun.ts",
|
||||
"node": "./src/runtime/import.node.ts",
|
||||
"default": "./src/runtime/import.bun.ts"
|
||||
|
||||
@@ -60,7 +60,10 @@ export function fileLogger(target = file(), id: string = runID()) {
|
||||
})
|
||||
}
|
||||
|
||||
const stderrLogger = Logger.make((options) => process.stderr.write(formatter().log(options) + "\n"))
|
||||
const stderrLogger = Logger.make((options) => {
|
||||
if (process.env.OPENCODE_PRINT_LOGS !== "1") return
|
||||
process.stderr.write(formatter().log(options) + "\n")
|
||||
})
|
||||
|
||||
export function minimumLogLevel() {
|
||||
const value = process.env.OPENCODE_LOG_LEVEL?.toUpperCase()
|
||||
@@ -74,8 +77,7 @@ export function minimumLogLevel() {
|
||||
}
|
||||
|
||||
export function loggers(local = true, channel = "local") {
|
||||
const logger = fileLogger(file(local, channel))
|
||||
return process.env.OPENCODE_PRINT_LOGS === "1" ? [logger, stderrLogger] : [logger]
|
||||
return [fileLogger(file(local, channel)), stderrLogger]
|
||||
}
|
||||
|
||||
export * as Logging from "./logging.js"
|
||||
|
||||
@@ -0,0 +1,9 @@
|
||||
const unavailable = () => new Error("Dynamic module loading is unavailable on workerd")
|
||||
|
||||
export function importModule(_specifier: string): Promise<unknown> {
|
||||
return Promise.reject(unavailable())
|
||||
}
|
||||
|
||||
export function resolveModule(_specifier: string, _directory: string): string {
|
||||
throw unavailable()
|
||||
}
|
||||
Reference in New Issue
Block a user