fix(core): reduce memory peaks in long sessions

This commit is contained in:
starptech 2026-06-03 16:26:06 +02:00
parent e707e416ed
commit fe28f3b1dd
21 changed files with 2331 additions and 104 deletions

View file

@ -0,0 +1,8 @@
ALTER TABLE `part` ADD `data_model` text;
--> statement-breakpoint
UPDATE part
SET data_model = json_remove(data, '$.state.metadata')
WHERE length(CAST(data AS BLOB)) > 65536
AND json_extract(data, '$.type') = 'tool'
AND json_extract(data, '$.state.status') = 'completed'
AND length(CAST(json_extract(data, '$.state.metadata') AS BLOB)) > 65536;

File diff suppressed because it is too large Load diff

View file

@ -8,12 +8,15 @@ import { Flag } from "../flag/flag"
import { isAbsolute, join } from "path"
import { DatabaseMigration } from "./migration"
import { InstallationChannel } from "../installation/version"
import { Sqlite } from "./sqlite"
const makeDatabase = EffectDrizzleSqlite.makeWithDefaults()
type DatabaseShape = Effect.Success<typeof makeDatabase>
export interface Interface {
db: DatabaseShape
// Lazy property getters cannot yield an Effect.
sync: Sqlite.DrizzleClient
}
export class Service extends Context.Service<Service, Interface>()("@opencode/v2/storage/Database") {}
@ -22,6 +25,7 @@ export const layer = Layer.effect(
Service,
Effect.gen(function* () {
const db = yield* makeDatabase
const sync = yield* Sqlite.Drizzle
yield* db.run("PRAGMA journal_mode = WAL")
yield* db.run("PRAGMA synchronous = NORMAL")
@ -31,7 +35,7 @@ export const layer = Layer.effect(
yield* db.run("PRAGMA wal_checkpoint(PASSIVE)")
yield* DatabaseMigration.apply(db)
return { db }
return { db, sync }
}).pipe(Effect.orDie),
)

View file

@ -27,5 +27,6 @@ export const migrations = (
import("./migration/20260601202201_amazing_prowler"),
import("./migration/20260602002951_lowly_union_jack"),
import("./migration/20260602182828_add_project_directories"),
import("./migration/20260603120017_warm_guardsmen"),
])
).map((module) => module.default) satisfies DatabaseMigration.Migration[]

View file

@ -0,0 +1,21 @@
import { Effect } from "effect"
import type { DatabaseMigration } from "../migration"
export default {
id: "20260603120017_warm_guardsmen",
up(tx) {
return Effect.gen(function* () {
yield* tx.run(`ALTER TABLE \`part\` ADD \`data_model\` text;`)
// Keep canonical history intact while avoiding prompt-time decoding. This
// one-time transactional backfill may briefly grow the WAL on large stores.
yield* tx.run(`
UPDATE part
SET data_model = json_remove(data, '$.state.metadata')
WHERE length(CAST(data AS BLOB)) > 65536
AND json_extract(data, '$.type') = 'tool'
AND json_extract(data, '$.state.status') = 'completed'
AND length(CAST(json_extract(data, '$.state.metadata') AS BLOB)) > 65536
`)
})
},
} satisfies DatabaseMigration.Migration

View file

@ -0,0 +1,22 @@
export * as SessionPartModelData from "./model-data"
import type { SessionV1 } from "../v1/session"
type V1PartData<Data extends SessionV1.Part = SessionV1.Part> = Data extends SessionV1.Part
? Omit<Data, "id" | "sessionID" | "messageID">
: never
export const THRESHOLD = 64 * 1024
// Strip UI-only metadata only when the stored prompt projection benefits.
export function create(data: unknown): V1PartData | null {
if (!data || typeof data !== "object") return null
if (!("type" in data) || data.type !== "tool") return null
if (!("state" in data) || !data.state || typeof data.state !== "object") return null
if (!("status" in data.state) || data.state.status !== "completed") return null
if (!("metadata" in data.state)) return null
const metadata = JSON.stringify(data.state.metadata)
if (!metadata || Buffer.byteLength(metadata) <= THRESHOLD) return null
const { metadata: _, ...state } = data.state
return { ...data, state } as V1PartData
}

View file

@ -11,6 +11,7 @@ import { SessionMessage } from "./message"
import { SessionMessageUpdater } from "./message-updater"
import { MessageTable, PartTable, SessionMessageTable, SessionTable } from "./sql"
import type { DeepMutable } from "../schema"
import { SessionPartModelData } from "./model-data"
type DatabaseService = Database.Interface["db"]
@ -309,7 +310,7 @@ export const layer = Layer.effectDiscard(
yield* events.project(SessionV1.Event.MessageRemoved, (event) =>
Effect.gen(function* () {
const rows = yield* db
.select()
.select({ session_id: PartTable.session_id, data: PartTable.data })
.from(PartTable)
.where(and(eq(PartTable.message_id, event.data.messageID), eq(PartTable.session_id, event.data.sessionID)))
.all()
@ -328,7 +329,7 @@ export const layer = Layer.effectDiscard(
yield* events.project(SessionV1.Event.PartRemoved, (event) =>
Effect.gen(function* () {
const row = yield* db
.select()
.select({ session_id: PartTable.session_id, data: PartTable.data })
.from(PartTable)
.where(and(eq(PartTable.id, event.data.partID), eq(PartTable.session_id, event.data.sessionID)))
.get()
@ -348,11 +349,17 @@ export const layer = Layer.effectDiscard(
const messageID = event.data.part.messageID
const sessionID = event.data.part.sessionID
const data = partData(event.data.part)
const row = yield* db.select().from(PartTable).where(eq(PartTable.id, id)).get().pipe(Effect.orDie)
const data_model = SessionPartModelData.create(data)
const row = yield* db
.select({ session_id: PartTable.session_id, data: PartTable.data })
.from(PartTable)
.where(eq(PartTable.id, id))
.get()
.pipe(Effect.orDie)
yield* db
.insert(PartTable)
.values({ id, message_id: messageID, session_id: sessionID, time_created: event.data.time, data })
.onConflictDoUpdate({ target: PartTable.id, set: { data } })
.values({ id, message_id: messageID, session_id: sessionID, time_created: event.data.time, data, data_model })
.onConflictDoUpdate({ target: PartTable.id, set: { data, data_model } })
.run()
.pipe(Effect.orDie)
const previous = row && usage(row.data)

View file

@ -85,6 +85,8 @@ export const PartTable = sqliteTable(
session_id: text().$type<SessionSchema.ID>().notNull(),
...Timestamps,
data: text({ mode: "json" }).notNull().$type<V1PartData>(),
// Derived prompt projection; data remains canonical.
data_model: text({ mode: "json" }).$type<V1PartData>(),
},
(table) => [
index("part_message_id_id_idx").on(table.message_id, table.id),

View file

@ -14,6 +14,7 @@ import { AbsolutePath } from "@opencode-ai/core/schema"
import { SessionSchema } from "@opencode-ai/core/session/schema"
import { SessionTable } from "@opencode-ai/core/session/sql"
import sessionMetadataMigration from "@opencode-ai/core/database/migration/20260511173437_session-metadata"
import partModelDataMigration from "@opencode-ai/core/database/migration/20260603120017_warm_guardsmen"
import type { SqlClient as SqlClientService } from "effect/unstable/sql/SqlClient"
const run = <A, E>(effect: Effect.Effect<A, E, SqlClientService>) =>
@ -43,7 +44,7 @@ describe("DatabaseMigration", () => {
expect(yield* db.get(sql`SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'session'`)).toEqual({
name: "session",
})
expect(yield* db.get(sql`SELECT count(*) as count FROM migration`)).toEqual({ count: 25 })
expect(yield* db.get(sql`SELECT count(*) as count FROM migration`)).toEqual({ count: 26 })
}),
)
})
@ -77,6 +78,30 @@ describe("DatabaseMigration", () => {
)
})
test("backfills lightweight model data for oversized completed tool metadata", async () => {
await run(
Effect.gen(function* () {
const db = yield* makeDb
yield* db.run(sql`CREATE TABLE part (id text PRIMARY KEY, data text NOT NULL)`)
const large = JSON.stringify({ type: "tool", state: { status: "completed", metadata: { diff: "x".repeat(70_000) } } })
const unicode = JSON.stringify({ type: "tool", state: { status: "completed", metadata: { diff: "😀".repeat(20_000) } } })
const small = JSON.stringify({ type: "tool", state: { status: "completed", metadata: { diff: "small" } } })
yield* db.run(sql`INSERT INTO part (id, data) VALUES (${"large"}, ${large}), (${"unicode"}, ${unicode}), (${"small"}, ${small})`)
yield* DatabaseMigration.applyOnly(db, [partModelDataMigration])
expect(yield* db.get(sql`SELECT data, data_model FROM part WHERE id = ${"large"}`)).toEqual({
data: large,
data_model: JSON.stringify({ type: "tool", state: { status: "completed" } }),
})
expect(yield* db.get(sql`SELECT data_model FROM part WHERE id = ${"small"}`)).toEqual({ data_model: null })
expect(yield* db.get(sql`SELECT data_model FROM part WHERE id = ${"unicode"}`)).toEqual({
data_model: JSON.stringify({ type: "tool", state: { status: "completed" } }),
})
}),
)
})
test("normalizes Windows storage paths and leaves POSIX paths untouched", async () => {
await run(
Effect.gen(function* () {

View file

@ -12,6 +12,7 @@ import path from "path"
import { FSUtil } from "@opencode-ai/core/fs-util"
import { Effect, Schema } from "effect"
import type { InstanceContext } from "@/project/instance-context"
import { SessionPartModelData } from "@opencode-ai/core/session/model-data"
const decodeMessageInfo = Schema.decodeUnknownSync(SessionV1.Info)
const decodePart = Schema.decodeUnknownSync(SessionV1.Part)
@ -212,6 +213,8 @@ const runImport = Effect.fn("Cli.import.body")(function* (file: string, ctx: Ins
message_id: messageID,
session_id: row.id,
data: partData,
// Keep imported history on the same prompt path as live updates.
data_model: SessionPartModelData.create(partData),
})
.onConflictDoNothing()
.run()

View file

@ -319,18 +319,20 @@ export function Prompt(props: PromptProps) {
let promptPartTypeId = 0
const event = useEvent()
event.on(TuiEvent.PromptAppend.type, (evt, { workspace }) => {
if (workspace !== project.workspace.current()) return
if (!input || input.isDestroyed) return
input.insertText(evt.properties.text)
setTimeout(() => {
// setTimeout is a workaround and needs to be addressed properly
onCleanup(
event.on(TuiEvent.PromptAppend.type, (evt, { workspace }) => {
if (workspace !== project.workspace.current()) return
if (!input || input.isDestroyed) return
input.getLayoutNode().markDirty()
input.gotoBufferEnd()
renderer.requestRender()
}, 0)
})
input.insertText(evt.properties.text)
setTimeout(() => {
// setTimeout is a workaround and needs to be addressed properly
if (!input || input.isDestroyed) return
input.getLayoutNode().markDirty()
input.gotoBufferEnd()
renderer.requestRender()
}, 0)
}),
)
createEffect(() => {
if (!input || input.isDestroyed) return

View file

@ -7,6 +7,7 @@ import {
For,
Match,
on,
onCleanup,
onMount,
Show,
Switch,
@ -291,21 +292,23 @@ export function Session() {
})
let lastSwitch: string | undefined = undefined
event.on("message.part.updated", (evt) => {
const part = evt.properties.part
if (part.type !== "tool") return
if (part.sessionID !== route.sessionID) return
if (part.state.status !== "completed") return
if (part.id === lastSwitch) return
onCleanup(
event.on("message.part.updated", (evt) => {
const part = evt.properties.part
if (part.type !== "tool") return
if (part.sessionID !== route.sessionID) return
if (part.state.status !== "completed") return
if (part.id === lastSwitch) return
if (part.tool === "plan_exit") {
local.agent.set("build")
lastSwitch = part.id
} else if (part.tool === "plan_enter") {
local.agent.set("plan")
lastSwitch = part.id
}
})
if (part.tool === "plan_exit") {
local.agent.set("build")
lastSwitch = part.id
} else if (part.tool === "plan_enter") {
local.agent.set("plan")
lastSwitch = part.id
}
}),
)
let seeded = false
let scroll: ScrollBoxRenderable
@ -321,25 +324,27 @@ export function Session() {
const dialog = useDialog()
const renderer = useRenderer()
event.on("session.status", (evt) => {
if (evt.properties.sessionID !== route.sessionID) return
if (evt.properties.status.type !== "retry") return
if (!evt.properties.status.action) return
if (dialog.stack.length > 0) return
onCleanup(
event.on("session.status", (evt) => {
if (evt.properties.sessionID !== route.sessionID) return
if (evt.properties.status.type !== "retry") return
if (!evt.properties.status.action) return
if (dialog.stack.length > 0) return
const keys = goUpsellKeys(evt.properties.status.action)
if (!keys) return
const keys = goUpsellKeys(evt.properties.status.action)
if (!keys) return
const seen = kv.get(keys.lastSeenAt)
if (typeof seen === "number" && Date.now() - seen < GO_UPSELL_WINDOW) return
const seen = kv.get(keys.lastSeenAt)
if (typeof seen === "number" && Date.now() - seen < GO_UPSELL_WINDOW) return
if (kv.get(keys.dontShow)) return
if (kv.get(keys.dontShow)) return
void DialogRetryAction.show(dialog, evt.properties.status.action).then((dontShowAgain) => {
if (dontShowAgain) kv.set(keys.dontShow, true)
kv.set(keys.lastSeenAt, Date.now())
})
})
void DialogRetryAction.show(dialog, evt.properties.status.action).then((dontShowAgain) => {
if (dontShowAgain) kv.set(keys.dontShow, true)
kv.set(keys.lastSeenAt, Date.now())
})
}),
)
const exit = useExit()

View file

@ -402,7 +402,7 @@ export const layer = Layer.effect(
{ context: [], prompt: undefined },
)
const nextPrompt = compacting.prompt ?? buildPrompt({ previousSummary, context: compacting.context })
const msgs = structuredClone(selected.head)
const msgs = MessageV2.cloneForTransform(selected.head)
yield* plugin.trigger("experimental.chat.messages.transform", {}, { messages: msgs })
const modelMessages = yield* MessageV2.toModelMessagesEffect(msgs, model, {
stripMedia: true,

View file

@ -29,6 +29,7 @@ import { eq } from "drizzle-orm"
import { inArray } from "drizzle-orm"
import { lt } from "drizzle-orm"
import { or } from "drizzle-orm"
import { sql } from "drizzle-orm"
import { MessageTable, PartTable, SessionTable } from "@opencode-ai/core/session/sql"
import { ProviderError } from "@/provider/error"
import { iife } from "@/util/iife"
@ -47,6 +48,7 @@ interface FetchDecompressionError extends Error {
}
export const SYNTHETIC_ATTACHMENT_PROMPT = "Attached media from tool result:"
const COMPACTED_TOOL_OUTPUT = "[Old tool result content cleared]"
export { isMedia }
function truncateToolOutput(text: string, maxChars?: number) {
@ -96,7 +98,7 @@ const info = (row: typeof MessageTable.$inferSelect) =>
sessionID: row.session_id,
}) as Info
const part = (row: typeof PartTable.$inferSelect) =>
const part = (row: Pick<typeof PartTable.$inferSelect, "id" | "session_id" | "message_id" | "data">) =>
({
...row.data,
id: row.id,
@ -107,20 +109,41 @@ const part = (row: typeof PartTable.$inferSelect) =>
const older = (row: Cursor) =>
or(lt(MessageTable.time_created, row.time), and(eq(MessageTable.time_created, row.time), lt(MessageTable.id, row.id)))
function hydrate(db: Database.Interface["db"], rows: (typeof MessageTable.$inferSelect)[]) {
function hydrate(
db: Database.Interface["db"],
rows: (typeof MessageTable.$inferSelect)[],
sync?: Database.Interface["sync"],
) {
const ids = rows.map((row) => row.id)
const partByMessage = new Map<string, Part[]>()
return Effect.gen(function* () {
if (ids.length > 0) {
const partRows = yield* db
.select()
.from(PartTable)
.where(inArray(PartTable.message_id, ids))
.orderBy(PartTable.message_id, PartTable.id)
.all()
.pipe(Effect.orDie)
const partRows = sync
? yield* db
.select({
id: PartTable.id,
message_id: PartTable.message_id,
session_id: PartTable.session_id,
// Keep oversized UI metadata available to extensions without decoding it for every prompt.
data: sql`coalesce(${PartTable.data_model}, ${PartTable.data})`
.mapWith(PartTable.data)
.as("data"),
lazy_metadata: sql<number>`${PartTable.data_model} IS NOT NULL`.as("lazy_metadata"),
})
.from(PartTable)
.where(inArray(PartTable.message_id, ids))
.orderBy(PartTable.message_id, PartTable.id)
.all()
.pipe(Effect.orDie)
: yield* db
.select({ id: PartTable.id, message_id: PartTable.message_id, session_id: PartTable.session_id, data: PartTable.data })
.from(PartTable)
.where(inArray(PartTable.message_id, ids))
.orderBy(PartTable.message_id, PartTable.id)
.all()
.pipe(Effect.orDie)
for (const row of partRows) {
const next = part(row)
const next = "lazy_metadata" in row && row.lazy_metadata && sync ? lazyMetadata(sync, part(row)) : part(row)
const list = partByMessage.get(row.message_id)
if (list) list.push(next)
else partByMessage.set(row.message_id, [next])
@ -134,6 +157,67 @@ function hydrate(db: Database.Interface["db"], rows: (typeof MessageTable.$infer
})
}
function lazyMetadata(db: Database.Interface["sync"], part: Part) {
if (part.type !== "tool" || part.state.status !== "completed") return part
defineLazyMetadata(part.state, () => {
// Prompt history is short-lived. Resolve against canonical storage only when
// an extension explicitly reads metadata instead of decoding it every turn.
const row = db
.select({ metadata: sql<string>`json_extract(${PartTable.data}, '$.state.metadata')` })
.from(PartTable)
.where(eq(PartTable.id, part.id))
.get()
return row?.metadata ? JSON.parse(row.metadata) : {}
})
return part
}
function defineLazyMetadata(state: SessionV1.ToolStateCompleted, load: () => Record<string, any>) {
let metadata: Record<string, any> | undefined
Object.defineProperty(state, "metadata", {
configurable: true,
enumerable: true,
get() {
metadata ??= load()
return metadata
},
set(value: Record<string, any>) {
metadata = value
},
})
}
export function cloneForTransform(input: WithParts[]) {
// structuredClone resolves getters, so mask and restore lazy metadata.
const lazy = new Map<string, () => Record<string, any>>()
const masked = input.map((msg) => ({
...msg,
parts: msg.parts.map((part) => {
if (part.type !== "tool" || part.state.status !== "completed") return part
const load = Object.getOwnPropertyDescriptor(part.state, "metadata")?.get
if (!load) return part
lazy.set(part.id, () => structuredClone(load.call(part.state)))
return {
...part,
state: Object.fromEntries(
Object.keys(part.state)
.filter((key) => key !== "metadata")
.map((key) => [key, Reflect.get(part.state, key)]),
),
} as Part
}),
}))
const result = structuredClone(masked) as WithParts[]
for (const msg of result) {
for (const part of msg.parts) {
if (part.type !== "tool" || part.state.status !== "completed") continue
const load = lazy.get(part.id)
if (load) defineLazyMetadata(part.state, load)
}
}
return result
}
function providerMeta(metadata: Record<string, any> | undefined) {
if (!metadata) return undefined
const { providerExecuted: _, ...rest } = metadata
@ -302,7 +386,7 @@ export const toModelMessagesEffect = Effect.fnUntraced(function* (
toolNames.add(part.tool)
if (part.state.status === "completed") {
const outputText = part.state.time.compacted
? "[Old tool result content cleared]"
? COMPACTED_TOOL_OUTPUT
: truncateToolOutput(part.state.output, options?.toolOutputMaxChars)
const attachments = part.state.time.compacted || options?.stripMedia ? [] : (part.state.attachments ?? [])
@ -433,12 +517,17 @@ export function toModelMessages(
return Effect.runPromise(toModelMessagesEffect(input, model, options).pipe(Effect.provide(EffectLogger.layer)))
}
export const page = Effect.fn("MessageV2.page")(function* (input: {
type PageInput = {
sessionID: SessionID
limit: number
before?: string
}) {
const { db } = yield* Database.Service
}
const pageWithOptions = Effect.fnUntraced(function* (
input: PageInput & { lazyCompletedToolMetadata?: boolean },
) {
const database = yield* Database.Service
const db = database.db
const before = input.before ? cursor.decode(input.before) : undefined
const where = before
? and(eq(MessageTable.session_id, input.sessionID), older(before))
@ -467,7 +556,7 @@ export const page = Effect.fn("MessageV2.page")(function* (input: {
const more = rows.length > input.limit
const slice = more ? rows.slice(0, input.limit) : rows
const items = yield* hydrate(db, slice)
const items = yield* hydrate(db, slice, input.lazyCompletedToolMetadata ? database.sync : undefined)
items.reverse()
const tail = slice.at(-1)
return {
@ -477,6 +566,10 @@ export const page = Effect.fn("MessageV2.page")(function* (input: {
}
})
export const page = Effect.fn("MessageV2.page")(function* (input: PageInput) {
return yield* pageWithOptions(input)
})
export function stream(sessionID: SessionID) {
const size = 50
return Effect.gen(function* () {
@ -504,7 +597,7 @@ export function parts(messageID: MessageID) {
return Effect.gen(function* () {
const { db } = yield* Database.Service
const rows = yield* db
.select()
.select({ id: PartTable.id, message_id: PartTable.message_id, session_id: PartTable.session_id, data: PartTable.data })
.from(PartTable)
.where(eq(PartTable.message_id, messageID))
.orderBy(PartTable.id)
@ -529,29 +622,57 @@ export const get = Effect.fn("MessageV2.get")(function* (input: { sessionID: Ses
}
})
export const related = Effect.fn("MessageV2.related")(function* (input: { sessionID: SessionID; messageID: MessageID }) {
const { db } = yield* Database.Service
return yield* db
.select()
.from(MessageTable)
.where(
and(
eq(MessageTable.session_id, input.sessionID),
or(
eq(MessageTable.id, input.messageID),
and(
sql`json_extract(${MessageTable.data}, '$.role') = 'assistant'`,
sql`json_extract(${MessageTable.data}, '$.parentID') = ${input.messageID}`,
),
),
),
)
.orderBy(MessageTable.time_created, MessageTable.id)
.all()
.pipe(Effect.flatMap((rows) => hydrate(db, rows)), Effect.orDie)
})
export function filterCompacted(msgs: Iterable<WithParts>) {
const result = [] as WithParts[]
const completed = new Set<string>()
let retain: MessageID | undefined
const state = compactedState()
for (const msg of msgs) {
result.push(msg)
if (retain) {
if (msg.info.id === retain) break
continue
}
if (msg.info.role === "user" && completed.has(msg.info.id)) {
const part = msg.parts.find((item): item is CompactionPart => item.type === "compaction")
if (!part) continue
if (!part.tail_start_id) break
retain = part.tail_start_id
if (msg.info.id === retain) break
continue
}
if (msg.info.role === "user" && completed.has(msg.info.id) && msg.parts.some((part) => part.type === "compaction"))
break
if (msg.info.role === "assistant" && msg.info.summary && msg.info.finish && !msg.info.error)
completed.add(msg.info.parentID)
if (reachedCompactedBoundary(state, msg)) break
}
return reorderCompacted(result)
}
function compactedState() {
return { completed: new Set<string>(), retain: undefined as MessageID | undefined }
}
function reachedCompactedBoundary(state: ReturnType<typeof compactedState>, msg: WithParts) {
if (state.retain) return msg.info.id === state.retain
if (msg.info.role === "user" && state.completed.has(msg.info.id)) {
const part = msg.parts.find((item): item is CompactionPart => item.type === "compaction")
if (!part) return false
if (!part.tail_start_id) return true
state.retain = part.tail_start_id
return msg.info.id === state.retain
}
if (msg.info.role === "assistant" && msg.info.summary && msg.info.finish && !msg.info.error)
state.completed.add(msg.info.parentID)
return false
}
function reorderCompacted(result: WithParts[]) {
result.reverse()
const compactionIndex = result.findLastIndex(
(msg) =>
@ -583,7 +704,28 @@ export function filterCompacted(msgs: Iterable<WithParts>) {
}
export const filterCompactedEffect = Effect.fnUntraced(function* (sessionID: SessionID) {
return filterCompacted(yield* stream(sessionID))
// Stop paging once older compacted history would be discarded anyway.
const size = 50
const result = [] as WithParts[]
const state = compactedState()
let before: string | undefined
while (true) {
const next = yield* pageWithOptions({ sessionID, limit: size, before, lazyCompletedToolMetadata: true }).pipe(
Effect.catchIf(NotFoundError.isInstance, () =>
Effect.succeed({ items: [] as WithParts[], more: false, cursor: undefined }),
),
)
if (next.items.length === 0) break
for (let i = next.items.length - 1; i >= 0; i--) {
const item = next.items[i]
if (!item) continue
result.push(item)
if (reachedCompactedBoundary(state, item)) return reorderCompacted(result)
}
if (!next.more || !next.cursor) break
before = next.cursor
}
return reorderCompacted(result)
})
// filterCompacted reorders messages for model consumption

View file

@ -720,7 +720,7 @@ export const layer: Layer.Layer<
const getPart: Interface["getPart"] = Effect.fn("Session.getPart")(function* (input) {
const row = yield* db
.select()
.select({ id: PartTable.id, session_id: PartTable.session_id, message_id: PartTable.message_id, data: PartTable.data })
.from(PartTable)
.where(
and(

View file

@ -5,6 +5,9 @@ import { Snapshot } from "@/snapshot"
import { Session } from "./session"
import { SessionID, MessageID } from "./schema"
import { Config } from "@/config/config"
import { MessageV2 } from "./message-v2"
import { Database } from "@opencode-ai/core/database/database"
import { NotFoundError } from "@/storage/storage"
function unquoteGitPath(input: string) {
if (!input.startsWith('"')) return input
@ -77,6 +80,7 @@ export const layer = Layer.effect(
const snapshot = yield* Snapshot.Service
const events = yield* EventV2Bridge.Service
const config = yield* Config.Service
const database = yield* Database.Service
const computeDiff = Effect.fn("SessionSummary.computeDiff")(function* (input: { messages: SessionV1.WithParts[] }) {
let from: string | undefined
@ -112,12 +116,7 @@ export const layer = Layer.effect(
})
yield* events.publish(Session.Event.Diff, { sessionID: input.sessionID, diff: [] })
if ((yield* config.get()).snapshot === false) return
const all = yield* sessions.messages({ sessionID: input.sessionID }).pipe(Effect.orDie)
if (!all.length) return
const messages = all.filter(
(m) => m.info.id === input.messageID || (m.info.role === "assistant" && m.info.parentID === input.messageID),
)
const messages = yield* MessageV2.related(input).pipe(Effect.provideService(Database.Service, database))
const target = messages.find((m) => m.info.id === input.messageID)
if (!target || target.info.role !== "user") return
const msgDiffs = yield* computeDiff({ messages })
@ -127,8 +126,11 @@ export const layer = Layer.effect(
const diff = Effect.fn("SessionSummary.diff")(function* (input: { sessionID: SessionID; messageID?: MessageID }) {
if (!input.messageID) return []
const message = (yield* sessions.messages({ sessionID: input.sessionID }).pipe(Effect.orDie)).find(
(item) => item.info.id === input.messageID,
const message = yield* MessageV2.get({ sessionID: input.sessionID, messageID: input.messageID }).pipe(
Effect.provideService(Database.Service, database),
Effect.catchIf(NotFoundError.isInstance, () =>
sessions.get(input.sessionID).pipe(Effect.orDie, Effect.as(undefined)),
),
)
if (!message || message.info.role !== "user") return []
const diffs = message.info.summary?.diffs ?? []
@ -150,6 +152,7 @@ export const defaultLayer = Layer.suspend(() =>
Layer.provide(Snapshot.defaultLayer),
Layer.provide(EventV2Bridge.defaultLayer),
Layer.provide(Config.defaultLayer),
Layer.provide(Database.defaultLayer),
),
)

View file

@ -170,14 +170,16 @@ export const layer = Layer.effect(
def: D,
fn: (data: EventV2.Data<D>) => Effect.Effect<void, unknown>,
) =>
events.listen((event) => {
if (event.type !== def.type || event.location?.directory !== _ctx.directory) return Effect.void
return fn(event.data as EventV2.Data<D>).pipe(
Effect.catchCause((cause) =>
Effect.sync(() => log.error("share subscriber failed", { type: def.type, cause })),
),
)
})
events
.listen((event) => {
if (event.type !== def.type || event.location?.directory !== _ctx.directory) return Effect.void
return fn(event.data as EventV2.Data<D>).pipe(
Effect.catchCause((cause) =>
Effect.sync(() => log.error("share subscriber failed", { type: def.type, cause })),
),
)
})
.pipe(Effect.tap((unsubscribe) => Effect.addFinalizer(() => unsubscribe)))
yield* watch(Session.Event.Updated, (data) =>
Effect.gen(function* () {

View file

@ -382,6 +382,20 @@ function autocontinue(enabled: boolean) {
})
}
function messagesTransform(inspect: (messages: SessionV1.WithParts[]) => void) {
return Layer.mock(Plugin.Service)({
trigger: <Name extends string, Input, Output>(name: Name, _input: Input, output: Output) => {
if (name !== "experimental.chat.messages.transform") return Effect.succeed(output)
return Effect.sync(() => {
inspect((output as { messages: SessionV1.WithParts[] }).messages)
return output
})
},
list: () => Effect.succeed([]),
init: () => Effect.void,
})
}
describe("session.compaction.isOverflow", () => {
it.live(
"returns true when token count exceeds usable context",
@ -682,8 +696,18 @@ describe("session.compaction.prune", () => {
status: "completed",
input: {},
output: "x".repeat(200_000),
attachments: [
{
id: PartID.ascending(),
messageID: b.id,
sessionID: info.id,
type: "file",
mime: "text/plain",
url: "data:text/plain;base64,eA==",
},
],
title: "done",
metadata: {},
metadata: { output: "x".repeat(200_000), description: "done" },
time: { start: Date.now(), end: Date.now() },
},
})
@ -713,9 +737,34 @@ describe("session.compaction.prune", () => {
expect(part?.state.status).toBe("completed")
if (part?.type === "tool" && part.state.status === "completed") {
expect(part.state.time.compacted).toBeNumber()
expect(part.state.output).toHaveLength(200_000)
expect(part.state.metadata.output).toHaveLength(200_000)
expect(part.state.attachments).toHaveLength(1)
const compacted = (yield* MessageV2.filterCompactedEffect(info.id))
.flatMap((msg) => msg.parts)
.find((part) => part.type === "tool")
expect(compacted?.type).toBe("tool")
if (compacted?.type === "tool" && compacted.state.status === "completed") {
expect(Object.getOwnPropertyDescriptor(compacted.state, "metadata")?.get).toBeFunction()
expect(compacted.state.output).toHaveLength(200_000)
expect(compacted.state.metadata.output).toHaveLength(200_000)
expect(compacted.state.attachments).toHaveLength(1)
expect(JSON.parse(JSON.stringify(compacted.state)).metadata.output).toHaveLength(200_000)
part.state.metadata = { description: "small" }
yield* ssn.updatePart(part)
const small = (yield* MessageV2.filterCompactedEffect(info.id))
.flatMap((msg) => msg.parts)
.find((item) => item.id === part.id)
expect(small?.type).toBe("tool")
if (small?.type === "tool" && small.state.status === "completed") {
expect(Object.getOwnPropertyDescriptor(small.state, "metadata")?.get).toBeUndefined()
expect(small.state.metadata).toEqual({ description: "small" })
}
}
}
}),
{
config: {
compaction: { prune: true },
@ -815,6 +864,52 @@ describe("session.compaction.prune", () => {
})
describe("session.compaction.process", () => {
itCompaction.instance(
"keeps oversized tool metadata lazy through the plugin transform clone",
() => {
let lazy = false
return Effect.gen(function* () {
const test = yield* TestInstance
const ssn = yield* SessionNs.Service
const session = yield* ssn.create({})
const first = yield* createUserMessage(session.id, "first")
const assistant = yield* createAssistantMessage(session.id, first.id, test.directory)
yield* ssn.updatePart({
id: PartID.ascending(),
messageID: assistant.id,
sessionID: session.id,
type: "tool",
tool: "apply_patch",
callID: "call_test",
state: {
status: "completed",
input: {},
output: "done",
title: "done",
metadata: { diff: "x".repeat(70_000) },
time: { start: Date.now(), end: Date.now() },
},
})
const parent = yield* createUserMessage(session.id, "compact")
const msgs = yield* MessageV2.filterCompactedEffect(session.id)
yield* SessionCompaction.use.process({ parentID: parent.id, messages: msgs, sessionID: session.id, auto: false })
expect(lazy).toBe(true)
}).pipe(
withCompaction({
config: cfg({ tail_turns: 0 }),
plugin: messagesTransform((messages) => {
const part = messages.flatMap((msg) => msg.parts).find((part) => part.type === "tool")
if (part?.type === "tool" && part.state.status === "completed") {
lazy = Object.getOwnPropertyDescriptor(part.state, "metadata")?.get !== undefined
}
}),
}),
)
},
)
it.instance(
"throws when parent is not a user message",
Effect.gen(function* () {

View file

@ -553,6 +553,32 @@ describe("MessageV2.get", () => {
)
})
describe("MessageV2.related", () => {
it.instance("returns one user turn and its assistant replies", () =>
withSession(({ session, sessionID }) =>
Effect.gen(function* () {
yield* addUser(sessionID, "older")
const user = yield* addUser(sessionID, "target")
const first = yield* addAssistant(sessionID, user)
const second = yield* addAssistant(sessionID, user)
yield* session.updatePart({
id: PartID.ascending(),
sessionID,
messageID: second,
type: "text",
text: "response",
})
yield* addUser(sessionID, "newer")
const result = yield* MessageV2.related({ sessionID, messageID: user })
expect(result.map((item) => item.info.id)).toEqual([user, first, second])
expect(result[2].parts).toHaveLength(1)
}),
),
)
})
describe("Session.messages", () => {
it.instance("returns all messages in chronological order across pages", () =>
withSession(({ session, sessionID }) =>
@ -645,6 +671,31 @@ describe("MessageV2.filterCompacted", () => {
),
)
it.instance("effect stops at compaction boundary beyond the first page", () =>
withSession(({ session, sessionID }) =>
Effect.gen(function* () {
yield* fill(sessionID, 55, (index: number) => index)
const compact = yield* addUser(sessionID)
yield* addCompactionPart(sessionID, compact)
const summary = yield* addAssistant(sessionID, compact, { summary: true, finish: "end_turn" })
yield* session.updatePart({
id: PartID.ascending(),
sessionID,
messageID: summary,
type: "text",
text: "summary",
})
yield* fill(sessionID, 55)
const expected = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
const result = yield* MessageV2.filterCompactedEffect(sessionID)
expect(result).toEqual(expected)
expect(result).toHaveLength(57)
}),
),
)
it.live("handles empty iterable", () =>
Effect.sync(() => {
const result = MessageV2.filterCompacted([])
@ -665,6 +716,22 @@ describe("MessageV2.filterCompacted", () => {
),
)
it.instance("does not break on summary without matching compaction part", () =>
withSession(({ sessionID }) =>
Effect.gen(function* () {
const user = yield* addUser(sessionID, "hello")
yield* addAssistant(sessionID, user, { summary: true, finish: "end_turn" })
yield* addUser(sessionID, "world")
const expected = MessageV2.filterCompacted(yield* MessageV2.stream(sessionID))
const result = yield* MessageV2.filterCompactedEffect(sessionID)
expect(result).toEqual(expected)
expect(result).toHaveLength(3)
}),
),
)
it.instance("skips assistant with error even if marked as summary", () =>
withSession(({ sessionID }) =>
Effect.gen(function* () {

View file

@ -19,6 +19,7 @@ import { eq } from "drizzle-orm"
import { provideTmpdirInstance } from "../fixture/fixture"
import { resetDatabase } from "../fixture/db"
import { testEffect } from "../lib/effect"
import { disposeInstance } from "@/effect/instance-registry"
const env = Layer.mergeAll(
Session.defaultLayer,
@ -40,10 +41,10 @@ const json = (req: Parameters<typeof HttpClientResponse.fromWeb>[0], body: unkno
const none = HttpClient.make(() => Effect.die("unexpected http call"))
function live(client: HttpClient.HttpClient) {
function live(client: HttpClient.HttpClient, events = EventV2Bridge.defaultLayer) {
const http = Layer.succeed(HttpClient.HttpClient, client)
return ShareNext.layer.pipe(
Layer.provide(EventV2Bridge.defaultLayer),
Layer.provide(events),
Layer.provide(Account.layer.pipe(Layer.provide(AccountRepo.defaultLayer), Layer.provide(http))),
Layer.provide(Config.defaultLayer),
Layer.provide(Database.defaultLayer),
@ -101,6 +102,34 @@ beforeEach(async () => {
})
describe("ShareNext", () => {
it.live("unsubscribes event listeners when the instance is disposed", () =>
provideTmpdirInstance((directory) => {
let active = 0
const events = Layer.mock(EventV2Bridge.Service, {
listen: () =>
Effect.sync(() => {
active++
return Effect.sync(() => {
active--
})
}),
})
return Effect.gen(function* () {
const share = yield* ShareNext.Service
let peak = 0
for (let index = 0; index < 20; index++) {
yield* share.init()
peak = Math.max(peak, active)
yield* Effect.promise(() => disposeInstance(directory))
}
expect(peak).toBe(5)
expect(active).toBe(0)
}).pipe(Effect.provide(live(none, events)))
}),
)
it.live("request uses legacy share API without active org account", () =>
provideTmpdirInstance(
() =>

34
script/install-local Executable file
View file

@ -0,0 +1,34 @@
#!/usr/bin/env bash
set -euo pipefail
ROOT=$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")/.." && pwd)
PACKAGE_DIR="$ROOT/packages/opencode"
INSTALLED=${OPENCODE_INSTALL_PATH:-$(command -v opencode || true)}
if [[ -z "$INSTALLED" ]]; then
echo "opencode is not installed or is not on PATH" >&2
echo "Set OPENCODE_INSTALL_PATH to the executable you want to replace." >&2
exit 1
fi
INSTALLED=$(realpath "$INSTALLED")
TARGET=$(bun --eval 'process.stdout.write(`opencode-${process.platform}-${process.arch}`)')
BINARY="$PACKAGE_DIR/dist/$TARGET/bin/opencode"
bun run --cwd "$PACKAGE_DIR" build --single "$@"
if [[ ! -x "$BINARY" ]]; then
echo "Expected build output does not exist: $BINARY" >&2
exit 1
fi
if [[ ! -e "$INSTALLED.backup" ]]; then
cp -p "$INSTALLED" "$INSTALLED.backup"
echo "Backed up $INSTALLED to $INSTALLED.backup"
fi
install -m 755 "$BINARY" "$INSTALLED"
echo "Installed local build at $INSTALLED"
"$INSTALLED" --version