remove batch abstraction
This commit is contained in:
+9
-52
@@ -7,63 +7,20 @@ import (
|
||||
type HandleRequestFunc func(ctx context.Context, req *TsdbQuery) (*Response, error)
|
||||
|
||||
func HandleRequest(ctx context.Context, req *TsdbQuery) (*Response, error) {
|
||||
tsdbQuery := &TsdbQuery{
|
||||
Queries: req.Queries,
|
||||
TimeRange: req.TimeRange,
|
||||
}
|
||||
|
||||
batches, err := getBatches(req)
|
||||
//TODO niceify
|
||||
endpoint, err := getTsdbQueryEndpointFor(req.Queries[0].DataSource)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
currentlyExecuting := 0
|
||||
resultsChan := make(chan *BatchResult)
|
||||
res := endpoint.Query(ctx, req)
|
||||
|
||||
for _, batch := range batches {
|
||||
if len(batch.Depends) == 0 {
|
||||
currentlyExecuting += 1
|
||||
batch.Started = true
|
||||
go batch.process(ctx, resultsChan, tsdbQuery)
|
||||
}
|
||||
if res.Error != nil {
|
||||
return nil, res.Error
|
||||
}
|
||||
|
||||
response := &Response{
|
||||
Results: make(map[string]*QueryResult),
|
||||
}
|
||||
|
||||
for currentlyExecuting != 0 {
|
||||
select {
|
||||
case batchResult := <-resultsChan:
|
||||
currentlyExecuting -= 1
|
||||
|
||||
response.BatchTimings = append(response.BatchTimings, batchResult.Timings)
|
||||
|
||||
if batchResult.Error != nil {
|
||||
return nil, batchResult.Error
|
||||
}
|
||||
|
||||
for refId, result := range batchResult.QueryResults {
|
||||
response.Results[refId] = result
|
||||
}
|
||||
|
||||
for _, batch := range batches {
|
||||
// not interested in started batches
|
||||
if batch.Started {
|
||||
continue
|
||||
}
|
||||
|
||||
if batch.allDependenciesAreIn(response) {
|
||||
currentlyExecuting += 1
|
||||
batch.Started = true
|
||||
go batch.process(ctx, resultsChan, tsdbQuery)
|
||||
}
|
||||
}
|
||||
case <-ctx.Done():
|
||||
return nil, ctx.Err()
|
||||
}
|
||||
}
|
||||
|
||||
//response.Results = tsdbQuery.Results
|
||||
return response, nil
|
||||
return &Response{
|
||||
Results: res.QueryResults,
|
||||
BatchTimings: []*BatchTiming{res.Timings},
|
||||
}, nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user