Files
db/engine/engine.go
T
what 21b80bdea4 feat: 完善扫描器、exec 及 schema 相关功能
- exec/scanner: 用 *interface{} 替换 **json.RawMessage 扫描目标,兼容 DuckDB 返回 map[string]interface{} 的场景;新增 toJSONRawMessage 转换函数
- exec/scanner: ScanVal 支持结构体指针,通过 JSON 中间层转换(DuckDB STRUCT 列)
- exec/scanner: 将 *sql.RawBytes 和 *[]byte 的处理从 ScanValContext 移入 scanner.ScanVal
- exec/query_executor: 简化 ScanValContext,移除私有 scan 方法
- exec: 补充 scanner 级别 ScanVal 测试用例
- internal/util/reflect: 重写 SafeSetVarValue,修复非指针 src 及 nil 指针字段的 panic
- internal/util/column_map: 恢复非匿名带标签结构体字段的展开逻辑
- schema: 新增 vector 列类型支持
- engine: 补充 DuckDB 相关配置
- dialect/sqlite3/vtab: 完善虚拟表适配器
- 各方言测试改用 sqlmock 虚拟连接
2026-05-20 17:52:28 +08:00

139 lines
3.2 KiB
Go

package engine
import (
"database/sql"
"fmt"
"log"
"git.fsdpf.net/go/db"
sqlite3dialect "git.fsdpf.net/go/db/dialect/sqlite3"
sqlite3vtab "git.fsdpf.net/go/db/dialect/sqlite3/vtab"
)
type Engine struct {
configs map[string]DBConfig
dbs map[string]*sql.DB
}
var _engine *Engine
func init() {
_engine = &Engine{
configs: make(map[string]DBConfig),
dbs: make(map[string]*sql.DB),
}
}
func (e Engine) Connection(name string) *db.Database {
cfg, ok := e.configs[name]
if !ok {
panic(fmt.Errorf("database connection %s not configured", name))
}
_db, ok := e.dbs[name]
if !ok {
_db = e.MakeConnection(cfg)
e.dbs[name] = _db
// vtable 连接:同时创建一个 :memory: 连接供虚拟表查询执行,
// 原文件连接仅用于 _vtab_cache 持久化,两者互不阻塞。
if cfg.Driver == "vtable" {
// 使用命名共享内存数据库,确保连接池中所有连接共享同一份内存数据,
// 避免匿名 :memory: 各连接独立导致虚表在其他连接不可见的问题。
sharedDSN := fmt.Sprintf("file:%s?mode=memory&cache=shared", name)
memDB, err := sql.Open(sqlite3vtab.DriverName, sharedDSN)
if err != nil {
panic(fmt.Sprintf("vtable: open memory connection for %s: %v", name, err))
}
e.dbs["__"+name] = memDB
e.configs["__"+name] = DBConfig{Driver: "vtable"}
}
}
return db.New(cfg.Driver, _db)
}
func (e Engine) MakeConnection(cfg DBConfig) (db *sql.DB) {
dsn := cfg.ToDSN()
driverName := cfg.Driver
switch cfg.Driver {
case "mysql":
case "sqlite3":
driverName = sqlite3dialect.DriverWithIF
case "vtable":
driverName = sqlite3vtab.DriverName
case "sqlserver":
case "postgres":
case "duckdb":
default:
panic(fmt.Sprintf("Unsupported driver: %s", cfg.Driver))
}
db, err := sql.Open(driverName, dsn)
if err != nil {
panic(err)
}
if err := db.Ping(); err != nil {
panic(err)
}
if cfg.Driver == "duckdb" {
for _, ext := range cfg.DuckDB.Extensions {
if _, err := db.Exec("INSTALL " + ext); err != nil {
panic(fmt.Sprintf("duckdb: install extension %q: %v", ext, err))
}
if _, err := db.Exec("LOAD " + ext); err != nil {
panic(fmt.Sprintf("duckdb: load extension %q: %v", ext, err))
}
}
}
// 应用连接池配置(对所有驱动生效)
if cfg.MaxOpenConns > 0 {
db.SetMaxOpenConns(cfg.MaxOpenConns)
}
if cfg.MaxIdleConns > 0 {
db.SetMaxIdleConns(cfg.MaxIdleConns)
}
if cfg.ConnMaxLifetime > 0 {
db.SetConnMaxLifetime(cfg.ConnMaxLifetime)
}
if cfg.ConnMaxIdleTime > 0 {
db.SetConnMaxIdleTime(cfg.ConnMaxIdleTime)
}
return db
}
// Shutdown 关闭所有已建立的数据库连接,使 DuckDB 等驱动得以完成 WAL checkpoint。
// 实现 do.Shutdownable 接口,由 DI 容器在关闭时调用。
func (e Engine) Shutdown() error {
for name, sqlDB := range e.dbs {
if err := sqlDB.Close(); err != nil {
log.Printf("engine: closing connection %q: %v", name, err)
}
}
return nil
}
func Open(cfgs map[string]DBConfig) Engine {
for n, cfg := range cfgs {
_engine.configs[n] = cfg
}
return *_engine
}
func Mock(cfgs map[string]MockDBConfig) Engine {
for k, cfg := range cfgs {
_engine.dbs[k] = cfg.Mock
_engine.configs[k] = DBConfig{Driver: cfg.Driver}
}
return *_engine
}