Files
query-orchestration/internal/query/result/set/set.go
T
Michael McGuinness 71f9802e1a Merged in feature/splitqueryrunning (pull request #57)
Split Query Running + Debugging Full Flow

* completedquerysyncrunner

* spliitinglogic

* synccomplete

* informdependents

* only push same collector

* deps

* livetesting

* foundissue

* some issues resolved

* activeupdate

* collectorupdatefixes

* fix dbquesries

* tests

* tests

* pollingdebug
2025-02-11 15:22:59 +00:00

95 lines
2.4 KiB
Go

package resultset
import (
"context"
"log/slog"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/query/result"
"queryorchestration/internal/serviceconfig/queue"
"github.com/google/uuid"
)
type Set struct {
DocumentID uuid.UUID `json:"document_id" validate:"required,uuid"`
QueryID uuid.UUID `json:"query_id" validate:"required,uuid"`
}
func (s *Service) Set(ctx context.Context, params *Set) error {
dbid := database.MustToDBUUID(params.DocumentID)
cleanVersion, err := s.cfg.GetDBQueries().GetDocumentCleanEntry(ctx, dbid)
if err != nil {
return err
}
textVersion, err := s.cfg.GetDBQueries().GetDocumentTextEntry(ctx, dbid)
if err != nil {
return err
}
query, err := s.svc.Query.Get(ctx, params.QueryID)
if err != nil {
return err
}
value, err := s.svc.Result.Process(ctx, &result.Process{
DocumentID: params.DocumentID,
QueryID: query.ID,
QueryVersion: query.ActiveVersion,
})
if err != nil {
return err
}
err = s.cfg.GetDBQueries().SetResult(ctx, &repository.SetResultParams{
Queryid: database.MustToDBUUID(params.QueryID),
Documentid: database.MustToDBUUID(params.DocumentID),
Value: value.GetStoreValue(),
Cleanversion: cleanVersion.Version,
Textversion: textVersion.Version,
Queryversion: query.ActiveVersion,
})
if err != nil {
return err
}
slog.Debug("set query result", "query_id", params.QueryID.String(), "document_id", params.DocumentID.String())
return s.informQueryDependents(ctx, params)
}
func (s *Service) informQueryDependents(ctx context.Context, params *Set) error {
ids, err := s.cfg.GetDBQueries().ListQueryDirectDependentsByDocumentID(ctx, &repository.ListQueryDirectDependentsByDocumentIDParams{
ID: database.MustToDBUUID(params.DocumentID),
Requiredids: database.MustToDBUUID(params.QueryID),
})
if err != nil {
return err
}
return s.TriggerQueriesSync(ctx, params.DocumentID, database.MustToUUIDArray(ids))
}
func (s *Service) TriggerQueriesSync(ctx context.Context, documentID uuid.UUID, queryIDs []uuid.UUID) error {
for _, id := range queryIDs {
err := s.TriggerQuerySync(ctx, &Set{
DocumentID: documentID,
QueryID: id,
})
if err != nil {
return err
}
}
return nil
}
func (s *Service) TriggerQuerySync(ctx context.Context, params *Set) error {
return s.cfg.SendToQueue(ctx, &queue.SendParams{
QueueURL: s.cfg.GetQueryURL(),
Body: params,
})
}