Compare commits

...

8 Commits

Author SHA1 Message Date
Shoubhit Dash 3e7efffcb6 fix(ai): harden websocket error contracts 2026-08-05 23:04:53 +05:30
Kit Langton 3e253c589e refactor(core): persist instruction updates as messages (#40679) 2026-08-05 12:55:29 -04:00
Kit Langton 0a0fc09533 refactor(core): remove unused file mutation methods (#40667) 2026-08-05 12:33:20 -04:00
Kit Langton 5aa0413fea refactor(core): remove dead fork pending copy (#40675) 2026-08-05 12:30:42 -04:00
Kit Langton 6f4c199629 fix(tui): wait for diff request in tests (#40670) 2026-08-05 16:02:22 +00:00
Kit Langton ed8e1f4654 fix(core): reconcile promoted prompt retries from messages (#40664) 2026-08-05 11:43:03 -04:00
Kit Langton 3b0195e045 fix(core): avoid eager directory snapshots (#40552) 2026-08-05 14:52:07 +00:00
Shoubhit Dash 143a776373 fix(acp): surface subagent activity (#40438) 2026-08-05 18:57:15 +05:30
32 changed files with 1048 additions and 698 deletions
+50 -5
View File
@@ -211,11 +211,43 @@ export type StreamItem = Schema.Schema.Type<typeof StreamItem>
// event-level `error` envelope, so accept all three shapes here.
// https://www.openresponses.org/specification
const OpenResponsesErrorPayload = Schema.Struct({
type: optionalNull(Schema.String),
code: optionalNull(Schema.String),
message: optionalNull(Schema.String),
param: optionalNull(Schema.String),
})
const WebSocketErrorHeader = Schema.Union([Schema.String, Schema.Number, Schema.Boolean])
export const WebSocketErrorEvent = Schema.StructWithRest(
Schema.Struct({
type: Schema.tag("error"),
status: Schema.optional(Schema.Number),
status_code: Schema.optional(Schema.Number),
code: optionalNull(Schema.String),
message: Schema.optional(Schema.String),
param: optionalNull(Schema.String),
error: optionalNull(OpenResponsesErrorPayload),
headers: Schema.optional(Schema.Record(Schema.String, WebSocketErrorHeader)),
}),
[Schema.Record(Schema.String, Schema.Unknown)],
)
const decodeWebSocketErrorEvent = Schema.decodeUnknownEffect(WebSocketErrorEvent)
const decodeKnownErrorEvent = (event: Event) =>
decodeWebSocketErrorEvent({
...event,
status: typeof event.status === "number" ? event.status : undefined,
status_code: typeof event.status_code === "number" ? event.status_code : undefined,
headers: ProviderShared.isRecord(event.headers)
? Object.fromEntries(
Object.entries(event.headers).filter(
(entry): entry is [string, string | number | boolean] =>
typeof entry[1] === "string" || typeof entry[1] === "number" || typeof entry[1] === "boolean",
),
)
: undefined,
})
export const Event = Schema.StructWithRest(
Schema.Struct({
type: Schema.String,
@@ -240,6 +272,9 @@ export const Event = Schema.StructWithRest(
message: Schema.optional(Schema.String),
param: optionalNull(Schema.String),
error: optionalNull(OpenResponsesErrorPayload),
status: Schema.optional(Schema.Unknown),
status_code: Schema.optional(Schema.Unknown),
headers: Schema.optional(Schema.Unknown),
}),
[Schema.Record(Schema.String, Schema.Unknown)],
)
@@ -632,9 +667,9 @@ export type StepResult = readonly [ParserState, ReadonlyArray<LLMEvent>]
const NO_EVENTS: StepResult["1"] = []
// `response.completed` / `response.incomplete` are clean finishes that emit a
// `finish` event; `response.failed` is a hard failure. All three end the stream,
// so keep this set aligned with `step` and the protocol's terminal predicate.
const TERMINAL_TYPES = new Set(["response.completed", "response.incomplete", "response.failed"])
// `finish` event; `response.failed` and `error` are hard failures. All four end
// the stream, so keep this set aligned with `step` and the protocol's terminal predicate.
const TERMINAL_TYPES = new Set(["error", "response.completed", "response.incomplete", "response.failed"])
export const terminal = (event: Event) => TERMINAL_TYPES.has(event.type)
const onOutputTextDelta = (state: ParserState, event: Event, id: string): StepResult => {
@@ -969,10 +1004,16 @@ const providerErrorMessage = (event: Event, fallback: string): string => {
const providerError = (state: ParserState, event: Event, fallback: string) => {
const code = event.code || event.error?.code || event.response?.error?.code || undefined
const message = providerErrorMessage(event, fallback)
const status =
typeof event.status === "number"
? event.status
: typeof event.status_code === "number"
? event.status_code
: undefined
return new AIError({
module: state.id,
method: "stream",
reason: classifyProviderFailure({ message, code }),
reason: classifyProviderFailure({ message, code, status }),
})
}
@@ -1015,7 +1056,11 @@ export const step = (state: ParserState, event: Event) => {
if (event.type === "response.completed" || event.type === "response.incomplete")
return Effect.succeed(onResponseFinish(state, event))
if (event.type === "response.failed") return providerError(state, event, `${state.name} response failed`)
if (event.type === "error") return providerError(state, event, `${state.name} stream error`)
if (event.type === "error")
return decodeKnownErrorEvent(event).pipe(
Effect.mapError(() => ProviderShared.eventError(state.id, `${state.name} returned a malformed error event`)),
Effect.flatMap(() => providerError(state, event, `${state.name} stream error`)),
)
return Effect.succeed<StepResult>([state, NO_EVENTS])
}
+1
View File
@@ -67,6 +67,7 @@ const SERVER_CODES = new Set([
"overloaded_error",
"server_error",
"server_is_overloaded",
"slow_down",
"serviceunavailableexception",
])
const INVALID_REQUEST_CODES = new Set(["invalid_prompt", "invalid_request_error", "validationexception"])
+84 -9
View File
@@ -29,14 +29,45 @@ export class Service extends Context.Service<Service, Interface>()("@opencode/AI
const transportError = (
method: string,
message: string,
input: { readonly url?: string; readonly kind?: string } = {},
input: {
readonly url?: string
readonly kind?: string
readonly phase?: TransportReason["phase"]
readonly delivery?: TransportReason["delivery"]
} = {},
) =>
new AIError({
module: "WebSocketExecutor",
method,
reason: new TransportReason({ message, url: input.url, kind: input.kind }),
reason: new TransportReason({
message,
url: input.url,
kind: input.kind,
phase: input.phase,
delivery: input.delivery,
}),
})
const annotateTransportError = (
error: AIError,
input: { readonly phase: TransportReason["phase"]; readonly delivery: TransportReason["delivery"] },
) =>
error.reason._tag === "Transport"
? new AIError({
module: error.module,
method: error.method,
reason: new TransportReason({
message: error.reason.message,
kind: error.reason.kind,
url: error.reason.url,
http: error.reason.http,
phase: input.phase,
delivery: input.delivery,
recovery: error.reason.recovery,
}),
})
: error
const eventMessage = (event: Event) => {
if ("message" in event && typeof event.message === "string") return event.message
return event.type
@@ -56,6 +87,8 @@ const waitOpen = (ws: globalThis.WebSocket, input: WebSocketRequest) => {
transportError("open", `WebSocket closed before opening (state ${ws.readyState})`, {
url: input.url,
kind: "open",
phase: "connect",
delivery: "not-sent",
}),
)
}
@@ -79,7 +112,12 @@ const waitOpen = (ws: globalThis.WebSocket, input: WebSocketRequest) => {
cleanup()
resume(
Effect.fail(
transportError("open", `Failed to open WebSocket: ${eventMessage(event)}`, { url: input.url, kind: "open" }),
transportError("open", `Failed to open WebSocket: ${eventMessage(event)}`, {
url: input.url,
kind: "open",
phase: "connect",
delivery: "not-sent",
}),
),
)
}
@@ -90,6 +128,8 @@ const waitOpen = (ws: globalThis.WebSocket, input: WebSocketRequest) => {
transportError("open", `WebSocket closed before opening with code ${event.code}`, {
url: input.url,
kind: "open",
phase: "connect",
delivery: "not-sent",
}),
),
)
@@ -119,6 +159,8 @@ const webSocketUrl = (value: string) =>
transportError("prepare", error instanceof Error ? error.message : "Invalid WebSocket URL", {
url: value,
kind: "websocket",
phase: "prepare",
delivery: "not-sent",
}),
})
@@ -130,6 +172,8 @@ export const open = (input: WebSocketRequest) =>
transportError("open", error instanceof Error ? error.message : "Failed to construct WebSocket", {
url: input.url,
kind: "open",
phase: "connect",
delivery: "not-sent",
}),
}).pipe(Effect.flatMap((ws) => fromWebSocket(ws, input)))
@@ -150,7 +194,11 @@ export const fromWebSocket = (
Queue.failCauseUnsafe(
messages,
Cause.fail(
transportError("message", "Unsupported WebSocket message payload", { url: input.url, kind: "message" }),
transportError("message", "Unsupported WebSocket message payload", {
url: input.url,
kind: "message",
phase: "receive",
}),
),
)
}
@@ -158,16 +206,23 @@ export const fromWebSocket = (
Queue.failCauseUnsafe(
messages,
Cause.fail(
transportError("message", `WebSocket error: ${eventMessage(event)}`, { url: input.url, kind: "message" }),
transportError("message", `WebSocket error: ${eventMessage(event)}`, {
url: input.url,
kind: "message",
phase: "receive",
}),
),
)
}
const onClose = (event: CloseEvent) => {
if (event.code === 1000 || event.code === 1005) return Queue.endUnsafe(messages)
Queue.failCauseUnsafe(
messages,
Cause.fail(
transportError("message", `WebSocket closed with code ${event.code}`, { url: input.url, kind: "close" }),
transportError("message", `WebSocket closed with code ${event.code}`, {
url: input.url,
kind: "close",
phase: "close",
}),
),
)
}
@@ -189,6 +244,8 @@ export const fromWebSocket = (
transportError("sendText", error instanceof Error ? error.message : "Failed to send WebSocket message", {
url: input.url,
kind: "write",
phase: "send",
delivery: "not-sent",
}),
}),
messages: Stream.fromQueue(messages),
@@ -244,6 +301,8 @@ export const json = <Body, Message>(input: JsonInput<Body, Message>): JsonTransp
transportError("json", "WebSocket JSON transport requires WebSocketExecutor.Service", {
url: prepared.url,
kind: "websocket",
phase: "prepare",
delivery: "not-sent",
}),
)
}
@@ -251,11 +310,27 @@ export const json = <Body, Message>(input: JsonInput<Body, Message>): JsonTransp
return Stream.unwrap(
Effect.gen(function* () {
const connection = yield* Effect.acquireRelease(
webSocket.open({ url: prepared.url, headers: prepared.headers }),
webSocket
.open({ url: prepared.url, headers: prepared.headers })
.pipe(
Effect.mapError((error) => annotateTransportError(error, { phase: "connect", delivery: "not-sent" })),
),
(connection) => connection.close,
)
yield* connection.sendText(prepared.message)
return connection.messages.pipe(Stream.map((message) => messageText(message, decoder)))
let observed = false
return connection.messages.pipe(
Stream.map((message) => {
observed = true
return messageText(message, decoder)
}),
Stream.mapError((error) =>
annotateTransportError(error, {
phase: error.reason._tag === "Transport" && error.reason.phase === "close" ? "close" : "receive",
delivery: observed ? "accepted" : "ambiguous",
}),
),
)
}),
)
},
+7
View File
@@ -98,6 +98,13 @@ export class TransportReason extends Schema.Class<TransportReason>("AI.Error.Tra
kind: Schema.optional(Schema.String),
url: Schema.optional(Schema.String),
http: Schema.optional(HttpContext),
phase: Schema.optional(
Schema.Literals(["prepare", "queue", "connect", "send", "receive", "decode", "complete", "fallback", "close"]),
),
delivery: Schema.optional(Schema.Literals(["not-sent", "rejected", "ambiguous", "accepted"])),
recovery: Schema.optional(
Schema.Literals(["retry-connect", "retry-full", "rotate-and-retry-full", "fallback-http", "fail"]),
),
}) {}
export class InvalidProviderOutputReason extends Schema.Class<InvalidProviderOutputReason>(
+2 -2
View File
@@ -69,10 +69,10 @@ describe("provider error classification", () => {
test("classifies V1 overloaded provider codes", () => {
expect(
['{"code":"resource_exhausted"}', '{"code":"service_unavailable"}'].map(
['{"code":"resource_exhausted"}', '{"code":"service_unavailable"}', '{"code":"slow_down"}'].map(
(message) => classifyProviderFailure({ message })._tag,
),
).toEqual(["ProviderInternal", "ProviderInternal"])
).toEqual(["ProviderInternal", "ProviderInternal", "ProviderInternal"])
})
test("classifies transient client statuses as provider internal", () => {
@@ -11,6 +11,7 @@ import {
ToolCallPart,
ToolDefinition,
ToolResultPart,
TransportReason,
Usage,
} from "../../src"
import { Auth, LLMClient, RequestExecutor, WebSocketExecutor } from "../../src/route"
@@ -288,6 +289,114 @@ describe("OpenAI Responses route", () => {
}),
)
it.effect("terminates WebSocket control events without waiting for the socket to close", () =>
Effect.gen(function* () {
const events = [
{ type: "error", error: { code: "slow_down", message: "Try later" } },
{
type: "error",
status_code: 429,
message: "Rate limited",
headers: { "retry-after": 1, "x-request-id": "request", cached: false, invalid: [] },
},
{
type: "response.failed",
response: { error: { code: "server_error", message: "Unavailable" } },
},
{ type: "error", status: "not-a-status", message: "Malformed status" },
]
const errors = yield* Effect.forEach(events, (event) =>
LLMClient.generate(
LLM.request({
model: OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responsesWebSocket(
"gpt-4.1-mini",
),
prompt: "Say hello.",
}),
).pipe(
Effect.provide(
LLMClient.layer.pipe(
Layer.provide(
Layer.mergeAll(
Layer.succeed(
RequestExecutor.Service,
RequestExecutor.Service.of({ execute: () => Effect.die("unexpected HTTP request") }),
),
Layer.succeed(
WebSocketExecutor.Service,
WebSocketExecutor.Service.of({
open: () =>
Effect.succeed({
sendText: () => Effect.void,
messages: Stream.make(ProviderShared.encodeJson(event)).pipe(Stream.concat(Stream.never)),
close: Effect.void,
}),
}),
),
),
),
),
),
Effect.flip,
),
)
expect(errors.map((error) => error.reason._tag)).toEqual([
"ProviderInternal",
"RateLimit",
"ProviderInternal",
"UnknownProvider",
])
}),
)
it.effect("marks post-send WebSocket failures with delivery state", () =>
Effect.gen(function* () {
const failure = new AIError({
module: "test",
method: "receive",
reason: new TransportReason({ message: "socket closed", phase: "close" }),
})
const streams = [
Stream.fail(failure),
Stream.make(ProviderShared.encodeJson({ type: "response.created" })).pipe(Stream.concat(Stream.fail(failure))),
]
const deps = Layer.mergeAll(
Layer.succeed(
RequestExecutor.Service,
RequestExecutor.Service.of({ execute: () => Effect.die("unexpected HTTP request") }),
),
Layer.succeed(
WebSocketExecutor.Service,
WebSocketExecutor.Service.of({
open: () =>
Effect.succeed({
sendText: () => Effect.void,
messages: streams.shift() ?? Stream.die("unexpected WebSocket open"),
close: Effect.void,
}),
}),
),
)
const model = OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responsesWebSocket(
"gpt-4.1-mini",
)
const errors = yield* Effect.forEach(["first", "second"], (prompt) =>
LLMClient.generate(LLM.request({ model, prompt })).pipe(
Effect.provide(LLMClient.layer.pipe(Layer.provide(deps))),
Effect.flip,
),
)
expect(errors.map((error) => error.reason)).toEqual([
expect.objectContaining({ _tag: "Transport", phase: "close", delivery: "ambiguous" }),
expect.objectContaining({ _tag: "Transport", phase: "close", delivery: "accepted" }),
])
}),
)
it.effect("fails immediately when WebSocket is already closed", () =>
Effect.gen(function* () {
const error = yield* WebSocketExecutor.fromWebSocket(
@@ -297,6 +406,7 @@ describe("OpenAI Responses route", () => {
).pipe(Effect.flip)
expect(error.message).toContain("closed before opening")
expect(error.reason).toMatchObject({ _tag: "Transport", phase: "connect", delivery: "not-sent" })
}),
)
+19
View File
@@ -11,6 +11,7 @@ import {
LanguageModel,
ModelID,
ProviderID,
TransportReason,
Usage,
} from "../src/schema"
import { ProviderShared } from "../src/protocols/shared"
@@ -108,3 +109,21 @@ test("AI errors expose the shared runtime tag", async () => {
await Effect.runPromise(Effect.fail(error).pipe(Effect.catchTag("AI.Error", () => Effect.succeed("caught")))),
).toBe("caught")
})
test("transport errors serialize execution facts", () => {
const reason = new TransportReason({
message: "connection closed",
phase: "receive",
delivery: "ambiguous",
recovery: "fail",
})
expect(Schema.encodeSync(TransportReason)(reason)).toEqual({
_tag: "Transport",
message: "connection closed",
phase: "receive",
delivery: "ambiguous",
recovery: "fail",
})
expect(Schema.decodeUnknownSync(TransportReason)(Schema.encodeSync(TransportReason)(reason))).toEqual(reason)
})
+193 -38
View File
@@ -1,4 +1,4 @@
import type { AgentSideConnection, PromptResponse } from "@agentclientprotocol/sdk"
import type { AgentSideConnection, PromptResponse, SessionUpdate } from "@agentclientprotocol/sdk"
import type {
EventSubscribeOutput,
OpenCodeClient,
@@ -37,6 +37,34 @@ export type TurnStart =
| { readonly type: "skill"; readonly id: string }
| { readonly type: "compaction"; readonly id: string }
export const ChildSessionUpdatesCapability = "opencode/child-session-updates"
export const ChildSessionUpdateMethod = "opencode/session/child_update"
type ChildSessionUpdateBase = {
readonly rootSessionId: string
readonly childSessionId: string
readonly parentSessionId: string
readonly depth: number
readonly title?: string
}
type ChildSessionEvent =
| { readonly type: "update"; readonly update: SessionUpdate }
| {
readonly type: "status"
readonly status: "created" | "running" | "completed" | "failed" | "interrupted"
readonly error?: { readonly type: string; readonly message: string }
}
export type ChildSessionUpdate = ChildSessionUpdateBase & ChildSessionEvent
type ChildSession = {
readonly id: string
readonly parentID: string
readonly depth: number
readonly title?: string
}
function emptyToolState(): ToolState {
return { name: "tool", input: {}, metadata: {}, content: [] }
}
@@ -50,8 +78,13 @@ export async function streamTurn(input: {
readonly writeTextFile: boolean
readonly submit: (signal: AbortSignal) => Promise<unknown>
readonly control: TurnControl
readonly childSessionUpdate?: (update: ChildSessionUpdate) => Promise<void>
readonly connectionSignal?: AbortSignal
readonly sessionSignal?: AbortSignal
}): Promise<PromptResponse> {
const streamController = new AbortController()
const connectionAbort = () => streamController.abort()
input.connectionSignal?.addEventListener("abort", connectionAbort, { once: true })
const stream = input.client.event.subscribe({ signal: streamController.signal })[Symbol.asyncIterator]()
const connected = await stream.next()
if (connected.done) throw new Error("event stream disconnected before prompt admission")
@@ -62,47 +95,101 @@ export async function streamTurn(input: {
let finish: SessionMessageAssistant["finish"]
let executionError: { readonly type: string; readonly message: string } | undefined
const tools = new Map<string, ToolState>()
const children = new Map<string, ChildSession>()
const openChildren = new Set<string>()
let handedOff = false
const update = (value: Parameters<Connection["sessionUpdate"]>[0]["update"]) =>
input.connection.sessionUpdate({ sessionId: input.sessionID, update: value })
const notifyChild = async (child: ChildSession, value: ChildSessionEvent) => {
if (!input.childSessionUpdate) return
await input
.childSessionUpdate({
rootSessionId: input.sessionID,
childSessionId: child.id,
parentSessionId: child.parentID,
depth: child.depth,
...(child.title ? { title: child.title } : {}),
...value,
})
.catch(() => {})
}
const consume = async () => {
const updateSession = async (value: SessionUpdate, child: ChildSession | undefined, mode: "turn" | "background") => {
const projected = child ? projectChildUpdate(value, child) : value
if (mode === "turn" && (!child || !input.childSessionUpdate)) {
await input.connection.sessionUpdate({ sessionId: input.sessionID, update: projected })
}
if (child) await notifyChild(child, { type: "update", update: projected })
}
const consume = async (mode: "turn" | "background") => {
while (!streamController.signal.aborted) {
const next = await stream.next()
if (next.done) throw new Error("event stream disconnected during prompt execution")
const event = next.value
if (event.type === "permission.asked" && event.data.sessionID === input.sessionID) {
const tool = event.data.source?.id ? tools.get(event.data.source.id) : undefined
if (event.type === "session.created") {
const parentID = event.data.info.parentID
if (!parentID) continue
const parent = parentID === input.sessionID ? undefined : children.get(parentID)
if ((mode === "turn" && parentID === input.sessionID) || parent) {
const child = {
id: event.data.sessionID,
parentID,
depth: parent ? parent.depth + 1 : 1,
title: event.data.info.title,
}
children.set(child.id, child)
openChildren.add(child.id)
await notifyChild(child, { type: "status", status: "created" })
}
continue
}
const eventSessionID = sessionIDFromEvent(event)
const child = eventSessionID ? children.get(eventSessionID) : undefined
const send = (update: SessionUpdate) => updateSession(update, child, mode)
if (mode === "background" && !child) continue
if (event.type === "permission.asked" && (event.data.sessionID === input.sessionID || child)) {
const tool = event.data.source?.id ? tools.get(toolKey(event.data.sessionID, event.data.source.id)) : undefined
await replyPermission({
client: input.client,
connection: input.connection,
event,
sessionID: input.sessionID,
sessionID: event.data.sessionID,
clientSessionID: input.sessionID,
cwd: input.cwd,
tool,
...(child ? { toolCallPrefix: child.id, titlePrefix: child.title } : {}),
})
continue
}
if (event.type === "form.created" && event.data.form.sessionID === input.sessionID) {
if (event.type === "form.created" && (event.data.form.sessionID === input.sessionID || child)) {
await input.client.form
.cancel({ sessionID: input.sessionID, formID: event.data.form.id })
.catch(() => input.client.session.interrupt({ sessionID: input.sessionID }).catch(() => {}))
.cancel({ sessionID: event.data.form.sessionID, formID: event.data.form.id })
.catch(() => input.client.session.interrupt({ sessionID: event.data.form.sessionID }).catch(() => {}))
continue
}
if (!("sessionID" in event.data) || event.data.sessionID !== input.sessionID) continue
if (!eventSessionID || (eventSessionID !== input.sessionID && !child)) continue
if (matchesStart(event, input.start)) {
started = true
continue
}
if (!started) continue
if (event.type === "session.execution.started") {
if (child) {
await notifyChild(child, { type: "status", status: "running" })
}
continue
}
if (event.type === "session.step.started") {
assistantMessageID = event.data.assistantMessageID
if (!child) assistantMessageID = event.data.assistantMessageID
continue
}
if (event.type === "session.text.delta") {
assistantMessageID = event.data.assistantMessageID
await update({
if (!child) assistantMessageID = event.data.assistantMessageID
await send({
sessionUpdate: "agent_message_chunk",
messageId: event.data.assistantMessageID,
content: { type: "text", text: event.data.delta },
@@ -110,8 +197,8 @@ export async function streamTurn(input: {
continue
}
if (event.type === "session.reasoning.delta") {
assistantMessageID = event.data.assistantMessageID
await update({
if (!child) assistantMessageID = event.data.assistantMessageID
await send({
sessionUpdate: "agent_thought_chunk",
messageId: event.data.assistantMessageID,
content: { type: "text", text: event.data.delta },
@@ -119,9 +206,14 @@ export async function streamTurn(input: {
continue
}
if (event.type === "session.tool.input.started") {
assistantMessageID = event.data.assistantMessageID
tools.set(event.data.id, { name: event.data.name, input: {}, metadata: {}, content: [] })
await update({
if (!child) assistantMessageID = event.data.assistantMessageID
tools.set(toolKey(event.data.sessionID, event.data.id), {
name: event.data.name,
input: {},
metadata: {},
content: [],
})
await send({
sessionUpdate: "tool_call",
...pendingToolCall({
toolCallId: event.data.id,
@@ -133,11 +225,12 @@ export async function streamTurn(input: {
continue
}
if (event.type === "session.tool.called") {
assistantMessageID = event.data.assistantMessageID
const current = tools.get(event.data.id) ?? emptyToolState()
if (!child) assistantMessageID = event.data.assistantMessageID
const key = toolKey(event.data.sessionID, event.data.id)
const current = tools.get(key) ?? emptyToolState()
current.input = event.data.input
tools.set(event.data.id, current)
await update({
tools.set(key, current)
await send({
sessionUpdate: "tool_call_update",
...runningToolUpdate({
toolCallId: event.data.id,
@@ -149,10 +242,10 @@ export async function streamTurn(input: {
continue
}
if (event.type === "session.tool.progress") {
const current = tools.get(event.data.id)
const current = tools.get(toolKey(event.data.sessionID, event.data.id))
if (!current) continue
current.metadata = event.data.metadata
await update({
await send({
sessionUpdate: "tool_call_update",
...runningToolUpdate({
toolCallId: event.data.id,
@@ -164,8 +257,9 @@ export async function streamTurn(input: {
continue
}
if (event.type === "session.tool.success") {
const current = tools.get(event.data.id) ?? emptyToolState()
tools.delete(event.data.id)
const key = toolKey(event.data.sessionID, event.data.id)
const current = tools.get(key) ?? emptyToolState()
tools.delete(key)
await syncEditedFiles({
connection: input.connection,
writeTextFile: input.writeTextFile,
@@ -175,7 +269,7 @@ export async function streamTurn(input: {
toolInput: current.input,
metadata: event.data.metadata ?? {},
}).catch(() => {})
await update({
await send({
sessionUpdate: "tool_call_update",
...completedToolUpdate({
toolCallId: event.data.id,
@@ -188,9 +282,10 @@ export async function streamTurn(input: {
continue
}
if (event.type === "session.tool.failed") {
const current = tools.get(event.data.id) ?? emptyToolState()
tools.delete(event.data.id)
await update({
const key = toolKey(event.data.sessionID, event.data.id)
const current = tools.get(key) ?? emptyToolState()
tools.delete(key)
await send({
sessionUpdate: "tool_call_update",
...errorToolUpdate({
toolCallId: event.data.id,
@@ -205,13 +300,33 @@ export async function streamTurn(input: {
continue
}
if (event.type === "session.step.ended") {
assistantMessageID = event.data.assistantMessageID
finish = event.data.finish
if (!child) {
assistantMessageID = event.data.assistantMessageID
finish = event.data.finish
}
continue
}
if (event.type === "session.execution.succeeded") {
if (!child) return "succeeded" as const
openChildren.delete(child.id)
await notifyChild(child, { type: "status", status: "completed" })
if (mode === "background" && openChildren.size === 0) return "succeeded" as const
continue
}
if (event.type === "session.execution.interrupted") {
if (!child) return "interrupted" as const
openChildren.delete(child.id)
await notifyChild(child, { type: "status", status: "interrupted" })
if (mode === "background" && openChildren.size === 0) return "interrupted" as const
continue
}
if (event.type === "session.execution.succeeded") return "succeeded" as const
if (event.type === "session.execution.interrupted") return "interrupted" as const
if (event.type === "session.execution.failed") {
if (child) {
openChildren.delete(child.id)
await notifyChild(child, { type: "status", status: "failed", error: event.data.error })
if (mode === "background" && openChildren.size === 0) return "failed" as const
continue
}
executionError = event.data.error
return "failed" as const
}
@@ -219,7 +334,13 @@ export async function streamTurn(input: {
return "interrupted" as const
}
const completed = consume()
const completed = consume("turn")
const closeStream = async () => {
streamController.abort()
input.connectionSignal?.removeEventListener("abort", connectionAbort)
input.sessionSignal?.removeEventListener("abort", connectionAbort)
await stream.return?.(undefined).catch(() => {})
}
try {
await input.submit(control.admission.signal).catch((error) => {
if (!control.cancelled) throw error
@@ -233,6 +354,13 @@ export async function streamTurn(input: {
}
}
const terminal = await completed
if (input.childSessionUpdate && openChildren.size > 0 && !input.sessionSignal?.aborted) {
handedOff = true
input.sessionSignal?.addEventListener("abort", connectionAbort, { once: true })
void consume("background")
.catch(() => {})
.finally(closeStream)
}
const assistant = assistantMessageID
? await input.client.session
.message({ sessionID: input.sessionID, messageID: assistantMessageID })
@@ -250,11 +378,38 @@ export async function streamTurn(input: {
await completed.catch(() => {})
throw error
} finally {
streamController.abort()
await stream.return?.(undefined).catch(() => {})
if (!handedOff) await closeStream()
}
}
function sessionIDFromEvent(event: EventSubscribeOutput) {
if ("sessionID" in event.data && typeof event.data.sessionID === "string") return event.data.sessionID
if (event.type === "form.created") return event.data.form.sessionID
return undefined
}
function toolKey(sessionID: string, id: string) {
return `${sessionID}:${id}`
}
function projectChildUpdate(update: SessionUpdate, child: ChildSession) {
const projected = { ...update }
projected._meta = {
...projected._meta,
"opencode/child-session": {
id: child.id,
parentID: child.parentID,
depth: child.depth,
...(child.title ? { title: child.title } : {}),
},
}
if (projected.sessionUpdate === "tool_call" || projected.sessionUpdate === "tool_call_update") {
projected.toolCallId = `${child.id}:${projected.toolCallId}`
if (projected.title && child.title) projected.title = `${child.title}: ${projected.title}`
}
return projected
}
export async function replayMessages(
connection: Pick<AgentSideConnection, "sessionUpdate">,
sessionID: string,
+17 -3
View File
@@ -20,20 +20,28 @@ export async function replyPermission(input: {
readonly connection: Connection
readonly event: PermissionEvent
readonly sessionID: string
readonly clientSessionID?: string
readonly cwd: string
readonly tool?: Tool
readonly toolCallPrefix?: string
readonly titlePrefix?: string
}) {
const toolName = input.tool?.name ?? input.event.data.action
const toolInput = { ...input.event.data.metadata, ...input.tool?.input }
const previews = await permissionPreviews(toolName, toolInput, input.cwd)
const toolCallID = input.event.data.source?.id ?? input.event.data.id
const title = permissionTitle(toolName, toolInput, previews)
const result = await input.connection
.requestPermission({
sessionId: input.sessionID,
sessionId: input.clientSessionID ?? input.sessionID,
toolCall: {
...pendingToolCall({
toolCallId: input.event.data.source?.id ?? input.event.data.id,
toolCallId: input.toolCallPrefix ? `${input.toolCallPrefix}:${toolCallID}` : toolCallID,
toolName,
state: { input: toolInput, title: permissionTitle(toolName, toolInput, previews) },
state: {
input: toolInput,
title: prefixedTitle(input.titlePrefix, title),
},
cwd: input.cwd,
}),
locations: permissionLocations(toolName, toolInput, input.event.data.resources, input.cwd, previews),
@@ -51,6 +59,12 @@ export async function replyPermission(input: {
})
}
function prefixedTitle(prefix: string | undefined, title: string | undefined) {
if (!prefix) return title
if (!title) return prefix
return `${prefix}: ${title}`
}
export async function syncEditedFiles(input: {
readonly connection: Partial<Pick<AgentSideConnection, "writeTextFile">>
readonly writeTextFile: boolean
+32 -7
View File
@@ -43,13 +43,21 @@ import { OPENCODE_VERSION } from "../version"
import { SessionMessage } from "@opencode-ai/schema/session-message"
import { buildConfigOptions, parseModelSelection, type ConfigOptionProvider } from "./config-option"
import { promptContentToParts } from "./content"
import { replayMessages, streamTurn, type TurnControl, type TurnStart } from "./event"
import {
ChildSessionUpdateMethod,
ChildSessionUpdatesCapability,
replayMessages,
streamTurn,
type ChildSessionUpdate,
type TurnControl,
type TurnStart,
} from "./event"
import { ACPError } from "./error"
export const AuthMethodID = "opencode-login"
type Connection = Pick<AgentSideConnection, "sessionUpdate" | "requestPermission"> &
Partial<Pick<AgentSideConnection, "writeTextFile">>
Partial<Pick<AgentSideConnection, "writeTextFile" | "extNotification" | "signal">>
type Catalog = {
readonly providers: ConfigOptionProvider[]
@@ -64,6 +72,7 @@ type Catalog = {
type Attached = {
readonly id: string
readonly cwd: string
readonly abort: AbortController
catalog: Catalog
model: ModelRef
modeID: string
@@ -100,7 +109,7 @@ export function make(input: { readonly client: OpenCodeClient; readonly connecti
const catalogs = new Map<string, Promise<Catalog>>()
const registeredMcp = new Map<string, Set<string>>()
const active = new Map<string, TurnControl>()
const capabilities = { writeTextFile: false }
const capabilities = { writeTextFile: false, childSessionUpdates: false }
const catalog = (cwd: string) => {
const cached = catalogs.get(cwd)
@@ -119,11 +128,19 @@ export function make(input: { readonly client: OpenCodeClient; readonly connecti
throw new ACPError.SessionNotFoundError({ sessionId: sessionID })
}
const detach = (sessionID: string) => {
sessions.get(sessionID)?.abort.abort()
sessions.delete(sessionID)
registeredMcp.delete(sessionID)
}
const attach = async (session: SessionInfo, cwd: string, mcpServers: readonly McpServer[]) => {
const currentCatalog = await catalog(cwd)
sessions.get(session.id)?.abort.abort()
const state: Attached = {
id: session.id,
cwd,
abort: new AbortController(),
catalog: currentCatalog,
model: session.model ?? currentCatalog.defaultModel,
modeID: session.agent ?? currentCatalog.defaultModeID,
@@ -161,6 +178,7 @@ export function make(input: { readonly client: OpenCodeClient; readonly connecti
return {
initialize: async (params) => {
capabilities.writeTextFile = params.clientCapabilities?.fs?.writeTextFile === true
capabilities.childSessionUpdates = params.clientCapabilities?._meta?.[ChildSessionUpdatesCapability] === true
const authMethod: AuthMethod = {
description: "Run `opencode auth login` in the terminal",
name: "Login with opencode",
@@ -178,6 +196,7 @@ export function make(input: { readonly client: OpenCodeClient; readonly connecti
mcpCapabilities: { http: true, sse: false },
promptCapabilities: { embeddedContext: true, image: true },
sessionCapabilities: { close: {}, delete: {}, fork: {}, list: {}, resume: {} },
_meta: { [ChildSessionUpdatesCapability]: true },
},
authMethods: [authMethod],
agentInfo: { name: "OpenCode", version: OPENCODE_VERSION },
@@ -224,8 +243,7 @@ export function make(input: { readonly client: OpenCodeClient; readonly connecti
await input.client.session.remove({ sessionID: params.sessionId }).catch((error) => {
if (!isSessionNotFoundError(error)) throw error
})
sessions.delete(params.sessionId)
registeredMcp.delete(params.sessionId)
detach(params.sessionId)
return {}
},
resumeSession: async (params) => {
@@ -234,8 +252,7 @@ export function make(input: { readonly client: OpenCodeClient; readonly connecti
return { configOptions: configOptions(state) }
},
closeSession: async (params) => {
sessions.delete(params.sessionId)
registeredMcp.delete(params.sessionId)
detach(params.sessionId)
const turn = active.get(params.sessionId)
if (turn) {
turn.cancelled = true
@@ -296,6 +313,11 @@ export function make(input: { readonly client: OpenCodeClient; readonly connecti
const messageID = SessionMessage.ID.create()
const prepared = preparePrompt(state.catalog, params.prompt, messageID)
const control: TurnControl = { cancelled: false, admission: new AbortController() }
const extNotification = input.connection.extNotification
const childSessionUpdate =
capabilities.childSessionUpdates && extNotification
? (update: ChildSessionUpdate) => extNotification(ChildSessionUpdateMethod, update).then(() => {})
: undefined
active.set(state.id, control)
const response = await streamTurn({
client: input.client,
@@ -305,7 +327,10 @@ export function make(input: { readonly client: OpenCodeClient; readonly connecti
start: prepared.start,
writeTextFile: capabilities.writeTextFile,
control,
connectionSignal: input.connection.signal,
sessionSignal: state.abort.signal,
submit: (signal) => submitPrompt(input.client, state, prepared, signal),
...(childSessionUpdate ? { childSessionUpdate } : {}),
}).finally(() => {
if (active.get(state.id) === control) active.delete(state.id)
})
+191 -1
View File
@@ -2,7 +2,7 @@ import { describe, expect, test } from "bun:test"
import type { AgentSideConnection } from "@agentclientprotocol/sdk"
import type { SessionMessageInfo } from "@opencode-ai/client/promise"
import { resolve } from "node:path"
import { replayMessages, streamTurn, type TurnControl } from "../../src/acp/event"
import { replayMessages, streamTurn, type ChildSessionUpdate, type TurnControl } from "../../src/acp/event"
import { createSseFixture, durableEvent, ephemeralEvent, withTimeout } from "./sse-fixture"
type SessionUpdateParams = Parameters<AgentSideConnection["sessionUpdate"]>[0]
@@ -191,6 +191,181 @@ describe("acp event behavior", () => {
}
})
test("projects foreground child session updates onto the parent turn", async () => {
const updates: SessionUpdateParams[] = []
const fixture = createSseFixture({
onPrompt({ id, send }) {
send(durableEvent("session.input.promoted", { sessionID: "ses_parent", inputID: id }))
send(
durableEvent("session.created", {
sessionID: "ses_child",
info: childSession("ses_child", "ses_parent", "Explore code"),
}),
)
send(durableEvent("session.execution.started", { sessionID: "ses_child" }))
send(
durableEvent("session.tool.input.started", {
sessionID: "ses_child",
assistantMessageID: "msg_child",
id: "call_read",
name: "read",
}),
)
send(
durableEvent("session.tool.called", {
sessionID: "ses_child",
assistantMessageID: "msg_child",
id: "call_read",
input: { path: "/workspace/src/index.ts" },
executed: false,
}),
)
send(
durableEvent("session.tool.success", {
sessionID: "ses_child",
assistantMessageID: "msg_child",
id: "call_read",
metadata: {},
content: [{ type: "text", text: "source" }],
executed: true,
}),
)
send(durableEvent("session.execution.succeeded", { sessionID: "ses_child" }))
send(durableEvent("session.execution.succeeded", { sessionID: "ses_parent" }))
},
})
try {
const response = await turn({
fixture,
connection: recordingConnection(updates),
sessionID: "ses_parent",
inputID: "input_parent",
})
expect(updates.map((item) => [item.sessionId, item.update.sessionUpdate])).toEqual([
["ses_parent", "tool_call"],
["ses_parent", "tool_call_update"],
["ses_parent", "tool_call_update"],
])
expect(updates.map((item) => ("toolCallId" in item.update ? item.update.toolCallId : undefined))).toEqual([
"ses_child:call_read",
"ses_child:call_read",
"ses_child:call_read",
])
expect(updates[0]?.update).toMatchObject({
title: "Explore code: read",
_meta: {
"opencode/child-session": {
id: "ses_child",
parentID: "ses_parent",
depth: 1,
title: "Explore code",
},
},
})
expect(response.stopReason).toBe("end_turn")
} finally {
await fixture.stop()
}
})
test("continues child extension updates after the parent turn ends", async () => {
const updates: SessionUpdateParams[] = []
const childUpdates: ChildSessionUpdate[] = []
const completed = Promise.withResolvers<void>()
const fixture = createSseFixture({
onPrompt({ id, send }) {
send(durableEvent("session.input.promoted", { sessionID: "ses_parent", inputID: id }))
send(
durableEvent("session.created", {
sessionID: "ses_background",
info: childSession("ses_background", "ses_parent", "Background research"),
}),
)
send(durableEvent("session.execution.succeeded", { sessionID: "ses_parent" }))
},
})
try {
const response = await turn({
fixture,
connection: recordingConnection(updates),
sessionID: "ses_parent",
inputID: "input_parent",
childSessionUpdate: async (update) => {
childUpdates.push(update)
if (update.type === "status" && update.status === "completed") completed.resolve()
},
})
expect(response.stopReason).toBe("end_turn")
fixture.send(
durableEvent("session.created", {
sessionID: "ses_future",
info: childSession("ses_future", "ses_parent", "Later turn child"),
}),
)
fixture.send(durableEvent("session.execution.started", { sessionID: "ses_future" }))
fixture.send(durableEvent("session.execution.started", { sessionID: "ses_background" }))
fixture.send(
durableEvent("session.tool.input.started", {
sessionID: "ses_background",
assistantMessageID: "msg_background",
id: "call_shell",
name: "shell",
}),
)
fixture.send(
durableEvent("session.tool.called", {
sessionID: "ses_background",
assistantMessageID: "msg_background",
id: "call_shell",
input: { command: "pwd" },
executed: false,
}),
)
fixture.send(
durableEvent("session.tool.success", {
sessionID: "ses_background",
assistantMessageID: "msg_background",
id: "call_shell",
metadata: { exit: 0 },
content: [{ type: "text", text: "/workspace" }],
executed: true,
}),
)
fixture.send(durableEvent("session.execution.succeeded", { sessionID: "ses_background" }))
await withTimeout(completed.promise, "background child completion was not delivered")
expect(updates).toEqual([])
expect(
childUpdates.map((update) =>
update.type === "status" ? [update.type, update.status] : [update.type, update.update.sessionUpdate],
),
).toEqual([
["status", "created"],
["status", "running"],
["update", "tool_call"],
["update", "tool_call_update"],
["update", "tool_call_update"],
["status", "completed"],
])
expect(childUpdates[2]).toMatchObject({
rootSessionId: "ses_parent",
childSessionId: "ses_background",
parentSessionId: "ses_parent",
depth: 1,
title: "Background research",
type: "update",
update: { toolCallId: "ses_background:call_shell" },
})
expect(childUpdates.some((update) => update.childSessionId === "ses_future")).toBe(false)
} finally {
await fixture.stop()
}
})
test("streams tool pending, progress, success, and failure updates", async () => {
const updates: SessionUpdateParams[] = []
const fixture = createSseFixture({
@@ -556,6 +731,7 @@ function turn(input: {
readonly connection: Connection
readonly sessionID: string
readonly inputID: string
readonly childSessionUpdate?: (update: ChildSessionUpdate) => Promise<void>
}) {
return streamTurn({
client: input.fixture.client,
@@ -565,11 +741,25 @@ function turn(input: {
start: { type: "input", id: input.inputID },
writeTextFile: false,
control: { cancelled: false, admission: new AbortController() },
childSessionUpdate: input.childSessionUpdate,
submit: (signal) =>
input.fixture.client.session.prompt({ sessionID: input.sessionID, id: input.inputID, text: "hello" }, { signal }),
})
}
function childSession(id: string, parentID: string, title: string) {
return {
id,
slug: id,
projectID: "project",
directory: "/workspace",
parentID,
title,
version: "test",
time: { created: 1, updated: 1 },
}
}
function tokens() {
return { input: 1, output: 1, reasoning: 0, cache: { read: 0, write: 0 } }
}
@@ -153,6 +153,68 @@ describe("acp permission behavior", () => {
}
})
test("routes foreground child permissions through the parent ACP session", async () => {
const permissionRequests: RequestPermissionRequest[] = []
const fixture = createSseFixture({
onPrompt({ id, send }) {
send(durableEvent("session.input.promoted", { sessionID: "ses_parent", inputID: id }))
send(
durableEvent("session.created", {
sessionID: "ses_child",
info: {
id: "ses_child",
slug: "ses_child",
projectID: "project",
directory: "/workspace",
parentID: "ses_parent",
title: "Review code",
version: "test",
time: { created: 1, updated: 1 },
},
}),
)
send(durableEvent("session.execution.started", { sessionID: "ses_child" }))
send(
permissionAsked("ses_child", "perm_child", {
action: "read",
metadata: { path: "/workspace/child.ts" },
source: { type: "tool", messageID: "msg_child", id: "call_child" },
}),
)
send(durableEvent("session.execution.succeeded", { sessionID: "ses_child" }))
send(durableEvent("session.execution.succeeded", { sessionID: "ses_parent" }))
},
})
const connection = {
sessionUpdate: async () => {},
requestPermission: async (request) => {
permissionRequests.push(request)
return { outcome: { outcome: "selected", optionId: "once" } } as const
},
} satisfies Connection
try {
await startTurn(fixture, connection, "ses_parent", "input_parent")
expect(permissionRequests).toHaveLength(1)
expect(permissionRequests[0]).toMatchObject({
sessionId: "ses_parent",
toolCall: {
toolCallId: "ses_child:call_child",
title: "Review code: /workspace/child.ts",
},
})
expect(fixture.requests).toContainEqual(
expect.objectContaining({
method: "POST",
path: "/api/session/ses_child/permission/perm_child/reply",
}),
)
} finally {
await fixture.stop()
}
})
test("previews edits during approval and syncs the completed file", async () => {
const cwd = await fs.mkdtemp(path.join(os.tmpdir(), "opencode-acp-permission-"))
const file = path.join(cwd, "file.ts")
+7
View File
@@ -2,6 +2,7 @@ import { describe, expect, test } from "bun:test"
import type { AgentSideConnection } from "@agentclientprotocol/sdk"
import { OpenCode } from "@opencode-ai/client/promise"
import { ACPService } from "../../src/acp/service"
import { ChildSessionUpdatesCapability } from "../../src/acp/event"
describe("acp service", () => {
test("creates a v2 session, registers mcp, and publishes commands", async () => {
@@ -39,11 +40,17 @@ describe("acp service", () => {
})
try {
const initialized = await service.initialize({
protocolVersion: 1,
clientCapabilities: { _meta: { [ChildSessionUpdatesCapability]: true } },
clientInfo: { name: "test", version: "1" },
})
const result = await service.newSession({
cwd: "/workspace",
mcpServers: [{ name: "docs", command: "bun", args: ["docs.ts"], env: [{ name: "TOKEN", value: "x" }] }],
})
expect(result.sessionId).toBe("ses_acp")
expect(initialized.agentCapabilities?._meta).toEqual({ [ChildSessionUpdatesCapability]: true })
expect(result.configOptions?.map((option) => option.id)).toEqual(["model", "effort", "mode"])
expect(requests).toContainEqual({
method: "PUT",
+1
View File
@@ -420,6 +420,7 @@ export type Endpoint5_26Output =
readonly data: {
readonly sessionID: Session.ID
readonly delta: { readonly [x: string]: (string & Brand.Brand<"Instruction.Hash">) | "removed" }
readonly text?: string | undefined
}
}
| {
@@ -676,7 +676,7 @@ export type SessionInstructionsUpdated = {
type: "session.instructions.updated"
durable: { aggregateID: string; seq: number; version: 2 }
location?: LocationRef
data: { sessionID: string; delta: { [x: string]: string | "removed" } }
data: { sessionID: string; delta: { [x: string]: string | "removed" }; text?: string }
}
export type SessionSynthetic = {
+2 -92
View File
@@ -1,8 +1,7 @@
export * as FileMutation from "./file-mutation"
import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
import { Context, Effect, Layer, Schema } from "effect"
import { dirname } from "path"
import { Context, Effect, Layer } from "effect"
import { KeyedMutex } from "./effect/keyed-mutex"
import { FSUtil } from "@opencode-ai/util/fs-util"
import { Bom } from "@opencode-ai/util/bom"
@@ -22,22 +21,6 @@ export interface TextWriteInput {
readonly content: string
}
export interface ConditionalWriteInput extends WriteInput {
readonly expected: Uint8Array
}
export interface RemoveInput {
readonly target: Target
}
export class StaleContentError extends Schema.TaggedErrorClass<StaleContentError>()("FileMutation.StaleContentError", {
path: Schema.String,
}) {}
export class TargetExistsError extends Schema.TaggedErrorClass<TargetExistsError>()("FileMutation.TargetExistsError", {
path: Schema.String,
}) {}
export interface WriteResult {
readonly operation: "write"
readonly target: string
@@ -45,24 +28,10 @@ export interface WriteResult {
readonly existed: boolean
}
export interface RemoveResult {
readonly operation: "remove"
readonly target: string
readonly resource: string
readonly existed: boolean
}
export interface Interface {
/** Create without replacing an existing target. */
readonly create: (input: WriteInput) => Effect.Effect<WriteResult, TargetExistsError | FSUtil.Error>
readonly write: (input: WriteInput) => Effect.Effect<WriteResult, FSUtil.Error>
/** Write text while retaining an existing UTF-8 BOM and emitting at most one BOM. */
readonly writeTextPreservingBom: (input: TextWriteInput) => Effect.Effect<WriteResult, FSUtil.Error>
/** Commit only if an existing target still has the expected bytes. */
readonly writeIfUnchanged: (
input: ConditionalWriteInput,
) => Effect.Effect<WriteResult, StaleContentError | FSUtil.Error>
readonly remove: (input: RemoveInput) => Effect.Effect<RemoveResult, FSUtil.Error>
}
export class Service extends Context.Service<Service, Interface>()("@opencode/FileMutation") {}
@@ -89,13 +58,6 @@ const layer = Layer.effect(
existed,
})
const removeResult = (target: Target, existed: boolean): RemoveResult => ({
operation: "remove",
target: target.canonical,
resource: target.resource,
existed,
})
const write = Effect.fn("FileMutation.write")((input: WriteInput) =>
withTargetLock(input.target)(
Effect.gen(function* () {
@@ -122,62 +84,10 @@ const layer = Layer.effect(
),
)
const create = Effect.fn("FileMutation.create")((input: WriteInput) =>
withTargetLock(input.target)(
Effect.gen(function* () {
const write =
typeof input.content === "string"
? fs.writeFileString(input.target.canonical, input.content, { flag: "wx" })
: fs.writeFile(input.target.canonical, input.content, { flag: "wx" })
yield* write.pipe(
Effect.catchReason("PlatformError", "NotFound", () =>
fs.ensureDir(dirname(input.target.canonical)).pipe(Effect.andThen(write)),
),
Effect.catchReason("PlatformError", "AlreadyExists", () =>
Effect.fail(new TargetExistsError({ path: input.target.canonical })),
),
)
return writeResult(input.target, false)
}),
),
)
const writeIfUnchanged = Effect.fn("FileMutation.writeIfUnchanged")((input: ConditionalWriteInput) =>
withTargetLock(input.target)(
Effect.gen(function* () {
const current = yield* fs.readFile(input.target.canonical)
if (!sameBytes(current, input.expected)) {
return yield* new StaleContentError({ path: input.target.canonical })
}
yield* typeof input.content === "string"
? fs.writeFileString(input.target.canonical, input.content)
: fs.writeFile(input.target.canonical, input.content)
return writeResult(input.target, true)
}),
),
)
const remove = Effect.fn("FileMutation.remove")((input: RemoveInput) =>
withTargetLock(input.target)(
Effect.gen(function* () {
const existed = yield* fs.remove(input.target.canonical).pipe(
Effect.as(true),
Effect.catchReason("PlatformError", "NotFound", () => Effect.succeed(false)),
)
return removeResult(input.target, existed)
}),
),
)
return Service.of({ create, write, writeTextPreservingBom, writeIfUnchanged, remove })
return Service.of({ write, writeTextPreservingBom })
}),
)
function sameBytes(left: Uint8Array, right: Uint8Array) {
if (left.length !== right.length) return false
return left.every((byte, index) => byte === right[index])
}
export const node = makeLocationNode({ service: Service, layer, deps: [FSUtil.node] })
/**
+5 -9
View File
@@ -31,10 +31,7 @@ export const ripgrepLayer = Layer.effect(
const location = yield* Location.Service
const ripgrep = yield* Ripgrep.Service
const scope = yield* Scope.Scope
const state = {
files: [] as string[],
directories: [] as string[],
}
const files: string[] = []
const directories = new Set<string>()
yield* ripgrep
.find({
@@ -43,10 +40,9 @@ export const ripgrepLayer = Layer.effect(
limit: location.vcs ? Number.MAX_SAFE_INTEGER : 100_000,
onEntry: (entry) =>
Effect.sync(() => {
state.files.push(entry.path)
files.push(entry.path)
const parts = entry.path.split("/")
parts.slice(0, -1).forEach((_, index) => directories.add(parts.slice(0, index + 1).join("/") + path.sep))
state.directories = Array.from(directories)
}),
})
.pipe(Effect.orDie, Effect.asVoid, Effect.forkIn(scope))
@@ -106,10 +102,10 @@ export const ripgrepLayer = Layer.effect(
Effect.gen(function* () {
const items =
input.type === "file"
? state.files
? files
: input.type === "directory"
? state.directories
: [...state.files, ...state.directories]
? Array.from(directories)
: [...files, ...directories]
return fuzzysort.go(input.query, items, { limit: input.limit ?? 50 }).map((item) => {
const relative = item.target
const type = relative.endsWith(path.sep) ? ("directory" as const) : ("file" as const)
+5 -4
View File
@@ -410,7 +410,7 @@ const layer = Layer.effect(
fork: Effect.fn("Session.fork")(function* (input) {
const parent = yield* result.get(input.sessionID)
const boundary = yield* db
.select({ id: SessionMessageTable.id, seq: SessionMessageTable.seq })
.select({ id: SessionMessageTable.id })
.from(SessionMessageTable)
.where(
and(
@@ -429,13 +429,14 @@ const layer = Layer.effect(
})
if (!boundary) return yield* new ForkEmptyError({ sessionID: input.sessionID })
const sessionID = SessionSchema.ID.create()
const instructionThrough =
input.boundary.type === "before" ? boundary.seq - 1 : yield* Bus.latestSequence(db, parent.id)
// The fork adopts the parent's newest instruction values rather than the
// values in effect at the boundary; copied history may contain frozen
// instruction-update text the initial baseline already reflects.
yield* bus.publish(SessionEvent.Forked, {
sessionID,
parentID: parent.id,
boundary: { ...input.boundary, messageID: boundary.id },
instructions: yield* InstructionState.valuesAt(db, parent.id, instructionThrough),
instructions: yield* InstructionState.current(db, parent.id),
})
return yield* result.get(sessionID).pipe(Effect.orDie)
}),
+3 -5
View File
@@ -80,10 +80,9 @@ export const entriesForRunner = Effect.fn("SessionHistory.entriesForRunner")(fun
.transaction(() =>
Effect.gen(function* () {
const messages = yield* messageEntries(db, sessionID)
const assembled = yield* InstructionState.assemble(db, sessionID, instructions)
return {
initial: assembled.initial,
entries: [...messages, ...assembled.updates].toSorted((a, b) => a.seq - b.seq),
initial: yield* InstructionState.initial(db, sessionID, instructions),
entries: messages,
}
}),
)
@@ -106,10 +105,9 @@ export const preview = Effect.fn("SessionHistory.preview")(function* (
)
const settled = unsettled === -1 ? messages : messages.slice(0, unsettled)
const assembled = yield* InstructionState.preview(db, sessionID, instructions, observed)
const entries = [...settled, ...assembled.updates].toSorted((a, b) => a.seq - b.seq)
return {
initial: assembled.initial,
messages: entries.map((entry) => entry.message),
messages: settled.map((entry) => entry.message),
instructionUpdate: assembled.update,
}
}),
+56 -230
View File
@@ -1,25 +1,20 @@
export * as InstructionState from "./instruction-state"
import { and, asc, desc, eq, gt, inArray, lte, sql } from "drizzle-orm"
import { DateTime, Effect, Option, Schema } from "effect"
import { eq, inArray, sql } from "drizzle-orm"
import { Effect, Option, Schema } from "effect"
import type { Database } from "../database/database"
import { Bus } from "../bus"
import { EventTable } from "../event/sql"
import type { Bus } from "../bus"
import { Instructions } from "../instructions/index"
import { SessionEvent } from "./event"
import { SessionMessage } from "./message"
import { Event } from "@opencode-ai/schema/event"
import { SessionSchema } from "./schema"
import { InstructionBlobTable, InstructionStateTable } from "./sql"
type DatabaseService = Database.Interface["db"]
const decodeInstructionsUpdated = Schema.decodeUnknownSync(SessionEvent.InstructionsUpdated.data)
const decodeForked = Schema.decodeUnknownSync(SessionEvent.Forked.data)
export interface Observation extends Instructions.Admission {
readonly sessionID: SessionSchema.ID
readonly initial: boolean
readonly previous: Instructions.Values
readonly current: Instructions.Values
}
@@ -28,13 +23,14 @@ export const observe = Effect.fn("InstructionState.observe")(function* (
instructions: Instructions.Instructions,
sessionID: SessionSchema.ID,
): Effect.fn.Return<Observation, Instructions.InitializationBlocked> {
const [observed, stored] = yield* Effect.all([Instructions.read(instructions), ensure(db, sessionID)], {
const [observed, stored] = yield* Effect.all([Instructions.read(instructions), find(db, sessionID)], {
concurrency: "unbounded",
})
const result = yield* observeAgainst(observed, stored?.current_values)
return {
sessionID,
initial: !stored,
previous: stored?.current_values ?? {},
...result,
}
})
@@ -42,12 +38,20 @@ export const observe = Effect.fn("InstructionState.observe")(function* (
export const commit = Effect.fn("InstructionState.commit")(function* (
db: DatabaseService,
bus: Bus.Interface,
instructions: Instructions.Instructions,
observation: Observation,
) {
if (!observation.initial && Object.keys(observation.delta).length === 0) return
// The rendered text is frozen into the durable event: replaying it later would
// require the Location-scoped registry that produced it.
const text = observation.initial ? "" : yield* renderUpdateText(db, instructions, observation)
yield* bus.publish(
SessionEvent.InstructionsUpdated,
{ sessionID: observation.sessionID, delta: observation.delta },
{
sessionID: observation.sessionID,
delta: observation.delta,
...(text.length > 0 ? { text } : {}),
},
{
// Initial sync establishes the baseline; unlike later deltas it is not chronological history.
...(observation.initial ? { metadata: { instructions: { initial: true } } } : {}),
@@ -56,13 +60,27 @@ export const commit = Effect.fn("InstructionState.commit")(function* (
)
})
const renderUpdateText = Effect.fnUntraced(function* (
db: DatabaseService,
instructions: Instructions.Instructions,
observation: Observation,
) {
const replaced = Object.entries(observation.previous).filter(([key]) => Object.hasOwn(observation.delta, key))
const blobs = yield* loadBlobs(db, replaced.map(([, hash]) => hash))
const previous = Object.fromEntries(replaced.map(([key, hash]) => [key, requireBlob(blobs, hash)]))
const admitted = new Map(
Object.entries(observation.blobs).map(([hash, value]) => [Instructions.Hash.make(hash), value]),
)
return Instructions.renderUpdate(instructions, previous, dereferenceDelta(observation.delta, admitted))
})
export const prepare = Effect.fn("InstructionState.prepare")(function* (
db: DatabaseService,
bus: Bus.Interface,
instructions: Instructions.Instructions,
sessionID: SessionSchema.ID,
) {
yield* commit(db, bus, yield* observe(db, instructions, sessionID))
yield* commit(db, bus, instructions, yield* observe(db, instructions, sessionID))
})
export const apply = Effect.fn("InstructionState.apply")(function* (
@@ -140,79 +158,24 @@ export const reset = Effect.fn("InstructionState.reset")(function* (db: Database
.pipe(Effect.orDie)
})
export const rebuild = Effect.fn("InstructionState.rebuild")(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
) {
const state = yield* stateFromEvents(db, sessionID)
if (!state) {
yield* reset(db, sessionID)
return undefined
}
yield* db
.insert(InstructionStateTable)
.values(state)
.onConflictDoUpdate({
target: InstructionStateTable.session_id,
set: {
epoch_start: state.epoch_start,
through_seq: state.through_seq,
initial_values: state.initial_values,
current_values: state.current_values,
},
})
.run()
.pipe(Effect.orDie)
return state
})
const assembleState = Effect.fnUntraced(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
instructions: Instructions.Instructions,
state: typeof InstructionStateTable.$inferSelect,
) {
const rows = yield* instructionUpdatesAfter(db, sessionID, state.epoch_start)
const updates = rows.map((row) => ({
row,
delta: decodeInstructionsUpdated(row.data).delta,
}))
const blobs = yield* loadBlobs(db, [
...Object.values(state.initial_values),
...updates.flatMap((update) =>
Object.values(update.delta).filter((hash): hash is Instructions.Hash => hash !== "removed"),
),
])
const valuesAtStart = dereference(state.initial_values, blobs)
let values = valuesAtStart
const result: Array<{ readonly seq: number; readonly message: SessionMessage.System }> = []
for (const update of updates) {
const delta = dereferenceDelta(update.delta, blobs)
const text = Instructions.renderUpdate(instructions, values, delta)
if (text.length > 0)
result.push({
seq: update.row.seq,
message: SessionMessage.System.make({
id: SessionMessage.ID.fromEvent(Event.ID.make(update.row.id)),
type: "system",
text,
time: { created: DateTime.makeUnsafe(update.row.created) },
}),
})
values = Instructions.applyDelta(values, delta)
}
return { initial: Instructions.renderInitial(instructions, valuesAtStart), updates: result, current: values }
})
export const assemble = Effect.fn("InstructionState.assemble")(function* (
/** Renders the epoch baseline shown at the start of every model request. */
export const initial = Effect.fn("InstructionState.initial")(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
instructions: Instructions.Instructions,
) {
const state = yield* find(db, sessionID)
if (!state) return yield* Effect.die(new Error(`Instruction state not found during assembly: ${sessionID}`))
const assembled = yield* assembleState(db, sessionID, instructions, state)
return { initial: assembled.initial, updates: assembled.updates }
const blobs = yield* loadBlobs(db, Object.values(state.initial_values))
return Instructions.renderInitial(instructions, dereference(state.initial_values, blobs))
})
/** The current instruction values, used to seed a fork's baseline. */
export const current = Effect.fn("InstructionState.current")(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
) {
return (yield* find(db, sessionID))?.current_values
})
export const preview = Effect.fn("InstructionState.preview")(function* (
@@ -221,20 +184,26 @@ export const preview = Effect.fn("InstructionState.preview")(function* (
instructions: Instructions.Instructions,
observed: Instructions.ReadResult,
) {
const state = yield* readState(db, sessionID)
const state = yield* find(db, sessionID)
const result = yield* observeAgainst(observed, state?.current_values)
const blobs = new Map<Instructions.Hash, Schema.Json>(
const observedBlobs = new Map<Instructions.Hash, Schema.Json>(
Object.entries(result.blobs).map(([hash, value]) => [Instructions.Hash.make(hash), value]),
)
if (!state) {
const values = dereference(result.current, blobs)
return { initial: Instructions.renderInitial(instructions, values), updates: [], update: "" }
const values = dereference(result.current, observedBlobs)
return { initial: Instructions.renderInitial(instructions, values), update: "" }
}
const assembled = yield* assembleState(db, sessionID, instructions, state)
const stored = yield* loadBlobs(db, [
...Object.values(state.initial_values),
...Object.values(state.current_values),
])
return {
initial: assembled.initial,
updates: assembled.updates,
update: Instructions.renderUpdate(instructions, assembled.current, dereferenceDelta(result.delta, blobs)),
initial: Instructions.renderInitial(instructions, dereference(state.initial_values, stored)),
update: Instructions.renderUpdate(
instructions,
dereference(state.current_values, stored),
dereferenceDelta(result.delta, new Map([...stored, ...observedBlobs])),
),
}
})
@@ -255,46 +224,6 @@ const find = Effect.fnUntraced(function* (db: DatabaseService, sessionID: Sessio
.pipe(Effect.orDie)
})
const ensure = Effect.fnUntraced(function* (db: DatabaseService, sessionID: SessionSchema.ID) {
const stored = yield* find(db, sessionID)
if (!stored) return yield* rebuild(db, sessionID)
const latest = yield* latestRelevantSequence(db, sessionID)
if (!latest || latest.seq <= stored.through_seq) return stored
return yield* rebuild(db, sessionID)
})
const readState = Effect.fnUntraced(function* (db: DatabaseService, sessionID: SessionSchema.ID) {
const stored = yield* find(db, sessionID)
if (!stored) return yield* stateFromEvents(db, sessionID)
const latest = yield* latestRelevantSequence(db, sessionID)
if (!latest || latest.seq <= stored.through_seq) return stored
return yield* stateFromEvents(db, sessionID)
})
const stateFromEvents = Effect.fnUntraced(function* (db: DatabaseService, sessionID: SessionSchema.ID) {
const folded = fold(yield* instructionEvents(db, sessionID))
return folded ? foldedState(sessionID, folded) : undefined
})
export const valuesAt = Effect.fn("InstructionState.valuesAt")(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
through: number,
) {
return fold(yield* instructionEvents(db, sessionID, through))?.current
})
const latestRelevantSequence = Effect.fnUntraced(function* (db: DatabaseService, sessionID: SessionSchema.ID) {
return yield* db
.select({ seq: EventTable.seq })
.from(EventTable)
.where(and(eq(EventTable.aggregate_id, sessionID), inArray(EventTable.type, relevantEventTypes)))
.orderBy(desc(EventTable.seq))
.limit(1)
.get()
.pipe(Effect.orDie)
})
const insertBlobs = Effect.fnUntraced(function* (db: DatabaseService, blobs: Readonly<Record<string, Schema.Json>>) {
const rows = Object.entries(blobs).map(([hash, value]) => ({ hash: Instructions.Hash.make(hash), value }))
if (rows.length === 0) return
@@ -339,106 +268,3 @@ function requireBlob(blobs: ReadonlyMap<Instructions.Hash, Schema.Json>, hash: I
if (value === undefined) throw new Error(`Instruction blob not found: ${hash}`)
return value
}
const instructionEventType = Bus.versionedType(
SessionEvent.InstructionsUpdated.type,
SessionEvent.InstructionsUpdated.durable.version,
)
const compactionEventType = Bus.versionedType(
SessionEvent.Compaction.Ended.type,
SessionEvent.Compaction.Ended.durable.version,
)
const movedEventType = Bus.versionedType(SessionEvent.Moved.type, SessionEvent.Moved.durable.version)
const revertedEventType = Bus.versionedType(
SessionEvent.RevertEvent.Committed.type,
SessionEvent.RevertEvent.Committed.durable.version,
)
const forkedEventType = Bus.versionedType(SessionEvent.Forked.type, SessionEvent.Forked.durable.version)
const relevantEventTypes = [
forkedEventType,
instructionEventType,
compactionEventType,
movedEventType,
revertedEventType,
]
type InstructionEventRow = typeof EventTable.$inferSelect
const instructionEvents = Effect.fnUntraced(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
through?: number,
): Effect.fn.Return<ReadonlyArray<InstructionEventRow>> {
return yield* eventRows(db, sessionID, relevantEventTypes, undefined, through)
})
const instructionUpdatesAfter = Effect.fnUntraced(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
after: number,
) {
return yield* eventRows(db, sessionID, [instructionEventType], after)
})
const eventRows = Effect.fnUntraced(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
types: ReadonlyArray<string>,
after?: number,
through?: number,
): Effect.fn.Return<ReadonlyArray<InstructionEventRow>> {
return yield* db
.select()
.from(EventTable)
.where(
and(
eq(EventTable.aggregate_id, sessionID),
inArray(EventTable.type, types),
after === undefined ? undefined : gt(EventTable.seq, after),
through === undefined ? undefined : lte(EventTable.seq, through),
),
)
.orderBy(asc(EventTable.seq))
.all()
.pipe(Effect.orDie)
})
function fold(rows: ReadonlyArray<InstructionEventRow>) {
return rows.reduce<
| {
readonly epochStart: number
readonly throughSeq: number
readonly initial: Instructions.Values
readonly current: Instructions.Values
}
| undefined
>((state, row) => {
if (row.type === forkedEventType) {
const instructions = decodeForked(row.data).instructions
return instructions
? { epochStart: row.seq, throughSeq: row.seq, initial: instructions, current: instructions }
: undefined
}
if (row.type === movedEventType || row.type === revertedEventType) return undefined
if (row.type === compactionEventType)
return state
? { epochStart: row.seq, throughSeq: row.seq, initial: state.current, current: state.current }
: undefined
if (row.type !== instructionEventType) return state
const delta = decodeInstructionsUpdated(row.data).delta
const current = Instructions.applyHashDelta(state?.current ?? {}, delta)
return state
? { ...state, throughSeq: row.seq, current }
: { epochStart: row.seq, throughSeq: row.seq, initial: current, current }
}, undefined)
}
function foldedState(sessionID: SessionSchema.ID, folded: NonNullable<ReturnType<typeof fold>>) {
return {
session_id: sessionID,
epoch_start: folded.epochStart,
through_seq: folded.throughSeq,
initial_values: folded.initial,
current_values: folded.current,
}
}
+12 -1
View File
@@ -179,7 +179,18 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
"session.execution.succeeded": () => clearCurrentRetry,
"session.execution.failed": () => clearCurrentRetry,
"session.execution.interrupted": () => clearCurrentRetry,
"session.instructions.updated": () => Effect.void,
"session.instructions.updated": (event) => {
if (event.data.text === undefined) return Effect.void
return adapter.appendMessage(
SessionMessage.System.make({
id: SessionMessage.ID.fromEvent(event.id),
type: "system",
text: event.data.text,
metadata: event.metadata,
time: { created: event.created },
}),
)
},
"session.synthetic": (event) => {
return adapter.appendMessage(
SessionMessage.Synthetic.make({
+22 -39
View File
@@ -12,10 +12,8 @@ import {
User,
UserData,
} from "@opencode-ai/schema/session-pending"
import { Event } from "@opencode-ai/schema/event"
import type { Database } from "../database/database"
import { Bus } from "../bus"
import { EventTable } from "../event/sql"
import { KeyedMutex } from "../effect/keyed-mutex"
import { SessionEvent } from "./event"
import { SessionMessage } from "./message"
@@ -37,11 +35,7 @@ const decodeUser = Schema.decodeUnknownSync(UserData)
const encodeUser = Schema.encodeSync(UserData)
const decodeSynthetic = Schema.decodeUnknownSync(SyntheticData)
const encodeSynthetic = Schema.encodeSync(SyntheticData)
const decodeAdmittedEvent = Schema.decodeUnknownOption(SessionEvent.InputAdmitted.data)
const admittedEventType = Bus.versionedType(
SessionEvent.InputAdmitted.type,
SessionEvent.InputAdmitted.durable.version,
)
const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Info)
const inboxLocks = KeyedMutex.makeUnsafe<SessionSchema.ID>()
export class LifecycleConflict extends Schema.TaggedErrorClass<LifecycleConflict>()(
@@ -103,46 +97,35 @@ export const compaction = Effect.fn("SessionPending.compaction")(function* (
return entry.type === "compaction" ? entry : undefined
})
/**
* Reconstruct the admitted record for a pending row that was already consumed
* by promotion. The projected `session_message` row proves promotion happened;
* the durable `session.input.admitted` event retains the exact admitted
* message, including delivery.
*/
const promotedFromHistory = Effect.fn("SessionPending.promotedFromHistory")(function* (
const promotedFromMessage = Effect.fn("SessionPending.promotedFromMessage")(function* (
db: DatabaseService,
sessionID: SessionSchema.ID,
id: SessionMessage.ID,
delivery: Delivery,
) {
const message = yield* db
const row = yield* db
.select()
.from(SessionMessageTable)
.where(eq(SessionMessageTable.id, id))
.get()
.pipe(Effect.orDie)
if (message === undefined) return undefined
if (message.session_id !== sessionID || (message.type !== "user" && message.type !== "synthetic"))
if (row === undefined) return undefined
if (row.session_id !== sessionID || (row.type !== "user" && row.type !== "synthetic"))
return yield* Effect.die(new LifecycleConflict({ id }))
const rows = yield* db
.select()
.from(EventTable)
.where(and(eq(EventTable.aggregate_id, sessionID), eq(EventTable.type, admittedEventType)))
.all()
.pipe(Effect.orDie)
for (const row of rows) {
const decoded = decodeAdmittedEvent(row.data)
if (decoded._tag !== "Some" || decoded.value.inputID !== id) continue
const base = {
id,
sessionID,
timeCreated: DateTime.makeUnsafe(row.created),
}
return decoded.value.input.type === "user"
? User.make({ ...base, ...decoded.value.input })
: Synthetic.make({ ...base, ...decoded.value.input })
}
// A projected message without an admitted event in this aggregate (for
// example fork-copied history) is not a retryable admission.
const message = decodeMessage({ ...row.data, id: row.id, type: row.type })
const base = { id, sessionID, timeCreated: message.time.created, delivery }
if (message.type === "user")
return User.make({
...base,
type: "user",
data: decodeUser(message),
})
if (message.type === "synthetic")
return Synthetic.make({
...base,
type: "synthetic",
data: decodeSynthetic(message),
})
return yield* Effect.die(new LifecycleConflict({ id }))
})
@@ -160,7 +143,7 @@ export const admit = Effect.fn("SessionPending.admit")(function* (
if (existing.type === "compaction") return yield* Effect.die(new LifecycleConflict({ id: request.id }))
return existing
}
const promoted = yield* promotedFromHistory(db, request.sessionID, request.id)
const promoted = yield* promotedFromMessage(db, request.sessionID, request.id, request.input.delivery)
if (promoted !== undefined) return promoted
return yield* bus
.publish(SessionEvent.InputAdmitted, {
@@ -426,7 +409,7 @@ const publish = Effect.fn("SessionPending.publish")(function* (
.pipe(
Effect.catchDefect((defect) =>
defect instanceof LifecycleConflict
? promotedFromHistory(db, sessionID, entry.id).pipe(
? promotedFromMessage(db, sessionID, entry.id, entry.delivery).pipe(
Effect.flatMap((stored) => (stored !== undefined ? Effect.void : Effect.die(defect))),
)
: Effect.die(defect),
+15 -59
View File
@@ -1,6 +1,6 @@
export * as SessionProjector from "./projector"
import { and, asc, desc, eq, gt, gte, inArray, lt, lte, sql } from "drizzle-orm"
import { and, asc, desc, eq, gt, gte, lt, lte, sql } from "drizzle-orm"
import { DateTime, Effect, Layer, Schema, Stream } from "effect"
import { Database } from "../database/database"
import { Bus } from "../bus"
@@ -21,10 +21,7 @@ import { Money } from "@opencode-ai/schema/money"
type DatabaseService = Database.Interface["db"]
type CurrentDurableEvent = Extract<SessionEvent.Event, { readonly durable: object }>
type MessageEvent = Exclude<
CurrentDurableEvent,
typeof SessionEvent.Forked.Type | typeof SessionEvent.Deleted.Type | typeof SessionEvent.InstructionsUpdated.Type
>
type MessageEvent = Exclude<CurrentDurableEvent, typeof SessionEvent.Forked.Type | typeof SessionEvent.Deleted.Type>
const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Info)
const encodeMessage = Schema.encodeSync(SessionMessage.Info)
@@ -255,66 +252,22 @@ const projectFork = Effect.fn("SessionProjector.projectFork")(function* (
.pipe(Effect.orDie)
if (rows.length === 0) break
const idMap = new Map(rows.map((row) => [row.id, SessionMessage.ID.create()]))
yield* db
.insert(SessionMessageTable)
.values(
rows.map((row) => {
const id = idMap.get(row.id)
if (!id) throw new Error(`Fork message ID mapping missing: ${row.id}`)
return {
id,
session_id: event.data.sessionID,
type: row.type,
seq: row.seq,
time_created: row.time_created,
time_updated: row.time_updated,
data: row.data,
}
}),
rows.map((row) => ({
id: SessionMessage.ID.create(),
session_id: event.data.sessionID,
type: row.type,
seq: row.seq,
time_created: row.time_created,
time_updated: row.time_updated,
data: row.data,
})),
)
.run()
.pipe(Effect.orDie)
const pendingRows = yield* db
.select()
.from(SessionPendingTable)
.where(
and(
eq(SessionPendingTable.session_id, event.data.parentID),
inArray(
SessionPendingTable.id,
rows.map((row) => row.id),
),
),
)
.all()
.pipe(Effect.orDie)
if (pendingRows.length > 0) {
yield* db
.insert(SessionPendingTable)
.values(
pendingRows.flatMap((row) => {
const id = idMap.get(row.id)
return id && row.type !== "compaction"
? [
{
id,
session_id: event.data.sessionID,
type: row.type,
data: row.data,
delivery: row.delivery,
admitted_seq: row.admitted_seq,
time_created: row.time_created,
},
]
: []
}),
)
.run()
.pipe(Effect.orDie)
}
cursor = rows.at(-1)!.seq
}
if (copiedSeq !== undefined) yield* Bus.reserveSequence(db, event.data.sessionID, copiedSeq)
@@ -682,7 +635,10 @@ const layer = Layer.effectDiscard(
yield* bus.project(SessionEvent.Execution.Failed, (event) => run(db, event))
yield* bus.project(SessionEvent.Execution.Interrupted, (event) => run(db, event))
yield* bus.project(SessionEvent.InstructionsUpdated, (event) =>
InstructionState.apply(db, event.data.sessionID, event.durable.seq, event.data.delta),
Effect.gen(function* () {
yield* run(db, event)
yield* InstructionState.apply(db, event.data.sessionID, event.durable.seq, event.data.delta)
}),
)
yield* bus.project(SessionEvent.Synthetic, (event) => run(db, event))
yield* bus.project(SessionEvent.Skill.Activated, (event) => run(db, event))
+2 -1
View File
@@ -18,8 +18,9 @@ export function isRetryable(error: AIError) {
switch (error.reason._tag) {
case "RateLimit":
case "ProviderInternal":
case "Transport":
return true
case "Transport":
return error.reason.delivery === undefined || error.reason.delivery === "not-sent"
case "InvalidProviderOutput":
return error.reason.classification === "incomplete-stream"
case "Authentication":
-162
View File
@@ -89,68 +89,6 @@ describe("FileMutation", () => {
),
)
it.live("rejects create when a prospective target appears after resolution", () =>
withTmp((directory) =>
Effect.gen(function* () {
const targetPath = path.join(directory, "appeared.txt")
const target = yield* (yield* LocationMutation.Service).resolve({ path: "appeared.txt" })
yield* Effect.promise(() => fs.writeFile(targetPath, "winner"))
expect(
yield* (yield* FileMutation.Service).create({ target, content: "replacement" }).pipe(Effect.flip),
).toMatchObject({
_tag: "FileMutation.TargetExistsError",
})
expect(yield* Effect.promise(() => fs.readFile(targetPath, "utf8"))).toBe("winner")
}).pipe(provide(directory)),
),
)
it.live("creates when an existing target disappears after resolution", () =>
withTmp((directory) =>
Effect.gen(function* () {
const targetPath = path.join(directory, "removed.txt")
yield* Effect.promise(() => fs.writeFile(targetPath, "before"))
const target = yield* (yield* LocationMutation.Service).resolve({ path: "removed.txt" })
yield* Effect.promise(() => fs.rm(targetPath))
expect(yield* (yield* FileMutation.Service).create({ target, content: "after" })).toEqual({
operation: "write",
target: target.canonical,
resource: "removed.txt",
existed: false,
})
expect(yield* Effect.promise(() => fs.readFile(targetPath, "utf8"))).toBe("after")
}).pipe(provide(directory)),
),
)
it.live("removes an existing internal file", () =>
withTmp((directory) =>
Effect.gen(function* () {
const targetPath = path.join(directory, "remove.txt")
yield* Effect.promise(() => fs.writeFile(targetPath, "remove"))
const target = yield* (yield* LocationMutation.Service).resolve({ path: "remove.txt" })
const result = yield* (yield* FileMutation.Service).remove({ target })
expect(result).toEqual({
operation: "remove",
target: target.canonical,
resource: "remove.txt",
existed: true,
})
expect(
yield* Effect.promise(() =>
fs.stat(targetPath).then(
() => true,
() => false,
),
),
).toBe(false)
}).pipe(provide(directory)),
),
)
it.live("writes an explicitly resolved external target", () =>
withTmp((directory) =>
withTmp((outside) =>
@@ -171,49 +109,6 @@ describe("FileMutation", () => {
),
)
it.live("removes an explicitly resolved external target", () =>
withTmp((directory) =>
withTmp((outside) =>
Effect.gen(function* () {
const targetPath = path.join(outside, "external.txt")
yield* Effect.promise(() => fs.writeFile(targetPath, "external"))
const target = yield* (yield* LocationMutation.Service).resolve({ path: targetPath })
const result = yield* (yield* FileMutation.Service).remove({ target })
expect(result).toEqual({
operation: "remove",
target: target.canonical,
resource: target.resource,
existed: true,
})
expect(
yield* Effect.promise(() =>
fs.stat(targetPath).then(
() => true,
() => false,
),
),
).toBe(false)
}).pipe(provide(directory)),
),
),
)
it.live("reports a missing target as not removed without checking existence first", () =>
withTmp((directory) =>
Effect.gen(function* () {
const target = yield* (yield* LocationMutation.Service).resolve({ path: "missing.txt" })
expect(yield* (yield* FileMutation.Service).remove({ target })).toEqual({
operation: "remove",
target: target.canonical,
resource: "missing.txt",
existed: false,
})
}).pipe(provide(directory)),
),
)
it.live("serializes concurrent writes to the same canonical target", () =>
withTmp((directory) =>
Effect.gen(function* () {
@@ -257,63 +152,6 @@ describe("FileMutation", () => {
),
)
it.live("allows only one concurrent conditional write based on the same bytes", () =>
withTmp((directory) =>
Effect.gen(function* () {
const targetPath = path.join(directory, "shared.txt")
yield* Effect.promise(() => fs.writeFile(targetPath, "initial"))
const firstStarted = yield* Deferred.make<void>()
const releaseFirst = yield* Deferred.make<void>()
let writes = 0
const filesystem = instrumentWrites((write) =>
Effect.gen(function* () {
writes++
if (writes === 1) {
yield* Deferred.succeed(firstStarted, undefined)
yield* Deferred.await(releaseFirst)
}
yield* write
}),
)
yield* Effect.gen(function* () {
const mutation = yield* LocationMutation.Service
const files = yield* FileMutation.Service
const target = yield* mutation.resolve({ path: "shared.txt" })
const expected = new TextEncoder().encode("initial")
const first = yield* files.writeIfUnchanged({ target, expected, content: "first" }).pipe(Effect.forkChild)
yield* Deferred.await(firstStarted)
const second = yield* files
.writeIfUnchanged({ target, expected, content: "second" })
.pipe(Effect.flip, Effect.forkChild)
yield* Deferred.succeed(releaseFirst, undefined)
yield* Fiber.join(first)
expect(yield* Fiber.join(second)).toMatchObject({ _tag: "FileMutation.StaleContentError" })
expect(yield* Effect.promise(() => fs.readFile(targetPath, "utf8"))).toBe("first")
expect(writes).toBe(1)
}).pipe(provide(directory, filesystem))
}),
),
)
it.live("rejects a conditional write when target content is already stale", () =>
withTmp((directory) =>
Effect.gen(function* () {
const targetPath = path.join(directory, "stale.txt")
yield* Effect.promise(() => fs.writeFile(targetPath, "current"))
const target = yield* (yield* LocationMutation.Service).resolve({ path: "stale.txt" })
expect(
yield* (yield* FileMutation.Service)
.writeIfUnchanged({ target, expected: new TextEncoder().encode("older"), content: "replacement" })
.pipe(Effect.flip),
).toMatchObject({ _tag: "FileMutation.StaleContentError", path: target.canonical })
expect(yield* Effect.promise(() => fs.readFile(targetPath, "utf8"))).toBe("current")
}).pipe(provide(directory)),
),
)
it.live("allows distinct canonical targets to proceed independently", () =>
withTmp((directory) =>
Effect.gen(function* () {
+55 -11
View File
@@ -14,7 +14,7 @@ import { AbsolutePath } from "@opencode-ai/core/schema"
import { InstructionState } from "@opencode-ai/core/session/instruction-state"
import { SessionProjector } from "@opencode-ai/core/session/projector"
import { SessionSchema } from "@opencode-ai/core/session/schema"
import { InstructionBlobTable, InstructionStateTable, SessionTable } from "@opencode-ai/core/session/sql"
import { InstructionBlobTable, InstructionStateTable, SessionMessageTable, SessionTable } from "@opencode-ai/core/session/sql"
import { testEffect } from "./lib/effect"
const it = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionProjector.node])))
@@ -105,6 +105,7 @@ describe("InstructionState", () => {
expect(observation).toEqual({
sessionID,
initial: true,
previous: {},
current: {
"test/first": Instructions.hash("first"),
"test/second": Instructions.hash("second"),
@@ -156,7 +157,7 @@ describe("InstructionState", () => {
const initial = yield* InstructionState.observe(db, instructions, sessionID)
expect(reads).toBe(2)
yield* InstructionState.commit(db, events, initial)
yield* InstructionState.commit(db, events, instructions, initial)
expect(reads).toBe(2)
current = "changed"
@@ -166,6 +167,10 @@ describe("InstructionState", () => {
expect(changed).toMatchObject({
sessionID,
initial: false,
previous: {
"test/current": Instructions.hash("initial"),
"test/retired": Instructions.hash("retired"),
},
current: { "test/current": Instructions.hash("changed") },
delta: {
"test/current": Instructions.hash("changed"),
@@ -173,7 +178,7 @@ describe("InstructionState", () => {
},
blobs: { [Instructions.hash("changed")]: "changed" },
})
yield* InstructionState.commit(db, events, changed)
yield* InstructionState.commit(db, events, instructions, changed)
expect(reads).toBe(4)
yield* unsubscribe
@@ -190,6 +195,11 @@ describe("InstructionState", () => {
"test/retired": "removed",
},
])
// The chronological update text is frozen into the event; the baseline has none.
expect((yield* instructionEvents(db, sessionID)).map((event) => event.data.text)).toEqual([
undefined,
"changed\n\nRemoved retired",
])
expect(yield* db.select().from(InstructionStateTable).get().pipe(Effect.orDie)).toMatchObject({
initial_values: {
"test/current": Instructions.hash("initial"),
@@ -222,18 +232,19 @@ describe("InstructionState", () => {
expect(observation).toEqual({
sessionID,
initial: false,
previous: { "test/context": Instructions.hash("unchanged") },
current: { "test/context": Instructions.hash("unchanged") },
delta: {},
blobs: {},
})
yield* InstructionState.commit(db, events, observation)
yield* InstructionState.commit(db, events, instructions, observation)
expect(yield* instructionEvents(db, sessionID)).toEqual(beforeEvents)
expect(yield* db.select().from(InstructionBlobTable).all().pipe(Effect.orDie)).toEqual(beforeBlobs)
}),
)
it.effect("assembles a fresh private update without repairing a missing cache", () =>
it.effect("treats a missing state row as a fresh baseline without repairing it", () =>
Effect.gen(function* () {
const sessionID = SessionSchema.ID.make("ses_instruction_generate")
const { db, events } = yield* setup(sessionID)
@@ -254,7 +265,7 @@ describe("InstructionState", () => {
const assembled = yield* preview(db, sessionID, instructions)
expect(assembled).toEqual({ initial: "Initial context", updates: [], update: "Changed context" })
expect(assembled).toEqual({ initial: "Changed context", update: "" })
expect(yield* instructionEvents(db, sessionID)).toEqual(beforeEvents)
expect(yield* db.select().from(InstructionBlobTable).all().pipe(Effect.orDie)).toEqual(beforeBlobs)
expect(
@@ -268,7 +279,7 @@ describe("InstructionState", () => {
}),
)
it.effect("reads through a stale cache without repairing it", () =>
it.effect("trusts the projected state without consulting durable events", () =>
Effect.gen(function* () {
const sessionID = SessionSchema.ID.make("ses_instruction_generate_stale")
const { db, events } = yield* setup(sessionID)
@@ -280,6 +291,7 @@ describe("InstructionState", () => {
yield* InstructionState.prepare(db, events, instructions, sessionID)
value = "Committed update"
yield* InstructionState.prepare(db, events, instructions, sessionID)
// Tamper with the projected state; the authoritative row wins over event history.
yield* db
.update(InstructionStateTable)
.set({ through_seq: 0, current_values: { "test/context": Instructions.hash("Initial context") } })
@@ -294,7 +306,6 @@ describe("InstructionState", () => {
const assembled = yield* preview(db, sessionID, instructions)
expect(assembled.initial).toBe("Initial context")
expect(assembled.updates.map((entry) => entry.message.text)).toEqual(["Committed update"])
expect(assembled.update).toBe("Private update")
expect(yield* instructionEvents(db, sessionID)).toEqual(beforeEvents)
expect(yield* db.select().from(InstructionBlobTable).all().pipe(Effect.orDie)).toEqual(beforeBlobs)
@@ -302,6 +313,41 @@ describe("InstructionState", () => {
}),
)
it.effect("persists chronological updates as system messages", () =>
Effect.gen(function* () {
const sessionID = SessionSchema.ID.make("ses_instruction_messages")
const { db, events } = yield* setup(sessionID)
let value = "Initial context"
const instructions = source(
"test/context",
Effect.sync(() => value),
)
const messages = () =>
db
.select()
.from(SessionMessageTable)
.where(and(eq(SessionMessageTable.session_id, sessionID), eq(SessionMessageTable.type, "system")))
.orderBy(asc(SessionMessageTable.seq))
.all()
.pipe(Effect.orDie)
// The initial baseline is not chronological history and produces no message.
yield* InstructionState.prepare(db, events, instructions, sessionID)
expect(yield* messages()).toEqual([])
value = "Changed context"
yield* InstructionState.prepare(db, events, instructions, sessionID)
const rows = yield* messages()
expect(rows).toHaveLength(1)
expect(rows[0]?.data).toMatchObject({ text: "Changed context" })
expect(rows.map((row) => row.seq)).toEqual([(yield* instructionEvents(db, sessionID)).at(-1)!.seq])
// A no-op observation adds nothing.
yield* InstructionState.prepare(db, events, instructions, sessionID)
expect(yield* messages()).toHaveLength(1)
}),
)
it.effect("assembles initial instructions without persisting a baseline", () =>
Effect.gen(function* () {
const sessionID = SessionSchema.ID.make("ses_instruction_generate_initial")
@@ -310,7 +356,6 @@ describe("InstructionState", () => {
expect(yield* preview(db, sessionID, instructions)).toEqual({
initial: "Initial context",
updates: [],
update: "",
})
expect(yield* instructionEvents(db, sessionID)).toEqual([])
@@ -336,7 +381,6 @@ describe("InstructionState", () => {
expect(yield* preview(db, sessionID, instructions)).toEqual({
initial: "Committed context",
updates: [],
update: "",
})
expect(yield* instructionEvents(db, sessionID)).toEqual(beforeEvents)
@@ -388,7 +432,7 @@ describe("InstructionState", () => {
for (const next of ["initial", "changed", "changed", Instructions.removed] as const) {
value = next
yield* InstructionState.observe(db, observedInstructions, observedSessionID).pipe(
Effect.flatMap((observation) => InstructionState.commit(db, events, observation)),
Effect.flatMap((observation) => InstructionState.commit(db, events, observedInstructions, observation)),
)
yield* InstructionState.prepare(db, events, preparedInstructions, preparedSessionID)
}
+2 -6
View File
@@ -284,13 +284,9 @@ describe("Session.create", () => {
})
expect(yield* SessionPending.find(db, forkContext[0].id)).toBeUndefined()
expect(yield* SessionPending.find(db, forkContext[1].id)).toBeUndefined()
// Fork-copied messages have no admitted event in the fork aggregate, so
// reusing their IDs as prompt IDs is conflicting reuse, not a retry.
expect(
yield* session
.prompt({ id: forkContext[0].id, sessionID: forked.id, text: "First", resume: false })
.pipe(Effect.flip),
).toMatchObject({ _tag: "Session.PromptConflictError", messageID: forkContext[0].id })
yield* session.prompt({ id: forkContext[0].id, sessionID: forked.id, text: "First", resume: false }),
).toMatchObject({ id: forkContext[0].id, type: "user", data: { text: "First" } })
yield* session.prompt({
sessionID: parent.id,
+22
View File
@@ -110,4 +110,26 @@ describe("toSessionError", () => {
expect(eligible.map(SessionRunnerRetry.isRetryable)).toEqual([true, true, true])
expect(ineligible.map(SessionRunnerRetry.isRetryable)).toEqual([false, false, false, false, false, false, false])
})
test("retries transport failures only when delivery is absent or not sent", () => {
const retryable = [
llm(new TransportReason({ message: "http transport" })),
llm(new TransportReason({ message: "connect failed", delivery: "not-sent", phase: "connect" })),
]
const ineligible = [
llm(new TransportReason({ message: "send uncertain", delivery: "ambiguous", phase: "send" })),
llm(new TransportReason({ message: "response interrupted", delivery: "accepted", phase: "receive" })),
llm(
new TransportReason({
message: "continuation rejected",
delivery: "rejected",
recovery: "retry-full",
phase: "receive",
}),
),
]
expect(retryable.map(SessionRunnerRetry.isRetryable)).toEqual([true, true])
expect(ineligible.map(SessionRunnerRetry.isRetryable)).toEqual([false, false, false])
})
})
+41
View File
@@ -553,6 +553,47 @@ describe("Session.prompt", () => {
}),
)
it.effect("reconciles an exact retry from the promoted message without admission history", () =>
Effect.gen(function* () {
yield* setup
const session = yield* Session.Service
const bus = yield* Bus.Service
const { db } = yield* Database.Service
const input = { sessionID, id: messageID, text: "Fix the failing tests", resume: false }
const first = yield* session.prompt(input)
yield* SessionPending.promote(db, bus, sessionID, "steer")
yield* db
.delete(EventTable)
.where(eq(EventTable.aggregate_id, sessionID))
.run()
.pipe(Effect.orDie)
const retried = yield* session.prompt(input)
expect(retried).toMatchObject({ id: first.id, type: "user", data: { text: first.data.text } })
expect(yield* session.messages({ sessionID })).toMatchObject([
{ id: messageID, type: "user", text: "Fix the failing tests" },
])
}),
)
it.effect("ignores delivery when retrying a promoted message", () =>
Effect.gen(function* () {
yield* setup
const session = yield* Session.Service
const bus = yield* Bus.Service
const { db } = yield* Database.Service
const input = { sessionID, id: messageID, text: "Fix the failing tests", resume: false }
yield* session.prompt(input)
yield* SessionPending.promote(db, bus, sessionID, "steer")
const retried = yield* session.prompt({ ...input, delivery: "queue" })
expect(retried).toMatchObject({ id: messageID, type: "user", data: { text: input.text } })
expect(yield* admitted(messageID)).toBeUndefined()
}),
)
it.effect("wakes execution when an exact prompt retry recovers a committed message", () =>
Effect.gen(function* () {
yield* setup
+22 -12
View File
@@ -1180,7 +1180,7 @@ describe("SessionRunnerLLM", () => {
}),
)
it.effect("forks instruction values at the selected message instead of the parent's latest state", () =>
it.effect("seeds a fork with the parent's newest instruction values", () =>
Effect.gen(function* () {
const session = yield* setup
yield* runPrompt(session, "First")
@@ -1197,14 +1197,16 @@ describe("SessionRunnerLLM", () => {
.where(eq(InstructionStateTable.session_id, forked.id))
.get(),
).toMatchObject({
initial_values: { "test/context": Instructions.hash("Changed context") },
current_values: { "test/context": Instructions.hash("Changed context") },
initial_values: { "test/context": Instructions.hash("Latest context") },
current_values: { "test/context": Instructions.hash("Latest context") },
})
yield* session.prompt({ sessionID: forked.id, text: "Forked", resume: false })
yield* session.resume(forked.id)
expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([defaultSystem, "Changed context"])
expect(systemTexts(requests.at(-1)!)).toContain("Latest context")
expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([defaultSystem, "Latest context"])
// Copied history keeps the frozen chronological update; no new update is emitted.
expect(systemTexts(requests.at(-1)!)).toContain("Changed context")
expect(systemTexts(requests.at(-1)!)).not.toContain("Latest context")
const { db } = yield* Database.Service
const bus = yield* Bus.Service
@@ -1263,7 +1265,7 @@ describe("SessionRunnerLLM", () => {
}),
)
it.effect("rebuilds a missing instruction cache without admitting another delta", () =>
it.effect("re-establishes a fresh baseline when instruction state is missing", () =>
Effect.gen(function* () {
const session = yield* setup
const { db } = yield* Database.Service
@@ -1277,13 +1279,15 @@ describe("SessionRunnerLLM", () => {
expect(requests).toHaveLength(1)
expect(requests[0]?.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
expect(messageRoles(requests[0])).toEqual(["user", "user"])
// The projected row is authoritative: a missing row admits a fresh baseline
// instead of rebuilding from durable events.
expect(
yield* db
.select({ id: EventTable.id })
.select({ data: EventTable.data })
.from(EventTable)
.where(eq(EventTable.type, "session.instructions.updated.2"))
.all(),
).toHaveLength(1)
).toHaveLength(2)
expect(yield* db.select().from(InstructionStateTable).get()).toMatchObject({
initial_values: { "test/context": Instructions.hash("Initial context") },
current_values: { "test/context": Instructions.hash("Initial context") },
@@ -1310,7 +1314,10 @@ describe("SessionRunnerLLM", () => {
])
expect(messageRoles(requests[1])).toEqual(["user", "system", "user"])
expect(requests[1]?.messages.at(1)?.content).toEqual([{ type: "text", text: "Changed context" }])
expect(yield* session.messages({ sessionID })).toHaveLength(2)
// The chronological update is a durable client-visible system message.
const messages = yield* session.messages({ sessionID })
expect(messages).toHaveLength(3)
expect(messages[1]).toMatchObject({ type: "system", text: "Changed context" })
const { db } = yield* Database.Service
const updates = yield* db
.select({ data: EventTable.data })
@@ -1327,9 +1334,10 @@ describe("SessionRunnerLLM", () => {
expect(updates[1]?.data).toEqual({
sessionID,
delta: { "test/context": Instructions.hash("Changed context") },
text: "Changed context",
})
yield* replaySessionProjection(sessionID)
expect(yield* session.messages({ sessionID })).toHaveLength(2)
expect(yield* session.messages({ sessionID })).toHaveLength(3)
}),
)
@@ -1596,7 +1604,7 @@ describe("SessionRunnerLLM", () => {
expect(requests[1]?.messages.at(1)?.content).toEqual([
{ type: "text", text: "System context source removed: test/context" },
])
expect(yield* session.messages({ sessionID })).toHaveLength(2)
expect(yield* session.messages({ sessionID })).toHaveLength(3)
}),
)
@@ -1708,12 +1716,14 @@ describe("SessionRunnerLLM", () => {
expect(requests[2]?.messages.filter((message) => message.role === "system")).toHaveLength(2)
expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([
"user",
"system",
"user",
"model-switched",
"system",
"user",
])
yield* replaySessionProjection(sessionID)
expect(yield* session.messages({ sessionID })).toHaveLength(4)
expect(yield* session.messages({ sessionID })).toHaveLength(6)
yield* runPrompt(session, "Fourth")
}),
)
+5
View File
@@ -186,6 +186,11 @@ export const InstructionsUpdated = Event.durable({
schema: {
...Base,
delta: Instruction.Delta,
/**
* The rendered chronological update shown to the model, frozen at emit time.
* Absent for the initial baseline observation and for deltas that render empty.
*/
text: Schema.String.pipe(optional),
},
})
export type InstructionsUpdated = typeof InstructionsUpdated.Type
@@ -147,12 +147,12 @@ async function renderDiffViewer(vcsDiff: unknown[], height = 20, initialRoute?:
const config = createTuiResolvedConfig()
const transport = createFetch((url) => {
if (url.pathname !== "/api/vcs/diff") return
if (fail) return json({ message: "boom" }, { status: 500 })
vcsDiffInput = {
location: { directory: url.searchParams.get("location[directory]") },
mode: url.searchParams.get("mode"),
context: url.searchParams.get("context"),
}
if (fail) return json({ message: "boom" }, { status: 500 })
return json({
location: { directory: "/repo/session", project: { id: "project-1", directory: "/repo/session" } },
data: vcsDiff,
@@ -238,6 +238,7 @@ async function renderDiffViewer(vcsDiff: unknown[], height = 20, initialRoute?:
const app = await testRender(() => <Harness />, { width: 80, height })
await waitForCommand(app, commands, "diff.close")
await app.waitFor(() => vcsDiffInput !== undefined)
return {
app,
commands,