+238
-58
@@ -10,6 +10,7 @@ import (
|
||||
log "gitlab.com/gdulai/simpleloglvl"
|
||||
)
|
||||
|
||||
// Created and executes the DDL created by [simpleorm.ORM]
|
||||
func ExecuteDDL(conn *simpleorm.DBConnection, orm *simpleorm.ORM) {
|
||||
schema, err := orm.CreateSchema()
|
||||
if err != nil {
|
||||
@@ -24,14 +25,143 @@ func ExecuteDDL(conn *simpleorm.DBConnection, orm *simpleorm.ORM) {
|
||||
}
|
||||
}
|
||||
|
||||
// Wraps and represents a DB transaction.
|
||||
// Allows for all at once or separate execution of [exec.Exec] implementations.
|
||||
type Transaction struct {
|
||||
conn *simpleorm.DBConnection
|
||||
executions []*Exec[any]
|
||||
tx *sql.Tx
|
||||
finished bool
|
||||
}
|
||||
|
||||
// Creates a new [exec.Transaction].
|
||||
// Can be initialized with a set of [exec.Exec]s.
|
||||
func NewTransaction(conn *simpleorm.DBConnection, executions ...*Exec[any]) *Transaction {
|
||||
return &Transaction{conn: conn, executions: executions, finished: false}
|
||||
}
|
||||
|
||||
// The executions assigned to this transaction.
|
||||
func (t *Transaction) Executions() []*Exec[any] {
|
||||
return t.executions
|
||||
}
|
||||
|
||||
// Executes all the [exec.Exec] implementations assigned to this transactions.
|
||||
// Begins, commits or rollbacks the transaction.
|
||||
// This is a terminal operation, the [exec.Transaction] is considered finished after calling this.
|
||||
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.Rollback()
|
||||
if rollbackErr != nil {
|
||||
return rollbackErr
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
err := t.tx.Commit()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
t.finished = true
|
||||
return nil
|
||||
}
|
||||
|
||||
// Executes and assigns the passed [exec.Exec]s to the Transaction.
|
||||
// Calls rollback in case of an error.
|
||||
// This is NOT a terminal operation, transactions still has to be committed.
|
||||
func (t *Transaction) Execute(execs ...Exec[any]) error {
|
||||
if t.finished {
|
||||
return errors.New("Transaction already finished.")
|
||||
}
|
||||
|
||||
for _, exec := range execs {
|
||||
execPtr := &exec
|
||||
t.executions = append(t.executions, execPtr)
|
||||
if t.tx == nil {
|
||||
tx, err := t.conn.Begin()
|
||||
t.tx = tx
|
||||
if err != nil {
|
||||
t.finished = true
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
err := (*execPtr).execute(nil, t.tx)
|
||||
if err != nil {
|
||||
rollbackErr := t.Rollback()
|
||||
if rollbackErr != nil {
|
||||
return rollbackErr
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// Finishes the transaction, commits the changes.
|
||||
// This is terminal operation.
|
||||
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
|
||||
}
|
||||
|
||||
// Rollbacks and aborts the transaction.
|
||||
// This is a terminal operation.
|
||||
func (t *Transaction) Rollback() error {
|
||||
if t.finished {
|
||||
return errors.New("Transaction already finished.")
|
||||
}
|
||||
|
||||
if t.tx == nil {
|
||||
return errors.New("Transaction is nil.")
|
||||
}
|
||||
|
||||
rollbackErr := t.tx.Rollback()
|
||||
if rollbackErr != nil {
|
||||
return rollbackErr
|
||||
}
|
||||
t.finished = true
|
||||
return nil
|
||||
}
|
||||
|
||||
type Exec[T any] interface {
|
||||
Execute(conn *simpleorm.DBConnection) ([]T, error)
|
||||
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) {
|
||||
@@ -43,10 +173,20 @@ func CreateSelect[T any](orm *simpleorm.ORM, whereStmt string, args ...any) (Sel
|
||||
return Select[T]{target: *table, whereStmt: whereStmt, args: args}, nil
|
||||
}
|
||||
|
||||
func (s Select[T]) Execute(conn *simpleorm.DBConnection) ([]T, error) {
|
||||
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 []T{}, err
|
||||
return err
|
||||
}
|
||||
|
||||
if s.whereStmt != "" {
|
||||
@@ -55,9 +195,15 @@ func (s Select[T]) Execute(conn *simpleorm.DBConnection) ([]T, error) {
|
||||
|
||||
log.LogDebug("Preparing sql: %s, with args: %s", dml, s.args)
|
||||
|
||||
stmt, err := conn.Prepare(dml)
|
||||
var stmt *sql.Stmt
|
||||
if tx != nil {
|
||||
stmt, err = tx.Prepare(dml)
|
||||
} else {
|
||||
stmt, err = conn.Prepare(dml)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return err
|
||||
}
|
||||
defer stmt.Close()
|
||||
|
||||
@@ -81,40 +227,21 @@ func (s Select[T]) Execute(conn *simpleorm.DBConnection) ([]T, error) {
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return err
|
||||
}
|
||||
|
||||
rowContainer := createSelectResultContainer(s.target)
|
||||
|
||||
var results []T
|
||||
|
||||
for rows.Next() {
|
||||
err = rows.Scan(rowContainer...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
targetType := s.target.Type
|
||||
parsedResult := reflect.New(targetType)
|
||||
for i, fieldVal := range rowContainer {
|
||||
col := s.target.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)
|
||||
s.results, err = readRows[T](s.target, rows)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return results, nil
|
||||
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) {
|
||||
@@ -126,10 +253,18 @@ func NewInsert[T any](orm *simpleorm.ORM, toInsert ...T) (Insert[T], error) {
|
||||
return Insert[T]{target: *table, toInsert: toInsert}, nil
|
||||
}
|
||||
|
||||
func (ins Insert[T]) Execute(conn *simpleorm.DBConnection) ([]T, error) {
|
||||
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 []T{}, err
|
||||
return err
|
||||
}
|
||||
|
||||
var params []any
|
||||
@@ -145,21 +280,26 @@ func (ins Insert[T]) Execute(conn *simpleorm.DBConnection) ([]T, error) {
|
||||
|
||||
log.LogInfo("%s [%s]", dml, params)
|
||||
|
||||
stmt, err := conn.Prepare(dml)
|
||||
if err != nil {
|
||||
return []T{}, err
|
||||
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...)
|
||||
rows, err := stmt.Query(params...)
|
||||
|
||||
if err != nil {
|
||||
} else {
|
||||
rowsAffected, _ := result.RowsAffected()
|
||||
log.LogInfo("Inserted %s row", rowsAffected)
|
||||
lastInsertId, _ := result.LastInsertId()
|
||||
log.LogInfo("Last ID: %s", lastInsertId)
|
||||
return err
|
||||
}
|
||||
return []T{}, nil
|
||||
|
||||
ins.results, err = readRows[T](ins.target, rows)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
type Update[T any] struct {
|
||||
@@ -176,17 +316,21 @@ func NewUpdate[T any](orm *simpleorm.ORM, toUpdate T) (Update[T], error) {
|
||||
return Update[T]{target: *table, toUpdate: toUpdate}, nil
|
||||
}
|
||||
|
||||
func (u Update[T]) Execute(conn *simpleorm.DBConnection) ([]T, error) {
|
||||
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 []T{}, err
|
||||
return err
|
||||
}
|
||||
|
||||
params := prepareParams(u.toUpdate, u.target)
|
||||
|
||||
pkCols, err := getPk(u.toUpdate, u.target)
|
||||
if err != nil {
|
||||
return []T{}, err
|
||||
return err
|
||||
}
|
||||
// Put the pks back at the end
|
||||
for _, pk := range pkCols {
|
||||
@@ -195,20 +339,22 @@ func (u Update[T]) Execute(conn *simpleorm.DBConnection) ([]T, error) {
|
||||
|
||||
log.LogInfo("%s [%s]", dml, params)
|
||||
|
||||
stmt, err := conn.Prepare(dml)
|
||||
if err != nil {
|
||||
return []T{}, err
|
||||
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 []T{}, err
|
||||
return err
|
||||
} else {
|
||||
rowsAffected, _ := result.RowsAffected()
|
||||
log.LogInfo("Updated %s row", rowsAffected)
|
||||
}
|
||||
return []T{u.toUpdate}, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
type Delete[T any] struct {
|
||||
@@ -225,11 +371,15 @@ func NewDelete[T any](orm *simpleorm.ORM, toDelete []T) (Delete[T], error) {
|
||||
return Delete[T]{target: *table, toDelete: toDelete}, nil
|
||||
}
|
||||
|
||||
func (d Delete[T]) Execute(conn *simpleorm.DBConnection) ([]T, error) {
|
||||
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 []T{}, err
|
||||
return err
|
||||
}
|
||||
|
||||
var params []any
|
||||
@@ -237,7 +387,7 @@ func (d Delete[T]) Execute(conn *simpleorm.DBConnection) ([]T, error) {
|
||||
for _, del := range d.toDelete {
|
||||
pkCols, err := getPk(del, d.target)
|
||||
if err != nil {
|
||||
return []T{}, err
|
||||
return err
|
||||
}
|
||||
|
||||
for _, pk := range pkCols {
|
||||
@@ -247,21 +397,23 @@ func (d Delete[T]) Execute(conn *simpleorm.DBConnection) ([]T, error) {
|
||||
|
||||
log.LogInfo("%s [%s]", dml, params)
|
||||
|
||||
stmt, err := conn.Prepare(dml)
|
||||
if err != nil {
|
||||
return []T{}, err
|
||||
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 []T{}, err
|
||||
return err
|
||||
} else {
|
||||
rowsAffected, _ := result.RowsAffected()
|
||||
log.LogInfo("Deleted %s row", rowsAffected)
|
||||
}
|
||||
|
||||
return d.toDelete, nil
|
||||
return nil
|
||||
|
||||
}
|
||||
|
||||
@@ -322,3 +474,31 @@ func getPk(src any, t schema.Table) ([]any, error) {
|
||||
|
||||
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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user