2022-05-19 07:49:32 +00:00
|
|
|
package txn
|
|
|
|
|
2022-08-11 06:14:57 +00:00
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"fmt"
|
|
|
|
)
|
2022-05-19 07:49:32 +00:00
|
|
|
|
|
|
|
type Manager interface {
|
2022-11-20 19:49:10 +00:00
|
|
|
Begin(ctx context.Context, exclusive bool) (context.Context, error)
|
2022-05-19 07:49:32 +00:00
|
|
|
Commit(ctx context.Context) error
|
|
|
|
Rollback(ctx context.Context) error
|
2022-07-13 06:30:54 +00:00
|
|
|
|
2022-08-11 06:14:57 +00:00
|
|
|
IsLocked(err error) bool
|
2022-07-13 06:30:54 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
type DatabaseProvider interface {
|
|
|
|
WithDatabase(ctx context.Context) (context.Context, error)
|
2022-05-19 07:49:32 +00:00
|
|
|
}
|
|
|
|
|
2022-11-20 19:49:10 +00:00
|
|
|
type DatabaseProviderManager interface {
|
|
|
|
DatabaseProvider
|
|
|
|
Manager
|
|
|
|
}
|
|
|
|
|
2022-05-19 07:49:32 +00:00
|
|
|
type TxnFunc func(ctx context.Context) error
|
|
|
|
|
2022-07-13 06:30:54 +00:00
|
|
|
// WithTxn executes fn in a transaction. If fn returns an error then
|
|
|
|
// the transaction is rolled back. Otherwise it is committed.
|
2022-11-20 19:49:10 +00:00
|
|
|
// Transaction is exclusive. Only one thread may run a transaction
|
|
|
|
// using this function at a time. This function will wait until the
|
|
|
|
// lock is available before executing.
|
|
|
|
// This function should be used for making changes to the database.
|
2022-05-19 07:49:32 +00:00
|
|
|
func WithTxn(ctx context.Context, m Manager, fn TxnFunc) error {
|
2022-11-20 19:49:10 +00:00
|
|
|
const (
|
|
|
|
execComplete = true
|
|
|
|
exclusive = true
|
|
|
|
)
|
|
|
|
return withTxn(ctx, m, fn, exclusive, execComplete)
|
|
|
|
}
|
|
|
|
|
|
|
|
// WithReadTxn executes fn in a transaction. If fn returns an error then
|
|
|
|
// the transaction is rolled back. Otherwise it is committed.
|
|
|
|
// Transaction is not exclusive and does not enforce read-only restrictions.
|
|
|
|
// Multiple threads can run transactions using this function concurrently,
|
|
|
|
// but concurrent writes may result in locked database error.
|
|
|
|
func WithReadTxn(ctx context.Context, m Manager, fn TxnFunc) error {
|
|
|
|
const (
|
|
|
|
execComplete = true
|
|
|
|
exclusive = false
|
|
|
|
)
|
|
|
|
return withTxn(ctx, m, fn, exclusive, execComplete)
|
2022-09-28 06:08:00 +00:00
|
|
|
}
|
|
|
|
|
2023-02-23 03:38:02 +00:00
|
|
|
func withTxn(outerCtx context.Context, m Manager, fn TxnFunc, exclusive bool, execCompleteOnLocked bool) error {
|
|
|
|
ctx, err := begin(outerCtx, m, exclusive)
|
2022-05-19 07:49:32 +00:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
defer func() {
|
|
|
|
if p := recover(); p != nil {
|
|
|
|
// a panic occurred, rollback and repanic
|
2023-02-23 03:38:02 +00:00
|
|
|
rollback(ctx, outerCtx, m)
|
2022-05-19 07:49:32 +00:00
|
|
|
panic(p)
|
|
|
|
}
|
|
|
|
|
|
|
|
if err != nil {
|
|
|
|
// something went wrong, rollback
|
2023-02-23 03:38:02 +00:00
|
|
|
rollback(ctx, outerCtx, m)
|
2022-09-28 06:08:00 +00:00
|
|
|
|
|
|
|
if execCompleteOnLocked || !m.IsLocked(err) {
|
2023-02-23 03:38:02 +00:00
|
|
|
executePostCompleteHooks(ctx, outerCtx)
|
2022-09-28 06:08:00 +00:00
|
|
|
}
|
2022-05-19 07:49:32 +00:00
|
|
|
} else {
|
|
|
|
// all good, commit
|
2023-02-23 03:38:02 +00:00
|
|
|
err = commit(ctx, outerCtx, m)
|
|
|
|
executePostCompleteHooks(ctx, outerCtx)
|
2022-05-19 07:49:32 +00:00
|
|
|
}
|
2022-09-28 06:08:00 +00:00
|
|
|
|
2022-05-19 07:49:32 +00:00
|
|
|
}()
|
|
|
|
|
|
|
|
err = fn(ctx)
|
|
|
|
return err
|
|
|
|
}
|
2022-07-13 06:30:54 +00:00
|
|
|
|
2022-11-20 19:49:10 +00:00
|
|
|
func begin(ctx context.Context, m Manager, exclusive bool) (context.Context, error) {
|
2022-09-19 04:53:06 +00:00
|
|
|
var err error
|
2022-11-20 19:49:10 +00:00
|
|
|
ctx, err = m.Begin(ctx, exclusive)
|
2022-09-19 04:53:06 +00:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
hm := hookManager{}
|
|
|
|
ctx = hm.register(ctx)
|
|
|
|
|
|
|
|
return ctx, nil
|
|
|
|
}
|
|
|
|
|
2023-02-23 03:38:02 +00:00
|
|
|
func commit(ctx context.Context, outerCtx context.Context, m Manager) error {
|
2022-09-19 04:53:06 +00:00
|
|
|
if err := m.Commit(ctx); err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
2023-02-23 03:38:02 +00:00
|
|
|
executePostCommitHooks(ctx, outerCtx)
|
2022-09-19 04:53:06 +00:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2023-02-23 03:38:02 +00:00
|
|
|
func rollback(ctx context.Context, outerCtx context.Context, m Manager) {
|
2022-09-19 04:53:06 +00:00
|
|
|
if err := m.Rollback(ctx); err != nil {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2023-02-23 03:38:02 +00:00
|
|
|
executePostRollbackHooks(ctx, outerCtx)
|
2022-09-19 04:53:06 +00:00
|
|
|
}
|
|
|
|
|
2022-07-13 06:30:54 +00:00
|
|
|
// WithDatabase executes fn with the context provided by p.WithDatabase.
|
|
|
|
// It does not run inside a transaction, so all database operations will be
|
|
|
|
// executed in their own transaction.
|
|
|
|
func WithDatabase(ctx context.Context, p DatabaseProvider, fn TxnFunc) error {
|
|
|
|
var err error
|
|
|
|
ctx, err = p.WithDatabase(ctx)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
return fn(ctx)
|
|
|
|
}
|
2022-08-11 06:14:57 +00:00
|
|
|
|
2022-11-20 19:49:10 +00:00
|
|
|
// Retryer is a provides WithTxn function that retries the transaction
|
|
|
|
// if it fails with a locked database error.
|
|
|
|
// Transactions are run in exclusive mode.
|
2022-08-11 06:14:57 +00:00
|
|
|
type Retryer struct {
|
|
|
|
Manager Manager
|
2022-09-01 07:54:34 +00:00
|
|
|
// use value < 0 to retry forever
|
2022-08-11 06:14:57 +00:00
|
|
|
Retries int
|
|
|
|
OnFail func(ctx context.Context, err error, attempt int) error
|
|
|
|
}
|
|
|
|
|
|
|
|
func (r Retryer) WithTxn(ctx context.Context, fn TxnFunc) error {
|
|
|
|
var attempt int
|
|
|
|
var err error
|
2022-09-01 07:54:34 +00:00
|
|
|
for attempt = 1; attempt <= r.Retries || r.Retries < 0; attempt++ {
|
2022-11-20 19:49:10 +00:00
|
|
|
const (
|
|
|
|
execComplete = false
|
|
|
|
exclusive = true
|
|
|
|
)
|
|
|
|
err = withTxn(ctx, r.Manager, fn, exclusive, execComplete)
|
2022-08-11 06:14:57 +00:00
|
|
|
|
|
|
|
if err == nil {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
if !r.Manager.IsLocked(err) {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
if r.OnFail != nil {
|
|
|
|
if err := r.OnFail(ctx, err, attempt); err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
return fmt.Errorf("failed after %d attempts: %w", attempt, err)
|
|
|
|
}
|