feat: 多租户数据权限 + 数据库连接修正
- adminPool 连服务器本地 5432 sbrain_admin (SCRAM 认证) - tenantPool/pool 连 FRP 隧道 15432 bill_query (trust 认证) - query() 自动注入 store_code 数据权限 (store/regional 角色) - 区域经理 region 字段存储逗号分隔的 store_code 列表 - /overview 和 /overview/daily 对 scoped 用户从 bill_fact 聚合 - tenant_users.region 扩展为 VARCHAR(500) - deploy.sh .env 增加 ADMIN_DB_* 配置 - run.md 更新数据库连接架构说明 - 新增海淀区区域经理用户 (21 家海淀区门店)
This commit is contained in:
@@ -0,0 +1,24 @@
|
||||
import pg from 'pg'
|
||||
|
||||
const { Pool } = pg
|
||||
|
||||
const adminPoolConfig: any = {
|
||||
host: process.env.ADMIN_DB_HOST || '127.0.0.1',
|
||||
port: parseInt(process.env.ADMIN_DB_PORT || '5432'),
|
||||
database: process.env.ADMIN_DB_NAME || 'sbrain_admin',
|
||||
user: process.env.ADMIN_DB_USER || 'sbrain_admin',
|
||||
max: 5,
|
||||
}
|
||||
|
||||
const adminDbPassword = process.env.ADMIN_DB_PASSWORD
|
||||
if (adminDbPassword) {
|
||||
adminPoolConfig.password = adminDbPassword
|
||||
}
|
||||
|
||||
const adminPool = new Pool(adminPoolConfig)
|
||||
|
||||
adminPool.on('error', (err) => {
|
||||
console.error('Unexpected error on admin pool', err)
|
||||
})
|
||||
|
||||
export default adminPool
|
||||
@@ -1,15 +1,19 @@
|
||||
import pg from 'pg'
|
||||
import { tenantContextStorage } from './tenant-db.js'
|
||||
|
||||
const { Pool } = pg
|
||||
|
||||
const pool = new Pool({
|
||||
const poolConfig: any = {
|
||||
host: process.env.DB_HOST || 'localhost',
|
||||
port: parseInt(process.env.DB_PORT || '5432'),
|
||||
database: process.env.DB_NAME || 'bill_query',
|
||||
user: process.env.DB_USER || 'freedak',
|
||||
password: process.env.DB_PASSWORD || '',
|
||||
max: parseInt(process.env.DB_POOL_MAX || '10'),
|
||||
})
|
||||
}
|
||||
if (process.env.DB_PASSWORD) {
|
||||
poolConfig.password = process.env.DB_PASSWORD
|
||||
}
|
||||
const pool = new Pool(poolConfig)
|
||||
|
||||
pool.on('error', (err) => {
|
||||
console.error('Unexpected error on idle client', err)
|
||||
@@ -20,18 +24,61 @@ export interface QueryResult<T = any> {
|
||||
rowCount: number | null
|
||||
}
|
||||
|
||||
export async function query<T = any>(text: string, params?: any[]): Promise<QueryResult<T>> {
|
||||
export async function query<T = any>(text: string, params?: any[], opts?: { skipScope?: boolean }): Promise<QueryResult<T>> {
|
||||
const ctx = tenantContextStorage.getStore()
|
||||
const usePool = ctx ? ctx.pool : pool
|
||||
const start = Date.now()
|
||||
const res = await pool.query(text, params)
|
||||
|
||||
let sql = text
|
||||
let sqlParams = params || []
|
||||
|
||||
if (!opts?.skipScope && ctx?.scope && (ctx.scope.role === 'store' || ctx.scope.role === 'regional')) {
|
||||
const sqlLower = sql.toLowerCase()
|
||||
const hasStoreRef = /store_code|store_name|v_store_|mv_store_|bill_fact|fact_bill|store_task|dim_store/.test(sqlLower)
|
||||
if (!hasStoreRef) {
|
||||
const res = await usePool.query(sql, sqlParams)
|
||||
const duration = Date.now() - start
|
||||
if (duration > 500) {
|
||||
console.warn(`Slow query (${duration}ms):`, sql.substring(0, 100))
|
||||
}
|
||||
return res
|
||||
}
|
||||
|
||||
const scopeClause = ctx.scope.role === 'store' && ctx.scope.storeCode
|
||||
? ` store_code = $${sqlParams.length + 1}`
|
||||
: ctx.scope.role === 'regional' && ctx.scope.region
|
||||
? ` store_code = ANY($${sqlParams.length + 1})`
|
||||
: null
|
||||
|
||||
if (scopeClause) {
|
||||
const scopeValue = ctx.scope.role === 'store'
|
||||
? ctx.scope.storeCode
|
||||
: ctx.scope.region!.split(',').map(s => s.trim()).filter(Boolean)
|
||||
sqlParams = [...sqlParams, scopeValue]
|
||||
if (sql.includes('WHERE')) {
|
||||
sql = sql.replace('WHERE', `WHERE${scopeClause} AND`)
|
||||
} else if (sql.includes('GROUP BY')) {
|
||||
sql = sql.replace('GROUP BY', `WHERE${scopeClause} GROUP BY`)
|
||||
} else if (sql.includes('ORDER BY')) {
|
||||
sql = sql.replace('ORDER BY', `WHERE${scopeClause} ORDER BY`)
|
||||
} else {
|
||||
sql = sql + ` WHERE${scopeClause}`
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const res = await usePool.query(sql, sqlParams)
|
||||
const duration = Date.now() - start
|
||||
if (duration > 500) {
|
||||
console.warn(`Slow query (${duration}ms):`, text.substring(0, 100))
|
||||
console.warn(`Slow query (${duration}ms):`, sql.substring(0, 100))
|
||||
}
|
||||
return res
|
||||
}
|
||||
|
||||
export async function withTransaction<T>(callback: (client: pg.PoolClient) => Promise<T>): Promise<T> {
|
||||
const client = await pool.connect()
|
||||
const ctx = tenantContextStorage.getStore()
|
||||
const usePool = ctx ? ctx.pool : pool
|
||||
const client = await usePool.connect()
|
||||
try {
|
||||
await client.query('BEGIN')
|
||||
const result = await callback(client)
|
||||
|
||||
@@ -0,0 +1,74 @@
|
||||
import { AsyncLocalStorage } from 'node:async_hooks'
|
||||
import pg from 'pg'
|
||||
|
||||
const { Pool } = pg
|
||||
|
||||
export interface TenantContext {
|
||||
tenantId: string
|
||||
pool: pg.Pool
|
||||
scope?: { role: string; storeCode?: string; region?: string } | null
|
||||
}
|
||||
|
||||
export const tenantContextStorage = new AsyncLocalStorage<TenantContext>()
|
||||
|
||||
const poolCache = new Map<string, { pool: pg.Pool; lastUsed: number }>()
|
||||
const IDLE_TIMEOUT = 10 * 60 * 1000
|
||||
|
||||
export function getTenantPool(config: {
|
||||
tenantId: string
|
||||
dbHost: string
|
||||
dbPort: number
|
||||
dbName: string
|
||||
dbUser: string
|
||||
dbPassword?: string
|
||||
}): pg.Pool {
|
||||
const key = config.tenantId
|
||||
const cached = poolCache.get(key)
|
||||
if (cached) {
|
||||
cached.lastUsed = Date.now()
|
||||
return cached.pool
|
||||
}
|
||||
const poolConfig: any = {
|
||||
host: config.dbHost,
|
||||
port: config.dbPort,
|
||||
database: config.dbName,
|
||||
user: config.dbUser,
|
||||
max: 5,
|
||||
idleTimeoutMillis: 30000,
|
||||
}
|
||||
if (config.dbPassword) {
|
||||
poolConfig.password = config.dbPassword
|
||||
}
|
||||
const pool = new Pool(poolConfig)
|
||||
pool.on('error', (err) => {
|
||||
console.error(`Tenant pool [${key}] error:`, err)
|
||||
closeTenantPool(key)
|
||||
})
|
||||
poolCache.set(key, { pool, lastUsed: Date.now() })
|
||||
return pool
|
||||
}
|
||||
|
||||
export function closeTenantPool(tenantId: string) {
|
||||
const cached = poolCache.get(tenantId)
|
||||
if (cached) {
|
||||
cached.pool.end()
|
||||
poolCache.delete(tenantId)
|
||||
}
|
||||
}
|
||||
|
||||
export function closeAllTenantPools() {
|
||||
for (const [, cached] of poolCache) {
|
||||
cached.pool.end()
|
||||
}
|
||||
poolCache.clear()
|
||||
}
|
||||
|
||||
setInterval(() => {
|
||||
const now = Date.now()
|
||||
for (const [id, cached] of poolCache) {
|
||||
if (now - cached.lastUsed > IDLE_TIMEOUT) {
|
||||
console.log(`Closing idle tenant pool [${id}]`)
|
||||
closeTenantPool(id)
|
||||
}
|
||||
}
|
||||
}, 5 * 60 * 1000)
|
||||
Reference in New Issue
Block a user