split query batches in the executor

We previously had some significant bugs in which our error handling for translation would cause an entire batch of queries to silently fail.
This change processes each call seperately, such that all executable queries are processed.
This commit is contained in:
Nia Weiss 2020-10-08 09:53:11 -04:00
parent 33af6c1ac7
commit 27289c5f35
No known key found for this signature in database
GPG key ID: 895E83409BFDA1BB

View file

@ -208,23 +208,6 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar
opt = &execOptions{}
}
// Translate query keys to ids, if necessary.
// No need to translate a remote call.
if !opt.Remote {
if err := e.translateCalls(ctx, index, q.Calls); err != nil {
if errors.Cause(err) == ErrTranslatingKeyNotFound {
// No error - return empty result
resp.Results = make([]interface{}, len(q.Calls))
for i, c := range q.Calls {
resp.Results[i] = emptyResult(c)
}
return resp, nil
}
return resp, err
} else if err := validateQueryContext(ctx); err != nil {
return resp, err
}
}
if opt.Profile {
var prof tracing.ProfiledSpan
prof, ctx = tracing.StartProfiledSpanFromContext(ctx, "Execute")
@ -235,24 +218,61 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar
return resp, fmt.Errorf("profiling execution failed: %T is not tracing.Profile", prof)
}
}
results, err := e.execute(ctx, index, q, shards, opt)
if err != nil {
return resp, err
} else if err := validateQueryContext(ctx); err != nil {
return resp, err
resp.Results = make([]interface{}, len(q.Calls))
for i, c := range q.Calls {
resp.Results[i] = emptyResult(c)
}
var columnAttrsRows []*Row
for i, c := range q.Calls {
// Translate query keys to ids, if necessary.
// No need to translate a remote call.
if !opt.Remote {
if err := e.translateCalls(ctx, index, []*pql.Call{c}); err != nil {
if errors.Cause(err) == ErrTranslatingKeyNotFound {
// No error - return empty result
continue
}
return resp, err
} else if err := validateQueryContext(ctx); err != nil {
return resp, err
}
}
resp.Results = results
results, err := e.execute(ctx, index, &pql.Query{Calls: []*pql.Call{c}}, shards, opt)
if err != nil {
return resp, err
} else if err := validateQueryContext(ctx); err != nil {
return resp, err
}
if opt.ColumnAttrs {
if resultRow, ok := results[0].(*Row); ok {
columnAttrsRows = append(columnAttrsRows, resultRow)
}
}
// Translate response objects from ids to keys, if necessary.
// No need to translate a remote call.
if !opt.Remote {
if err := e.translateResults(ctx, index, idx, []*pql.Call{c}, results); err != nil {
if errors.Cause(err) == ErrTranslatingKeyNotFound {
// No error - return empty result
continue
}
return resp, err
} else if err := validateQueryContext(ctx); err != nil {
return resp, err
}
}
resp.Results[i] = results[0]
}
// Fill column attributes if requested.
if opt.ColumnAttrs {
// Consolidate all column ids across all calls.
var columnIDs []uint64
for _, result := range results {
bm, ok := result.(*Row)
if !ok {
continue
}
for _, bm := range columnAttrsRows {
columnIDs = uint64Slice(columnIDs).merge(bm.Columns())
}
@ -282,24 +302,6 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar
resp.ColumnAttrSets = columnAttrSets
}
// Translate response objects from ids to keys, if necessary.
// No need to translate a remote call.
if !opt.Remote {
if err := e.translateResults(ctx, index, idx, q.Calls, results); err != nil {
if errors.Cause(err) == ErrTranslatingKeyNotFound {
// No error - return empty result
resp.Results = make([]interface{}, len(q.Calls))
for i, c := range q.Calls {
resp.Results[i] = emptyResult(c)
}
return resp, nil
}
return resp, err
} else if err := validateQueryContext(ctx); err != nil {
return resp, err
}
}
return resp, nil
}