Implement the DenoSqliteDriver correctly
This commit is contained in:
parent
7c2f290775
commit
781ca741dd
|
@ -1,17 +1,17 @@
|
|||
export {
|
||||
CompiledQuery,
|
||||
type DatabaseConnection,
|
||||
type DatabaseIntrospector,
|
||||
type Dialect,
|
||||
type DialectAdapter,
|
||||
type Driver,
|
||||
Kysely,
|
||||
type QueryCompiler,
|
||||
type QueryResult,
|
||||
SqliteAdapter,
|
||||
type SqliteDatabase,
|
||||
type SqliteDialectConfig,
|
||||
SqliteDriver,
|
||||
SqliteIntrospector,
|
||||
SqliteQueryCompiler,
|
||||
type SqliteStatement,
|
||||
} from 'npm:kysely@^0.25.0';
|
||||
|
||||
export type { DB as DenoSqlite } from 'https://deno.land/x/sqlite@v3.7.3/mod.ts';
|
||||
|
|
|
@ -1,50 +1,107 @@
|
|||
import { type DenoSqlite, type SqliteDatabase, SqliteDriver, type SqliteStatement } from '../deps.ts';
|
||||
import { CompiledQuery, type DatabaseConnection, type DenoSqlite, type Driver, type QueryResult } from '../deps.ts';
|
||||
|
||||
import type { DenoSqliteDialectConfig } from './deno-sqlite-dialect-config.ts';
|
||||
|
||||
class DenoSqliteDriver extends SqliteDriver {
|
||||
class DenoSqliteDriver implements Driver {
|
||||
readonly #config: DenoSqliteDialectConfig;
|
||||
readonly #connectionMutex = new ConnectionMutex();
|
||||
|
||||
#db?: DenoSqlite;
|
||||
#connection?: DatabaseConnection;
|
||||
|
||||
constructor(config: DenoSqliteDialectConfig) {
|
||||
super({
|
||||
...config,
|
||||
database: async () =>
|
||||
new DenoSqliteDatabase(
|
||||
typeof config.database === 'function' ? await config.database() : config.database,
|
||||
),
|
||||
});
|
||||
this.#config = Object.freeze({ ...config });
|
||||
}
|
||||
|
||||
async init(): Promise<void> {
|
||||
this.#db = typeof this.#config.database === 'function' ? await this.#config.database() : this.#config.database;
|
||||
|
||||
this.#connection = new DenoSqliteConnection(this.#db);
|
||||
|
||||
if (this.#config.onCreateConnection) {
|
||||
await this.#config.onCreateConnection(this.#connection);
|
||||
}
|
||||
}
|
||||
|
||||
async acquireConnection(): Promise<DatabaseConnection> {
|
||||
// SQLite only has one single connection. We use a mutex here to wait
|
||||
// until the single connection has been released.
|
||||
await this.#connectionMutex.lock();
|
||||
return this.#connection!;
|
||||
}
|
||||
|
||||
async beginTransaction(connection: DatabaseConnection): Promise<void> {
|
||||
await connection.executeQuery(CompiledQuery.raw('begin'));
|
||||
}
|
||||
|
||||
async commitTransaction(connection: DatabaseConnection): Promise<void> {
|
||||
await connection.executeQuery(CompiledQuery.raw('commit'));
|
||||
}
|
||||
|
||||
async rollbackTransaction(connection: DatabaseConnection): Promise<void> {
|
||||
await connection.executeQuery(CompiledQuery.raw('rollback'));
|
||||
}
|
||||
|
||||
// deno-lint-ignore require-await
|
||||
async releaseConnection(): Promise<void> {
|
||||
this.#connectionMutex.unlock();
|
||||
}
|
||||
|
||||
// deno-lint-ignore require-await
|
||||
async destroy(): Promise<void> {
|
||||
this.#db?.close();
|
||||
}
|
||||
}
|
||||
|
||||
/** HACK: This is an adapter class. */
|
||||
class DenoSqliteDatabase implements SqliteDatabase {
|
||||
#db: DenoSqlite;
|
||||
class DenoSqliteConnection implements DatabaseConnection {
|
||||
readonly #db: DenoSqlite;
|
||||
|
||||
constructor(db: DenoSqlite) {
|
||||
this.#db = db;
|
||||
}
|
||||
|
||||
close(): void {
|
||||
this.#db.close();
|
||||
executeQuery<O>({ sql, parameters }: CompiledQuery): Promise<QueryResult<O>> {
|
||||
const query = this.#db.prepareQuery(sql);
|
||||
|
||||
const rows = query.allEntries(parameters as any);
|
||||
const { changes, lastInsertRowId } = this.#db;
|
||||
|
||||
query.finalize();
|
||||
|
||||
return Promise.resolve({
|
||||
rows: rows as O[],
|
||||
numAffectedRows: BigInt(changes),
|
||||
insertId: BigInt(lastInsertRowId),
|
||||
});
|
||||
}
|
||||
|
||||
prepare(sql: string): SqliteStatement {
|
||||
const query = this.#db.prepareQuery(sql);
|
||||
return {
|
||||
// HACK: implement an actual driver to fix this.
|
||||
reader: true,
|
||||
all: (parameters: ReadonlyArray<unknown>) => {
|
||||
const result = query.allEntries(parameters as any);
|
||||
query.finalize();
|
||||
return result;
|
||||
},
|
||||
run: (parameters: ReadonlyArray<unknown>) => {
|
||||
query.execute(parameters as any);
|
||||
query.finalize();
|
||||
return {
|
||||
changes: this.#db.changes,
|
||||
lastInsertRowid: this.#db.lastInsertRowId,
|
||||
};
|
||||
},
|
||||
};
|
||||
// deno-lint-ignore require-yield
|
||||
async *streamQuery<R>(): AsyncIterableIterator<QueryResult<R>> {
|
||||
throw new Error('Sqlite driver doesn\'t support streaming');
|
||||
}
|
||||
}
|
||||
|
||||
class ConnectionMutex {
|
||||
#promise?: Promise<void>;
|
||||
#resolve?: () => void;
|
||||
|
||||
async lock(): Promise<void> {
|
||||
while (this.#promise) {
|
||||
await this.#promise;
|
||||
}
|
||||
|
||||
this.#promise = new Promise((resolve) => {
|
||||
this.#resolve = resolve;
|
||||
});
|
||||
}
|
||||
|
||||
unlock(): void {
|
||||
const resolve = this.#resolve;
|
||||
|
||||
this.#promise = undefined;
|
||||
this.#resolve = undefined;
|
||||
|
||||
resolve?.();
|
||||
}
|
||||
}
|
||||
|
||||
|
|
Loading…
Reference in New Issue