Compare commits
2 commits
dev
...
fix/core-v
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e1ccbf5eb6 | ||
|
|
bfe2da2072 |
7 changed files with 2029 additions and 13 deletions
|
|
@ -0,0 +1,7 @@
|
|||
DROP INDEX IF EXISTS `event_aggregate_seq_idx`;--> statement-breakpoint
|
||||
DROP INDEX IF EXISTS `event_aggregate_type_seq_idx`;--> statement-breakpoint
|
||||
DROP INDEX IF EXISTS `session_message_session_seq_idx`;--> statement-breakpoint
|
||||
DROP INDEX IF EXISTS `session_message_session_time_created_id_idx`;--> statement-breakpoint
|
||||
DROP INDEX IF EXISTS `session_message_time_created_idx`;--> statement-breakpoint
|
||||
CREATE UNIQUE INDEX `event_aggregate_seq_idx` ON `event` (`aggregate_id`,`seq`);--> statement-breakpoint
|
||||
CREATE UNIQUE INDEX `session_message_session_seq_idx` ON `session_message` (`session_id`,`seq`);
|
||||
1900
packages/core/migration/20260604150931_harden_v2_sequence_indexes/snapshot.json
generated
Normal file
1900
packages/core/migration/20260604150931_harden_v2_sequence_indexes/snapshot.json
generated
Normal file
File diff suppressed because it is too large
Load diff
1
packages/core/src/database/migration.gen.ts
generated
1
packages/core/src/database/migration.gen.ts
generated
|
|
@ -31,5 +31,6 @@ export const migrations = (
|
|||
import("./migration/20260603040000_session_message_projection_order"),
|
||||
import("./migration/20260603141458_session_input_inbox"),
|
||||
import("./migration/20260603160727_jittery_ezekiel_stane"),
|
||||
import("./migration/20260604150931_harden_v2_sequence_indexes"),
|
||||
])
|
||||
).map((module) => module.default) satisfies DatabaseMigration.Migration[]
|
||||
|
|
|
|||
|
|
@ -0,0 +1,19 @@
|
|||
import { Effect } from "effect"
|
||||
import type { DatabaseMigration } from "../migration"
|
||||
|
||||
export default {
|
||||
id: "20260604150931_harden_v2_sequence_indexes",
|
||||
up(tx) {
|
||||
return Effect.gen(function* () {
|
||||
yield* tx.run(`DROP INDEX IF EXISTS \`event_aggregate_seq_idx\`;`)
|
||||
yield* tx.run(`DROP INDEX IF EXISTS \`event_aggregate_type_seq_idx\`;`)
|
||||
yield* tx.run(`DROP INDEX IF EXISTS \`session_message_session_seq_idx\`;`)
|
||||
yield* tx.run(`DROP INDEX IF EXISTS \`session_message_session_time_created_id_idx\`;`)
|
||||
yield* tx.run(`DROP INDEX IF EXISTS \`session_message_time_created_idx\`;`)
|
||||
yield* tx.run(`CREATE UNIQUE INDEX \`event_aggregate_seq_idx\` ON \`event\` (\`aggregate_id\`,\`seq\`);`)
|
||||
yield* tx.run(
|
||||
`CREATE UNIQUE INDEX \`session_message_session_seq_idx\` ON \`session_message\` (\`session_id\`,\`seq\`);`,
|
||||
)
|
||||
})
|
||||
},
|
||||
} satisfies DatabaseMigration.Migration
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
import { sqliteTable, text, integer, index } from "drizzle-orm/sqlite-core"
|
||||
import { sqliteTable, text, integer, uniqueIndex } from "drizzle-orm/sqlite-core"
|
||||
import type { EventV2 } from "../event"
|
||||
|
||||
export const EventSequenceTable = sqliteTable("event_sequence", {
|
||||
|
|
@ -18,8 +18,5 @@ export const EventTable = sqliteTable(
|
|||
type: text().notNull(),
|
||||
data: text({ mode: "json" }).$type<Record<string, unknown>>().notNull(),
|
||||
},
|
||||
(table) => [
|
||||
index("event_aggregate_seq_idx").on(table.aggregate_id, table.seq),
|
||||
index("event_aggregate_type_seq_idx").on(table.aggregate_id, table.type, table.seq),
|
||||
],
|
||||
(table) => [uniqueIndex("event_aggregate_seq_idx").on(table.aggregate_id, table.seq)],
|
||||
)
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
import { sqliteTable, text, integer, index, primaryKey, real } from "drizzle-orm/sqlite-core"
|
||||
import { sqliteTable, text, integer, index, primaryKey, real, uniqueIndex } from "drizzle-orm/sqlite-core"
|
||||
import * as DatabasePath from "../database/path"
|
||||
import { ProjectTable } from "../project/sql"
|
||||
import type { SessionMessage } from "./message"
|
||||
|
|
@ -122,15 +122,14 @@ export const SessionMessageTable = sqliteTable(
|
|||
.notNull()
|
||||
.references(() => SessionTable.id, { onDelete: "cascade" }),
|
||||
type: text().$type<SessionMessage.Type>().notNull(),
|
||||
/** Immutable sequence of the durable event that first projected this timeline row. */
|
||||
seq: integer().notNull(),
|
||||
...Timestamps,
|
||||
data: text({ mode: "json" }).notNull().$type<SessionMessageData>(),
|
||||
},
|
||||
(table) => [
|
||||
index("session_message_session_seq_idx").on(table.session_id, table.seq),
|
||||
uniqueIndex("session_message_session_seq_idx").on(table.session_id, table.seq),
|
||||
index("session_message_session_type_seq_idx").on(table.session_id, table.type, table.seq),
|
||||
index("session_message_session_time_created_id_idx").on(table.session_id, table.time_created, table.id),
|
||||
index("session_message_time_created_idx").on(table.time_created),
|
||||
],
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -10,6 +10,7 @@ import { DatabaseMigration } from "@opencode-ai/core/database/migration"
|
|||
import sessionUsageMigration from "@opencode-ai/core/database/migration/20260510033149_session_usage"
|
||||
import normalizeStoragePathsMigration from "@opencode-ai/core/database/migration/20260601010001_normalize_storage_paths"
|
||||
import sessionMessageProjectionOrderMigration from "@opencode-ai/core/database/migration/20260603040000_session_message_projection_order"
|
||||
import hardenV2SequenceIndexesMigration from "@opencode-ai/core/database/migration/20260604150931_harden_v2_sequence_indexes"
|
||||
import { ProjectV2 } from "@opencode-ai/core/project"
|
||||
import { ProjectTable } from "@opencode-ai/core/project/sql"
|
||||
import { AbsolutePath } from "@opencode-ai/core/schema"
|
||||
|
|
@ -62,23 +63,115 @@ describe("DatabaseMigration", () => {
|
|||
expect(
|
||||
yield* db.get(sql`SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'session_input'`),
|
||||
).toEqual({ name: "session_input" })
|
||||
expect(yield* db.get(sql`SELECT count(*) as count FROM migration`)).toEqual({ count: 29 })
|
||||
expect(yield* db.get(sql`SELECT count(*) as count FROM migration`)).toEqual({ count: 30 })
|
||||
expect(
|
||||
yield* db.all(
|
||||
sql`SELECT name FROM sqlite_master WHERE type = 'index' AND name IN ('event_aggregate_seq_idx', 'event_aggregate_type_seq_idx', 'session_input_session_pending_seq_idx', 'session_input_session_pending_delivery_seq_idx', 'session_message_session_idx', 'session_message_session_type_idx', 'session_message_session_seq_idx', 'session_message_session_type_seq_idx', 'session_message_session_time_created_id_idx') ORDER BY name`,
|
||||
sql`SELECT name FROM sqlite_master WHERE type = 'index' AND name IN ('event_aggregate_seq_idx', 'event_aggregate_type_seq_idx', 'session_input_session_pending_seq_idx', 'session_input_session_pending_delivery_seq_idx', 'session_message_session_idx', 'session_message_session_type_idx', 'session_message_session_seq_idx', 'session_message_session_type_seq_idx', 'session_message_session_time_created_id_idx', 'session_message_time_created_idx') ORDER BY name`,
|
||||
),
|
||||
).toEqual([
|
||||
{ name: "event_aggregate_seq_idx" },
|
||||
{ name: "event_aggregate_type_seq_idx" },
|
||||
{ name: "session_input_session_pending_delivery_seq_idx" },
|
||||
{ name: "session_message_session_seq_idx" },
|
||||
{ name: "session_message_session_time_created_id_idx" },
|
||||
{ name: "session_message_session_type_seq_idx" },
|
||||
])
|
||||
}),
|
||||
)
|
||||
})
|
||||
|
||||
test("enforces unique durable and timeline sequence positions", async () => {
|
||||
await run(
|
||||
Effect.gen(function* () {
|
||||
const db = yield* makeDb
|
||||
yield* db.run(
|
||||
sql`CREATE TABLE event (id text PRIMARY KEY, aggregate_id text NOT NULL, seq integer NOT NULL, type text NOT NULL)`,
|
||||
)
|
||||
yield* db.run(sql`CREATE INDEX event_aggregate_seq_idx ON event (aggregate_id, seq)`)
|
||||
yield* db.run(sql`CREATE INDEX event_aggregate_type_seq_idx ON event (aggregate_id, type, seq)`)
|
||||
yield* db.run(
|
||||
sql`CREATE TABLE session_message (id text PRIMARY KEY, session_id text NOT NULL, seq integer NOT NULL)`,
|
||||
)
|
||||
yield* db.run(sql`CREATE INDEX session_message_session_seq_idx ON session_message (session_id, seq)`)
|
||||
|
||||
yield* DatabaseMigration.applyOnly(db, [hardenV2SequenceIndexesMigration])
|
||||
|
||||
yield* db.run(sql`INSERT INTO event (id, aggregate_id, seq, type) VALUES ('event_1', 'session_1', 1, 'one')`)
|
||||
yield* db.run(sql`INSERT INTO session_message (id, session_id, seq) VALUES ('message_1', 'session_1', 1)`)
|
||||
expect(() =>
|
||||
Effect.runSync(
|
||||
db.run(sql`INSERT INTO event (id, aggregate_id, seq, type) VALUES ('event_2', 'session_1', 1, 'two')`),
|
||||
),
|
||||
).toThrow()
|
||||
expect(() =>
|
||||
Effect.runSync(
|
||||
db.run(sql`INSERT INTO session_message (id, session_id, seq) VALUES ('message_2', 'session_1', 1)`),
|
||||
),
|
||||
).toThrow()
|
||||
}),
|
||||
)
|
||||
})
|
||||
|
||||
test.each(["event", "session_message"] as const)(
|
||||
"rejects duplicate existing %s sequence positions without dropping old indexes",
|
||||
async (duplicate) => {
|
||||
await run(
|
||||
Effect.gen(function* () {
|
||||
const db = yield* makeDb
|
||||
yield* db.run(
|
||||
sql`CREATE TABLE event (id text PRIMARY KEY, aggregate_id text NOT NULL, seq integer NOT NULL, type text NOT NULL)`,
|
||||
)
|
||||
yield* db.run(sql`CREATE INDEX event_aggregate_seq_idx ON event (aggregate_id, seq)`)
|
||||
yield* db.run(sql`CREATE INDEX event_aggregate_type_seq_idx ON event (aggregate_id, type, seq)`)
|
||||
yield* db.run(
|
||||
sql`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)`,
|
||||
)
|
||||
yield* db.run(sql`CREATE INDEX session_message_session_seq_idx ON session_message (session_id, seq)`)
|
||||
yield* db.run(
|
||||
sql`CREATE INDEX session_message_session_type_seq_idx ON session_message (session_id, type, seq)`,
|
||||
)
|
||||
yield* db.run(
|
||||
sql`CREATE INDEX session_message_session_time_created_id_idx ON session_message (session_id, time_created, id)`,
|
||||
)
|
||||
yield* db.run(sql`CREATE INDEX session_message_time_created_idx ON session_message (time_created)`)
|
||||
|
||||
yield* db.run(sql`INSERT INTO event (id, aggregate_id, seq, type) VALUES ('event_1', 'session_1', 1, 'one')`)
|
||||
yield* db.run(
|
||||
sql`INSERT INTO session_message (id, session_id, type, seq, time_created) VALUES ('message_1', 'session_1', 'user', 1, 1)`,
|
||||
)
|
||||
if (duplicate === "event") {
|
||||
yield* db.run(
|
||||
sql`INSERT INTO event (id, aggregate_id, seq, type) VALUES ('event_2', 'session_1', 1, 'two')`,
|
||||
)
|
||||
}
|
||||
if (duplicate === "session_message") {
|
||||
yield* db.run(
|
||||
sql`INSERT INTO session_message (id, session_id, type, seq, time_created) VALUES ('message_2', 'session_1', 'assistant', 1, 2)`,
|
||||
)
|
||||
}
|
||||
|
||||
expect(
|
||||
(yield* DatabaseMigration.applyOnly(db, [hardenV2SequenceIndexesMigration]).pipe(Effect.exit))._tag,
|
||||
).toBe("Failure")
|
||||
expect(
|
||||
yield* db.all(sql`
|
||||
SELECT name, [unique] AS is_unique FROM pragma_index_list('event') WHERE origin = 'c'
|
||||
UNION ALL
|
||||
SELECT name, [unique] AS is_unique FROM pragma_index_list('session_message') WHERE origin = 'c'
|
||||
ORDER BY name
|
||||
`),
|
||||
).toEqual([
|
||||
{ name: "event_aggregate_seq_idx", is_unique: 0 },
|
||||
{ name: "event_aggregate_type_seq_idx", is_unique: 0 },
|
||||
{ name: "session_message_session_seq_idx", is_unique: 0 },
|
||||
{ name: "session_message_session_time_created_id_idx", is_unique: 0 },
|
||||
{ name: "session_message_session_type_seq_idx", is_unique: 0 },
|
||||
{ name: "session_message_time_created_idx", is_unique: 0 },
|
||||
])
|
||||
expect(yield* db.all(sql`SELECT id FROM migration`)).toEqual([])
|
||||
}),
|
||||
)
|
||||
},
|
||||
)
|
||||
|
||||
test("resets incompatible projected Session messages before adding sequence order", async () => {
|
||||
await run(
|
||||
Effect.gen(function* () {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue