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 }