-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
fix: support multiple instances of niledatabase
- Loading branch information
Showing
7 changed files
with
149 additions
and
132 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,42 @@ | ||
import { Config } from '../utils/Config'; | ||
import { watchEvictPool } from '../utils/Event'; | ||
|
||
import NileDatabase, { NileDatabaseI } from './NileInstance'; | ||
|
||
export default class DBManager { | ||
connections: Map<string, NileDatabase>; | ||
|
||
private makeId( | ||
tenantId?: string | undefined | null, | ||
userId?: string | undefined | null | ||
) { | ||
if (tenantId && userId) { | ||
return `${tenantId}:${userId}`; | ||
} | ||
if (tenantId) { | ||
return `${tenantId}`; | ||
} | ||
return 'base'; | ||
} | ||
constructor(config: Config) { | ||
this.connections = new Map(); | ||
// add the base one, so you can at least query | ||
const id = this.makeId(); | ||
this.connections.set(id, new NileDatabase(new Config(config), id)); | ||
watchEvictPool((id) => { | ||
if (id && this.connections.has(id)) { | ||
this.connections.delete(id); | ||
} | ||
}); | ||
} | ||
|
||
getConnection(config: Config): NileDatabaseI { | ||
const id = this.makeId(config.tenantId, config.userId); | ||
const existing = this.connections.get(id); | ||
if (existing) { | ||
return existing as unknown as NileDatabaseI; | ||
} | ||
this.connections.set(id, new NileDatabase(new Config(config), id)); | ||
return this.connections.get(id) as unknown as NileDatabaseI; | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,87 @@ | ||
/* eslint-disable @typescript-eslint/no-explicit-any */ | ||
import knex, { Knex } from 'knex'; | ||
|
||
import { Config } from '../utils/Config'; | ||
import { evictPool } from '../utils/Event'; | ||
|
||
// doing this now, to provide flexibility later | ||
class NileDatabase { | ||
knex: Knex; | ||
// db: Knex; | ||
tenantId?: undefined | null | string; | ||
userId?: undefined | null | string; | ||
id: string; | ||
config: any; | ||
|
||
constructor(config: Config, id: string) { | ||
this.id = id; | ||
let poolConfig = {}; | ||
const afterCreate = ( | ||
conn: { | ||
on: any; | ||
query: (query: string, cb: (err: unknown) => void) => void; | ||
}, | ||
done: (err: unknown, conn: unknown) => void | ||
) => { | ||
const query = [`SET nile.tenant_id = '${config.tenantId}'`]; | ||
if (config.userId) { | ||
if (!config.tenantId) { | ||
// eslint-disable-next-line no-console | ||
console.warn( | ||
'A user id cannot be set in context without a tenant id' | ||
); | ||
} | ||
query.push(`SET nile.user_id = '${config.userId}'`); | ||
} | ||
// in this example we use pg driver's connection API | ||
conn.query(query.join(';'), function (err: unknown) { | ||
done(err, conn); | ||
}); | ||
}; | ||
if (config.tenantId) { | ||
if (config.db.pool?.afterCreate) { | ||
// eslint-disable-next-line no-console | ||
console.log( | ||
'Providing an pool configuration will stop automatic tenant context setting.' | ||
); | ||
} else if (config.db.pool) { | ||
poolConfig = { | ||
...config.db.pool, | ||
afterCreate, | ||
}; | ||
} else if (!config.db.pool) { | ||
poolConfig = { | ||
afterCreate, | ||
}; | ||
} | ||
} | ||
|
||
this.config = { | ||
...config, | ||
db: { | ||
...config.db, | ||
connection: { | ||
...config.db.connection, | ||
database: config.db.connection.database ?? config.database, | ||
}, | ||
pool: poolConfig, | ||
}, | ||
}; | ||
const knexConfig = { ...this.config.db, client: 'pg' }; | ||
|
||
// start the timer for cleanup | ||
this.startTimeout(); | ||
|
||
this.knex = knex(knexConfig); | ||
} | ||
|
||
startTimeout() { | ||
setTimeout(() => { | ||
this.knex.destroy(); | ||
evictPool(this.id); | ||
}, this.config.db.pool.idleTimeoutMillis ?? 30000); | ||
} | ||
} | ||
|
||
export type NileDatabaseI = (table?: string) => Knex; | ||
export default NileDatabase; |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,102 +1,2 @@ | ||
/* eslint-disable @typescript-eslint/no-explicit-any */ | ||
import knex, { Knex } from 'knex'; | ||
|
||
import { Config } from '../utils/Config'; | ||
|
||
// doing this now, to provide flexibility later | ||
class NileDatabase { | ||
knex: Knex; | ||
db: Knex; | ||
tenantId?: undefined | null | string; | ||
userId?: undefined | null | string; | ||
config: any; | ||
|
||
constructor(config: Config) { | ||
this.config = { ...config, client: 'pg' }; | ||
this.knex = knex(this.config); | ||
// Create a proxy to intercept method calls | ||
// @ts-expect-error - proxy, but knex | ||
this.db = new Proxy(this, { | ||
get: (target, method) => { | ||
if (method === 'tenantId') { | ||
return this.tenantId; | ||
} | ||
if (method === 'userId') { | ||
return this.userId; | ||
} | ||
if (method === 'db') { | ||
return (...args: any) => { | ||
//@ts-expect-error - its a string | ||
return target.knex.table(...args); | ||
}; | ||
} | ||
//@ts-expect-error - its a string | ||
if (typeof target.knex[method] === 'function') { | ||
return (...args: any) => { | ||
//@ts-expect-error - its a string | ||
return target.knex[method](...args); | ||
}; | ||
} else { | ||
return target.knex; | ||
} | ||
}, | ||
}); | ||
} | ||
ensureUpToDate() { | ||
// Close the existing pool connections and update the Knex instance with the latest config | ||
this.knex.destroy(); | ||
this.knex = knex({ ...this.config.db, client: 'pg' }); | ||
} | ||
setConfig(newConfig: Config) { | ||
const { tenantId, userId } = newConfig; | ||
this.tenantId = tenantId; | ||
this.userId = userId; | ||
let poolConfig = {}; | ||
const afterCreate = ( | ||
conn: { | ||
on: any; | ||
query: (query: string, cb: (err: unknown) => void) => void; | ||
}, | ||
done: (err: unknown, conn: unknown) => void | ||
) => { | ||
// console.log(this.tenantId, this.userId, 'in create'); | ||
const query = [`SET nile.tenant_id = '${this.tenantId}'`]; | ||
if (this.userId) { | ||
if (!this.tenantId) { | ||
// eslint-disable-next-line no-console | ||
console.warn( | ||
'A user id cannot be set in context without a tenant id' | ||
); | ||
} | ||
query.push(`SET nile.user_id = '${this.userId}'`); | ||
} | ||
// in this example we use pg driver's connection API | ||
conn.query(query.join(';'), function (err: unknown) { | ||
done(err, conn); | ||
}); | ||
}; | ||
if (this.tenantId) { | ||
if (newConfig.db.pool?.afterCreate) { | ||
// eslint-disable-next-line no-console | ||
console.log( | ||
'Providing an pool configuration will stop automatic tenant context setting.' | ||
); | ||
} else if (newConfig.db.pool) { | ||
poolConfig = { | ||
...newConfig.db.pool, | ||
afterCreate, | ||
}; | ||
} else if (!newConfig.db.pool) { | ||
poolConfig = { | ||
afterCreate, | ||
}; | ||
} | ||
} | ||
|
||
this.config = { ...newConfig, db: { ...newConfig.db, pool: poolConfig } }; | ||
this.ensureUpToDate(); | ||
} | ||
} | ||
|
||
export type NileDatabaseI = (table?: string) => Knex; | ||
export default NileDatabase; | ||
export { default } from './DBManager'; | ||
export { NileDatabaseI } from './NileInstance'; |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters