Files
query-orchestration/internal/query/sync.go
T

164 lines
3.1 KiB
Go
Raw Normal View History

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
}