470 lines
9.4 KiB
Go
470 lines
9.4 KiB
Go
package exec
|
|
|
|
import (
|
|
"database/sql"
|
|
"errors"
|
|
"reflect"
|
|
|
|
"git.gdulai.com/gdulai/simpleorm"
|
|
"git.gdulai.com/gdulai/simpleorm/schema"
|
|
log "gitlab.com/gdulai/simpleloglvl"
|
|
)
|
|
|
|
func ExecuteDDL(conn *simpleorm.DBConnection, orm *simpleorm.ORM) {
|
|
schema, err := orm.CreateSchema()
|
|
if err != nil {
|
|
log.LogError("Failed to run DDL: %s", schema)
|
|
return
|
|
}
|
|
|
|
log.LogInfo("Executing DDL:\n%s", schema)
|
|
|
|
if _, err := conn.Exec(schema); err != nil {
|
|
log.LogFatalError("%", err)
|
|
}
|
|
}
|
|
|
|
type Transaction struct {
|
|
conn *simpleorm.DBConnection
|
|
executions []*Exec[any]
|
|
tx *sql.Tx
|
|
finished bool
|
|
}
|
|
|
|
func (t *Transaction) NewTransaction(conn *simpleorm.DBConnection, executions ...*Exec[any]) *Transaction {
|
|
return &Transaction{conn: conn, executions: executions, finished: false}
|
|
}
|
|
|
|
func (t *Transaction) Executions() []*Exec[any] {
|
|
return t.executions
|
|
}
|
|
|
|
func (t *Transaction) ExecuteAtOnce() error {
|
|
if t.finished {
|
|
return errors.New("Transaction already finished.")
|
|
}
|
|
|
|
if t.tx == nil {
|
|
tx, err := t.conn.Begin()
|
|
t.tx = tx
|
|
if err != nil {
|
|
t.finished = true
|
|
return err
|
|
}
|
|
}
|
|
|
|
for _, exec := range t.executions {
|
|
err := (*exec).execute(nil, t.tx)
|
|
if err != nil {
|
|
rollbackErr := t.tx.Rollback()
|
|
if rollbackErr != nil {
|
|
return rollbackErr
|
|
}
|
|
t.finished = true
|
|
return err
|
|
}
|
|
|
|
}
|
|
|
|
err := t.tx.Commit()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
t.finished = true
|
|
return nil
|
|
}
|
|
|
|
func (t *Transaction) Execute(exec *Exec[any]) error {
|
|
if t.finished {
|
|
return errors.New("Transaction already finished.")
|
|
}
|
|
|
|
t.executions = append(t.executions, exec)
|
|
if t.tx == nil {
|
|
tx, err := t.conn.Begin()
|
|
t.tx = tx
|
|
if err != nil {
|
|
t.finished = true
|
|
return err
|
|
}
|
|
}
|
|
|
|
err := (*exec).execute(nil, t.tx)
|
|
if err != nil {
|
|
rollbackErr := t.tx.Rollback()
|
|
if rollbackErr != nil {
|
|
return rollbackErr
|
|
}
|
|
t.finished = true
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (t *Transaction) Finish() error {
|
|
if t.finished {
|
|
return errors.New("Transaction already finished.")
|
|
}
|
|
|
|
if t.tx == nil {
|
|
return errors.New("Transaction is nil.")
|
|
}
|
|
|
|
t.tx.Commit()
|
|
t.finished = true
|
|
return nil
|
|
}
|
|
|
|
type Exec[T any] interface {
|
|
Execute(conn *simpleorm.DBConnection) error
|
|
execute(conn *simpleorm.DBConnection, tx *sql.Tx) error
|
|
}
|
|
|
|
type Select[T any] struct {
|
|
target schema.Table
|
|
whereStmt string
|
|
args []any
|
|
results []T
|
|
}
|
|
|
|
func CreateSelect[T any](orm *simpleorm.ORM, whereStmt string, args ...any) (Select[T], error) {
|
|
table, ok := orm.Cache().Get(reflect.TypeFor[T]().Name())
|
|
if !ok {
|
|
return Select[T]{}, errors.New("Failed to get table from schema cache")
|
|
}
|
|
|
|
return Select[T]{target: *table, whereStmt: whereStmt, args: args}, nil
|
|
}
|
|
|
|
func (s *Select[T]) Results() []T {
|
|
return s.results
|
|
}
|
|
|
|
func (s *Select[T]) Execute(conn *simpleorm.DBConnection) error {
|
|
return s.execute(conn, nil)
|
|
}
|
|
|
|
func (s *Select[T]) execute(conn *simpleorm.DBConnection, tx *sql.Tx) error {
|
|
// Reinit the results, new execution
|
|
s.results = []T{}
|
|
dml, err := s.target.GetSelectDML()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if s.whereStmt != "" {
|
|
dml += " WHERE " + s.whereStmt
|
|
}
|
|
|
|
log.LogDebug("Preparing sql: %s, with args: %s", dml, s.args)
|
|
|
|
var stmt *sql.Stmt
|
|
if tx != nil {
|
|
stmt, err = tx.Prepare(dml)
|
|
} else {
|
|
stmt, err = conn.Prepare(dml)
|
|
}
|
|
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer stmt.Close()
|
|
|
|
log.LogDebug("Executing statement: %s", stmt)
|
|
|
|
var rows *sql.Rows
|
|
if len(s.args) == 0 {
|
|
rows, err = stmt.Query()
|
|
} else {
|
|
// Flattent args to make sure it can be parsed correctly
|
|
var flatArgs []any
|
|
for _, a := range s.args {
|
|
if s, ok := a.([]any); ok {
|
|
flatArgs = append(flatArgs, s...)
|
|
} else {
|
|
flatArgs = append(flatArgs, a)
|
|
}
|
|
}
|
|
rows, err = stmt.Query(flatArgs...)
|
|
|
|
}
|
|
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
s.results, err = readRows[T](s.target, rows)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
type Insert[T any] struct {
|
|
target schema.Table
|
|
toInsert []T
|
|
results []T
|
|
}
|
|
|
|
func NewInsert[T any](orm *simpleorm.ORM, toInsert ...T) (Insert[T], error) {
|
|
table, ok := orm.Cache().Get(reflect.TypeFor[T]().Name())
|
|
if !ok {
|
|
return Insert[T]{}, errors.New("Failed to get table from schema cache")
|
|
}
|
|
|
|
return Insert[T]{target: *table, toInsert: toInsert}, nil
|
|
}
|
|
|
|
func (ins *Insert[T]) Results() []T {
|
|
return ins.results
|
|
}
|
|
|
|
func (ins *Insert[T]) Execute(conn *simpleorm.DBConnection) error {
|
|
return ins.execute(conn, nil)
|
|
}
|
|
|
|
func (ins *Insert[T]) execute(conn *simpleorm.DBConnection, tx *sql.Tx) error {
|
|
dml, err := ins.target.GetInsertDML(len(ins.toInsert))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
var params []any
|
|
for i := range len(ins.toInsert) {
|
|
actualParams := prepareParams(ins.toInsert[i], ins.target)
|
|
if len(params) == 0 {
|
|
params = make([]any, len(ins.toInsert)*len(actualParams))
|
|
}
|
|
for j := range actualParams {
|
|
params[(i*len(actualParams))+j] = actualParams[j]
|
|
}
|
|
}
|
|
|
|
log.LogInfo("%s [%s]", dml, params)
|
|
|
|
var stmt *sql.Stmt
|
|
if tx != nil {
|
|
stmt, err = tx.Prepare(dml)
|
|
} else {
|
|
stmt, err = conn.Prepare(dml)
|
|
}
|
|
defer stmt.Close()
|
|
|
|
rows, err := stmt.Query(params...)
|
|
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
ins.results, err = readRows[T](ins.target, rows)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
type Update[T any] struct {
|
|
target schema.Table
|
|
toUpdate T
|
|
}
|
|
|
|
func NewUpdate[T any](orm *simpleorm.ORM, toUpdate T) (Update[T], error) {
|
|
table, ok := orm.Cache().Get(reflect.TypeFor[T]().Name())
|
|
if !ok {
|
|
return Update[T]{}, errors.New("Failed to get table from schema cache")
|
|
}
|
|
|
|
return Update[T]{target: *table, toUpdate: toUpdate}, nil
|
|
}
|
|
|
|
func (u Update[T]) Execute(conn *simpleorm.DBConnection) error {
|
|
return u.execute(conn, nil)
|
|
}
|
|
|
|
func (u Update[T]) execute(conn *simpleorm.DBConnection, tx *sql.Tx) error {
|
|
dml, err := u.target.GetUpdateDML()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
params := prepareParams(u.toUpdate, u.target)
|
|
|
|
pkCols, err := getPk(u.toUpdate, u.target)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// Put the pks back at the end
|
|
for _, pk := range pkCols {
|
|
params = append(params, pk)
|
|
}
|
|
|
|
log.LogInfo("%s [%s]", dml, params)
|
|
|
|
var stmt *sql.Stmt
|
|
if tx != nil {
|
|
stmt, err = tx.Prepare(dml)
|
|
} else {
|
|
stmt, err = conn.Prepare(dml)
|
|
}
|
|
defer stmt.Close()
|
|
|
|
result, err := stmt.Exec(params...)
|
|
if err != nil {
|
|
return err
|
|
} else {
|
|
rowsAffected, _ := result.RowsAffected()
|
|
log.LogInfo("Updated %s row", rowsAffected)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
type Delete[T any] struct {
|
|
target schema.Table
|
|
toDelete []T
|
|
}
|
|
|
|
func NewDelete[T any](orm *simpleorm.ORM, toDelete []T) (Delete[T], error) {
|
|
table, ok := orm.Cache().Get(reflect.TypeFor[T]().Name())
|
|
if !ok {
|
|
return Delete[T]{}, errors.New("Failed to get table from schema cache")
|
|
}
|
|
|
|
return Delete[T]{target: *table, toDelete: toDelete}, nil
|
|
}
|
|
|
|
func (d *Delete[T]) Execute(conn *simpleorm.DBConnection) error {
|
|
return d.execute(conn, nil)
|
|
}
|
|
|
|
func (d *Delete[T]) execute(conn *simpleorm.DBConnection, tx *sql.Tx) error {
|
|
count := len(d.toDelete)
|
|
dml, err := d.target.GetDeleteDML(count)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
var params []any
|
|
|
|
for _, del := range d.toDelete {
|
|
pkCols, err := getPk(del, d.target)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
for _, pk := range pkCols {
|
|
params = append(params, pk)
|
|
}
|
|
}
|
|
|
|
log.LogInfo("%s [%s]", dml, params)
|
|
|
|
var stmt *sql.Stmt
|
|
if tx != nil {
|
|
stmt, err = tx.Prepare(dml)
|
|
} else {
|
|
stmt, err = conn.Prepare(dml)
|
|
}
|
|
defer stmt.Close()
|
|
|
|
result, err := stmt.Exec(params...)
|
|
if err != nil {
|
|
return err
|
|
} else {
|
|
rowsAffected, _ := result.RowsAffected()
|
|
log.LogInfo("Deleted %s row", rowsAffected)
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
func createSelectResultContainer(t schema.Table) []any {
|
|
vals := make([]any, t.Type.NumField())
|
|
for i := range vals {
|
|
switch t.Type.Field(i).Type.Kind().String() {
|
|
case "string":
|
|
var fieldContainer string
|
|
vals[i] = &fieldContainer
|
|
case "int", "bool":
|
|
var fieldContainer int
|
|
vals[i] = &fieldContainer
|
|
case "int64", "time.Time":
|
|
var fieldContainer int64
|
|
vals[i] = &fieldContainer
|
|
}
|
|
}
|
|
return vals
|
|
}
|
|
|
|
func prepareParams(src any, t schema.Table) []any {
|
|
var params []any
|
|
for _, col := range t.Columns {
|
|
_, ok := col.Modifiers["pk"]
|
|
if ok && t.IsPkAuto() {
|
|
continue
|
|
}
|
|
field, ok := t.Type.FieldByName(col.FieldName)
|
|
if !ok {
|
|
continue
|
|
}
|
|
fieldValue := reflect.ValueOf(src).FieldByIndex(field.Index)
|
|
|
|
params = append(params, col.Encode(fieldValue))
|
|
}
|
|
return params
|
|
}
|
|
|
|
func getPk(src any, t schema.Table) ([]any, error) {
|
|
var values []any
|
|
for _, constraint := range t.Constraints {
|
|
if constraint.Type == "pk" {
|
|
for _, col := range constraint.Columns {
|
|
field, ok := t.Type.FieldByName(col.FieldName)
|
|
if !ok {
|
|
continue
|
|
|
|
}
|
|
fieldValue := reflect.ValueOf(src).FieldByIndex(field.Index)
|
|
values = append(values, fieldValue.Interface())
|
|
}
|
|
}
|
|
}
|
|
if len(values) == 0 {
|
|
return nil, errors.New("Could not determine pk column!")
|
|
}
|
|
|
|
return values, nil
|
|
}
|
|
|
|
func readRows[T any](table schema.Table, rows *sql.Rows) ([]T, error) {
|
|
rowContainer := createSelectResultContainer(table)
|
|
|
|
var results []T
|
|
|
|
for rows.Next() {
|
|
err := rows.Scan(rowContainer...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
targetType := table.Type
|
|
parsedResult := reflect.New(targetType)
|
|
for i, fieldVal := range rowContainer {
|
|
col := table.Columns[i]
|
|
targetField := parsedResult.Elem().Field(i)
|
|
|
|
rawValue := reflect.Indirect(reflect.ValueOf(fieldVal))
|
|
decoded := reflect.ValueOf(col.Decode(targetField.Type().Name(), rawValue))
|
|
|
|
targetField.Set(decoded)
|
|
}
|
|
parsedObj := reflect.Indirect(parsedResult).Interface().(T)
|
|
results = append(results, parsedObj)
|
|
}
|
|
return results, nil
|
|
}
|