diff --git a/packages/core/src/session/projector.ts b/packages/core/src/session/projector.ts index 792067017d14..5d1b48ca7f21 100644 --- a/packages/core/src/session/projector.ts +++ b/packages/core/src/session/projector.ts @@ -76,8 +76,30 @@ function sessionRow(info: SessionV1.SessionInfo): typeof SessionTable.$inferInse function messageData( info: (typeof SessionV1.Event.MessageUpdated.Type)["data"]["info"], + current?: typeof MessageTable.$inferSelect.data, ): typeof MessageTable.$inferInsert.data { const { id: _, sessionID: __, ...rest } = info + const summary = current?.summary + if ( + info.role === "user" && + info.summary?.diffs !== undefined && + current?.role === "user" && + typeof summary === "object" && + summary.diffs !== undefined && + info.summary.diffs.length === summary.diffs.length && + info.summary.diffs.every((item, index) => { + const stored = summary.diffs[index] + return ( + item.patch === undefined && + item.file === stored?.file && + item.additions === stored?.additions && + item.deletions === stored?.deletions && + item.status === stored?.status + ) + }) + ) { + return { ...rest, summary: { ...info.summary, diffs: summary.diffs } } as DeepMutable + } return rest as DeepMutable } @@ -262,7 +284,16 @@ const layer = Layer.effectDiscard( const time_created = event.data.info.time.created const id = event.data.info.id const sessionID = event.data.info.sessionID - const data = messageData(event.data.info) + const current = + event.data.info.role === "user" && event.data.info.summary?.diffs?.some((item) => item.patch === undefined) + ? yield* db + .select({ data: MessageTable.data }) + .from(MessageTable) + .where(and(eq(MessageTable.id, id), eq(MessageTable.session_id, sessionID))) + .get() + .pipe(Effect.orDie) + : undefined + const data = messageData(event.data.info, current?.data) yield* db .insert(MessageTable) .values({ id, session_id: sessionID, time_created, data }) diff --git a/packages/opencode/src/session/session.ts b/packages/opencode/src/session/session.ts index a2a91cd47b5e..c4417e7459d4 100644 --- a/packages/opencode/src/session/session.ts +++ b/packages/opencode/src/session/session.ts @@ -26,7 +26,7 @@ import { inArray } from "drizzle-orm" import { lt } from "drizzle-orm" import { or } from "drizzle-orm" import type { SQL } from "drizzle-orm" -import { PartTable, SessionTable } from "@opencode-ai/core/session/sql" +import { MessageTable, PartTable, SessionTable } from "@opencode-ai/core/session/sql" import { ProjectTable } from "@opencode-ai/core/project/sql" import { MessageV2 } from "./message-v2" import type { InstanceContext } from "../project/instance-context" @@ -628,7 +628,28 @@ const layer: Layer.Layer< const updateMessage = (msg: T): Effect.Effect => Effect.gen(function* () { - yield* events.publish(SessionV1.Event.MessageUpdated, { sessionID: msg.sessionID, info: msg }) + const existing = + msg.role === "user" && msg.summary?.diffs !== undefined + ? yield* db + .select({ + matches: sql`json_extract(${MessageTable.data}, '$.summary.diffs') = json(${JSON.stringify(msg.summary.diffs)})`, + }) + .from(MessageTable) + .where(and(eq(MessageTable.id, msg.id), eq(MessageTable.session_id, msg.sessionID))) + .get() + .pipe(Effect.orDie) + : undefined + const info = + existing?.matches === 1 && msg.role === "user" && msg.summary?.diffs !== undefined + ? { + ...msg, + summary: { + ...msg.summary, + diffs: msg.summary.diffs.map(({ patch: _, ...item }) => item), + }, + } + : msg + yield* events.publish(SessionV1.Event.MessageUpdated, { sessionID: msg.sessionID, info }) return msg }).pipe(Effect.withSpan("Session.updateMessage")) diff --git a/packages/opencode/test/server/session-diff-missing-patch.test.ts b/packages/opencode/test/server/session-diff-missing-patch.test.ts index d2a4211ff1a9..d75f40c19591 100644 --- a/packages/opencode/test/server/session-diff-missing-patch.test.ts +++ b/packages/opencode/test/server/session-diff-missing-patch.test.ts @@ -13,19 +13,34 @@ import { afterEach, describe, expect } from "bun:test" import { LayerNode } from "@opencode-ai/core/effect/layer-node" import { Effect, Layer } from "effect" +import path from "path" import { SessionPaths } from "@/server/routes/instance/httpapi/groups/session" import { Session } from "@/session/session" import { Storage } from "@/storage/storage" import { SessionV1 } from "@opencode-ai/core/v1/session" -import { MessageID } from "@/session/schema" +import { MessageID, PartID } from "@/session/schema" +import { SessionSummary } from "@/session/summary" +import { Snapshot } from "@/snapshot" import { ProviderV2 } from "@opencode-ai/core/provider" import { ModelV2 } from "@opencode-ai/core/model" +import { Database } from "@opencode-ai/core/database/database" +import { EventV2 } from "@opencode-ai/core/event" +import { EventSequenceTable, EventTable } from "@opencode-ai/core/event/sql" +import { MessageTable, SessionTable } from "@opencode-ai/core/session/sql" +import { and, eq, sql } from "drizzle-orm" import { resetDatabase } from "../fixture/db" import { disposeAllInstances, TestInstance } from "../fixture/fixture" import { testEffect } from "../lib/effect" import { httpApiLayer, requestInDirectory } from "./httpapi-layer" -const it = testEffect(Layer.mergeAll(LayerNode.compile(LayerNode.group([Session.node, Storage.node])), httpApiLayer)) +const it = testEffect( + Layer.mergeAll( + LayerNode.compile( + LayerNode.group([Database.node, EventV2.node, Session.node, SessionSummary.node, Snapshot.node, Storage.node]), + ), + httpApiLayer, + ), +) afterEach(async () => { await disposeAllInstances() @@ -94,4 +109,179 @@ describe("session diff with missing patch (#26574)", () => { }), { git: true, config: { formatter: false, lsp: false } }, ) + + it.instance( + "keeps stored turn diffs while compacting later message update events", + () => + Effect.gen(function* () { + const test = yield* TestInstance + const session = yield* withSession({ title: "compact-turn-diff" }) + const messageID = MessageID.ascending() + const diff = { + file: "turn.ts", + additions: 1, + deletions: 0, + patch: "x".repeat(262_144), + status: "modified" as const, + } + const message = { + id: messageID, + sessionID: session.id, + role: "user" as const, + time: { created: Date.now() }, + agent: "build", + model: { providerID: ProviderV2.ID.make("test"), modelID: ModelV2.ID.make("model") }, + summary: { diffs: [diff] }, + } satisfies SessionV1.User + yield* Session.use.updateMessage(message) + yield* Session.use.updateMessage({ ...message, tools: { read: true } }) + + const { db } = yield* Database.Service + const events = yield* db + .select({ bytes: sql`length(${EventTable.data})` }) + .from(EventTable) + .where( + and( + eq(EventTable.aggregate_id, session.id), + eq(EventTable.type, "message.updated.1"), + ), + ) + .orderBy(EventTable.seq) + .all() + .pipe(Effect.orDie) + + expect(events).toHaveLength(2) + expect(events[0]?.bytes).toBeGreaterThan(diff.patch.length) + expect(events[1]?.bytes).toBeLessThan(1_000) + + const response = yield* requestInDirectory( + `${pathFor(SessionPaths.diff, { sessionID: session.id })}?messageID=${messageID}`, + test.directory, + ) + + expect(response.status).toBe(200) + expect(yield* response.json).toEqual([diff]) + }), + { git: true, config: { formatter: false, lsp: false } }, + ) + + it.instance( + "persists changed turn diffs through a fresh event replay", + () => + Effect.gen(function* () { + const test = yield* TestInstance + const session = yield* withSession({ title: "changed-turn-diff" }) + const messageID = MessageID.ascending() + const message = { + id: messageID, + sessionID: session.id, + role: "user" as const, + time: { created: Date.now() }, + agent: "build", + model: { providerID: ProviderV2.ID.make("test"), modelID: ModelV2.ID.make("model") }, + } satisfies SessionV1.User + yield* Session.use.updateMessage(message) + const assistant = yield* Session.use.updateMessage({ + id: MessageID.ascending(), + sessionID: session.id, + role: "assistant", + time: { created: Date.now() }, + parentID: messageID, + agent: "build", + modelID: ModelV2.ID.make("model"), + providerID: ProviderV2.ID.make("test"), + mode: "build", + path: { cwd: test.directory, root: test.directory }, + cost: 0, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + } satisfies SessionV1.Assistant) + const snapshot = yield* Snapshot.Service + const summary = yield* SessionSummary.Service + const start = yield* snapshot.track() + if (!start) return yield* Effect.die("expected initial snapshot") + yield* Effect.promise(() => Bun.write(path.join(test.directory, "first.ts"), "first")) + const first = yield* snapshot.track() + if (!first) return yield* Effect.die("expected first snapshot") + yield* Session.use.updatePart({ + id: PartID.ascending(), + messageID: assistant.id, + sessionID: session.id, + type: "step-start", + snapshot: start, + }) + yield* Session.use.updatePart({ + id: PartID.ascending(), + messageID: assistant.id, + sessionID: session.id, + type: "step-finish", + reason: "stop", + snapshot: first, + cost: 0, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + }) + yield* summary.summarize({ sessionID: session.id, messageID }) + + yield* Effect.promise(() => Bun.write(path.join(test.directory, "second.ts"), "second")) + const second = yield* snapshot.track() + if (!second) return yield* Effect.die("expected second snapshot") + yield* Session.use.updatePart({ + id: PartID.ascending(), + messageID: assistant.id, + sessionID: session.id, + type: "step-finish", + reason: "stop", + snapshot: second, + cost: 0, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + }) + yield* summary.summarize({ sessionID: session.id, messageID }) + + const { db } = yield* Database.Service + const before = yield* db + .select({ data: MessageTable.data }) + .from(MessageTable) + .where(eq(MessageTable.id, messageID)) + .get() + .pipe(Effect.orDie) + expect(before?.data).toMatchObject({ + role: "user", + summary: { diffs: [{ file: "first.ts" }, { file: "second.ts" }] }, + }) + + const events = yield* db + .select() + .from(EventTable) + .where(eq(EventTable.aggregate_id, session.id)) + .orderBy(EventTable.seq) + .all() + .pipe(Effect.orDie) + yield* db.delete(MessageTable).where(eq(MessageTable.session_id, session.id)).run().pipe(Effect.orDie) + yield* db.delete(EventTable).where(eq(EventTable.aggregate_id, session.id)).run().pipe(Effect.orDie) + yield* db.delete(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, session.id)).run().pipe(Effect.orDie) + yield* db.delete(SessionTable).where(eq(SessionTable.id, session.id)).run().pipe(Effect.orDie) + + const event = yield* EventV2.Service + yield* event.replayAll( + events.map((item) => ({ + id: item.id, + type: item.type, + data: item.data, + seq: item.seq, + aggregateID: item.aggregate_id, + })), + ) + + const after = yield* db + .select({ data: MessageTable.data }) + .from(MessageTable) + .where(eq(MessageTable.id, messageID)) + .get() + .pipe(Effect.orDie) + expect(after?.data).toMatchObject({ + role: "user", + summary: { diffs: [{ file: "first.ts" }, { file: "second.ts" }] }, + }) + }), + { git: true, config: { formatter: false, lsp: false } }, + ) })