Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 14 additions & 6 deletions .scratch/batch-b/issues/01-u1-fork-rollback.md
Original file line number Diff line number Diff line change
Expand Up @@ -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;现有生产实现满足契约,无生产代码修改。
82 changes: 76 additions & 6 deletions packages/opencode/test/session/fork-batch.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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
Expand All @@ -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)
}
}
Expand All @@ -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<Record<string, unknown>>)
return Effect.succeed(
(statement.all(...(params as SQLQueryBindings[])) ?? []) as Array<Record<string, unknown>>,
)
} catch (cause) {
return Effect.fail(
new SqlError({
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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),
Expand All @@ -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)
}),
)
})
Loading