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`.
This commit is contained in:
Mariell Hoversholm
2025-03-26 09:02:25 +01:00
parent 03d6d8f854
commit 094c7d860e
2 changed files with 80 additions and 0 deletions
+49
View File
@@ -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()
+31
View File
@@ -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