From 568ea044c7d3ba3fb4444aa445443c1eaee5e31f Mon Sep 17 00:00:00 2001 From: "grafana-delivery-bot[bot]" <132647405+grafana-delivery-bot[bot]@users.noreply.github.com> Date: Thu, 28 Dec 2023 11:57:43 +0100 Subject: [PATCH] [v10.2.x] ServerLock: Rework serverlock to use raw SQL and not depend on id (#79877) * ServerLock: Rework serverlock to use raw SQL and not depend on id (#79859) * rework SQL to use raw sql and more resistant to DBs that do not return ID * rework SQL to return ID in all DBs. Avoid using ID as operator (cherry picked from commit a595353d579002f66763e34750b8f2a8e9244703) * remove print statement * fix missing return ID for postgres (cherry picked from commit eb9c7fea072a2df15fa238d56aa385ee0d848cef) --------- Co-authored-by: Jo Co-authored-by: jguer --- pkg/infra/serverlock/serverlock.go | 132 +++++++++++++++++++---------- 1 file changed, 89 insertions(+), 43 deletions(-) diff --git a/pkg/infra/serverlock/serverlock.go b/pkg/infra/serverlock/serverlock.go index 9dffb1c0871..3b44a4ebc78 100644 --- a/pkg/infra/serverlock/serverlock.go +++ b/pkg/infra/serverlock/serverlock.go @@ -11,6 +11,8 @@ import ( "github.com/grafana/grafana/pkg/infra/db" "github.com/grafana/grafana/pkg/infra/log" "github.com/grafana/grafana/pkg/infra/tracing" + "github.com/grafana/grafana/pkg/services/sqlstore" + "github.com/grafana/grafana/pkg/services/sqlstore/migrator" ) func ProvideService(sqlStore db.DB, tracer tracing.Tracer) *ServerLockService { @@ -81,9 +83,10 @@ func (sl *ServerLockService) acquireLock(ctx context.Context, serverLock *server version = ?, last_execution = ? WHERE - id = ? AND version = ?` + operation_uid = ? AND version = ?` - res, err := dbSession.Exec(sql, newVersion, time.Now().Unix(), serverLock.Id, serverLock.Version) + res, err := dbSession.Exec(sql, newVersion, time.Now().Unix(), + serverLock.OperationUID, serverLock.Version) if err != nil { return err } @@ -102,16 +105,16 @@ func (sl *ServerLockService) getOrCreate(ctx context.Context, actionName string) defer span.End() var result *serverLock - err := sl.SQLStore.WithTransactionalDbSession(ctx, func(dbSession *db.Session) error { - lockRows := []*serverLock{} - err := dbSession.Where("operation_uid = ?", actionName).Find(&lockRows) + sqlRes := &serverLock{} + has, err := dbSession.SQL("SELECT * FROM server_lock WHERE operation_uid = ?", + actionName).Get(sqlRes) if err != nil { return err } - if len(lockRows) > 0 { - result = lockRows[0] + if has { + result = sqlRes return nil } @@ -119,14 +122,8 @@ func (sl *ServerLockService) getOrCreate(ctx context.Context, actionName string) OperationUID: actionName, LastExecution: 0, } - - _, err = dbSession.Insert(lockRow) - if err != nil { - return err - } - - result = lockRow - return nil + result, err = sl.createLock(ctx, lockRow, dbSession) + return err }) return result, err @@ -243,50 +240,54 @@ func (sl *ServerLockService) acquireForRelease(ctx context.Context, actionName s // getting the lock - as the action name has a Unique constraint, this will fail if the lock is already on the database err := sl.SQLStore.WithTransactionalDbSession(ctx, func(dbSession *db.Session) error { // we need to find if the lock is in the database - lockRows := []*serverLock{} - err := dbSession.Where("operation_uid = ?", actionName).Find(&lockRows) + result := &serverLock{} + sqlRaw := `SELECT * FROM server_lock WHERE operation_uid = ?` + if sl.SQLStore.GetDBType() == migrator.MySQL || sl.SQLStore.GetDBType() == migrator.Postgres { + sqlRaw += ` FOR UPDATE` + } + + has, err := dbSession.SQL( + sqlRaw, + actionName).Get(result) if err != nil { return err } ctxLogger := sl.log.FromContext(ctx) - if len(lockRows) > 0 { - result := lockRows[0] + if has { if sl.isLockWithinInterval(result, maxInterval) { return &ServerLockExistsError{actionName: actionName} - } else { - // lock has timeouted, so we update the timestamp - result.LastExecution = time.Now().Unix() - cond := &serverLock{OperationUID: actionName} - affected, err := dbSession.Update(result, cond) - if err != nil { - return err - } - if affected != 1 { - ctxLogger.Error("Expected rows affected to be 1 if there was no error", "actionName", actionName, "rowsAffected", affected) - } - return nil } - } else { - // lock not found, creating it - lockRow := &serverLock{ - OperationUID: actionName, - LastExecution: time.Now().Unix(), - } - - affected, err := dbSession.Insert(lockRow) + // lock has timed out, so we update the timestamp + result.LastExecution = time.Now().Unix() + res, err := dbSession.Exec("UPDATE server_lock SET last_execution = ? WHERE operation_uid = ?", + result.LastExecution, actionName) if err != nil { return err } + affected, err := res.RowsAffected() + if err != nil { + ctxLogger.Error("Error getting rows affected", "actionName", actionName, "error", err) + } + if affected != 1 { - // this means that there was no error but there is something not working correctly ctxLogger.Error("Expected rows affected to be 1 if there was no error", "actionName", actionName, "rowsAffected", affected) } + + return nil } - return nil + + // lock not found, creating it + lock := &serverLock{ + OperationUID: actionName, + LastExecution: time.Now().Unix(), + } + _, err = sl.createLock(ctx, lock, dbSession) + return err }) + return err } @@ -304,10 +305,14 @@ func (sl *ServerLockService) releaseLock(ctx context.Context, actionName string) return err } affected, err := res.RowsAffected() + if err != nil { + sl.log.FromContext(ctx).Debug("Error getting rows affected", "actionName", actionName, "error", err) + } + if affected != 1 { sl.log.FromContext(ctx).Debug("Error releasing lock", "actionName", actionName, "rowsAffected", affected) } - return err + return nil }) return err @@ -323,7 +328,7 @@ func (sl *ServerLockService) isLockWithinInterval(lock *serverLock, maxInterval return false } -func (sl ServerLockService) executeFunc(ctx context.Context, actionName string, fn func(ctx context.Context)) { +func (sl *ServerLockService) executeFunc(ctx context.Context, actionName string, fn func(ctx context.Context)) { start := time.Now() ctx, span := sl.tracer.Start(ctx, "ServerLockService.executeFunc") defer span.End() @@ -335,3 +340,44 @@ func (sl ServerLockService) executeFunc(ctx context.Context, actionName string, ctxLogger.Debug("Execution finished", "actionName", actionName, "duration", time.Since(start)) } + +func (sl *ServerLockService) createLock(ctx context.Context, + lockRow *serverLock, dbSession *sqlstore.DBSession) (*serverLock, error) { + affected := int64(1) + rawSQL := `INSERT INTO server_lock (operation_uid, last_execution, version) VALUES (?, ?, ?)` + if sl.SQLStore.GetDBType() == migrator.Postgres { + rawSQL += ` RETURNING id` + var id int64 + _, err := dbSession.SQL(rawSQL, lockRow.OperationUID, lockRow.LastExecution, 0).Get(&id) + if err != nil { + return nil, err + } + lockRow.Id = id + } else { + res, err := dbSession.Exec( + rawSQL, + lockRow.OperationUID, lockRow.LastExecution, 0) + if err != nil { + return nil, err + } + lastID, err := res.LastInsertId() + if err != nil { + sl.log.FromContext(ctx).Error("Error getting last insert id", "actionName", lockRow.OperationUID, "error", err) + } + lockRow.Id = lastID + + affected, err = res.RowsAffected() + if err != nil { + sl.log.FromContext(ctx).Error("Error getting rows affected", "actionName", lockRow.OperationUID, "error", err) + } + } + + if affected != 1 || lockRow.Id == 0 { + sl.log.FromContext(ctx).Error("Expected rows affected to be 1 if there was no error", + "actionName", lockRow.OperationUID, + "rowsAffected", affected, + "lockRow ID", lockRow.Id) + } + + return lockRow, nil +}