mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-15 07:48:24 -04:00
refactor(core): trim sqlite adapter paths (#42457)
This commit is contained in:
@@ -239,58 +239,5 @@ export class SQLiteEffectDatabase<
|
||||
) => Effect.Effect<A, E | SqlError, R> = (tx, config) => this.session.transaction(tx, config)
|
||||
}
|
||||
|
||||
export type SQLiteEffectWithReplicas<Q> = Q & { $primary: Q; $replicas: Q[] }
|
||||
|
||||
export const withReplicas = <
|
||||
TEffectHKT extends QueryEffectHKTBase,
|
||||
TRunResult,
|
||||
TRelations extends AnyRelations,
|
||||
Q extends SQLiteEffectDatabase<TEffectHKT, TRunResult, TRelations>,
|
||||
>(
|
||||
primary: Q,
|
||||
replicas: [Q, ...Q[]],
|
||||
getReplica: (replicas: Q[]) => Q = () => replicas[Math.floor(Math.random() * replicas.length)]!,
|
||||
): SQLiteEffectWithReplicas<Q> => {
|
||||
const select: Q["select"] = (...args: []) => getReplica(replicas).select(...args)
|
||||
const selectDistinct: Q["selectDistinct"] = (...args: []) => getReplica(replicas).selectDistinct(...args)
|
||||
const $count: Q["$count"] = (...args: [any]) => getReplica(replicas).$count(...args)
|
||||
const _with: Q["with"] = (...args: []) => getReplica(replicas).with(...args)
|
||||
const $with = ((...args: [string] | [string, ColumnsSelection]) =>
|
||||
args.length === 1
|
||||
? getReplica(replicas).$with(args[0])
|
||||
: getReplica(replicas).$with(args[0], args[1])) as Q["$with"]
|
||||
|
||||
const update: Q["update"] = (...args: [any]) => primary.update(...args)
|
||||
const insert: Q["insert"] = (...args: [any]) => primary.insert(...args)
|
||||
const $delete: Q["delete"] = (...args: [any]) => primary.delete(...args)
|
||||
const run: Q["run"] = (...args: [any]) => primary.run(...args)
|
||||
const all: Q["all"] = (...args: [any]) => primary.all(...args)
|
||||
const get: Q["get"] = (...args: [any]) => primary.get(...args)
|
||||
const values: Q["values"] = (...args: [any]) => primary.values(...args)
|
||||
const transaction: Q["transaction"] = (...args: [any]) => primary.transaction(...args)
|
||||
|
||||
return {
|
||||
...primary,
|
||||
update,
|
||||
insert,
|
||||
delete: $delete,
|
||||
run,
|
||||
all,
|
||||
get,
|
||||
values,
|
||||
transaction,
|
||||
$primary: primary,
|
||||
$replicas: replicas,
|
||||
select,
|
||||
selectDistinct,
|
||||
$count,
|
||||
$with,
|
||||
with: _with,
|
||||
get query() {
|
||||
return getReplica(replicas).query
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
export type AnySQLiteEffectDatabase = SQLiteEffectDatabase<any, any, any>
|
||||
export type AnySQLiteEffectSelectBase = SQLiteEffectSelectBase<any, any, any, any, any, any, any, any, any, any>
|
||||
|
||||
@@ -1,22 +1,10 @@
|
||||
/* oxlint-disable */
|
||||
import type { TablesRelationalConfig } from "drizzle-orm/_relations"
|
||||
import type { MigrationMeta } from "drizzle-orm/migrator"
|
||||
import type { AnyRelations } from "drizzle-orm/relations"
|
||||
import { type SQL, sql } from "drizzle-orm/sql/sql"
|
||||
import type { BaseSQLiteDatabase } from "drizzle-orm/sqlite-core"
|
||||
import type { SQLiteSession } from "drizzle-orm/sqlite-core/session"
|
||||
import { GET_VERSION_FOR, MIGRATIONS_TABLE_VERSIONS, type UpgradeResult } from "./utils.js"
|
||||
|
||||
/** @internal */
|
||||
export type SQLiteMigrationTableRow = { id: number | null; hash: string; created_at: number }
|
||||
|
||||
type AsyncSQLiteDatabaseWithSession = BaseSQLiteDatabase<"async", unknown, Record<string, unknown>> & {
|
||||
session: {
|
||||
all<T>(query: SQL): Promise<T[]>
|
||||
}
|
||||
transaction<T>(transaction: (tx: { run(query: SQL): Promise<unknown> }) => Promise<T>): Promise<T>
|
||||
}
|
||||
|
||||
type SQLiteMigrationBackfillEntry = {
|
||||
name: string
|
||||
selector:
|
||||
@@ -115,139 +103,3 @@ export function buildSQLiteMigrationBackfillStatements(
|
||||
|
||||
return statements
|
||||
}
|
||||
|
||||
/**
|
||||
* Detects the current version of the migrations table schema and upgrades it if needed.
|
||||
*
|
||||
* Version 0: Original schema (id, hash, created_at)
|
||||
* Version 1: Extended schema (id, hash, created_at, name, applied_at)
|
||||
*/
|
||||
export function upgradeSyncIfNeeded(
|
||||
migrationsTable: string,
|
||||
session: SQLiteSession<"sync", unknown, Record<string, unknown>, AnyRelations, TablesRelationalConfig>,
|
||||
localMigrations: MigrationMeta[],
|
||||
): UpgradeResult {
|
||||
const tableExists = session.all(sql`SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ${migrationsTable}`)
|
||||
|
||||
if (tableExists.length === 0) {
|
||||
return { newDb: true }
|
||||
}
|
||||
|
||||
// Table exists, check table shape
|
||||
const rows = session.all<{ column_name: string }>(
|
||||
sql`SELECT name as column_name FROM pragma_table_info(${migrationsTable})`,
|
||||
)
|
||||
|
||||
const version = GET_VERSION_FOR.sqlite(rows.map((r) => r.column_name))
|
||||
|
||||
for (let v = version; v < MIGRATIONS_TABLE_VERSIONS.sqlite; v++) {
|
||||
const upgradeFn = upgradeSyncFunctions[v]
|
||||
if (!upgradeFn) {
|
||||
throw new Error(`No upgrade path from migration table version ${v} to ${v + 1}`)
|
||||
}
|
||||
upgradeFn(migrationsTable, session, localMigrations)
|
||||
}
|
||||
|
||||
return { newDb: false }
|
||||
}
|
||||
|
||||
const upgradeSyncFunctions: Record<
|
||||
number,
|
||||
(
|
||||
migrationsTable: string,
|
||||
session: SQLiteSession<"sync", unknown, Record<string, unknown>, AnyRelations, TablesRelationalConfig>,
|
||||
localMigrations: MigrationMeta[],
|
||||
) => void
|
||||
> = {
|
||||
/**
|
||||
* Upgrade from version 0 to version 1:
|
||||
* 1. Read all existing DB migrations
|
||||
* 2. Sort localMigrations ASC by millis and if the same - sort by name
|
||||
* 3. Match each DB row to a local migration
|
||||
* If multiple migrations share the same second, use hash matching as a tiebreaker
|
||||
* Not implemented for now -> If hash matching fails, fall back to serial id ordering
|
||||
* 5. Create extra column and backfill names for matched migrations
|
||||
*/
|
||||
0: (migrationsTable, session, localMigrations) => {
|
||||
const table = sql`${sql.identifier(migrationsTable)}`
|
||||
const dbRows = session.all<SQLiteMigrationTableRow>(sql`SELECT id, hash, created_at FROM ${table} ORDER BY id ASC`)
|
||||
const statements = buildSQLiteMigrationBackfillStatements(
|
||||
migrationsTable,
|
||||
prepareSQLiteMigrationBackfill(dbRows, localMigrations),
|
||||
)
|
||||
|
||||
session.transaction((tx) => {
|
||||
for (const statement of statements) {
|
||||
tx.run(statement)
|
||||
}
|
||||
})
|
||||
},
|
||||
}
|
||||
|
||||
/**
|
||||
* Detects the current version of the migrations table schema and upgrades it if needed.
|
||||
*
|
||||
* Version 0: Original schema (id, hash, created_at)
|
||||
* Version 1: Extended schema (id, hash, created_at, name, applied_at)
|
||||
*/
|
||||
export async function upgradeAsyncIfNeeded(
|
||||
migrationsTable: string,
|
||||
db: AsyncSQLiteDatabaseWithSession,
|
||||
localMigrations: MigrationMeta[],
|
||||
): Promise<UpgradeResult> {
|
||||
// Check if the table exists at all
|
||||
const tableExists = await db.session.all(
|
||||
sql`SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ${migrationsTable}`,
|
||||
)
|
||||
|
||||
if (tableExists.length === 0) {
|
||||
return { newDb: true }
|
||||
}
|
||||
|
||||
const rows = await db.session.all<{ column_name: string }>(
|
||||
sql`SELECT name as column_name FROM pragma_table_info(${migrationsTable})`,
|
||||
)
|
||||
|
||||
const version = GET_VERSION_FOR.sqlite(rows.map((r) => r.column_name))
|
||||
|
||||
for (let v = version; v < MIGRATIONS_TABLE_VERSIONS.sqlite; v++) {
|
||||
const upgradeFn = upgradeAsyncFunctions[v]
|
||||
if (!upgradeFn) {
|
||||
throw new Error(`No upgrade path from migration table version ${v} to ${v + 1}`)
|
||||
}
|
||||
await upgradeFn(migrationsTable, db, localMigrations)
|
||||
}
|
||||
|
||||
return { newDb: false }
|
||||
}
|
||||
|
||||
const upgradeAsyncFunctions: Record<
|
||||
number,
|
||||
(migrationsTable: string, db: AsyncSQLiteDatabaseWithSession, localMigrations: MigrationMeta[]) => Promise<void>
|
||||
> = {
|
||||
/**
|
||||
* Upgrade from version 0 to version 1:
|
||||
* 1. Read all existing DB migrations
|
||||
* 2. Sort localMigrations ASC by millis and if the same - sort by name
|
||||
* 3. Match each DB row to a local migration
|
||||
* If multiple migrations share the same second, use hash matching as a tiebreaker
|
||||
* Not implemented for now -> If hash matching fails, fall back to serial id ordering
|
||||
* 5. Create extra column and backfill names for matched migrations
|
||||
*/
|
||||
0: async (migrationsTable, db, localMigrations) => {
|
||||
const table = sql`${sql.identifier(migrationsTable)}`
|
||||
const dbRows = await db.session.all<SQLiteMigrationTableRow>(
|
||||
sql`SELECT id, hash, created_at FROM ${table} ORDER BY id ASC`,
|
||||
)
|
||||
const statements = buildSQLiteMigrationBackfillStatements(
|
||||
migrationsTable,
|
||||
prepareSQLiteMigrationBackfill(dbRows, localMigrations),
|
||||
)
|
||||
|
||||
await db.transaction(async (tx) => {
|
||||
for (const statement of statements) {
|
||||
await tx.run(statement)
|
||||
}
|
||||
})
|
||||
},
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user