Compare commits

..

1 Commits

Author SHA1 Message Date
Kit Langton 23e4e35ae8 fix(core): tolerate older migration schemas 2026-08-11 12:16:59 -04:00
16 changed files with 275 additions and 325 deletions
@@ -348,10 +348,7 @@ const lowerMessages = Effect.fn("BedrockConverse.lowerMessages")(function* (
continue
}
}
const previous = messages.at(-1)
if (previous?.role === "user")
messages[messages.length - 1] = { role: "user", content: [...previous.content, ...content] }
else messages.push({ role: "user", content })
messages.push({ role: "user", content })
continue
}
@@ -395,10 +392,7 @@ const lowerMessages = Effect.fn("BedrockConverse.lowerMessages")(function* (
const cachePoint = BedrockCache.block(breakpoints, part.cache)
if (cachePoint) content.push(cachePoint)
}
const previous = messages.at(-1)
if (previous?.role === "user")
messages[messages.length - 1] = { role: "user", content: [...previous.content, ...content] }
else messages.push({ role: "user", content })
messages.push({ role: "user", content })
}
return messages
@@ -1,36 +0,0 @@
{
"version": 1,
"metadata": {
"tags": [
"prefix:bedrock-converse",
"provider:amazon-bedrock",
"protocol:bedrock-converse",
"tool",
"tool-loop",
"parallel"
],
"name": "bedrock-converse/continues-after-parallel-tool-results",
"recordedAt": "2026-08-11T16:41:46.482Z"
},
"interactions": [
{
"transport": "http",
"request": {
"method": "POST",
"url": "https://bedrock-runtime.us-east-1.amazonaws.com/model/us.amazon.nova-micro-v1%3A0/converse-stream",
"headers": {
"content-type": "application/json"
},
"body": "{\"modelId\":\"us.amazon.nova-micro-v1:0\",\"messages\":[{\"role\":\"user\",\"content\":[{\"text\":\"Compare the weather in Paris and London.\"}]},{\"role\":\"assistant\",\"content\":[{\"toolUse\":{\"toolUseId\":\"weather_paris\",\"name\":\"get_weather\",\"input\":{\"city\":\"Paris\"}}},{\"toolUse\":{\"toolUseId\":\"weather_london\",\"name\":\"get_weather\",\"input\":{\"city\":\"London\"}}}]},{\"role\":\"user\",\"content\":[{\"toolResult\":{\"toolUseId\":\"weather_paris\",\"content\":[{\"json\":{\"temperature\":22,\"condition\":\"sunny\"}}],\"status\":\"success\"}},{\"toolResult\":{\"toolUseId\":\"weather_london\",\"content\":[{\"json\":{\"temperature\":14,\"condition\":\"rainy\"}}],\"status\":\"success\"}}]}],\"system\":[{\"text\":\"After receiving both tool results, reply exactly: Paris is sunny; London is rainy.\"}],\"inferenceConfig\":{\"maxTokens\":40,\"temperature\":0},\"toolConfig\":{\"tools\":[{\"toolSpec\":{\"name\":\"get_weather\",\"description\":\"Get current weather for a city.\",\"inputSchema\":{\"json\":{\"type\":\"object\",\"properties\":{\"city\":{\"type\":\"string\"}},\"required\":[\"city\"],\"additionalProperties\":false}}}}]}}"
},
"response": {
"status": 200,
"headers": {
"content-type": "application/vnd.amazon.eventstream"
},
"body": "AAAAqAAAAFKgEDvmCzpldmVudC10eXBlBwAMbWVzc2FnZVN0YXJ0DTpjb250ZW50LXR5cGUHABBhcHBsaWNhdGlvbi9qc29uDTptZXNzYWdlLXR5cGUHAAVldmVudHsicCI6ImFiY2RlZmdoaWprbG1ub3BxcnN0dXZ3eHl6QUJDREVGR0hJSktMTU5PUFEiLCJyb2xlIjoiYXNzaXN0YW50In189ig4AAAAzgAAAFfGCE2ECzpldmVudC10eXBlBwARY29udGVudEJsb2NrRGVsdGENOmNvbnRlbnQtdHlwZQcAEGFwcGxpY2F0aW9uL2pzb24NOm1lc3NhZ2UtdHlwZQcABWV2ZW50eyJjb250ZW50QmxvY2tJbmRleCI6MCwiZGVsdGEiOnsidGV4dCI6IlBhcmlzIn0sInAiOiJhYmNkZWZnaGlqa2xtbm9wcXJzdHV2d3h5ekFCQ0RFRkdISUpLTE1OT1BRUlNUVVYifWttETIAAADDAAAAVz6YiTULOmV2ZW50LXR5cGUHABFjb250ZW50QmxvY2tEZWx0YQ06Y29udGVudC10eXBlBwAQYXBwbGljYXRpb24vanNvbg06bWVzc2FnZS10eXBlBwAFZXZlbnR7ImNvbnRlbnRCbG9ja0luZGV4IjowLCJkZWx0YSI6eyJ0ZXh0IjoiIGlzIn0sInAiOiJhYmNkZWZnaGlqa2xtbm9wcXJzdHV2d3h5ekFCQ0RFRkdISUpLTE0ifRHJ8Q0AAACqAAAAV6q6nAkLOmV2ZW50LXR5cGUHABFjb250ZW50QmxvY2tEZWx0YQ06Y29udGVudC10eXBlBwAQYXBwbGljYXRpb24vanNvbg06bWVzc2FnZS10eXBlBwAFZXZlbnR7ImNvbnRlbnRCbG9ja0luZGV4IjowLCJkZWx0YSI6eyJ0ZXh0IjoiIHN1bm55In0sInAiOiJhYmNkZWZnaGlqayJ98ZCy6gAAALAAAABXgOoTKgs6ZXZlbnQtdHlwZQcAEWNvbnRlbnRCbG9ja0RlbHRhDTpjb250ZW50LXR5cGUHABBhcHBsaWNhdGlvbi9qc29uDTptZXNzYWdlLXR5cGUHAAVldmVudHsiY29udGVudEJsb2NrSW5kZXgiOjAsImRlbHRhIjp7InRleHQiOiI7In0sInAiOiJhYmNkZWZnaGlqa2xtbm9wcXJzdHV2In020bBKAAAAygAAAFcziOtECzpldmVudC10eXBlBwARY29udGVudEJsb2NrRGVsdGENOmNvbnRlbnQtdHlwZQcAEGFwcGxpY2F0aW9uL2pzb24NOm1lc3NhZ2UtdHlwZQcABWV2ZW50eyJjb250ZW50QmxvY2tJbmRleCI6MCwiZGVsdGEiOnsidGV4dCI6IiBMb25kb24ifSwicCI6ImFiY2RlZmdoaWprbG1ub3BxcnN0dXZ3eHl6QUJDREVGR0hJSktMTU5PUCJ9ew04hAAAAK4AAABXXzo6yQs6ZXZlbnQtdHlwZQcAEWNvbnRlbnRCbG9ja0RlbHRhDTpjb250ZW50LXR5cGUHABBhcHBsaWNhdGlvbi9qc29uDTptZXNzYWdlLXR5cGUHAAVldmVudHsiY29udGVudEJsb2NrSW5kZXgiOjAsImRlbHRhIjp7InRleHQiOiIgaXMifSwicCI6ImFiY2RlZmdoaWprbG1ub3BxciJ9yK3bdAAAALcAAABXMsrPOgs6ZXZlbnQtdHlwZQcAEWNvbnRlbnRCbG9ja0RlbHRhDTpjb250ZW50LXR5cGUHABBhcHBsaWNhdGlvbi9qc29uDTptZXNzYWdlLXR5cGUHAAVldmVudHsiY29udGVudEJsb2NrSW5kZXgiOjAsImRlbHRhIjp7InRleHQiOiIgcmFpbnkifSwicCI6ImFiY2RlZmdoaWprbG1ub3BxcnN0dXZ3eCJ9JoCYVwAAALoAAABXyloLiws6ZXZlbnQtdHlwZQcAEWNvbnRlbnRCbG9ja0RlbHRhDTpjb250ZW50LXR5cGUHABBhcHBsaWNhdGlvbi9qc29uDTptZXNzYWdlLXR5cGUHAAVldmVudHsiY29udGVudEJsb2NrSW5kZXgiOjAsImRlbHRhIjp7InRleHQiOiIuIn0sInAiOiJhYmNkZWZnaGlqa2xtbm9wcXJzdHV2d3h5ekFCQ0RFRiJ9WJwR8wAAAMEAAABXRFjaVQs6ZXZlbnQtdHlwZQcAEWNvbnRlbnRCbG9ja0RlbHRhDTpjb250ZW50LXR5cGUHABBhcHBsaWNhdGlvbi9qc29uDTptZXNzYWdlLXR5cGUHAAVldmVudHsiY29udGVudEJsb2NrSW5kZXgiOjAsImRlbHRhIjp7InRleHQiOiIifSwicCI6ImFiY2RlZmdoaWprbG1ub3BxcnN0dXZ3eHl6QUJDREVGR0hJSktMTU4ifYzp4V0AAAChAAAAVqptnY4LOmV2ZW50LXR5cGUHABBjb250ZW50QmxvY2tTdG9wDTpjb250ZW50LXR5cGUHABBhcHBsaWNhdGlvbi9qc29uDTptZXNzYWdlLXR5cGUHAAVldmVudHsiY29udGVudEJsb2NrSW5kZXgiOjAsInAiOiJhYmNkZWZnaGlqa2xtbm9wcXJzdHV2d3h5ekFCQyJ9AHyeLwAAAJUAAABRYKgWaws6ZXZlbnQtdHlwZQcAC21lc3NhZ2VTdG9wDTpjb250ZW50LXR5cGUHABBhcHBsaWNhdGlvbi9qc29uDTptZXNzYWdlLXR5cGUHAAVldmVudHsicCI6ImFiY2RlZmdoaWprbG1ub3BxcnN0Iiwic3RvcFJlYXNvbiI6ImVuZF90dXJuIn2HCXz0AAABBgAAAE6wWpX7CzpldmVudC10eXBlBwAIbWV0YWRhdGENOmNvbnRlbnQtdHlwZQcAEGFwcGxpY2F0aW9uL2pzb24NOm1lc3NhZ2UtdHlwZQcABWV2ZW50eyJtZXRyaWNzIjp7ImxhdGVuY3lNcyI6MTA5MH0sInAiOiJhYmNkZWZnaGlqa2xtbm9wcXJzdHV2d3h5ekFCQ0RFRkdISUpLTE1OT1BRUlNUVSIsInVzYWdlIjp7ImlucHV0VG9rZW5zIjo1MjEsIm91dHB1dFRva2VucyI6OSwic2VydmVyVG9vbFVzYWdlIjp7fSwidG90YWxUb2tlbnMiOjUzMH19Uwfxiw==",
"bodyEncoding": "base64"
}
}
]
}
@@ -255,57 +255,6 @@ describe("Bedrock Converse route", () => {
}),
)
it.effect("merges parallel tool results into one user message", () =>
Effect.gen(function* () {
const prepared = yield* compileRequest(
LLM.request({
id: "req_parallel_history",
model,
messages: [
Message.user("Compare the weather."),
Message.assistant([
ToolCallPart.make({ id: "tool_paris", name: "lookup", input: { city: "Paris" } }),
ToolCallPart.make({ id: "tool_london", name: "lookup", input: { city: "London" } }),
]),
Message.tool({ id: "tool_paris", name: "lookup", result: { forecast: "sunny" } }),
Message.tool({ id: "tool_london", name: "lookup", result: { forecast: "rainy" } }),
],
cache: "none",
}),
)
expect(prepared.body.messages).toEqual([
{ role: "user", content: [{ text: "Compare the weather." }] },
{
role: "assistant",
content: [
{ toolUse: { toolUseId: "tool_paris", name: "lookup", input: { city: "Paris" } } },
{ toolUse: { toolUseId: "tool_london", name: "lookup", input: { city: "London" } } },
],
},
{
role: "user",
content: [
{
toolResult: {
toolUseId: "tool_paris",
content: [{ json: { forecast: "sunny" } }],
status: "success",
},
},
{
toolResult: {
toolUseId: "tool_london",
content: [{ json: { forecast: "rainy" } }],
status: "success",
},
},
],
},
])
}),
)
it.effect("lowers image content in tool-result messages", () =>
Effect.gen(function* () {
const prepared = yield* compileRequest(
@@ -1216,39 +1165,4 @@ describe("Bedrock Converse recorded", () => {
)
}),
)
recorded.effect.with("continues after parallel tool results", { tags: ["tool", "tool-loop", "parallel"] }, () =>
Effect.gen(function* () {
const response = yield* LLMClient.generate(
LLM.request({
id: "recorded_bedrock_parallel_tool_results",
model: recordedModel(),
system: "After receiving both tool results, reply exactly: Paris is sunny; London is rainy.",
messages: [
Message.user("Compare the weather in Paris and London."),
Message.assistant([
ToolCallPart.make({ id: "weather_paris", name: weatherToolName, input: { city: "Paris" } }),
ToolCallPart.make({ id: "weather_london", name: weatherToolName, input: { city: "London" } }),
]),
Message.tool({
id: "weather_paris",
name: weatherToolName,
result: { temperature: 22, condition: "sunny" },
}),
Message.tool({
id: "weather_london",
name: weatherToolName,
result: { temperature: 14, condition: "rainy" },
}),
],
tools: [weatherTool],
cache: "none",
generation: { maxTokens: 40, temperature: 0 },
}),
)
expect(response.text.trim()).toBe("Paris is sunny; London is rainy.")
expect(response.finishReason?.normalized).toBe("stop")
}),
)
})
+9 -14
View File
@@ -39,8 +39,6 @@ export const run = Effect.fnUntraced(function* (options: Options) {
})
const processEffect = Effect.fnUntraced(function* (options: Options) {
const serviceErrorFormat = process.env.OPENCODE_SERVICE_ERROR_FORMAT
delete process.env.OPENCODE_SERVICE_ERROR_FORMAT
const global = yield* Global.Service
if (options.mode === "service") yield* Effect.sync(() => process.chdir(global.home))
return yield* Effect.scoped(
@@ -129,7 +127,15 @@ const processEffect = Effect.fnUntraced(function* (options: Options) {
if (serviceOptions === undefined || port === undefined || !addressInUse(error)) return Effect.fail(error)
return recognizeIncumbent(serviceOptions, hostname, port).pipe(
Effect.flatMap((found) =>
found ? Effect.void : managedPortInUse(hostname, port, error, serviceErrorFormat),
found
? Effect.void
: Effect.fail(
new Error(
`Managed service port ${port} on ${hostname} is already in use by another process. ` +
"Configure another port with `opencode service set port <port>` and start the service again.",
{ cause: error },
),
),
),
)
}),
@@ -208,17 +214,6 @@ function serviceURL(hostname: string, port: number) {
return `http://${hostname.includes(":") ? `[${hostname}]` : hostname}:${port}`
}
function managedPortInUse(hostname: string, port: number, cause: unknown, format?: string) {
const message =
`Managed service port ${port} on ${hostname} is already in use by another process. ` +
"Configure another port with `opencode service set port <port>` and start the service again."
const failure = new Error(message, { cause })
if (format !== "plain") return Effect.fail(failure)
return Effect.sync(() => process.stderr.write(`OPENCODE_SERVICE_ERROR:${message}\n`)).pipe(
Effect.andThen(Effect.fail(failure)),
)
}
function truthy(value?: string) {
return value === "1" || value?.toLowerCase() === "true"
}
+33 -8
View File
@@ -1,9 +1,9 @@
import { ServiceStatus } from "@opencode-ai/protocol/groups/health"
import { Effect, FileSystem, Option, Schedule, Schema } from "effect"
import { spawn, type ChildProcess } from "node:child_process"
import { homedir } from "node:os"
import { join } from "node:path"
import type { DiscoverOptions, Endpoint, EnsureOptions, StopOptions } from "../service.js"
import { ServiceProcess } from "../service-process.js"
export * from "../service.js"
/** Contents of the local service registration file. */
@@ -17,6 +17,11 @@ export type Info = import("../service.js").Info
// is all a client needs to connect. The daemon's own configuration (port,
// persisted password) is CLI-owned and never read here.
type Contender = {
readonly child: ChildProcess
readonly error: () => Error | undefined
}
// Read-only lookup: registration file plus health check and version gate.
// Never spawns; escalation to ensure() is the caller's policy.
/** Discover a healthy, compatible local service without starting one. */
@@ -47,12 +52,11 @@ const discoverLocal = Effect.fnUntraced(function* (options: DiscoverOptions) {
// becomes discoverable. A contender is never killed merely for slow startup.
/** Ensure a healthy, compatible local service is running. */
export const ensure = Effect.fn("service.ensure")(function* (options: EnsureOptions = {}) {
const contenders = new Set<ServiceProcess.Contender>()
const contenders = new Set<Contender>()
let timeouts: { readonly info: Info; readonly count: number } | undefined
let announced = false
let lastSpawn = 0
let spawnDelay = 5_000
let lastFailure: Error | undefined
const announce = (reason: "missing" | "version-mismatch", previousVersion?: string) =>
Effect.sync(() => {
if (announced) return
@@ -63,7 +67,15 @@ export const ensure = Effect.fn("service.ensure")(function* (options: EnsureOpti
const [command, ...args] = options.command ?? ["opencode", "serve", "--service"]
if (command === undefined) return yield* Effect.fail(new Error("Missing service command"))
return yield* Effect.try({
try: () => ServiceProcess.start(command, args),
try: () => {
const child = spawn(command, args, { detached: true, stdio: "ignore" })
let error: Error | undefined
child.once("error", (cause) => {
error = new Error("Failed to start server", { cause })
})
child.unref()
return { child, error: () => error }
},
catch: (cause) => new Error("Failed to start server", { cause }),
})
})
@@ -96,9 +108,8 @@ export const ensure = Effect.fn("service.ensure")(function* (options: EnsureOpti
return Option.none<LocalService>()
} else if (lastSpawn === 0 && info !== undefined) lastSpawn = Date.now()
const finished = [...contenders].filter(ServiceProcess.finished)
const failure = finished.map(ServiceProcess.failure).find((error): error is Error => error !== undefined)
if (failure !== undefined) lastFailure = failure
const finished = [...contenders].filter(contenderFinished)
const failure = finished.map(contenderFailure).find((error): error is Error => error !== undefined)
if (finished.some((item) => item.child.exitCode === 0)) {
spawnDelay = Math.min(spawnDelay * 2, 30_000)
}
@@ -118,10 +129,24 @@ export const ensure = Effect.fn("service.ensure")(function* (options: EnsureOpti
}),
)
if (Option.isNone(found))
return yield* Effect.fail(lastFailure ?? new Error("Timed out waiting for the background service to start"))
return yield* Effect.fail(new Error("Timed out waiting for the background service to start"))
return found.value.endpoint
})
function contenderFailure(contender: Contender) {
const error = contender.error()
if (error !== undefined) return error
if (contender.child.exitCode !== null && contender.child.exitCode !== 0)
return new Error(`Server process exited with code ${contender.child.exitCode}`)
if (contender.child.signalCode !== null)
return new Error(`Server process terminated by ${contender.child.signalCode}`)
return undefined
}
function contenderFinished(contender: Contender) {
return contender.error() !== undefined || contender.child.exitCode !== null || contender.child.signalCode !== null
}
/** Stop the registered local service. */
export const stop = Effect.fn("service.stop")(function* (options: StopOptions = {}) {
const existing = yield* find(options)
+35 -8
View File
@@ -1,8 +1,8 @@
import { readFile } from "node:fs/promises"
import { spawn, type ChildProcess } from "node:child_process"
import { homedir } from "node:os"
import { join } from "node:path"
import type { DiscoverOptions, Endpoint, Info, EnsureOptions, StopOptions } from "../service.js"
import { ServiceProcess } from "../service-process.js"
import type { ServiceHealth, ServiceStopResponse } from "./generated/types.js"
export * from "../service.js"
@@ -13,6 +13,11 @@ export * from "../service.js"
// intentionally implemented with Node APIs so Promise clients do not need
// Effect or @effect/platform-node at runtime.
type Contender = {
readonly child: ChildProcess
readonly error: () => Error | undefined
}
/** Discover a healthy, compatible local service without starting one. */
export async function discover(options: DiscoverOptions = {}) {
return (await discoverLocal(options))?.endpoint
@@ -28,12 +33,11 @@ async function discoverLocal(options: DiscoverOptions) {
/** Ensure a healthy, compatible local service is running. */
export async function ensure(options: EnsureOptions = {}): Promise<Endpoint> {
const deadline = Date.now() + 120_000
const contenders = new Set<ServiceProcess.Contender>()
const contenders = new Set<Contender>()
let timeouts: { readonly info: Info; readonly count: number } | undefined
let announced = false
let lastSpawn = 0
let spawnDelay = 5_000
let lastFailure: Error | undefined
const announce = (reason: "missing" | "version-mismatch", previousVersion?: string) => {
if (announced) return
@@ -43,11 +47,21 @@ export async function ensure(options: EnsureOptions = {}): Promise<Endpoint> {
const spawnContender = () => {
const [command, ...args] = options.command ?? ["opencode", "serve", "--service"]
if (command === undefined) throw new Error("Missing service command")
return ServiceProcess.start(command, args)
try {
const child = spawn(command, args, { detached: true, stdio: "ignore" })
let error: Error | undefined
child.once("error", (cause) => {
error = new Error("Failed to start server", { cause })
})
child.unref()
return { child, error: () => error }
} catch (cause) {
throw new Error("Failed to start server", { cause })
}
}
while (true) {
if (Date.now() >= deadline) throw lastFailure ?? new Error("Timed out waiting for the background service to start")
if (Date.now() >= deadline) throw new Error("Timed out waiting for the background service to start")
const registration = await registered(options.file, true)
if (registration.timedOut && registration.info !== undefined) {
timeouts = {
@@ -75,9 +89,8 @@ export async function ensure(options: EnsureOptions = {}): Promise<Endpoint> {
}
} else {
if (lastSpawn === 0 && registration.info !== undefined) lastSpawn = Date.now()
const finished = [...contenders].filter(ServiceProcess.finished)
const failure = finished.map(ServiceProcess.failure).find((error) => error !== undefined)
if (failure !== undefined) lastFailure = failure
const finished = [...contenders].filter(contenderFinished)
const failure = finished.map(contenderFailure).find((error) => error !== undefined)
if (finished.some((item) => item.child.exitCode === 0)) {
spawnDelay = Math.min(spawnDelay * 2, 30_000)
}
@@ -94,6 +107,20 @@ export async function ensure(options: EnsureOptions = {}): Promise<Endpoint> {
}
}
function contenderFailure(contender: Contender) {
const error = contender.error()
if (error !== undefined) return error
if (contender.child.exitCode !== null && contender.child.exitCode !== 0)
return new Error(`Server process exited with code ${contender.child.exitCode}`)
if (contender.child.signalCode !== null)
return new Error(`Server process terminated by ${contender.child.signalCode}`)
return undefined
}
function contenderFinished(contender: Contender) {
return contender.error() !== undefined || contender.child.exitCode !== null || contender.child.signalCode !== null
}
/** Stop the registered local service. */
export async function stop(options: StopOptions = {}) {
const existing = await find(options)
-57
View File
@@ -1,57 +0,0 @@
export * as ServiceProcess from "./service-process"
import { spawn, type ChildProcess } from "node:child_process"
const errorPrefix = "OPENCODE_SERVICE_ERROR:"
export type Contender = {
readonly child: ChildProcess
readonly error: () => Error | undefined
readonly startupError: () => string
}
export function start(command: string, args: ReadonlyArray<string>) {
try {
const child = spawn(command, args, {
detached: true,
stdio: ["ignore", "ignore", "pipe"],
env: { ...process.env, OPENCODE_SERVICE_ERROR_FORMAT: "plain" },
})
let error: Error | undefined
let pending = ""
let startupError = ""
child.once("error", (cause) => {
error = new Error("Failed to start server", { cause })
})
child.stderr?.on("data", (chunk) => {
const lines = (pending + chunk.toString()).split(/\r?\n/)
pending = lines.pop()?.slice(-64 * 1024) ?? ""
const message = lines.findLast((line) => line.startsWith(errorPrefix))
if (message !== undefined) startupError = message.slice(errorPrefix.length)
})
unref(child.stderr)
child.unref()
return { child, error: () => error, startupError: () => startupError } satisfies Contender
} catch (cause) {
throw new Error("Failed to start server", { cause })
}
}
export function failure(contender: Contender) {
const error = contender.error()
if (error !== undefined) return error
if (contender.child.exitCode !== null && contender.child.exitCode !== 0)
return new Error(contender.startupError() || `Server process exited with code ${contender.child.exitCode}`)
if (contender.child.signalCode !== null)
return new Error(`Server process terminated by ${contender.child.signalCode}`)
return undefined
}
export function finished(contender: Contender) {
return contender.error() !== undefined || contender.child.exitCode !== null || contender.child.signalCode !== null
}
function unref(stream: ChildProcess["stderr"]) {
if (!stream || !("unref" in stream) || typeof stream.unref !== "function") return
stream.unref()
}
-5
View File
@@ -3,11 +3,6 @@ import { appendFile, rename, writeFile } from "node:fs/promises"
const [registration, mode, delay] = process.argv.slice(2)
if (registration === undefined || mode === undefined) throw new Error("Missing service fixture arguments")
if (mode === "failed") process.exit(1)
if (mode === "failed-message") {
console.error("sensitive startup detail")
console.error("OPENCODE_SERVICE_ERROR:Managed service port is already in use")
process.exit(1)
}
if (mode === "record-start") {
await writeFile(registration + ".started", "")
process.exit(1)
@@ -70,19 +70,6 @@ test("reports a failed registered service", async () => {
)
})
test("reports the native contender's startup error", async () => {
const directory = await temp()
const registration = join(directory, "service.json")
await expect(
Service.ensure({
file: registration,
version: "test",
command: [process.execPath, fixture, registration, "failed-message"],
}),
).rejects.toThrow(/^Managed service port is already in use$/)
}, 10_000)
test("evicts an unresponsive registered service before starting its replacement", async () => {
const directory = await temp()
const registration = join(directory, "service.json")
-14
View File
@@ -197,20 +197,6 @@ test("reports a contender that fails to start", async () => {
).rejects.toThrow("Server process exited with code 1")
}, 10_000)
test("reports the contender's startup error", async () => {
const directory = await temp()
const registration = join(directory, "service.json")
await expect(
run(
Service.ensure({
file: registration,
version: "test",
command: [process.execPath, fixture, registration, "failed-message"],
}),
),
).rejects.toThrow(/^Managed service port is already in use$/)
}, 10_000)
test("reports a contender terminated by a signal", async () => {
const directory = await temp()
const registration = join(directory, "service.json")
+81 -5
View File
@@ -159,6 +159,58 @@ type NextMessage = {
readonly data: string
}
type SourceColumns<Row> = Readonly<Record<keyof Row, true | null>>
const nextProjectColumns = {
id: true,
worktree: true,
vcs: null,
name: null,
icon_url: null,
icon_url_override: null,
icon_color: null,
time_created: true,
time_updated: true,
time_initialized: null,
sandboxes: true,
commands: null,
} satisfies SourceColumns<NextProject>
const nextSessionColumns = {
id: true,
project_id: true,
workspace_id: null,
parent_id: null,
fork_session_id: null,
fork_boundary: null,
slug: true,
directory: true,
path: null,
title: null,
version: true,
share_url: null,
summary_additions: null,
summary_deletions: null,
summary_files: null,
summary_diffs: null,
metadata: null,
cost: true,
tokens_input: true,
tokens_output: true,
tokens_reasoning: true,
tokens_cache_read: true,
tokens_cache_write: true,
revert: null,
permission: null,
agent: null,
model: null,
time_created: true,
time_updated: true,
time_compacting: null,
time_archived: null,
time_suspended: null,
} satisfies SourceColumns<NextSession>
const lock = Semaphore.makeUnsafe(1)
const MIGRATION_STATE_KEY = "migration.v1-v2"
const EVENT_DELETE_BATCH_SIZE = 1_000
@@ -678,12 +730,9 @@ function importNextDatabase(
}),
)
const projects = new Map(
source
.query<NextProject, []>("SELECT * FROM project")
.all()
.map((project) => [project.id, project]),
selectSourceRows<NextProject>(source, "project", nextProjectColumns).map((project) => [project.id, project]),
)
const sessions = source.query<NextSession, []>("SELECT * FROM session ORDER BY id DESC").all()
const sessions = selectSourceRows<NextSession>(source, "session", nextSessionColumns, "id")
for (const [index, session] of sessions.entries()) {
const project = projects.get(session.project_id)
const projectID = project ? session.project_id : Project.ID.global
@@ -776,6 +825,33 @@ function isNextDatabase(source: SQLiteDatabase) {
return tables.has("project") && tables.has("session") && tables.has("session_message")
}
function selectSourceRows<Row>(
source: SQLiteDatabase,
table: "project" | "session",
columns: SourceColumns<Row>,
orderBy?: keyof Row,
) {
const available = new Set(
source
.query<{ name: string }, []>(`PRAGMA table_info("${table}")`)
.all()
.map((column) => column.name),
)
const selection = Object.entries(columns)
.map(([name, required]) => {
if (available.has(name)) return `"${name}"`
if (!required) return `NULL AS "${name}"`
throw new Error(`Previous V2 database ${table} table is missing required column ${name}`)
})
.join(", ")
return source
.query<
Row,
[]
>(`SELECT ${selection} FROM "${table}"${orderBy === undefined ? "" : ` ORDER BY "${String(orderBy)}" DESC`}`)
.all()
}
function row(
source: SourceMessage,
message: {
+3 -1
View File
@@ -21,7 +21,9 @@ Usage notes:
- 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 Input = Schema.Struct({
questions: Schema.Array(Question.Prompt).check(Schema.isNonEmpty()).annotate({ description: "Questions to ask" }),
questions: Schema.Array(Question.Prompt)
.check(Schema.isNonEmpty())
.annotate({ description: "Questions to ask" }),
})
export const Output = Schema.Struct({
@@ -0,0 +1,74 @@
CREATE TABLE project (
id text PRIMARY KEY,
worktree text NOT NULL,
vcs text,
name text,
icon_url text,
time_created integer NOT NULL,
time_updated integer NOT NULL,
time_initialized integer,
sandboxes text NOT NULL
);
CREATE TABLE session (
id text PRIMARY KEY,
project_id text NOT NULL,
workspace_id text,
parent_id text,
fork_session_id text,
fork_message_id text,
fork_seq integer,
slug text NOT NULL,
directory text NOT NULL,
path text,
title text,
version text NOT NULL,
share_url text,
summary_additions integer,
summary_deletions integer,
summary_files integer,
summary_diffs text,
metadata text,
cost real DEFAULT 0 NOT NULL,
tokens_input integer DEFAULT 0 NOT NULL,
tokens_output integer DEFAULT 0 NOT NULL,
tokens_reasoning integer DEFAULT 0 NOT NULL,
tokens_cache_read integer DEFAULT 0 NOT NULL,
tokens_cache_write integer DEFAULT 0 NOT NULL,
revert text,
permission text,
agent text,
model text,
time_created integer NOT NULL,
time_updated integer NOT NULL,
time_compacting integer,
time_archived integer
);
CREATE TABLE session_message (
id text PRIMARY KEY,
session_id text NOT NULL,
type text NOT NULL,
seq integer NOT NULL,
time_created integer NOT NULL,
time_updated integer NOT NULL,
data text NOT NULL
);
INSERT INTO project (
id, worktree, vcs, name, icon_url, time_created, time_updated, time_initialized, sandboxes
) VALUES (
'old-project', '/tmp/old-next', 'git', 'Old project', 'https://example.test/icon.png', 1, 2, 3, '[]'
);
INSERT INTO session (
id, project_id, fork_session_id, fork_message_id, fork_seq, slug, directory, title, version,
time_created, time_updated
) VALUES (
'ses_old_next', 'old-project', 'ses_parent', 'msg_parent', 4, 'old-next', '/tmp/old-next',
'Old imported session', '2', 10, 20
);
INSERT INTO session_message VALUES (
'msg_old_next', 'ses_old_next', 'user', 0, 12, 13, '{"text":"from old next","time":{"created":12}}'
);
+38
View File
@@ -945,6 +945,44 @@ describe("V1Migration database workflow", () => {
)
})
test("imports previous V2 sessions from an older source schema", async () => {
await using tmp = await tmpdir()
const filename = path.join(tmp.path, "opencode-next.db")
const sqlite = await import("bun:sqlite")
const source = new sqlite.Database(filename)
source.exec(await Bun.file(path.join(import.meta.dir, "fixture/v1-migration-old-next.sql")).text())
source.close()
await database(
Effect.gen(function* () {
const { db } = yield* Database.Service
expect(yield* V1Migration.run({ nextDatabasePath: filename })).toEqual({ status: "completed" })
expect(
yield* db.get(
sql`SELECT fork_session_id, fork_boundary, time_suspended FROM session_v2 WHERE id = 'ses_old_next'`,
),
).toEqual({ fork_session_id: "ses_parent", fork_boundary: null, time_suspended: null })
expect(
yield* db.get(
sql`SELECT icon_url, icon_url_override, icon_color, commands FROM project WHERE id = 'old-project'`,
),
).toEqual({
icon_url: "https://example.test/icon.png",
icon_url_override: null,
icon_color: null,
commands: null,
})
expect(yield* db.all(sql`SELECT id, seq FROM session_message WHERE session_id = 'ses_old_next'`)).toEqual([
{ id: "msg_old_next", seq: 0 },
])
expect(yield* db.get(sql`SELECT seq FROM event_sequence WHERE aggregate_id = 'ses_old_next'`)).toEqual({
seq: 0,
})
}),
)
})
test("derives required status from the durable cursor", async () => {
await database(
Effect.gen(function* () {
-13
View File
@@ -12,19 +12,6 @@ export const ModelHandler = HttpApiBuilder.group(Api, "server.model", (handlers)
.handle(
"model.list",
Effect.fn(function* () {
const plugins = yield* PluginSupervisor.Service
yield* plugins.flush.pipe(
Effect.timeoutOrElse({
duration: "5 seconds",
orElse: () =>
Effect.fail(
new ServiceUnavailableError({
message: "Model catalog initialization timed out",
service: "model.catalog",
}),
),
}),
)
const catalog = yield* Catalog.Service
return yield* response(catalog.model.available())
}),
-57
View File
@@ -1,57 +0,0 @@
import fs from "node:fs/promises"
import path from "node:path"
import { expect } from "bun:test"
import { Effect } from "effect"
import { HttpServer } from "effect/unstable/http"
import { tmpdir } from "../../core/test/fixture/tmpdir"
import { it } from "../../core/test/lib/effect"
import { ServerProcess } from "../src/process"
it.live("waits for plugin initialization before listing models", () =>
Effect.acquireUseRelease(
Effect.promise(() => tmpdir("opencode-model-endpoint-")),
(tmp) =>
Effect.gen(function* () {
yield* Effect.promise(() =>
fs.writeFile(
path.join(tmp.path, "opencode.json"),
JSON.stringify({
providers: {
custom: {
package: "aisdk:@ai-sdk/openai-compatible",
settings: { apiKey: "secret" },
models: { chat: {} },
},
},
}),
),
)
const server = yield* ServerProcess.start<never, never>({
hostname: "127.0.0.1",
port: 0,
password: "secret",
app: { version: "test-version" },
database: { path: ":memory:" },
config: { directory: tmp.path },
fs: { filewatcher: false },
})
const url = new URL("/api/model", HttpServer.formatAddress(server.address))
url.searchParams.set("location[directory]", tmp.path)
const response = yield* Effect.promise(() =>
fetch(url, { headers: { authorization: `Basic ${btoa("opencode:secret")}` } }),
)
expect(response.status).toBe(200)
const body: unknown = yield* Effect.promise(() => response.json())
if (!isRecord(body) || !Array.isArray(body["data"])) throw new Error("Expected a model list response")
expect(
body["data"].some((model) => isRecord(model) && model["providerID"] === "custom" && model["id"] === "chat"),
).toBeTrue()
}),
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
),
)
function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === "object" && value !== null && !Array.isArray(value)
}