From 9e67f1dddd0c4ed7beee89eaddae7483af884d18 Mon Sep 17 00:00:00 2001 From: David Kagan <102766847+DKagan07@users.noreply.github.com> Date: Mon, 3 Apr 2023 11:13:08 -0400 Subject: [PATCH] Cluster nodes for serverless (#2336) * slowly making a serverless systemAPI for ClusterNodes() * implemented some methods for fb_database_info * fixed linting * fixed comments --- dax/queryer/queryer.go | 10 +++---- dax/queryer/system_api.go | 58 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 63 insertions(+), 5 deletions(-) create mode 100644 dax/queryer/system_api.go diff --git a/dax/queryer/queryer.go b/dax/queryer/queryer.go index 033db940a..07e86e84c 100644 --- a/dax/queryer/queryer.go +++ b/dax/queryer/queryer.go @@ -39,6 +39,8 @@ type Queryer struct { controller dax.Controller + systemLayer *systemlayer.SystemLayer + logger logger.Logger } @@ -47,6 +49,7 @@ func New(cfg Config) *Queryer { q := &Queryer{ controller: dax.NewNopController(), orchestrators: make(map[dax.QualifiedDatabaseID]*qualifiedOrchestrator), + systemLayer: systemlayer.NewSystemLayer(), logger: logger.NopLogger, } @@ -198,17 +201,14 @@ func (q *Queryer) QuerySQL(ctx context.Context, qdbid dax.QualifiedDatabaseID, s // Importer imp := idkserverless.NewImporter(q.controller, qdbid, nil) - // TODO(tlt): We need a serverless-compatible implementation of the // SystemAPI. - sysapi := &featurebase.NopSystemAPI{} - - systemLayer := systemlayer.NewSystemLayer() + sysapi := newSystemAPI(q.controller, qdbid) // We intentionally don't pass the sql argument here because we're working // with an io.Reader rather than a string and it's just not necessary to // send it as a string to this method. Also, what happens if the sql is a // large BULK INSERT? - pl := planner.NewExecutionPlanner(q.Orchestrator(qdbid), sapi, sysapi, systemLayer, imp, q.logger, "") + pl := planner.NewExecutionPlanner(q.Orchestrator(qdbid), sapi, sysapi, q.systemLayer, imp, q.logger, "") planOp, err := pl.CompilePlan(ctx, st) if err != nil { diff --git a/dax/queryer/system_api.go b/dax/queryer/system_api.go new file mode 100644 index 000000000..ed94bc9a8 --- /dev/null +++ b/dax/queryer/system_api.go @@ -0,0 +1,58 @@ +package queryer + +import ( + "context" + + featurebase "github.com/featurebasedb/featurebase/v3" + "github.com/featurebasedb/featurebase/v3/dax" +) + +// systemAPI is an implementation of the systemAPI. +type systemAPI struct { + featurebase.NopSystemAPI + controller dax.Controller + qdbid dax.QualifiedDatabaseID +} + +func newSystemAPI(c dax.Controller, qdbid dax.QualifiedDatabaseID) *systemAPI { + return &systemAPI{ + controller: c, + qdbid: qdbid, + } +} + +// ClusterNodes returns a list of featurebase.ClusterNodes +// with length of the minimum number of workers +func (s *systemAPI) ClusterNodes() []featurebase.ClusterNode { + + ctx := context.Background() + + qdb, err := s.controller.DatabaseByID(ctx, s.qdbid) + if err != nil { + return []featurebase.ClusterNode{} + } + out := make([]featurebase.ClusterNode, qdb.Options.WorkersMin) + return out +} + +func (s *systemAPI) PlatformDescription() string { + return "Serverless" +} + +func (s *systemAPI) ClusterName() string { + return "Serverless" +} + +func (s *systemAPI) ClusterNodeCount() int { + ctx := context.Background() + + qdb, err := s.controller.DatabaseByID(ctx, s.qdbid) + if err != nil { + return -1 + } + return qdb.Options.WorkersMin +} + +func (s *systemAPI) ClusterState() string { + return "NORMAL" +}