0ac5ff9e15
Test Query * depstextandclean * startedcleaningresult * resulttidyup * roundone * cleaning * unsyncedquery * startedtestsandsimplification * api * querytests * resultprocessortests * unittests * cleanup
164 lines
3.1 KiB
Go
164 lines
3.1 KiB
Go
package query
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
documenttext "queryorchestration/internal/document/text"
|
|
"queryorchestration/internal/query/result"
|
|
resultprocessor "queryorchestration/internal/query/result/processor"
|
|
"sync"
|
|
|
|
"github.com/google/uuid"
|
|
)
|
|
|
|
func (s *Service) Sync(ctx context.Context, doc *Document) error {
|
|
err := s.svc.Text.IsExtracted(&documenttext.IsExtractedParams{
|
|
DocumentID: doc.ID,
|
|
MinCleanVersion: doc.CleanVersion,
|
|
MinTextVersion: doc.TextVersion,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
unsyncedQueries, err := s.svc.Result.ListUnsyncedQueriesByDocId(ctx, doc.ID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
batchedQueries := s.batchQueries(unsyncedQueries)
|
|
|
|
for _, queries := range batchedQueries {
|
|
err := s.processBatch(ctx, doc, queries)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (s *Service) batchQueries(queries []*resultprocessor.Query) [][]*resultprocessor.Query {
|
|
n := len(queries)
|
|
if n == 0 {
|
|
return nil
|
|
}
|
|
|
|
idToIndex := make(map[uuid.UUID]int, n)
|
|
for i, q := range queries {
|
|
idToIndex[q.ID] = i
|
|
}
|
|
|
|
result := make([][]*resultprocessor.Query, 0, n)
|
|
assigned := make([]bool, n)
|
|
deps := make([][]int, n)
|
|
remaining := n
|
|
|
|
for i, q := range queries {
|
|
toAppend := false
|
|
if q.RequiredQueryIDs == nil || len(*q.RequiredQueryIDs) == 0 {
|
|
toAppend = true
|
|
} else {
|
|
deps[i] = make([]int, 0, len(*q.RequiredQueryIDs))
|
|
|
|
for _, reqID := range *q.RequiredQueryIDs {
|
|
if idx, exists := idToIndex[reqID]; exists {
|
|
deps[i] = append(deps[i], idx)
|
|
}
|
|
}
|
|
|
|
toAppend = len(deps[i]) == 0
|
|
}
|
|
|
|
if toAppend {
|
|
if len(result) == 0 {
|
|
result = [][]*resultprocessor.Query{{q}}
|
|
} else {
|
|
result[0] = append(result[0], q)
|
|
}
|
|
assigned[idToIndex[q.ID]] = true
|
|
remaining--
|
|
}
|
|
}
|
|
|
|
for remaining > 0 {
|
|
currentLayer := make([]*resultprocessor.Query, 0, remaining)
|
|
|
|
for i, q := range queries {
|
|
if assigned[i] {
|
|
continue
|
|
}
|
|
|
|
allSatisfied := true
|
|
for _, depIdx := range deps[i] {
|
|
if !assigned[depIdx] {
|
|
allSatisfied = false
|
|
break
|
|
}
|
|
}
|
|
|
|
if allSatisfied {
|
|
currentLayer = append(currentLayer, q)
|
|
}
|
|
}
|
|
|
|
for _, v := range currentLayer {
|
|
assigned[idToIndex[v.ID]] = true
|
|
remaining--
|
|
}
|
|
|
|
result = append(result, currentLayer)
|
|
}
|
|
|
|
return result
|
|
}
|
|
|
|
func (s *Service) processBatch(ctx context.Context, doc *Document, queries []*resultprocessor.Query) error {
|
|
if doc == nil {
|
|
return errors.New("document required")
|
|
}
|
|
|
|
errChan := make(chan error, len(queries))
|
|
var wg sync.WaitGroup
|
|
sem := make(chan struct{}, 10)
|
|
|
|
for _, query := range queries {
|
|
wg.Add(1)
|
|
go func(q *resultprocessor.Query) {
|
|
defer wg.Done()
|
|
|
|
sem <- struct{}{}
|
|
defer func() {
|
|
<-sem
|
|
}()
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
errChan <- ctx.Err()
|
|
return
|
|
default:
|
|
}
|
|
|
|
_, err := s.svc.Result.Set(ctx, &result.Set{
|
|
DocumentID: doc.ID,
|
|
CleanVersion: doc.CleanVersion,
|
|
TextVersion: doc.TextVersion,
|
|
Query: query,
|
|
})
|
|
if err != nil {
|
|
errChan <- err
|
|
}
|
|
}(query)
|
|
}
|
|
|
|
wg.Wait()
|
|
close(errChan)
|
|
for err := range errChan {
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|