From 094c7d860e7f584428173f7970a235a173b5c613 Mon Sep 17 00:00:00 2001 From: Mariell Hoversholm Date: Wed, 26 Mar 2025 09:02:25 +0100 Subject: [PATCH] Unified Storage: Return an already exists error When inserting a resource that already exists (i.e. race condition), we can safely catch UNIQUE constraint violations and transform them into a `k8s.io/apimachinery/pkg/api/errors` error that stands the test of `errors.IsAlreadyExists`. --- pkg/storage/unified/sql/backend.go | 49 +++++++++++++++++++++++++ pkg/storage/unified/sql/backend_test.go | 31 ++++++++++++++++ 2 files changed, 80 insertions(+) diff --git a/pkg/storage/unified/sql/backend.go b/pkg/storage/unified/sql/backend.go index 6d2e0b4c01e..d0c0f27daf6 100644 --- a/pkg/storage/unified/sql/backend.go +++ b/pkg/storage/unified/sql/backend.go @@ -6,15 +6,22 @@ import ( "errors" "fmt" "math" + "net/http" "sync" "time" + "cloud.google.com/go/spanner" + "github.com/go-sql-driver/mysql" "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgconn" + "github.com/mattn/go-sqlite3" "github.com/prometheus/client_golang/prometheus" "go.opentelemetry.io/otel/trace" "go.opentelemetry.io/otel/trace/noop" + "google.golang.org/grpc/codes" "google.golang.org/protobuf/proto" apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "github.com/grafana/grafana/pkg/infra/log" "github.com/grafana/grafana/pkg/storage/unified/resource" @@ -24,6 +31,15 @@ import ( "github.com/grafana/grafana/pkg/util/debouncer" ) +var ErrResourceAlreadyExists error = &apierrors.StatusError{ + ErrStatus: metav1.Status{ + Status: metav1.StatusFailure, + Reason: metav1.StatusReasonAlreadyExists, + Message: "the resource already exists", + Code: http.StatusConflict, + }, +} + const tracePrefix = "sql.resource." const defaultPollingInterval = 100 * time.Millisecond const defaultWatchBufferSize = 100 // number of events to buffer in the watch stream @@ -337,6 +353,9 @@ func (b *backend) create(ctx context.Context, event resource.WriteEvent) (int64, Folder: folder, GUID: guid, }); err != nil { + if isRowAlreadyExistsError(err) { + return guid, ErrResourceAlreadyExists + } return guid, fmt.Errorf("insert into resource: %w", err) } @@ -377,6 +396,36 @@ func (b *backend) create(ctx context.Context, event resource.WriteEvent) (int64, return rv, nil } +// isRowAlreadyExistsError checks if the error is the result of the row inserted already existing. +// +// On SQLite and Postgres, this is known as a UNIQUE constraint violation. On MySQL, it's known as a duplicate entry. +// On Spanner, it's a gRPC ALREADY_EXISTS error. +func isRowAlreadyExistsError(err error) bool { + var sqlite sqlite3.Error + if errors.As(err, &sqlite) { + return sqlite.ExtendedCode == sqlite3.ErrConstraintUnique + } + + var pg *pgconn.PgError + if errors.As(err, &pg) { + // https://www.postgresql.org/docs/current/errcodes-appendix.html + return pg.Code == "23505" // unique_violation + } + + var mysqlerr *mysql.MySQLError + if errors.As(err, &mysqlerr) { + // https://dev.mysql.com/doc/mysql-errors/8.0/en/server-error-reference.html + return mysqlerr.Number == 1062 // ER_DUP_ENTRY + } + + // ErrCode returns Unknown for non-gRPC errors. + if spanner.ErrCode(err) == codes.AlreadyExists { + return true + } + + return false +} + func (b *backend) update(ctx context.Context, event resource.WriteEvent) (int64, error) { ctx, span := b.tracer.Start(ctx, tracePrefix+"Update") defer span.End() diff --git a/pkg/storage/unified/sql/backend_test.go b/pkg/storage/unified/sql/backend_test.go index f275b2b78e5..de6a4028782 100644 --- a/pkg/storage/unified/sql/backend_test.go +++ b/pkg/storage/unified/sql/backend_test.go @@ -8,7 +8,9 @@ import ( "testing" sqlmock "github.com/DATA-DOG/go-sqlmock" + "github.com/mattn/go-sqlite3" "github.com/stretchr/testify/require" + apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "github.com/grafana/grafana/pkg/apimachinery/utils" @@ -229,6 +231,29 @@ func TestBackend_create(t *testing.T) { require.Equal(t, int64(200), v) }) + t.Run("resource already exists", func(t *testing.T) { + t.Parallel() + b, ctx := setupBackendTest(t) + b.SQLMock.ExpectBegin() + expectSuccessfulResourceVersionExec(t, b.TestDBProvider, + func() { b.ExecWithResult("insert resource", 0, 1) }, + func() { b.ExecWithResult("insert resource_history", 0, 1) }, + ) + b.SQLMock.ExpectCommit() + b.SQLMock.ExpectBegin() + b.SQLMock.ExpectExec("insert resource").WillReturnError(sqlite3.Error{Code: sqlite3.ErrConstraint, ExtendedCode: sqlite3.ErrConstraintUnique}) + b.SQLMock.ExpectRollback() + + // First we insert the resource successfully. This is what the happy path test does as well. + v, err := b.create(ctx, event) + require.NoError(t, err) + require.Equal(t, int64(200), v) + + // Then we try to insert the same resource again. This should fail. + _, err = b.create(ctx, event) + require.ErrorIs(t, err, ErrResourceAlreadyExists) + }) + t.Run("error inserting into resource", func(t *testing.T) { t.Parallel() b, ctx := setupBackendTest(t) @@ -642,6 +667,12 @@ func TestBackend_getHistoryPagination(t *testing.T) { }) } +func TestErrResourceAlreadyExistsIsRecognisable(t *testing.T) { + t.Parallel() + + require.True(t, apierrors.IsAlreadyExists(ErrResourceAlreadyExists), "ErrResourceAlreadyExists should be recognised as an AlreadyExists error") +} + // setupHistoryTest creates the necessary mock expectations for a history test func setupHistoryTest(b testBackend, resourceVersions []int64, latestRV int64) *sqlmock.Rows { // Expect fetch latest RV call - set to the highest resource version