Этот коммит содержится в:
Ibrahim Serdar Acikgoz
2022-04-04 14:03:39 +03:00
коммит произвёл GitHub
родитель b1b9774efa
Коммит 97d3368a98
16 изменённых файлов: 173 добавлений и 188 удалений

67
vendor/github.com/mattermost/morph/drivers/postgres/lock.go сгенерированный поставляемый
Просмотреть файл

@@ -8,6 +8,7 @@ import (
"sync"
"time"
"github.com/lib/pq"
"github.com/mattermost/morph/drivers"
)
@@ -26,12 +27,14 @@ type Mutex struct {
stopRefresh chan bool
refreshDone chan bool
conn *sql.Conn
logger drivers.Logger
}
// NewMutex creates a mutex with the given key name.
//
// returns error if key is empty.
func NewMutex(key string, driver drivers.Driver) (*Mutex, error) {
func NewMutex(key string, driver drivers.Driver, logger drivers.Logger) (*Mutex, error) {
key, err := drivers.MakeLockKey(key)
if err != nil {
return nil, err
@@ -56,8 +59,9 @@ func NewMutex(key string, driver drivers.Driver) (*Mutex, error) {
}
return &Mutex{
key: key,
conn: conn,
key: key,
conn: conn,
logger: logger,
}, nil
}
@@ -68,10 +72,16 @@ func (m *Mutex) tryLock(ctx context.Context) (bool, error) {
if err != nil {
return false, err
}
defer m.finalizeTx(tx)
query := fmt.Sprintf("INSERT INTO %s (id, expireat) VALUES ($1, $2)", drivers.MutexTableName)
if _, err := tx.Exec(query, m.key, now.Add(drivers.TTL).Unix()); err != nil {
err2 := m.releaseLock(tx, now)
if pqErr, ok := err.(*pq.Error); ok && pqErr.Code == "23505" {
m.logger.Println("DB is locked, going to try acquire the lock if it is expired.")
}
m.finalizeTx(tx)
err2 := m.releaseLock(ctx, now)
if err2 == nil { // lock has been released due to expiration
return true, nil
}
@@ -81,27 +91,25 @@ func (m *Mutex) tryLock(ctx context.Context) (bool, error) {
err = tx.Commit()
if err != nil {
if txErr := tx.Rollback(); txErr != nil {
return false, txErr
}
return false, err
}
return true, nil
}
func (m *Mutex) releaseLock(tx *sql.Tx, t time.Time) error {
func (m *Mutex) releaseLock(ctx context.Context, t time.Time) error {
tx, err := m.conn.BeginTx(ctx, nil)
if err != nil {
return err
}
defer m.finalizeTx(tx)
e, err := m.getExpireAt(tx)
if err != nil {
return err
}
if t.Unix() < e {
if txErr := tx.Rollback(); txErr != nil {
return fmt.Errorf("could not rollback: %w", txErr)
}
return errors.New("could not release the lock")
}
@@ -112,10 +120,6 @@ func (m *Mutex) releaseLock(tx *sql.Tx, t time.Time) error {
err = tx.Commit()
if err != nil {
if txErr := tx.Rollback(); txErr != nil {
return fmt.Errorf("could not rollback transaction: %w", txErr)
}
return fmt.Errorf("unable to set new expireat for mutex: %w", err)
}
@@ -127,10 +131,6 @@ func (m *Mutex) getExpireAt(tx *sql.Tx) (int64, error) {
query := fmt.Sprintf("SELECT expireat FROM %s WHERE id = $1", drivers.MutexTableName)
err := tx.QueryRow(query, m.key).Scan(&expireAt)
if err != nil {
if txErr := tx.Rollback(); txErr != nil {
return -1, fmt.Errorf("could not rollback: %w", txErr)
}
return -1, fmt.Errorf("failed to fetch mutex from db: %w", err)
}
@@ -143,6 +143,7 @@ func (m *Mutex) refreshLock(ctx context.Context) error {
if err != nil {
return err
}
defer m.finalizeTx(tx)
e, err := m.getExpireAt(tx)
if err != nil {
@@ -157,28 +158,18 @@ func (m *Mutex) refreshLock(ctx context.Context) error {
err = tx.Commit()
if err != nil {
if txErr := tx.Rollback(); txErr != nil {
return fmt.Errorf("could not rollback: %w", txErr)
}
return fmt.Errorf("unable to refresh expireat for mutex: %w", err)
}
return nil
}
// Lock locks m. If the mutex is already locked by any other morph instance, including the current one,
// the calling goroutine blocks until the mutex can be locked.
func (m *Mutex) Lock() error {
return m.LockWithContext(context.Background())
}
// LockWithContext locks m unless the context is canceled. If the mutex is already locked by any other
// Lock locks m unless the context is canceled. If the mutex is already locked by any other
// instance, including the current one, the calling goroutine blocks until the mutex can be locked,
// or the context is canceled.
//
// The mutex is locked only if a nil error is returned.
func (m *Mutex) LockWithContext(ctx context.Context) error {
func (m *Mutex) Lock(ctx context.Context) error {
var waitInterval time.Duration
for {
@@ -248,16 +239,18 @@ func (m *Mutex) Unlock() error {
func executeTx(tx *sql.Tx, query string, args ...interface{}) error {
if _, err := tx.Exec(query, args...); err != nil {
if txErr := tx.Rollback(); txErr != nil {
return fmt.Errorf("could not rollback transaction: %w", txErr)
}
return err
}
return nil
}
func (m *Mutex) finalizeTx(tx *sql.Tx) {
if err := tx.Rollback(); err != nil && err != sql.ErrTxDone {
m.logger.Printf("failed to rollback transaction: %s", err)
}
}
// noCopy may be embedded into structs which must not be copied
// after the first use.
//