featurebase/sql3/planner/opquery.go
Travis Turner 92f64f42e9
Stub out the queryer/systemAPI implementation
This is used to fix the fb_exec_requests query; by default, it was doing
"fanout", but that doesn't make sense in serverless. So we needed a way
to prevent fanout when running Serverless. This gets that distinction
from the `SystemAPI.ClusterName()` method.
2023-03-15 11:14:23 -05:00

144 lines
3.4 KiB
Go

package planner
import (
"context"
"encoding/json"
"fmt"
"time"
pilosa "github.com/featurebasedb/featurebase/v3"
fbcontext "github.com/featurebasedb/featurebase/v3/context"
"github.com/featurebasedb/featurebase/v3/sql3"
"github.com/featurebasedb/featurebase/v3/sql3/planner/types"
"github.com/pkg/errors"
)
// PlanOpQuery is a query - this is the root node of an execution plan
type PlanOpQuery struct {
planner *ExecutionPlanner
ChildOp types.PlanOperator
sql string
warnings []string
}
var _ types.PlanOperator = (*PlanOpQuery)(nil)
func NewPlanOpQuery(p *ExecutionPlanner, child types.PlanOperator, sql string) *PlanOpQuery {
return &PlanOpQuery{
planner: p,
ChildOp: child,
warnings: make([]string, 0),
sql: sql,
}
}
func (p *PlanOpQuery) Schema() types.Schema {
return p.ChildOp.Schema()
}
func (p *PlanOpQuery) Child() types.PlanOperator {
return p.ChildOp
}
func (p *PlanOpQuery) Iterator(ctx context.Context, row types.Row) (types.RowIterator, error) {
iter, err := p.ChildOp.Iterator(ctx, row)
if err != nil {
return nil, err
}
return newQueryIterator(p.planner.systemLayerAPI.ExecutionRequests(), p, iter), nil
}
func (p *PlanOpQuery) Children() []types.PlanOperator {
return []types.PlanOperator{
p.ChildOp,
}
}
func (p *PlanOpQuery) WithChildren(children ...types.PlanOperator) (types.PlanOperator, error) {
if len(children) != 1 {
return nil, sql3.NewErrInternalf("unexpected number of children '%d'", len(children))
}
op := NewPlanOpQuery(p.planner, children[0], p.sql)
op.warnings = append(op.warnings, p.warnings...)
return op, nil
}
func (p *PlanOpQuery) Plan() map[string]interface{} {
result := make(map[string]interface{})
result["_op"] = fmt.Sprintf("%T", p)
result["_schema"] = p.Schema().Plan()
result["sql"] = p.sql
result["warnings"] = p.warnings
result["child"] = p.ChildOp.Plan()
return result
}
func (p *PlanOpQuery) AddWarning(warning string) {
p.warnings = append(p.warnings, warning)
}
func (p *PlanOpQuery) Warnings() []string {
var w []string
w = append(w, p.warnings...)
if p.ChildOp != nil {
w = append(w, p.ChildOp.Warnings()...)
}
return w
}
func (p *PlanOpQuery) String() string {
return ""
}
type queryIterator struct {
requests pilosa.ExecutionRequestsAPI
query *PlanOpQuery
child types.RowIterator
hasStarted *struct{}
}
func newQueryIterator(requests pilosa.ExecutionRequestsAPI, query *PlanOpQuery, child types.RowIterator) *queryIterator {
return &queryIterator{
requests: requests,
query: query,
child: child,
}
}
func (i *queryIterator) Next(ctx context.Context) (types.Row, error) {
if i.hasStarted == nil {
requestId, ok := fbcontext.RequestID(ctx)
if !ok {
return nil, sql3.NewErrInternalf("unable to get request id from context")
}
userId := ""
userId, _ = fbcontext.UserID(ctx)
i.requests.AddRequest(requestId, userId, time.Now(), i.query.sql)
i.hasStarted = &struct{}{}
}
row, err := i.child.Next(ctx)
if err != nil {
// either error or no more rows, either way update the request
requestId, ok := fbcontext.RequestID(ctx)
if !ok {
return nil, errors.Wrapf(sql3.NewErrInternalf("unable to get request id from context"), "next on child: %s", err)
}
plan, err := json.MarshalIndent(i.query.Plan(), "", " ")
if err != nil {
i.query.planner.logger.Infof("marshal indent: %s", err)
}
i.requests.UpdateRequest(requestId, time.Now(), "complete", "", 0, "", 0, 0, 0, 0, 0, string(plan))
}
return row, err
}