| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213 |
- import { describe, expect, test } from "bun:test"
- import { SqliteClient } from "@effect/sql-sqlite-bun"
- import { EffectDrizzleSqlite } from "@opencode-ai/effect-drizzle-sqlite"
- import { Database } from "@opencode-ai/core/database/database"
- import { DatabaseMigration } from "@opencode-ai/core/database/migration"
- import { V1Migration } from "@opencode-ai/core/database/v1-migration"
- import { SessionMessage } from "@opencode-ai/core/session/message"
- import { SessionSchema } from "@opencode-ai/core/session/schema"
- import { SessionTable } from "@opencode-ai/core/session/sql"
- import { Project } from "@opencode-ai/core/project"
- import { ProjectTable } from "@opencode-ai/core/project/sql"
- import { AbsolutePath } from "@opencode-ai/core/schema"
- import { Global } from "@opencode-ai/util/global"
- import { Effect, Layer, Logger, Schedule, Schema, Scope } from "effect"
- import { eq, sql } from "drizzle-orm"
- import type { SqlClient } from "effect/unstable/sql/SqlClient"
- import { tmpdir } from "./fixture/tmpdir"
- import path from "path"
- const makeDb = EffectDrizzleSqlite.makeWithDefaults()
- const run = <A, E>(effect: Effect.Effect<A, E, SqlClient | Scope.Scope | Global.Service>) =>
- Effect.runPromise(
- Effect.scoped(
- effect.pipe(
- Effect.provideService(Global.Service, Global.make({ data: path.join(process.cwd(), ".test-data") })),
- Effect.provide(SqliteClient.layer({ filename: ":memory:", disableWAL: true })),
- ),
- ),
- )
- const session = (
- overrides: Partial<V1Migration.TransformInput["session"]> = {},
- ): V1Migration.TransformInput["session"] => ({
- id: SessionSchema.ID.make("ses_test"),
- project_id: Project.ID.global,
- workspace_id: null,
- parent_id: null,
- fork_session_id: null,
- fork_boundary: null,
- slug: "test",
- directory: "/tmp/test",
- path: null,
- title: "Test",
- version: "1",
- share_url: null,
- summary_additions: null,
- summary_deletions: null,
- summary_files: null,
- summary_diffs: null,
- metadata: null,
- cost: 99,
- tokens_input: 99,
- tokens_output: 99,
- tokens_reasoning: 99,
- tokens_cache_read: 99,
- tokens_cache_write: 99,
- revert: null,
- permission: null,
- agent: null,
- model: null,
- time_created: 1,
- time_updated: 2,
- time_compacting: 3,
- time_archived: null,
- time_suspended: null,
- ...overrides,
- })
- const user = (id: string, overrides: Record<string, unknown> = {}, time = 10): V1Migration.SourceMessage => ({
- id,
- session_id: "ses_test",
- time_created: time,
- time_updated: time + 1,
- data: JSON.stringify({
- role: "user",
- time: { created: time },
- agent: "build",
- model: { providerID: "provider", modelID: "model" },
- ...overrides,
- }),
- })
- const assistant = (
- id: string,
- parentID: string,
- overrides: Record<string, unknown> = {},
- time = 20,
- ): V1Migration.SourceMessage => ({
- id,
- session_id: "ses_test",
- time_created: time,
- time_updated: time + 1,
- data: JSON.stringify({
- role: "assistant",
- time: { created: time, completed: time + 5 },
- parentID,
- modelID: "model",
- providerID: "provider",
- mode: "build",
- agent: "build",
- path: { cwd: "/tmp/test", root: "/tmp/test" },
- cost: 1,
- tokens: { total: 10, input: 2, output: 3, reasoning: 4, cache: { read: 5, write: 6 } },
- ...overrides,
- }),
- })
- const part = (id: string, messageID: string, data: Record<string, unknown> | string): V1Migration.SourcePart => ({
- id,
- message_id: messageID,
- session_id: "ses_test",
- time_created: 1,
- time_updated: 2,
- data: typeof data === "string" ? data : JSON.stringify(data),
- })
- const transform = (messages: V1Migration.SourceMessage[], parts: V1Migration.SourcePart[], info = session()) => {
- const result = V1Migration.transformSession({ session: info, messages, parts })
- result.messages.forEach((row) =>
- Schema.decodeUnknownSync(SessionMessage.Info)({ id: row.id, type: row.type, ...row.data }),
- )
- return result
- }
- describe("V1Migration.transformSession", () => {
- test("maps ordinary user text, agents, ignored fields, order, and timestamps", () => {
- const message = user("msg_000000000001aaaaaaaaaaaaaa", {
- system: "discard",
- tools: { read: false },
- format: { type: "text" },
- summary: { title: "discard", diffs: [] },
- })
- const result = transform(
- [message],
- [
- part("prt_4", message.id, { type: "agent", name: "review" }),
- part("prt_2", message.id, { type: "text", text: "ignored", ignored: true }),
- part("prt_3", message.id, { type: "agent", name: "build", source: { value: "@build", start: 2, end: 8 } }),
- part("prt_1", message.id, { type: "text", text: "first" }),
- part("prt_5", message.id, { type: "text", text: "second" }),
- ],
- )
- expect(result.messages).toEqual([
- {
- id: message.id,
- session_id: "ses_test",
- type: "user",
- seq: 0,
- time_created: 10,
- time_updated: 11,
- data: {
- text: "first\n\nsecond",
- agents: [{ name: "build", mention: { text: "@build", start: 2, end: 8 } }, { name: "review" }],
- time: { created: 10 },
- },
- },
- ])
- expect(result.watermark).toBe(0)
- })
- test("maps embedded files and deterministic placeholders without external IO", () => {
- const message = user("msg_000000000002aaaaaaaaaaaaaa")
- const result = transform(
- [message],
- [
- part("prt_1", message.id, { type: "text", text: "prompt" }),
- part("prt_2", message.id, {
- type: "file",
- mime: "text/plain",
- filename: "inline.txt",
- url: "data:text/plain,hello%20world",
- }),
- part("prt_3", message.id, {
- type: "file",
- mime: "text/plain",
- url: "data:text/plain;base64,aGk=",
- source: {
- type: "resource",
- clientName: "mcp",
- uri: "resource://item",
- text: { value: "item", start: 1, end: 5 },
- },
- }),
- part("prt_4", message.id, {
- type: "file",
- mime: "text/plain",
- filename: "named.txt",
- url: "file:///tmp/named.txt",
- }),
- part("prt_5", message.id, { type: "file", mime: "application/octet-stream", url: "https://example.test/raw" }),
- ],
- )
- expect(result.messages[0].data).toEqual({
- text: "prompt\n\n[Attachment unavailable after migration: named.txt (text/plain)]\n\n[Attachment unavailable after migration: https://example.test/raw (application/octet-stream)]",
- files: [
- { data: "aGVsbG8gd29ybGQ=", mime: "text/plain", source: { type: "inline" }, name: "inline.txt" },
- {
- data: "aGk=",
- mime: "text/plain",
- source: { type: "uri", uri: "resource://item" },
- mention: { text: "item", start: 1, end: 5 },
- },
- ],
- time: { created: 10 },
- })
- })
- test("maps attachment-only source variants and uses resource URIs as unavailable labels", () => {
- const message = user("msg_000000000049aaaaaaaaaaaaaa")
- const result = transform(
- [message],
- [
- part("prt_1", message.id, {
- type: "file",
- mime: "text/plain",
- url: "data:text/plain,file",
- source: { type: "file", path: "/tmp/a", text: { value: "a", start: 0, end: 1 } },
- }),
- part("prt_2", message.id, {
- type: "file",
- mime: "text/plain",
- url: "data:text/plain,symbol",
- source: {
- type: "symbol",
- path: "/tmp/a",
- range: { start: { line: 0, character: 0 }, end: { line: 0, character: 1 } },
- name: "a",
- kind: 1,
- text: { value: "symbol", start: 2, end: 8 },
- },
- }),
- part("prt_3", message.id, {
- type: "file",
- mime: "application/json",
- url: "https://example.test/resource",
- source: {
- type: "resource",
- clientName: "mcp",
- uri: "resource://fallback",
- text: { value: "resource", start: 0, end: 8 },
- },
- }),
- ],
- )
- expect(result.messages[0].data).toEqual({
- text: "[Attachment unavailable after migration: resource://fallback (application/json)]",
- files: [
- {
- data: "ZmlsZQ==",
- mime: "text/plain",
- source: { type: "inline" },
- mention: { text: "a", start: 0, end: 1 },
- },
- {
- data: "c3ltYm9s",
- mime: "text/plain",
- source: { type: "inline" },
- mention: { text: "symbol", start: 2, end: 8 },
- },
- ],
- time: { created: 10 },
- })
- })
- test("splits mixed synthetic content deterministically and preserves adjacency", () => {
- const all = user("msg_000000000003aaaaaaaaaaaaaa", {}, 1)
- const mixed = user("msg_000000000004aaaaaaaaaaaaaa", {}, 2)
- const later = user("msg_000000000005aaaaaaaaaaaaaa", {}, 3)
- const parts = [
- part("prt_1", all.id, { type: "text", text: "context", synthetic: true }),
- part("prt_2", mixed.id, { type: "text", text: "hello" }),
- part("prt_3", mixed.id, { type: "text", text: "hidden", synthetic: true }),
- part("prt_4", mixed.id, { type: "text", text: "ignored", synthetic: true, ignored: true }),
- part("prt_5", later.id, { type: "text", text: "later" }),
- ]
- const first = transform([later, mixed, all], parts)
- const second = transform([later, mixed, all], parts)
- expect(first.messages.map((row) => row.type)).toEqual(["synthetic", "user", "synthetic", "user"])
- expect(first.messages[0].id).toBe(all.id)
- expect(first.messages[2].id).not.toBe(mixed.id)
- expect(first.messages[2].id).toBe("msg_000000000004ST20Jh98kGJtjL")
- expect(first.messages[2].id.slice(0, 16)).toBe(mixed.id.slice(0, 16))
- expect(first.messages[2].id).toMatch(/^msg_[0-9A-Za-z]{26}$/)
- expect(first.messages[2].id).toBe(second.messages[2].id)
- expect(first.messages[2].data).toEqual({ text: "hidden", time: { created: 2 } })
- const collision = user(first.messages[2].id, {}, 4)
- const collided = transform(
- [all, mixed, later, collision],
- [...parts, part("prt_6", collision.id, { type: "text", text: "collision" })],
- )
- expect(collided.messages[2].id).not.toBe(first.messages[2].id)
- expect(collided.messages[2].id.slice(0, 16)).toBe(mixed.id.slice(0, 16))
- })
- test("preserves assistant content, model, usage, finish, snapshots, and marker filtering", () => {
- const parent = user("msg_000000000006aaaaaaaaaaaaaa")
- const message = assistant("msg_000000000007aaaaaaaaaaaaaa", parent.id, {
- variant: "fast",
- structured: { discard: true },
- finish: "stop",
- cost: 2.5,
- })
- const result = transform(
- [parent, message],
- [
- part("prt_1", message.id, { type: "text", text: "", metadata: { separator: true } }),
- part("prt_2", message.id, {
- type: "reasoning",
- text: "think",
- metadata: { provider: 1 },
- time: { start: 21, end: 22 },
- }),
- part("prt_3", message.id, { type: "step-start", snapshot: "snap_start" }),
- part("prt_4", message.id, { type: "snapshot", snapshot: "snap_ignored" }),
- part("prt_5", message.id, { type: "patch", hash: "snap_patch", files: ["a.ts", "b.ts"] }),
- part("prt_6", message.id, { type: "patch", hash: "snap_patch_2", files: ["b.ts", "c.ts"] }),
- part("prt_7", message.id, {
- type: "step-finish",
- reason: "stop",
- snapshot: "snap_end",
- cost: 99,
- tokens: { input: 99, output: 99, reasoning: 99, cache: { read: 99, write: 99 } },
- }),
- part("prt_8", message.id, {
- type: "retry",
- attempt: 1,
- error: { name: "APIError", data: { message: "retry", isRetryable: true } },
- time: { created: 1 },
- }),
- ],
- )
- expect(result.messages[1].data).toEqual({
- agent: "build",
- model: { id: "model", providerID: "provider", variant: "fast" },
- content: [
- { type: "text", text: "", state: { separator: true } },
- { type: "reasoning", text: "think", state: { provider: 1 }, time: { created: 21, completed: 22 } },
- ],
- snapshot: { start: "snap_start", end: "snap_end", files: ["a.ts", "b.ts", "c.ts"] },
- finish: "stop",
- cost: 2.5,
- tokens: { input: 2, output: 3, reasoning: 4, cache: { read: 5, write: 6 } },
- time: { created: 20, completed: 21 },
- })
- expect(result.messages[1]).toMatchObject({ time_created: 20, time_updated: 21 })
- })
- test("normalizes every tool state", () => {
- const parent = user("msg_000000000008aaaaaaaaaaaaaa")
- const message = assistant("msg_000000000009aaaaaaaaaaaaaa", parent.id)
- const tool = (id: string, callID: string, state: Record<string, unknown>, metadata?: Record<string, unknown>) =>
- part(id, message.id, { type: "tool", callID, tool: "read", state, ...(metadata ? { metadata } : {}) })
- const result = transform(
- [parent, message],
- [
- tool("prt_1", "pending", { status: "pending", input: { a: 1 }, raw: "{}" }),
- tool("prt_2", "running", {
- status: "running",
- input: { b: 2 },
- metadata: { phase: "read" },
- time: { start: 30 },
- }),
- tool(
- "prt_3",
- "completed",
- {
- status: "completed",
- input: { c: 3 },
- output: "done",
- title: "Read",
- metadata: { result: true },
- time: { start: 31, end: 32 },
- attachments: [
- {
- id: "prt_attachment",
- sessionID: "ses_test",
- messageID: message.id,
- type: "file",
- mime: "text/plain",
- filename: "out.txt",
- url: "file:///out.txt",
- },
- ],
- },
- { provider: true },
- ),
- tool("prt_4", "compacted", {
- status: "completed",
- input: {},
- output: "secret",
- title: "Read",
- metadata: {},
- time: { start: 33, end: 34, compacted: 35 },
- attachments: [],
- }),
- tool("prt_5", "failed", {
- status: "error",
- input: { e: 5 },
- error: "boom",
- metadata: { output: "partial" },
- time: { start: 36, end: 37 },
- }),
- tool("prt_6", "failed-object-output", {
- status: "error",
- input: { f: 6 },
- error: "object output",
- metadata: { output: { nested: true } },
- time: { start: 38, end: 39 },
- }),
- ],
- )
- const content = result.messages[1].data.content
- if (!Array.isArray(content)) throw new Error("Expected assistant content")
- expect(content[0]).toMatchObject({
- id: "pending",
- state: {
- status: "error",
- error: { type: "tool.interrupted", message: "Tool execution was interrupted before V2 migration" },
- },
- time: { created: 20 },
- })
- expect(content[1]).toMatchObject({
- id: "running",
- state: { status: "error", metadata: { phase: "read" } },
- time: { created: 30 },
- })
- expect(content[2]).toMatchObject({
- id: "completed",
- providerState: { provider: true },
- state: {
- status: "completed",
- content: [
- { type: "text", text: "done" },
- { type: "file", uri: "file:///out.txt", mime: "text/plain", name: "out.txt" },
- ],
- metadata: { result: true },
- },
- time: { created: 31, completed: 32 },
- })
- expect(content[3]).toMatchObject({
- id: "compacted",
- state: { content: [{ type: "text", text: "[Old tool result content cleared]" }] },
- })
- expect(content[4]).toMatchObject({
- id: "failed",
- state: {
- status: "error",
- error: { type: "tool.execution", message: "boom" },
- content: [{ type: "text", text: "partial" }],
- metadata: { output: "partial" },
- },
- time: { created: 36, completed: 37 },
- })
- expect(content[5]).toEqual({
- type: "tool",
- id: "failed-object-output",
- name: "read",
- state: {
- status: "error",
- input: { f: 6 },
- error: { type: "tool.execution", message: "object output" },
- metadata: { output: { nested: true } },
- },
- time: { created: 38, completed: 39 },
- })
- })
- test("normalizes assistant errors and finish reasons", () => {
- const parent = user("msg_000000000010aaaaaaaaaaaaaa")
- const cases = [
- ["ProviderAuthError", { providerID: "provider", message: "auth" }, "provider.auth", "auth"],
- ["ContentFilterError", { message: "filtered" }, "provider.content-filter", "filtered"],
- ["ContextOverflowError", { message: "overflow" }, "provider.invalid-request", "overflow"],
- ["StructuredOutputError", { message: "shape", retries: 2 }, "provider.invalid-output", "shape"],
- ["MessageOutputLengthError", {}, "provider.invalid-output", "The model exceeded its output limit"],
- ["MessageAbortedError", { message: "stopped" }, "aborted", "stopped"],
- [
- "APIError",
- { message: "api", statusCode: 503, isRetryable: true, responseBody: "discard" },
- "provider.error",
- "api",
- ],
- ["UnknownError", { message: "unknown", ref: "discard" }, "unknown", "unknown"],
- ] as const
- const messages = cases.map(([name, data], index) =>
- assistant(
- `msg_00000000001${index}aaaaaaaaaaaaaa`,
- parent.id,
- { error: { name, data }, finish: index === 0 ? "new-provider-value" : undefined },
- 20 + index,
- ),
- )
- const result = transform([parent, ...messages], [])
- cases.forEach((entry, index) => {
- expect(result.messages[index + 1].data.error).toMatchObject({ type: entry[2], message: entry[3] })
- const error = result.messages[index + 1].data.error
- if (!error || typeof error !== "object") throw new Error("Expected assistant error")
- expect(Object.keys(error).sort()).toEqual(["message", "type"])
- })
- expect(result.messages[1].data.finish).toBe("unknown")
- })
- test("preserves every supported finish reason and omits an absent finish", () => {
- const parent = user("msg_000000000050aaaaaaaaaaaaaa")
- const finishes = ["stop", "length", "tool-calls", "content-filter", "error", "unknown", undefined] as const
- const result = transform(
- [
- parent,
- ...finishes.map((finish, index) =>
- assistant(`msg_00000000005${index + 1}aaaaaaaaaaaaaa`, parent.id, finish ? { finish } : {}, 20 + index),
- ),
- ],
- [],
- )
- expect(result.messages.slice(1).map((row) => row.data.finish)).toEqual([...finishes])
- })
- test("filters subtasks, collapses compactions, and keeps contiguous sequences", () => {
- const subtask = user("msg_000000000020aaaaaaaaaaaaaa", {}, 1)
- const taskAssistant = assistant("msg_000000000021aaaaaaaaaaaaaa", subtask.id, {}, 2)
- const compact = user("msg_000000000022aaaaaaaaaaaaaa", {}, 3)
- const unrelated = assistant("msg_000000000023aaaaaaaaaaaaaa", subtask.id, {}, 4)
- const summary = assistant("msg_000000000024aaaaaaaaaaaaaa", compact.id, { summary: true }, 5)
- const result = transform(
- [summary, compact, unrelated, taskAssistant, subtask],
- [
- part("prt_1", subtask.id, { type: "subtask", prompt: "work", description: "work", agent: "build" }),
- part("prt_2", taskAssistant.id, {
- type: "tool",
- callID: "task",
- tool: "task",
- state: {
- status: "completed",
- input: {},
- output: "done",
- title: "Task",
- metadata: {},
- time: { start: 1, end: 2 },
- },
- }),
- part("prt_3", unrelated.id, { type: "text", text: "keep" }),
- part("prt_4", compact.id, { type: "compaction", auto: false }),
- part("prt_5", summary.id, { type: "text", text: "summary" }),
- part("prt_6", summary.id, { type: "text", text: "" }),
- ],
- )
- expect(result.messages.map((row) => [row.type, row.seq])).toEqual([
- ["compaction", 0],
- ["assistant", 1],
- ])
- expect(result.messages[0]).toMatchObject({
- id: compact.id,
- time_created: 3,
- time_updated: 6,
- data: { status: "completed", reason: "manual", summary: "summary", recent: "" },
- })
- })
- test("serializes the retained compaction tail and preserves existing session selections", () => {
- const tailUser = user("msg_000000000025aaaaaaaaaaaaaa", {}, 1)
- const tailAssistant = assistant("msg_000000000026aaaaaaaaaaaaaa", tailUser.id, {}, 2)
- const compact = user("msg_000000000027aaaaaaaaaaaaaa", {}, 3)
- const summary = assistant("msg_000000000028aaaaaaaaaaaaaa", compact.id, { summary: true }, 4)
- const existing = session({
- agent: "existing",
- model: { id: "existing-model", providerID: "existing-provider", variant: "existing" },
- })
- const result = transform(
- [summary, compact, tailAssistant, tailUser],
- [
- part("prt_1", tailUser.id, { type: "text", text: "question" }),
- part("prt_2", tailAssistant.id, { type: "text", text: "answer" }),
- part("prt_3", compact.id, { type: "compaction", auto: true, tail_start_id: tailUser.id }),
- part("prt_4", summary.id, { type: "text", text: "summary" }),
- ],
- existing,
- )
- expect(result.messages[2].data).toMatchObject({
- reason: "auto",
- summary: "summary",
- recent: "[User]: question\n\n[Assistant]: answer",
- })
- expect(result.session.agent).toBe("existing")
- expect(result.session.model).toEqual(existing.model)
- })
- test("skips malformed rows, reports exact identifiers, and still backfills session aggregates", () => {
- const good = user(
- "msg_000000000030aaaaaaaaaaaaaa",
- { agent: "review", model: { providerID: "p2", modelID: "m2" } },
- 2,
- )
- const badJson = { ...user("msg_000000000031aaaaaaaaaaaaaa"), data: "{" }
- const badSchema = { ...user("msg_000000000032aaaaaaaaaaaaaa"), data: JSON.stringify({ role: "user" }) }
- const internal = assistant(
- "msg_000000000033aaaaaaaaaaaaaa",
- good.id,
- { cost: 7, tokens: { input: 8, output: 9, reasoning: 10, cache: { read: 11, write: 12 } } },
- 3,
- )
- const source = [good, badJson, badSchema, internal]
- const result = transform(source, [
- part("prt_1", good.id, { type: "text", text: "before" }),
- part("prt_2", good.id, "{"),
- part("prt_3", good.id, { type: "future", value: true }),
- part("prt_4", "msg_missing", { type: "text", text: "orphan" }),
- part("prt_5", good.id, { type: "text", text: "after" }),
- ])
- expect(result.messages[0].data.text).toBe("before\n\nafter")
- expect(result.messages.map((row) => row.seq)).toEqual([0, 1])
- expect(result.warnings).toEqual([
- { reason: "invalid-message", sessionID: "ses_test", messageID: badJson.id },
- { reason: "invalid-message", sessionID: "ses_test", messageID: badSchema.id },
- { reason: "invalid-part", sessionID: "ses_test", messageID: good.id, partID: "prt_2", observedType: undefined },
- { reason: "invalid-part", sessionID: "ses_test", messageID: good.id, partID: "prt_3", observedType: "future" },
- {
- reason: "orphan-part",
- sessionID: "ses_test",
- messageID: "msg_missing",
- partID: "prt_4",
- observedType: "text",
- },
- ])
- expect(result.session).toEqual({
- agent: "review",
- model: { id: "m2", providerID: "p2", variant: "default" },
- cost: 7,
- tokens_input: 8,
- tokens_output: 9,
- tokens_reasoning: 10,
- tokens_cache_read: 11,
- tokens_cache_write: 12,
- revert: null,
- time_compacting: null,
- })
- expect(source[0]).toBe(good)
- })
- test("retains empty ordinary messages and omits failed compactions", () => {
- const empty = user("msg_000000000034aaaaaaaaaaaaaa", {}, 1)
- const assistantMessage = assistant("msg_000000000035aaaaaaaaaaaaaa", empty.id, {}, 2)
- const compact = user("msg_000000000036aaaaaaaaaaaaaa", {}, 3)
- const failed = assistant(
- "msg_000000000037aaaaaaaaaaaaaa",
- compact.id,
- { summary: true, error: { name: "UnknownError", data: { message: "failed" } } },
- 4,
- )
- const result = transform(
- [empty, assistantMessage, compact, failed],
- [
- part("prt_1", empty.id, { type: "text", text: "ignored", ignored: true }),
- part("prt_2", assistantMessage.id, { type: "snapshot", snapshot: "standalone" }),
- part("prt_3", compact.id, { type: "compaction", auto: true }),
- part("prt_4", failed.id, { type: "text", text: "not committed" }),
- ],
- )
- expect(result.messages.map((row) => row.type)).toEqual(["user", "assistant"])
- expect(result.messages[0].data).toEqual({ text: "", time: { created: 1 } })
- expect(result.messages[1].data).toMatchObject({ content: [], snapshot: { start: "standalone" } })
- expect(result.watermark).toBe(1)
- })
- test("orders equal-time rows by ID and keeps SQL and payload timestamps consistent", () => {
- const first = user("msg_000000000041aaaaaaaaaaaaaa", { time: { created: 999 } }, 10)
- const second = assistant("msg_000000000042aaaaaaaaaaaaaa", first.id, { time: { created: 998, completed: 997 } }, 10)
- const result = transform(
- [second, first],
- [
- part("prt_1", first.id, { type: "text", text: "first" }),
- part("prt_2", second.id, { type: "text", text: "second" }),
- ],
- )
- expect(result.messages.map((row) => row.id)).toEqual([first.id, second.id])
- expect(result.messages.map((row) => [row.time_created, row.time_updated, row.data.time])).toEqual([
- [10, 11, { created: 10 }],
- [10, 11, { created: 10, completed: 11 }],
- ])
- })
- test("omits incomplete compactions and subtask assistants while retaining their aggregate usage", () => {
- const mixed = user("msg_000000000043aaaaaaaaaaaaaa", {}, 1)
- const task = assistant(
- "msg_000000000044aaaaaaaaaaaaaa",
- mixed.id,
- { cost: 4, tokens: { input: 5, output: 6, reasoning: 7, cache: { read: 8, write: 9 } } },
- 2,
- )
- const compact = user("msg_000000000045aaaaaaaaaaaaaa", {}, 3)
- const unfinished = assistant(
- "msg_000000000046aaaaaaaaaaaaaa",
- compact.id,
- { summary: true, time: { created: 4 } },
- 4,
- )
- const result = transform(
- [unfinished, compact, task, mixed],
- [
- part("prt_1", mixed.id, { type: "text", text: "keep" }),
- part("prt_2", mixed.id, { type: "subtask", prompt: "work", description: "work", agent: "build" }),
- part("prt_3", task.id, {
- type: "tool",
- callID: "task",
- tool: "task",
- state: { status: "pending", input: {}, raw: "{}" },
- }),
- part("prt_4", compact.id, { type: "compaction", auto: true }),
- part("prt_5", unfinished.id, { type: "text", text: "unfinished" }),
- ],
- )
- expect(result.messages.map((row) => [row.id, row.type])).toEqual([[mixed.id, "user"]])
- expect(result.watermark).toBe(0)
- expect(result.session).toMatchObject({
- cost: 5,
- tokens_input: 7,
- tokens_output: 9,
- tokens_reasoning: 11,
- tokens_cache_read: 13,
- tokens_cache_write: 15,
- })
- })
- })
- describe("V1Migration database workflow", () => {
- const createLegacyTables = Effect.fnUntraced(function* (db: Effect.Success<typeof makeDb>) {
- yield* db.run(sql`
- CREATE TABLE session (
- id text PRIMARY KEY,
- project_id text NOT NULL,
- workspace_id text,
- parent_id text,
- slug text NOT NULL,
- directory text NOT NULL,
- path text,
- title text NOT NULL,
- version text NOT NULL,
- share_url text,
- summary_additions integer,
- summary_deletions integer,
- summary_files integer,
- summary_diffs text,
- metadata text,
- cost real DEFAULT 0 NOT NULL,
- tokens_input integer DEFAULT 0 NOT NULL,
- tokens_output integer DEFAULT 0 NOT NULL,
- tokens_reasoning integer DEFAULT 0 NOT NULL,
- tokens_cache_read integer DEFAULT 0 NOT NULL,
- tokens_cache_write integer DEFAULT 0 NOT NULL,
- revert text,
- permission text,
- agent text,
- model text,
- time_created integer NOT NULL,
- time_updated integer NOT NULL,
- time_compacting integer,
- time_archived integer
- )
- `)
- yield* db.run(sql`
- CREATE TABLE message (
- id text PRIMARY KEY,
- session_id text NOT NULL,
- time_created integer NOT NULL,
- time_updated integer NOT NULL,
- data text NOT NULL
- )
- `)
- yield* db.run(sql`
- CREATE TABLE part (
- id text PRIMARY KEY,
- message_id text NOT NULL,
- session_id text NOT NULL,
- time_created integer NOT NULL,
- time_updated integer NOT NULL,
- data text NOT NULL
- )
- `)
- })
- const database = <A, E>(effect: Effect.Effect<A, E, Database.Service | Global.Service | Scope.Scope>) =>
- run(
- Effect.gen(function* () {
- const db = yield* makeDb
- yield* DatabaseMigration.apply(db)
- yield* createLegacyTables(db)
- return yield* effect.pipe(Effect.provideService(Database.Service, { db }))
- }),
- )
- test("reports required and completed status and completes an empty database idempotently", async () => {
- await database(
- Effect.gen(function* () {
- expect(yield* V1Migration.status()).toEqual({ status: "required" })
- expect(yield* V1Migration.run()).toEqual({ status: "completed" })
- expect(yield* V1Migration.status()).toEqual({ status: "completed" })
- expect(yield* V1Migration.run()).toEqual({ status: "completed" })
- }),
- )
- })
- test("imports previous V2 sessions and messages as part of the migration", async () => {
- await using tmp = await tmpdir()
- const filename = path.join(tmp.path, "opencode-next.db")
- const sqlite = await import("bun:sqlite")
- const source = new sqlite.Database(filename)
- source.run(`
- CREATE TABLE project (
- id text PRIMARY KEY, worktree text NOT NULL, vcs text, name text, icon_url text, icon_url_override text,
- icon_color text, time_created integer NOT NULL, time_updated integer NOT NULL, time_initialized integer,
- sandboxes text NOT NULL, commands text
- );
- CREATE TABLE session (
- id text PRIMARY KEY, project_id text NOT NULL, workspace_id text, parent_id text, fork_session_id text,
- fork_boundary text, slug text NOT NULL, directory text NOT NULL, path text, title text, version text NOT NULL,
- share_url text, summary_additions integer, summary_deletions integer, summary_files integer, summary_diffs text,
- metadata text, cost real DEFAULT 0 NOT NULL, tokens_input integer DEFAULT 0 NOT NULL,
- tokens_output integer DEFAULT 0 NOT NULL, tokens_reasoning integer DEFAULT 0 NOT NULL,
- tokens_cache_read integer DEFAULT 0 NOT NULL, tokens_cache_write integer DEFAULT 0 NOT NULL, revert text,
- permission text, agent text, model text, time_created integer NOT NULL, time_updated integer NOT NULL,
- time_compacting integer, time_archived integer, time_suspended integer
- );
- CREATE TABLE session_message (
- id text PRIMARY KEY, session_id text NOT NULL, type text NOT NULL, seq integer NOT NULL,
- time_created integer NOT NULL, time_updated integer NOT NULL, data text NOT NULL
- );
- INSERT INTO project VALUES (
- 'next-project', 'C:/Users/sewer', 'git', 'Source project', NULL, NULL, NULL, 1, 2, NULL, '[]', NULL
- );
- INSERT INTO session (
- id, project_id, slug, directory, title, version, agent, model, time_created, time_updated
- ) VALUES
- ('ses_next', 'next-project', 'next', 'C:/Users/sewer', 'Imported', '2', 'build',
- '{"id":"model","providerID":"provider"}', 10, 20),
- ('ses_existing', 'next-project', 'source-existing', '/tmp/next', 'Source existing', '2', NULL, NULL, 11, 21),
- ('ses_orphan', 'missing-project', 'orphan', '/tmp/orphan', 'Orphan', '2', NULL, NULL, 12, 22);
- INSERT INTO session_message VALUES
- ('msg_next', 'ses_next', 'user', 4, 12, 13, '{"text":"from next","time":{"created":12}}'),
- ('msg_source_existing', 'ses_existing', 'user', 2, 12, 13, '{"text":"source","time":{"created":12}}'),
- ('msg_orphan', 'ses_orphan', 'user', 0, 12, 13, '{"text":"orphan","time":{"created":12}}');
- `)
- source.close()
- await database(
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db.run(sql`
- INSERT INTO project (id, worktree, name, time_created, time_updated, sandboxes)
- VALUES ('next-project', '/tmp/current', 'Current project', 1, 2, '[]')
- `)
- yield* db.run(sql`
- INSERT INTO session_v2 (id, project_id, slug, directory, title, version, time_created, time_updated)
- VALUES ('ses_existing', 'next-project', 'current-existing', '/tmp/current', 'Current existing', '2', 1, 2)
- `)
- yield* db.run(sql`
- INSERT INTO session_message (id, session_id, type, seq, time_created, time_updated, data)
- VALUES ('msg_current_existing', 'ses_existing', 'user', 0, 1, 2, '{"text":"current","time":{"created":1}}')
- `)
- expect(yield* V1Migration.status()).toEqual({
- status: "required",
- })
- expect(yield* V1Migration.run({ nextDatabasePath: filename })).toEqual({ status: "completed" })
- expect(yield* V1Migration.status()).toEqual({
- status: "completed",
- })
- expect(yield* db.get(sql`SELECT title, agent, model FROM session_v2 WHERE id = 'ses_next'`)).toEqual({
- title: "Imported",
- agent: "build",
- model: '{"id":"model","providerID":"provider"}',
- })
- expect(
- yield* db
- .select({ directory: SessionTable.directory })
- .from(SessionTable)
- .where(eq(SessionTable.id, SessionSchema.ID.make("ses_next")))
- .get(),
- ).toEqual({ directory: process.platform === "win32" ? "C:\\Users\\sewer" : "C:/Users/sewer" })
- expect(yield* db.all(sql`SELECT id, seq, data FROM session_message WHERE session_id = 'ses_next'`)).toEqual([
- {
- id: "msg_next",
- seq: 4,
- data: '{"text":"from next","time":{"created":12}}',
- },
- ])
- expect(yield* db.get(sql`SELECT seq, owner_id FROM event_sequence WHERE aggregate_id = 'ses_next'`)).toEqual({
- seq: 4,
- owner_id: null,
- })
- expect(yield* db.get(sql`SELECT title FROM session_v2 WHERE id = 'ses_existing'`)).toEqual({
- title: "Current existing",
- })
- expect(yield* db.all(sql`SELECT id FROM session_message WHERE session_id = 'ses_existing'`)).toEqual([
- { id: "msg_current_existing" },
- ])
- expect(yield* db.get(sql`SELECT project_id FROM session_v2 WHERE id = 'ses_orphan'`)).toEqual({
- project_id: "global",
- })
- expect(yield* db.get(sql`SELECT name, worktree FROM project WHERE id = 'next-project'`)).toEqual({
- name: "Current project",
- worktree: "/tmp/current",
- })
- yield* db.run(sql`UPDATE project SET worktree = 'C:/Users/sewer' WHERE id = 'next-project'`)
- expect(
- yield* db
- .select({ worktree: ProjectTable.worktree })
- .from(ProjectTable)
- .where(eq(ProjectTable.id, Project.ID.make("next-project")))
- .get(),
- ).toEqual({
- worktree: AbsolutePath.make(process.platform === "win32" ? "C:\\Users\\sewer" : "C:/Users/sewer"),
- })
- expect(yield* db.get(sql`SELECT value FROM kv WHERE key = 'migration.v1-v2'`)).toEqual({
- value: '{"phase":"completed"}',
- })
- }),
- )
- })
- test("derives required status from the durable cursor", async () => {
- await database(
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db.run(
- sql`INSERT INTO project (id, worktree, time_created, time_updated, sandboxes) VALUES ('global', '/tmp/test', 1, 2, '[]')`,
- )
- yield* Effect.forEach(["ses_c", "ses_a", "ses_b"], (id) =>
- db.run(
- sql`INSERT INTO session (id, project_id, slug, directory, title, version, time_created, time_updated) VALUES (${id}, 'global', ${id}, '/tmp/test', 'Test', '1', 1, 2)`,
- ),
- )
- yield* db.run(
- sql`INSERT INTO kv (key, value, time_created, time_updated) VALUES ('migration.v1-v2', '{"phase":"sessions","cursor":"ses_b"}', 1, 1)`,
- )
- expect(yield* V1Migration.status()).toEqual({ status: "required" })
- }),
- )
- })
- test("reassigns V1 sessions whose projects are missing to the global project", async () => {
- await database(
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- const global = yield* Global.Service
- yield* db.run(
- sql`INSERT INTO session (id, project_id, slug, directory, title, version, time_created, time_updated) VALUES ('ses_orphan', 'missing-project', 'orphan', '/tmp/orphan', 'Orphan', '1', 1, 2)`,
- )
- expect(yield* V1Migration.run()).toEqual({ status: "completed" })
- expect(yield* db.get(sql`SELECT project_id FROM session_v2 WHERE id = 'ses_orphan'`)).toEqual({
- project_id: "global",
- })
- expect(yield* db.get(sql`SELECT worktree FROM project WHERE id = 'global'`)).toEqual({
- worktree: path.parse(global.data).root,
- })
- expect(yield* db.get(sql`SELECT value FROM kv WHERE key = 'migration.v1-v2'`)).toEqual({
- value: '{"phase":"completed"}',
- })
- }),
- )
- })
- test("replaces projections, updates sessions, preserves V1 rows, and checkpoints completion", async () => {
- await database(
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db.run(
- sql`INSERT INTO project (id, worktree, time_created, time_updated, sandboxes) VALUES ('global', '/tmp/test', 1, 2, '[]')`,
- )
- yield* db.run(sql`INSERT INTO session (
- id, project_id, slug, directory, title, version, cost, tokens_input, tokens_output, tokens_reasoning,
- tokens_cache_read, tokens_cache_write, revert, agent, model, metadata, time_created, time_updated,
- time_compacting, time_archived
- ) VALUES (
- 'ses_test', 'global', 'test', '/tmp/test', 'Test', '1', 99, 99, 99, 99, 99, 99, '{}', 'preserved',
- '{"id":"selected","providerID":"selected-provider","variant":"selected-variant"}', '{"keep":true}',
- 1, 2, 3, 4
- )`)
- const source = user("msg_000000000040aaaaaaaaaaaaaa")
- const sourcePart = part("prt_1", source.id, { type: "text", text: "hello" })
- yield* db.run(
- sql`INSERT INTO message (id, session_id, time_created, time_updated, data) VALUES (${source.id}, 'ses_test', 10, 11, ${source.data})`,
- )
- yield* db.run(
- sql`INSERT INTO part (id, message_id, session_id, time_created, time_updated, data) VALUES ('prt_1', ${source.id}, 'ses_test', 1, 2, ${sourcePart.data})`,
- )
- yield* db.run(
- sql`INSERT INTO session_message (id, session_id, type, seq, time_created, time_updated, data) VALUES ('msg_stale', 'ses_test', 'user', 0, 1, 1, '{"text":"stale","time":{"created":1}}')`,
- )
- yield* db.run(sql`INSERT INTO event_sequence (aggregate_id, seq) VALUES ('ses_test', 9)`)
- yield* db.run(
- sql`INSERT INTO event (id, aggregate_id, seq, created, type, data) VALUES ('event_stale', 'ses_test', 9, 1, 'session.renamed.1', '{}')`,
- )
- expect(yield* V1Migration.run()).toEqual({ status: "completed" })
- expect(yield* db.all(sql`SELECT id, type, seq, time_created, time_updated, data FROM session_message`)).toEqual(
- [
- {
- id: source.id,
- type: "user",
- seq: 0,
- time_created: 10,
- time_updated: 11,
- data: '{"text":"hello","time":{"created":10}}',
- },
- ],
- )
- expect(yield* db.all(sql`SELECT id, data FROM message`)).toEqual([{ id: source.id, data: source.data }])
- expect(yield* db.all(sql`SELECT id, data FROM part`)).toEqual([{ id: "prt_1", data: sourcePart.data }])
- expect(yield* db.get(sql`SELECT seq FROM event_sequence WHERE aggregate_id = 'ses_test'`)).toEqual({ seq: 0 })
- expect(yield* db.all(sql`SELECT id FROM event`)).toEqual([])
- expect(
- yield* db.get(
- sql`SELECT agent, model, metadata, cost, tokens_input, tokens_output, tokens_reasoning, tokens_cache_read, tokens_cache_write, revert, time_created, time_updated, time_compacting, time_archived FROM session_v2 WHERE id = 'ses_test'`,
- ),
- ).toEqual({
- agent: "preserved",
- model: '{"id":"selected","providerID":"selected-provider","variant":"selected-variant"}',
- metadata: '{"keep":true}',
- cost: 0,
- tokens_input: 0,
- tokens_output: 0,
- tokens_reasoning: 0,
- tokens_cache_read: 0,
- tokens_cache_write: 0,
- revert: null,
- time_created: 1,
- time_updated: 2,
- time_compacting: null,
- time_archived: 4,
- })
- expect(yield* db.get(sql`SELECT cost, revert, time_compacting FROM session WHERE id = 'ses_test'`)).toEqual({
- cost: 99,
- revert: "{}",
- time_compacting: 3,
- })
- expect(yield* db.get(sql`SELECT value FROM kv WHERE key = 'migration.v1-v2'`)).toEqual({
- value: '{"phase":"completed"}',
- })
- }),
- )
- })
- test("rolls back one session atomically and resumes from the committed cursor", async () => {
- await database(
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db.run(
- sql`INSERT INTO project (id, worktree, time_created, time_updated, sandboxes) VALUES ('global', '/tmp/test', 1, 2, '[]')`,
- )
- yield* Effect.forEach(["ses_a", "ses_b", "ses_c"], (id) =>
- db.run(
- sql`INSERT INTO session (id, project_id, slug, directory, title, version, cost, time_created, time_updated) VALUES (${id}, 'global', ${id}, '/tmp/test', 'Test', '1', 99, 1, 2)`,
- ),
- )
- yield* db.run(
- sql`CREATE TRIGGER fail_b BEFORE UPDATE ON session_v2 WHEN NEW.id = 'ses_b' BEGIN SELECT RAISE(ABORT, 'stop'); END`,
- )
- yield* db.run(
- sql`INSERT INTO session_message (id, session_id, type, seq, time_created, time_updated, data) VALUES ('msg_stale_b', 'ses_b', 'user', 0, 7, 8, '{"text":"stale","time":{"created":7}}')`,
- )
- yield* db.run(sql`INSERT INTO event_sequence (aggregate_id, seq, owner_id) VALUES ('ses_b', 7, 'owner')`)
- yield* db.run(
- sql`INSERT INTO event (id, aggregate_id, seq, created, type, data) VALUES ('event_stale_b', 'ses_b', 7, 1, 'session.renamed.1', '{}')`,
- )
- yield* Layer.launch(V1Migration.layer).pipe(Effect.forkScoped)
- const failed = yield* V1Migration.status().pipe(
- Effect.filterOrFail((status) => status.status === "error"),
- Effect.retry(Schedule.spaced("10 millis")),
- )
- expect(failed.status).toBe("error")
- if (failed.status === "error") expect(failed.error).toContain("stop")
- expect(yield* db.get(sql`SELECT value FROM kv WHERE key = 'migration.v1-v2'`)).toEqual({
- value: '{"phase":"sessions","cursor":"ses_c"}',
- })
- expect(yield* db.get(sql`SELECT cost FROM session_v2 WHERE id = 'ses_c'`)).toEqual({ cost: 0 })
- expect(yield* db.get(sql`SELECT cost FROM session_v2 WHERE id = 'ses_b'`)).toBeUndefined()
- expect(yield* db.get(sql`SELECT cost FROM session WHERE id = 'ses_b'`)).toEqual({ cost: 99 })
- expect(
- yield* db.all(
- sql`SELECT id, seq, time_created, time_updated, data FROM session_message WHERE session_id = 'ses_b'`,
- ),
- ).toEqual([
- {
- id: "msg_stale_b",
- seq: 0,
- time_created: 7,
- time_updated: 8,
- data: '{"text":"stale","time":{"created":7}}',
- },
- ])
- expect(yield* db.get(sql`SELECT seq, owner_id FROM event_sequence WHERE aggregate_id = 'ses_b'`)).toEqual({
- seq: 7,
- owner_id: "owner",
- })
- expect(yield* db.all(sql`SELECT id FROM event WHERE aggregate_id = 'ses_b'`)).toEqual([])
- expect(yield* db.get(sql`SELECT value FROM kv WHERE key = 'migration.v1-v2'`)).toEqual({
- value: '{"phase":"sessions","cursor":"ses_c"}',
- })
- yield* db.run(
- sql`INSERT INTO event (id, aggregate_id, seq, created, type, data) VALUES ('event_after_clear', 'ses_c', 0, 2, 'session.renamed.1', '{}')`,
- )
- yield* db.run(sql`DROP TRIGGER fail_b`)
- yield* Layer.launch(V1Migration.layer).pipe(Effect.forkScoped)
- yield* V1Migration.status().pipe(
- Effect.filterOrFail((status) => status.status === "completed"),
- Effect.retry(Schedule.spaced("10 millis")),
- )
- expect(yield* db.get(sql`SELECT cost FROM session_v2 WHERE id = 'ses_b'`)).toEqual({ cost: 0 })
- expect(yield* db.all(sql`SELECT id FROM session_message WHERE session_id = 'ses_b'`)).toEqual([])
- expect(yield* db.get(sql`SELECT seq, owner_id FROM event_sequence WHERE aggregate_id = 'ses_b'`)).toEqual({
- seq: -1,
- owner_id: null,
- })
- expect(yield* db.all(sql`SELECT id FROM event`)).toEqual([{ id: "event_after_clear" }])
- }),
- )
- })
- test("processes root, child, archived, empty, malformed-only, subtask-only, and incomplete-compaction sessions", async () => {
- const output = new Array<ReturnType<typeof Logger.formatStructured.log>>()
- await database(
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db.run(
- sql`INSERT INTO project (id, worktree, time_created, time_updated, sandboxes) VALUES ('global', '/tmp/test', 1, 2, '[]')`,
- )
- yield* db.run(
- sql`INSERT INTO session (id, project_id, slug, directory, title, version, time_created, time_updated) VALUES ('ses_root', 'global', 'root', '/tmp/test', 'Root', '1', 1, 2)`,
- )
- yield* db.run(
- sql`INSERT INTO session (id, project_id, parent_id, slug, directory, title, version, time_created, time_updated) VALUES ('ses_child', 'global', 'ses_root', 'child', '/tmp/test', 'Child', '1', 1, 2)`,
- )
- yield* db.run(
- sql`INSERT INTO session (id, project_id, slug, directory, title, version, time_created, time_updated, time_archived) VALUES ('ses_archived', 'global', 'archived', '/tmp/test', 'Archived', '1', 1, 2, 3)`,
- )
- yield* Effect.forEach(["ses_empty", "ses_malformed", "ses_subtask", "ses_compaction"], (id) =>
- db.run(
- sql`INSERT INTO session (id, project_id, slug, directory, title, version, time_created, time_updated) VALUES (${id}, 'global', ${id}, '/tmp/test', ${id}, '1', 1, 2)`,
- ),
- )
- yield* db.run(
- sql`INSERT INTO message (id, session_id, time_created, time_updated, data) VALUES ('msg_bad', 'ses_malformed', 1, 2, '{')`,
- )
- const subtask = user("msg_000000000047aaaaaaaaaaaaaa")
- yield* db.run(
- sql`INSERT INTO message (id, session_id, time_created, time_updated, data) VALUES (${subtask.id}, 'ses_subtask', 10, 11, ${subtask.data})`,
- )
- yield* db.run(
- sql`INSERT INTO part (id, message_id, session_id, time_created, time_updated, data) VALUES ('prt_subtask', ${subtask.id}, 'ses_subtask', 1, 2, '{"type":"subtask","prompt":"work","description":"work","agent":"build"}')`,
- )
- const compact = user("msg_000000000048aaaaaaaaaaaaaa")
- yield* db.run(
- sql`INSERT INTO message (id, session_id, time_created, time_updated, data) VALUES (${compact.id}, 'ses_compaction', 10, 11, ${compact.data})`,
- )
- yield* db.run(
- sql`INSERT INTO part (id, message_id, session_id, time_created, time_updated, data) VALUES ('prt_compaction', ${compact.id}, 'ses_compaction', 1, 2, '{"type":"compaction","auto":true}')`,
- )
- expect(yield* V1Migration.run()).toEqual({ status: "completed" })
- expect(yield* db.all(sql`SELECT aggregate_id, seq FROM event_sequence ORDER BY aggregate_id`)).toEqual([
- { aggregate_id: "ses_archived", seq: -1 },
- { aggregate_id: "ses_child", seq: -1 },
- { aggregate_id: "ses_compaction", seq: -1 },
- { aggregate_id: "ses_empty", seq: -1 },
- { aggregate_id: "ses_malformed", seq: -1 },
- { aggregate_id: "ses_root", seq: -1 },
- { aggregate_id: "ses_subtask", seq: -1 },
- ])
- expect(yield* db.all(sql`SELECT id FROM session_message`)).toEqual([])
- expect(yield* V1Migration.status()).toEqual({ status: "completed" })
- expect(output.map((entry) => entry.message)).toContainEqual([
- "Skipped V1 migration row",
- {
- reason: "invalid-message",
- sessionID: "ses_malformed",
- messageID: "msg_bad",
- },
- ])
- }).pipe(
- Effect.provide(
- Logger.layer([
- Logger.map(Logger.formatStructured, (entry) => {
- output.push(entry)
- }),
- ]),
- ),
- ),
- )
- })
- test("serializes concurrent callers and migrates each session once", async () => {
- await database(
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db.run(
- sql`INSERT INTO project (id, worktree, time_created, time_updated, sandboxes) VALUES ('global', '/tmp/test', 1, 2, '[]')`,
- )
- yield* db.run(
- sql`INSERT INTO session (id, project_id, slug, directory, title, version, time_created, time_updated) VALUES ('ses_test', 'global', 'test', '/tmp/test', 'Test', '1', 1, 2)`,
- )
- yield* db.run(sql`CREATE TABLE audit (count integer NOT NULL)`)
- yield* db.run(sql`INSERT INTO audit VALUES (0)`)
- yield* db.run(
- sql`CREATE TRIGGER audit_session AFTER UPDATE ON session_v2 BEGIN UPDATE audit SET count = count + 1; END`,
- )
- expect(yield* Effect.all([V1Migration.run(), V1Migration.run()], { concurrency: "unbounded" })).toEqual([
- { status: "completed" },
- { status: "completed" },
- ])
- expect(yield* db.get(sql`SELECT count FROM audit`)).toEqual({ count: 1 })
- }),
- )
- })
- })
|