From ddb5f125f04d3eb05c18bad96480acbff49081c6 Mon Sep 17 00:00:00 2001 From: Ryan McKinley Date: Tue, 2 Jul 2024 18:21:59 -0700 Subject: [PATCH] add sqlobj basic implementation --- pkg/services/apiserver/service.go | 10 +- .../sqlstore/migrations/migrations.go | 2 + .../sqlstore/migrations/object_mig.go | 28 +++ pkg/storage/unified/sqlnext/sql_resources.go | 44 ----- pkg/storage/unified/sqlobj/sql_resources.go | 181 ++++++++++++++++++ 5 files changed, 219 insertions(+), 46 deletions(-) create mode 100644 pkg/services/sqlstore/migrations/object_mig.go create mode 100644 pkg/storage/unified/sqlobj/sql_resources.go diff --git a/pkg/services/apiserver/service.go b/pkg/services/apiserver/service.go index 074d2f5590f..6a772a7c771 100644 --- a/pkg/services/apiserver/service.go +++ b/pkg/services/apiserver/service.go @@ -45,8 +45,8 @@ import ( "github.com/grafana/grafana/pkg/services/store/entity/sqlstash" "github.com/grafana/grafana/pkg/setting" "github.com/grafana/grafana/pkg/storage/unified/apistore" - "github.com/grafana/grafana/pkg/storage/unified/entitybridge" "github.com/grafana/grafana/pkg/storage/unified/resource" + "github.com/grafana/grafana/pkg/storage/unified/sqlobj" ) var ( @@ -267,7 +267,13 @@ func (s *service) start(ctx context.Context) error { return fmt.Errorf("unified storage requires the unifiedStorage feature flag") } - resourceServer, err := entitybridge.ProvideResourceServer(s.db, s.cfg, s.features, s.tracing) + // resourceServer, err := entitybridge.ProvideResourceServer(s.db, s.cfg, s.features, s.tracing) + // if err != nil { + // return err + // } + + // HACK... for now + resourceServer, err := sqlobj.ProvideSQLResourceServer(s.db, s.tracing) if err != nil { return err } diff --git a/pkg/services/sqlstore/migrations/migrations.go b/pkg/services/sqlstore/migrations/migrations.go index ae88297da96..f867f2f0de6 100644 --- a/pkg/services/sqlstore/migrations/migrations.go +++ b/pkg/services/sqlstore/migrations/migrations.go @@ -123,6 +123,8 @@ func (oss *OSSMigrations) AddMigration(mg *Migrator) { accesscontrol.AddManagedFolderAlertingSilencesActionsMigrator(mg) ualert.AddRecordingRuleColumns(mg) + + addObjectMigrations(mg) } func addStarMigrations(mg *Migrator) { diff --git a/pkg/services/sqlstore/migrations/object_mig.go b/pkg/services/sqlstore/migrations/object_mig.go new file mode 100644 index 00000000000..3bdaf1751a2 --- /dev/null +++ b/pkg/services/sqlstore/migrations/object_mig.go @@ -0,0 +1,28 @@ +package migrations + +import "github.com/grafana/grafana/pkg/services/sqlstore/migrator" + +// Add SQL table for a simple unified storage object backend +// NOTE: it would be nice to have this defined in the unified storage package, however that +// introduces a circular dependency. +func addObjectMigrations(mg *migrator.Migrator) { + mg.AddMigration("create unified storage object table", migrator.NewAddTableMigration(migrator.Table{ + Name: "object", + Columns: []*migrator.Column{ + // Sequential resource version + {Name: "rv", Type: migrator.DB_BigInt, Nullable: false, IsPrimaryKey: true, IsAutoIncrement: true}, + + // Properties that exist in path/key (and duplicated in the json value) + {Name: "group", Type: migrator.DB_NVarchar, Length: 190, Nullable: false}, + {Name: "namespace", Type: migrator.DB_NVarchar, Length: 63, Nullable: true}, // namespace is not required (cluster scope) + {Name: "resource", Type: migrator.DB_NVarchar, Length: 190, Nullable: false}, + {Name: "name", Type: migrator.DB_NVarchar, Length: 190, Nullable: false}, + + // The k8s resource JSON text (without the resourceVersion populated) + {Name: "value", Type: migrator.DB_MediumText, Nullable: false}, + }, + Indices: []*migrator.Index{ + {Cols: []string{"group", "namespace", "resource", "name"}, Type: migrator.UniqueIndex}, + }, + })) +} diff --git a/pkg/storage/unified/sqlnext/sql_resources.go b/pkg/storage/unified/sqlnext/sql_resources.go index 6f5cda1e7e3..01f5f72c267 100644 --- a/pkg/storage/unified/sqlnext/sql_resources.go +++ b/pkg/storage/unified/sqlnext/sql_resources.go @@ -4,7 +4,6 @@ import ( "context" "errors" "fmt" - "net/http" "strings" "github.com/prometheus/client_golang/prometheus" @@ -18,7 +17,6 @@ import ( "github.com/grafana/grafana/pkg/services/store/entity/sqlstash" "github.com/grafana/grafana/pkg/services/store/entity/sqlstash/sqltemplate" "github.com/grafana/grafana/pkg/storage/unified/resource" - "github.com/grafana/grafana/pkg/util" ) // Package-level errors. @@ -160,45 +158,3 @@ func (s *sqlResourceStore) PrepareList(ctx context.Context, req *resource.ListRe return nil, ErrNotImplementedYet } - -func (s *sqlResourceStore) PutBlob(ctx context.Context, req *resource.PutBlobRequest) (*resource.PutBlobResponse, error) { - if req.Method == resource.PutBlobRequest_HTTP { - return &resource.PutBlobResponse{ - Status: &resource.StatusResult{ - Status: "Failure", - Message: "http upload not supported", - Code: http.StatusNotImplemented, - }, - }, nil - } - - uid := util.GenerateShortUID() - - fmt.Printf("TODO, UPLOAD: %s // %+v", uid, req) - - return nil, ErrNotImplementedYet -} - -func (s *sqlResourceStore) GetBlob(ctx context.Context, uid string, mustProxy bool) (*resource.GetBlobResponse, error) { - return nil, ErrNotImplementedYet -} - -// Show resource history (and trash) -func (s *sqlResourceStore) History(ctx context.Context, req *resource.HistoryRequest) (*resource.HistoryResponse, error) { - _, span := s.tracer.Start(ctx, "storage_server.History") - defer span.End() - - fmt.Printf("TODO, GET History: %+v", req.Key) - - return nil, ErrNotImplementedYet -} - -// Used for efficient provisioning -func (s *sqlResourceStore) Origin(ctx context.Context, req *resource.OriginRequest) (*resource.OriginResponse, error) { - _, span := s.tracer.Start(ctx, "storage_server.History") - defer span.End() - - fmt.Printf("TODO, GET History: %+v", req.Key) - - return nil, ErrNotImplementedYet -} diff --git a/pkg/storage/unified/sqlobj/sql_resources.go b/pkg/storage/unified/sqlobj/sql_resources.go new file mode 100644 index 00000000000..5274fab449e --- /dev/null +++ b/pkg/storage/unified/sqlobj/sql_resources.go @@ -0,0 +1,181 @@ +package sqlobj + +import ( + "context" + "database/sql" + "errors" + "fmt" + "log/slog" + "time" + + "go.opentelemetry.io/otel/trace" + + "github.com/grafana/grafana/pkg/infra/db" + "github.com/grafana/grafana/pkg/infra/tracing" + "github.com/grafana/grafana/pkg/services/sqlstore/session" + "github.com/grafana/grafana/pkg/storage/unified/resource" +) + +// Package-level errors. +var ( + ErrNotImplementedYet = errors.New("not implemented yet (sqlobj)") +) + +func ProvideSQLResourceServer(db db.DB, tracer tracing.Tracer) (resource.ResourceServer, error) { + store := &sqlResourceStore{ + db: db, + log: slog.Default().With("logger", "unistore-sql-objects"), + tracer: tracer, + } + + return resource.NewResourceServer(resource.ResourceServerOptions{ + Tracer: tracer, + Backend: store, + Diagnostics: store, + Lifecycle: store, + }) +} + +type sqlResourceStore struct { + log *slog.Logger + db db.DB + tracer trace.Tracer + + broadcaster resource.Broadcaster[*resource.WrittenEvent] + + // Simple watch stream -- NOTE, this only works for single tenant! + stream chan<- *resource.WrittenEvent +} + +func (s *sqlResourceStore) Init() (err error) { + s.broadcaster, err = resource.NewBroadcaster(context.Background(), func(c chan<- *resource.WrittenEvent) error { + s.stream = c + return nil + }) + return +} + +func (s *sqlResourceStore) IsHealthy(ctx context.Context, r *resource.HealthCheckRequest) (*resource.HealthCheckResponse, error) { + return &resource.HealthCheckResponse{Status: resource.HealthCheckResponse_SERVING}, nil +} + +func (s *sqlResourceStore) Stop() { + if s.stream != nil { + close(s.stream) + } +} + +func (s *sqlResourceStore) WriteEvent(ctx context.Context, event resource.WriteEvent) (rv int64, err error) { + _, span := s.tracer.Start(ctx, "sql_resource.WriteEvent") + defer span.End() + + key := event.Key + + // This delegates resource version creation to auto-increment + // At scale, this is not a great strategy since everything is locked across all resources while this executes + appender := func(tx *session.SessionTx) (int64, error) { + return tx.ExecWithReturningId(ctx, + `INSERT INTO "object" ("group","namespace","resource","name","value") VALUES($1,$2,$3,$4,$5)`, + key.Group, key.Namespace, key.Resource, key.Name, event.Value) + } + + wiper := func(tx *session.SessionTx) (sql.Result, error) { + return tx.Exec(ctx, `DELETE FROM "object" WHERE `+ + `"group"=$1 AND `+ + `"namespace"=$2 AND `+ + `"resource"=$3 AND `+ + `"name"=$4`, + key.Group, key.Namespace, key.Resource, key.Name) + } + + err = s.db.GetSqlxSession().WithTransaction(ctx, func(tx *session.SessionTx) error { + switch event.Type { + case resource.WatchEvent_ADDED: + rv, err = appender(tx) + + case resource.WatchEvent_MODIFIED: + _, err = wiper(tx) + if err == nil { + rv, err = appender(tx) + } + case resource.WatchEvent_DELETED: + _, err = wiper(tx) + default: + return fmt.Errorf("unsupported event type") + } + return err + }) + + // Async notify all subscribers + if s.stream != nil { + go func() { + write := &resource.WrittenEvent{ + WriteEvent: event, + Timestamp: time.Now().UnixMilli(), + ResourceVersion: rv, + } + s.stream <- write + }() + } + return +} + +func (s *sqlResourceStore) WatchWriteEvents(ctx context.Context) (<-chan *resource.WrittenEvent, error) { + return s.broadcaster.Subscribe(ctx) +} + +func (s *sqlResourceStore) Read(ctx context.Context, req *resource.ReadRequest) (*resource.ReadResponse, error) { + _, span := s.tracer.Start(ctx, "storage_server.GetResource") + defer span.End() + + key := req.Key + rows, err := s.db.GetSqlxSession().Query(ctx, "SELECT rv,value FROM object WHERE group=$1 AND namespace=$2 AND resource=$3 AND name=$4", + key.Group, key.Namespace, key.Resource, key.Name) + if err != nil { + return nil, err + } + if rows.Next() { + rsp := &resource.ReadResponse{} + err = rows.Scan(&rsp.ResourceVersion, &rsp.Value) + if err == nil && rows.Next() { + return nil, fmt.Errorf("unexpected multiple results found") // should not be possible with the index strategy + } + return rsp, err + } + return nil, fmt.Errorf("NOT FOUND ERROR") +} + +// This implementation is only ever called from inside single tenant grafana, so there is no need to decode +// the value and try filtering first -- that will happen one layer up anyway +func (s *sqlResourceStore) PrepareList(ctx context.Context, req *resource.ListRequest) (*resource.ListResponse, error) { + _, span := s.tracer.Start(ctx, "storage_server.List") + defer span.End() + + if req.NextPageToken != "" { + return nil, fmt.Errorf("This storage backend does not support paging") + } + + max := 250 + key := req.Options.Key + rsp := &resource.ListResponse{} + rows, err := s.db.GetSqlxSession().Query(ctx, + "SELECT rv,value FROM object \n"+ + ` WHERE "group"=$1 AND namespace=$2 AND resource=$3 `+ + " ORDER BY name asc LIMIT $4", + key.Group, key.Namespace, key.Resource, max+1) + if err != nil { + return nil, err + } + for rows.Next() { + wrapper := &resource.ResourceWrapper{} + err = rows.Scan(&wrapper.ResourceVersion, &wrapper.Value) + if err != nil { + break + } + rsp.Items = append(rsp.Items, wrapper) + } + if len(rsp.Items) > max { + err = fmt.Errorf("more values that are supported by this storage engine") + } + return rsp, err +}