Skip to content
Open
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
124 changes: 124 additions & 0 deletions apps/server/src/persistence/Migrations.history.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,124 @@
import { assert, it } from "@effect/vitest";
import * as Effect from "effect/Effect";
import * as SqlClient from "effect/unstable/sql/SqlClient";
import * as NodeSqliteClient from "@t3tools/shared/nodeSqliteClient";

import { migrationEntries, migrationManifest, runMigrations } from "./Migrations.ts";

const seedHistorical = Effect.fn("seedHistorical")(function* (base: number, count: number) {
const sql = yield* SqlClient.SqlClient;
yield* runMigrations({ toMigrationInclusive: base });
for (const [id, name, migration] of migrationEntries.filter(
([id]) => id >= 48 && id < 48 + count,
)) {
yield* migration;
yield* sql`INSERT INTO effect_sql_migrations (migration_id, name) VALUES (${base + 1 + id - 48}, ${name})`;
}
});

for (const [base, count] of [
[43, 9],
[44, 9],
[44, 11],
[43, 1],
[44, 5],
] as const) {
it.effect(`upgrades historical V2 ${base + 1}–${base + count} without replaying its DDL`, () =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
yield* seedHistorical(base, count);
yield* runMigrations();
const rows = yield* sql<{
migration_id: number;
name: string;
}>`SELECT migration_id, name FROM effect_sql_migrations ORDER BY migration_id`;
assert.deepStrictEqual(
rows.map(({ migration_id, name }) => [migration_id, name] as const),
migrationManifest,
);
const columns = yield* sql<{ name: string }>`PRAGMA table_info(projection_projects)`;
assert.ok(columns.some(({ name }) => name === "auto_pull"));
assert.ok(columns.some(({ name }) => name === "project_icon_json"));
assert.deepStrictEqual(yield* runMigrations(), []);
}).pipe(Effect.provide(NodeSqliteClient.layerMemory())),
);
}

it.effect("preserves the historical manifest when a missing main migration fails", () =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
yield* seedHistorical(44, 9);
yield* sql`ALTER TABLE projection_thread_messages RENAME TO unavailable_messages`;
const before = yield* sql`SELECT * FROM effect_sql_migrations ORDER BY migration_id`;
const result = yield* Effect.exit(runMigrations());
assert.strictEqual(result._tag, "Failure");
assert.deepStrictEqual(
yield* sql`SELECT * FROM effect_sql_migrations ORDER BY migration_id`,
before,
);
const columns = yield* sql<{ name: string }>`PRAGMA table_info(projection_projects)`;
assert.ok(!columns.some(({ name }) => name === "auto_pull"));
}).pipe(Effect.provide(NodeSqliteClient.layerMemory())),
);

it.effect("rejects an unknown migration in a historical cohort without changing it", () =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
yield* seedHistorical(44, 9);
yield* sql`UPDATE effect_sql_migrations SET name = 'UnknownMigration' WHERE migration_id = 50`;
const before = yield* sql`SELECT * FROM effect_sql_migrations ORDER BY migration_id`;
assert.strictEqual((yield* Effect.exit(runMigrations()))._tag, "Failure");
assert.deepStrictEqual(
yield* sql`SELECT * FROM effect_sql_migrations ORDER BY migration_id`,
before,
);
}).pipe(Effect.provide(NodeSqliteClient.layerMemory())),
);

it.effect("rolls back reconciliation when a later V2 migration fails", () =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
yield* seedHistorical(44, 11);
yield* sql`CREATE INDEX orchestration_v2_projection_turn_items_thread_run_idx ON orchestration_v2_projection_turn_items(thread_id, run_id)`;
const before = yield* sql`SELECT * FROM effect_sql_migrations ORDER BY migration_id`;
assert.strictEqual((yield* Effect.exit(runMigrations()))._tag, "Failure");
assert.deepStrictEqual(
yield* sql`SELECT * FROM effect_sql_migrations ORDER BY migration_id`,
before,
);
const columns = yield* sql<{ name: string }>`PRAGMA table_info(projection_projects)`;
assert.ok(!columns.some(({ name }) => name === "auto_pull"));
}).pipe(Effect.provide(NodeSqliteClient.layerMemory())),
);

it.effect("preserves V2 import progress and original migration timestamps", () =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
yield* seedHistorical(43, 9);
yield* sql`INSERT INTO orchestration_v2_legacy_imports (thread_id, source_updated_at, shell_imported_at, imported_message_count) VALUES ('thread-1', '2026-09-01', '2026-09-02', 123)`;
yield* sql`UPDATE effect_sql_migrations SET created_at = '2026-09-01 00:00:00' WHERE migration_id = 44`;
const before = yield* sql`SELECT * FROM orchestration_v2_legacy_imports`;
yield* runMigrations();
assert.deepStrictEqual(yield* sql`SELECT * FROM orchestration_v2_legacy_imports`, before);
assert.deepStrictEqual(
yield* sql`SELECT created_at FROM effect_sql_migrations WHERE migration_id = 48`,
[{ created_at: "2026-09-01 00:00:00" }],
);
}).pipe(Effect.provide(NodeSqliteClient.layerMemory())),
);

it.effect("rejects a historical migration ceiling below the required main schema", () =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
yield* seedHistorical(43, 9);
const before = yield* sql`SELECT * FROM effect_sql_migrations ORDER BY migration_id`;
assert.strictEqual(
(yield* Effect.exit(runMigrations({ toMigrationInclusive: 46 })))._tag,
"Failure",
);
assert.deepStrictEqual(
yield* sql`SELECT * FROM effect_sql_migrations ORDER BY migration_id`,
before,
);
}).pipe(Effect.provide(NodeSqliteClient.layerMemory())),
);
70 changes: 69 additions & 1 deletion apps/server/src/persistence/Migrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@

import * as Migrator from "effect/unstable/sql/Migrator";
import * as Effect from "effect/Effect";
import * as SqlClient from "effect/unstable/sql/SqlClient";

// Import all migrations statically
import Migration0001 from "./Migrations/001_OrchestrationEvents.ts";
Expand Down Expand Up @@ -161,6 +162,66 @@ export const makeMigrationLoader = (throughId?: number) =>
*/
const run = Migrator.make({});

// Early V2 builds numbered this same migration sequence from 44 or 45.
// Match the complete recorded prefix before moving IDs; names alone must not
// cause an unknown or partially applied schema to be accepted as current.
const reconcileHistoricalV2 = Effect.fn("reconcileHistoricalV2")(function* (
toMigrationInclusive?: number,
) {
const sql = yield* SqlClient.SqlClient;
const tables =
yield* sql`SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'effect_sql_migrations'`;
if (tables.length === 0) return [];

const rows = yield* sql<{ migration_id: number; name: string }>`
SELECT migration_id, name FROM effect_sql_migrations ORDER BY migration_id
`;
const firstV2 = rows.find(({ name }) => name === "OrchestrationV2");
if (!firstV2 || firstV2.migration_id === 48) return [];

const base = firstV2.migration_id - 1;
const valid =
(base === 43 || base === 44) &&
rows.every((row, index) => {
const entry = migrationEntries[index < base ? index : index + 47 - base];
return row.migration_id === index + 1 && entry?.[1] === row.name;
});
if (!valid) {
return yield* new Migrator.MigrationError({
kind: "BadState",
message: "Unrecognized historical Orchestrator V2 migration manifest",
});
}

if (toMigrationInclusive !== undefined && toMigrationInclusive < 47) {
return yield* new Migrator.MigrationError({
kind: "BadState",
message: "Historical V2 reconciliation requires a migration ceiling of at least 47",
});
}

// Descending updates leave room for each lower ID and retain original dates.
for (const row of rows.slice(base).toReversed()) {
yield* sql`UPDATE effect_sql_migrations SET migration_id = ${row.migration_id + 47 - base}
WHERE migration_id = ${row.migration_id}`;
}
const executed: Array<readonly [number, string]> = [];
for (const [id, name, migration] of migrationEntries.filter(([id]) => id > base && id < 48)) {
yield* Effect.mapError(
migration,
(cause: unknown) =>
new Migrator.MigrationError({
kind: "Failed",
message: `Migration "${id}_${name}" failed during historical V2 reconciliation`,
cause,
}),
);
yield* sql`INSERT INTO effect_sql_migrations (migration_id, name) VALUES (${id}, ${name})`;
executed.push([id, name]);
}
return executed;
});

export interface RunMigrationsOptions {
readonly toMigrationInclusive?: number | undefined;
}
Expand All @@ -178,7 +239,14 @@ export interface RunMigrationsOptions {
export const runMigrations = Effect.fn("runMigrations")(function* ({
toMigrationInclusive,
}: RunMigrationsOptions = {}) {
const executedMigrations = yield* run({ loader: makeMigrationLoader(toMigrationInclusive) });
const sql = yield* SqlClient.SqlClient;
const executedMigrations = yield* sql.withTransaction(
Effect.gen(function* () {
const reconciled = yield* reconcileHistoricalV2(toMigrationInclusive);
const pending = yield* run({ loader: makeMigrationLoader(toMigrationInclusive) });
return [...reconciled, ...pending];
}),
);
const migrations = executedMigrations.map(([id, name]) => `${id}_${name}`);
yield* migrations.length === 0
? Effect.logDebug("Database schema is current")
Expand Down
Loading