diff --git a/packages/core/migration/20260604150931_harden_v2_sequence_indexes/migration.sql b/packages/core/migration/20260604150931_harden_v2_sequence_indexes/migration.sql index 02cbb2997ae..fdbb41ac6f2 100644 --- a/packages/core/migration/20260604150931_harden_v2_sequence_indexes/migration.sql +++ b/packages/core/migration/20260604150931_harden_v2_sequence_indexes/migration.sql @@ -1,5 +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 -CREATE UNIQUE INDEX `event_aggregate_seq_uidx` ON `event` (`aggregate_id`,`seq`);--> statement-breakpoint -CREATE UNIQUE INDEX `session_message_session_seq_uidx` ON `session_message` (`session_id`,`seq`); \ No newline at end of file +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`); diff --git a/packages/core/migration/20260604150931_harden_v2_sequence_indexes/snapshot.json b/packages/core/migration/20260604150931_harden_v2_sequence_indexes/snapshot.json index d0f513343aa..a066459c7e3 100644 --- a/packages/core/migration/20260604150931_harden_v2_sequence_indexes/snapshot.json +++ b/packages/core/migration/20260604150931_harden_v2_sequence_indexes/snapshot.json @@ -1,7 +1,7 @@ { "version": "7", "dialect": "sqlite", - "id": "3ab6b23c-f36f-4bd5-9f54-2f78af422812", + "id": "4a2f3a76-7952-44cc-abb4-122fb76e17a3", "prevIds": [ "fc92fa34-8074-44c3-88f0-a5417f7fd92d" ], @@ -1684,7 +1684,7 @@ "isUnique": true, "where": null, "origin": "manual", - "name": "event_aggregate_seq_uidx", + "name": "event_aggregate_seq_idx", "entityType": "indexes", "table": "event" }, @@ -1804,7 +1804,7 @@ "isUnique": true, "where": null, "origin": "manual", - "name": "session_message_session_seq_uidx", + "name": "session_message_session_seq_idx", "entityType": "indexes", "table": "session_message" }, @@ -1830,42 +1830,6 @@ "entityType": "indexes", "table": "session_message" }, - { - "columns": [ - { - "value": "session_id", - "isExpression": false - }, - { - "value": "time_created", - "isExpression": false - }, - { - "value": "id", - "isExpression": false - } - ], - "isUnique": false, - "where": null, - "origin": "manual", - "name": "session_message_session_time_created_id_idx", - "entityType": "indexes", - "table": "session_message" - }, - { - "columns": [ - { - "value": "time_created", - "isExpression": false - } - ], - "isUnique": false, - "where": null, - "origin": "manual", - "name": "session_message_time_created_idx", - "entityType": "indexes", - "table": "session_message" - }, { "columns": [ { diff --git a/packages/core/src/database/migration/20260604150931_harden_v2_sequence_indexes.ts b/packages/core/src/database/migration/20260604150931_harden_v2_sequence_indexes.ts index 3cd56dbfe25..f6f7b831785 100644 --- a/packages/core/src/database/migration/20260604150931_harden_v2_sequence_indexes.ts +++ b/packages/core/src/database/migration/20260604150931_harden_v2_sequence_indexes.ts @@ -8,9 +8,11 @@ export default { 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(`CREATE UNIQUE INDEX \`event_aggregate_seq_uidx\` ON \`event\` (\`aggregate_id\`,\`seq\`);`) + 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_uidx\` ON \`session_message\` (\`session_id\`,\`seq\`);`, + `CREATE UNIQUE INDEX \`session_message_session_seq_idx\` ON \`session_message\` (\`session_id\`,\`seq\`);`, ) }) }, diff --git a/packages/core/src/event/sql.ts b/packages/core/src/event/sql.ts index 7553f056029..ef95e0b7941 100644 --- a/packages/core/src/event/sql.ts +++ b/packages/core/src/event/sql.ts @@ -18,5 +18,5 @@ export const EventTable = sqliteTable( type: text().notNull(), data: text({ mode: "json" }).$type>().notNull(), }, - (table) => [uniqueIndex("event_aggregate_seq_uidx").on(table.aggregate_id, table.seq)], + (table) => [uniqueIndex("event_aggregate_seq_idx").on(table.aggregate_id, table.seq)], ) diff --git a/packages/core/src/session/sql.ts b/packages/core/src/session/sql.ts index 9f3dfae14f3..c0865c637f1 100644 --- a/packages/core/src/session/sql.ts +++ b/packages/core/src/session/sql.ts @@ -128,10 +128,8 @@ export const SessionMessageTable = sqliteTable( data: text({ mode: "json" }).notNull().$type(), }, (table) => [ - uniqueIndex("session_message_session_seq_uidx").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), ], ) diff --git a/packages/core/test/database-migration.test.ts b/packages/core/test/database-migration.test.ts index 6554460e393..0e96f5c1b49 100644 --- a/packages/core/test/database-migration.test.ts +++ b/packages/core/test/database-migration.test.ts @@ -66,13 +66,12 @@ describe("DatabaseMigration", () => { 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_seq_uidx', '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_seq_uidx', '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_uidx" }, + { name: "event_aggregate_seq_idx" }, { name: "session_input_session_pending_delivery_seq_idx" }, - { name: "session_message_session_seq_uidx" }, - { name: "session_message_session_time_created_id_idx" }, + { name: "session_message_session_seq_idx" }, { name: "session_message_session_type_seq_idx" }, ]) }), @@ -111,6 +110,68 @@ describe("DatabaseMigration", () => { ) }) + 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* () {