diff --git a/.scratch/batch-b/issues/01-u1-fork-rollback.md b/.scratch/batch-b/issues/01-u1-fork-rollback.md index fa0ae38611..6c5c8af49d 100644 --- a/.scratch/batch-b/issues/01-u1-fork-rollback.md +++ b/.scratch/batch-b/issues/01-u1-fork-rollback.md @@ -6,10 +6,18 @@ **Evidence:** `.scratch/batch-b/evidence.md#u-1--session-fork-嵌套事务回滚` **Branch:** `test/fork-rollback` **Blocked by:** None -**Status:** ready-for-agent +**Status:** done -- [ ] 复用 `packages/opencode/test/session/fork-batch.test.ts` 的真实 SQLite fixture;不写 adapter-only 替代测试 -- [ ] 至少一个 message/part 发布完成后再确定性失败,旧实现若违约时测试能红 -- [ ] 目标 session 无复制出的 durable events 与 projections,源 session 不变;Session Created 可保留 -- [ ] 若红灯暴露生产缺陷,只做本契约所需的最小修复 -- [ ] 在 `packages/opencode` 运行目标测试与 `bun typecheck`,结果附入票据 +- [x] 复用 `packages/opencode/test/session/fork-batch.test.ts` 的真实 SQLite fixture;不写 adapter-only 替代测试 +- [x] 至少一个 message/part 发布完成后再确定性失败,旧实现若违约时测试能红 +- [x] 目标 session 无复制出的 durable events 与 projections,源 session 不变;Session Created 可保留 +- [x] 若红灯暴露生产缺陷,只做本契约所需的最小修复 +- [x] 在 `packages/opencode` 运行目标测试与 `bun typecheck`,结果附入票据 + +## 验证证据 + +- 基线:`dev@8f8465753b6517b3deeb6ad37002263d3da287fe`;分支:`test/fork-rollback`。 +- 失败注入:真实 SQLite trigger 在第二条复制 part 的 durable event insert 上执行 `RAISE(ABORT)`;此前 3 个嵌套 publication 已释放 savepoint。 +- mutation 红灯:临时移除外层复制事务后,目标 projection 残留 2 条 message(第一条含已复制 part),新增用例 0 pass / 1 fail;mutation 未保留。 +- `cd packages/opencode && bun test test/session/fork-batch.test.ts`:3 pass,0 fail,52 expect。 +- `cd packages/opencode && bun typecheck`:`tsgo --noEmit`,exit 0;现有生产实现满足契约,无生产代码修改。 diff --git a/packages/opencode/test/session/fork-batch.test.ts b/packages/opencode/test/session/fork-batch.test.ts index 51633bd9dd..ab7225d4a9 100644 --- a/packages/opencode/test/session/fork-batch.test.ts +++ b/packages/opencode/test/session/fork-batch.test.ts @@ -2,15 +2,17 @@ import { describe, expect } from "bun:test" import { Database as BunDatabase, type SQLQueryBindings } from "bun:sqlite" import { Database } from "@opencode-ai/core/database/database" import { EventV2 } from "@opencode-ai/core/event" +import { EventTable } from "@opencode-ai/core/event/sql" import { SessionProjector } from "@opencode-ai/core/session/projector" import { SessionV1 } from "@opencode-ai/core/v1/session" import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner" -import { Context, Effect, Fiber, Layer, Scope, Semaphore, Stream } from "effect" +import { Context, Effect, Exit, Fiber, Layer, Scope, Semaphore, Stream } from "effect" import * as Client from "effect/unstable/sql/SqlClient" import type { Connection } from "effect/unstable/sql/SqlConnection" import { SqlError, classifySqliteError } from "effect/unstable/sql/SqlError" import * as Statement from "effect/unstable/sql/Statement" import * as Reactivity from "effect/unstable/reactivity/Reactivity" +import { eq, sql } from "drizzle-orm" import { Session as SessionNs } from "@/session/session" import { MessageID, PartID } from "../../src/session/schema" import { testInstanceStoreLayer } from "../fixture/fixture" @@ -24,9 +26,11 @@ interface SqlCounter { begins: number commits: number savepoints: number + releases: number + savepointRollbacks: number } -const counter: SqlCounter = { begins: 0, commits: 0, savepoints: 0 } +const counter: SqlCounter = { begins: 0, commits: 0, savepoints: 0, releases: 0, savepointRollbacks: 0 } // The Database layer's sqlite client is a closed graph (its native provider // cannot be overridden from outside), so this test builds its own SqlClient @@ -51,6 +55,8 @@ const countingClientLayer = Layer.effect( if (/^\s*begin\b/i.test(sql)) counter.begins++ else if (/^\s*commit\b/i.test(sql)) counter.commits++ else if (/^\s*savepoint\b/i.test(sql)) counter.savepoints++ + else if (/^\s*release\b/i.test(sql)) counter.releases++ + else if (/^\s*rollback to\b/i.test(sql)) counter.savepointRollbacks++ return target.query(sql) } } @@ -65,7 +71,9 @@ const countingClientLayer = Layer.effect( // @ts-ignore bun-types missing safeIntegers method statement.safeIntegers(Context.get(fiber.context, Client.SafeIntegers)) try { - return Effect.succeed((statement.all(...(params as SQLQueryBindings[])) ?? []) as Array>) + return Effect.succeed( + (statement.all(...(params as SQLQueryBindings[])) ?? []) as Array>, + ) } catch (cause) { return Effect.fail( new SqlError({ @@ -125,15 +133,14 @@ const countingClientLayer = Layer.effect( }), ) -const dbLayer = Database.layer.pipe( - Layer.provide(countingClientLayer.pipe(Layer.provide(Reactivity.layer))), -) +const dbLayer = Database.layer.pipe(Layer.provide(countingClientLayer.pipe(Layer.provide(Reactivity.layer)))) const eventV2Layer = EventV2.layer.pipe(Layer.provide(dbLayer)) const eventV2BridgeLayer = EventV2Bridge.layer.pipe(Layer.provide(eventV2Layer)) const projectorLayer = SessionProjector.layer.pipe(Layer.provide(eventV2Layer), Layer.provide(dbLayer)) const it = testEffect( Layer.mergeAll( + dbLayer, SessionNs.layer.pipe( Layer.provide(Storage.defaultLayer), Layer.provide(dbLayer), @@ -250,6 +257,8 @@ describe("Session.fork", () => { counter.begins = 0 counter.commits = 0 counter.savepoints = 0 + counter.releases = 0 + counter.savepointRollbacks = 0 const fork = yield* Effect.acquireRelease(session.fork({ sessionID: original.id }), (info) => session.remove(info.id).pipe(Effect.ignore), @@ -268,4 +277,65 @@ describe("Session.fork", () => { for (const msg of target) expect(msg.parts.length).toBe(2) }), ) + + it.instance("fork rolls back copied durable events and projections when a nested publication fails", () => + Effect.gen(function* () { + const database = yield* Database.Service + const session = yield* SessionNs.Service + const original = yield* Effect.acquireRelease(session.create({ title: "fork-source" }), (info) => + session.remove(info.id).pipe(Effect.ignore), + ) + + const firstMessageID = MessageID.ascending() + const failingMessageID = MessageID.ascending() + yield* session.updateMessage(userInfo(original.id, firstMessageID)) + yield* session.updatePart(textPart(original.id, firstMessageID, "copied before failure")) + yield* session.updateMessage(userInfo(original.id, failingMessageID)) + yield* session.updatePart(textPart(original.id, failingMessageID, "force fork copy failure")) + const sourceBefore = yield* session.messages({ sessionID: original.id }) + + yield* database.db + .run( + sql` + CREATE TRIGGER reject_fork_copy + BEFORE INSERT ON event + WHEN NEW.type = 'message.part.updated.1' + AND json_extract(NEW.data, '$.part.text') = 'force fork copy failure' + BEGIN + SELECT RAISE(ABORT, 'forced fork copy failure'); + END + `, + ) + .pipe(Effect.orDie) + + counter.begins = 0 + counter.commits = 0 + counter.savepoints = 0 + counter.releases = 0 + counter.savepointRollbacks = 0 + + const forkExit = yield* Effect.exit(session.fork({ sessionID: original.id })) + + expect(Exit.isFailure(forkExit)).toBe(true) + + const forkedSessions = (yield* session.list()).filter((info) => info.id !== original.id) + expect(forkedSessions).toHaveLength(1) + const forked = forkedSessions[0] + if (!forked) return + yield* Effect.addFinalizer(() => session.remove(forked.id).pipe(Effect.ignore)) + + expect(yield* session.messages({ sessionID: forked.id })).toEqual([]) + expect(yield* session.messages({ sessionID: original.id })).toEqual(sourceBefore) + + const durableEvents = yield* database.db + .select({ type: EventTable.type }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, forked.id)) + expect(durableEvents).toEqual([{ type: "session.created.1" }]) + + expect(counter.savepoints).toBe(4) + expect(counter.releases).toBe(4) + expect(counter.savepointRollbacks).toBe(1) + }), + ) })