Merged in feature/testwithlogs (pull request #65)

Query Version Sync Runner

* testing

* queryversiosyncworking

* update

* tests

* fixtests
This commit is contained in:
Michael McGuinness
2025-02-14 10:56:24 +00:00
parent 477518e5eb
commit 0df3d16976
91 changed files with 1737 additions and 624 deletions
+7 -1
View File
@@ -6,6 +6,7 @@ import (
resultset "queryorchestration/internal/query/result/set"
"github.com/go-playground/validator/v10"
"github.com/google/uuid"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
)
@@ -28,8 +29,13 @@ func New(validator *validator.Validate, svc *Services) Runner {
}
}
type Body struct {
DocumentID uuid.UUID `json:"document_id" validate:"required,uuid"`
QueryID uuid.UUID `json:"query_id" validate:"required,uuid"`
}
func (s *Runner) Process(ctx context.Context, req *types.Message) error {
var body resultset.Set
var body Body
err := json.Unmarshal([]byte(*req.Body), &body)
if err != nil {
return err
+3 -1
View File
@@ -11,6 +11,7 @@ import (
"queryorchestration/internal/query/result"
resultprocessor "queryorchestration/internal/query/result/processor"
resultset "queryorchestration/internal/query/result/set"
resultsync "queryorchestration/internal/query/result/sync"
"queryorchestration/internal/server/runner"
queryc "queryorchestration/internal/serviceconfig/queue/query"
queuemock "queryorchestration/mocks/queue"
@@ -53,6 +54,7 @@ func TestQueryRunner(t *testing.T) {
Result: result.New(cfg, &result.Services{
Query: que,
}),
Sync: resultsync.New(cfg),
}),
})
assert.NotNil(t, runner)
@@ -93,7 +95,7 @@ func TestQueryRunner(t *testing.T) {
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
AddRow(database.MustToDBUUID(query.ID), repository.QuerytypeJsonExtractor, query.Version, query.Version, []byte(*query.Config), database.MustToDBUUIDArray(*query.RequiredQueryIDs)),
)
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(query.ID), query.Version).WillReturnRows(
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(query.ID), &query.Version).WillReturnRows(
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
AddRow(database.MustToDBUUID(query.ID), repository.QuerytypeJsonExtractor, query.Version, query.Version, []byte(*query.Config), database.MustToDBUUIDArray(*query.RequiredQueryIDs)),
)
+3 -1
View File
@@ -8,6 +8,7 @@ import (
collectorupdate "queryorchestration/internal/job/collector/update"
"queryorchestration/internal/query"
querytest "queryorchestration/internal/query/test"
queryupdate "queryorchestration/internal/query/update"
"github.com/go-playground/validator/v10"
)
@@ -19,9 +20,10 @@ type Services struct {
Collector *collector.Service
CollectorUpdate *collectorupdate.Service
Query *query.Service
QueryUpdate *queryupdate.Service
QueryTest *querytest.Service
Client *client.Service
Job *job.Service
QueryTest *querytest.Service
}
type Controllers struct {
+1 -1
View File
@@ -71,7 +71,7 @@ func (s *Controllers) UpdateQuery(ctx echo.Context, id types.UUID) error {
return echo.NewHTTPError(http.StatusBadRequest, err)
}
err := s.svc.Query.Update(ctx.Request().Context(), &resultprocessor.Update{
err := s.svc.QueryUpdate.Update(ctx.Request().Context(), &resultprocessor.Update{
ActiveVersion: req.ActiveVersion,
Config: req.Config,
RequiredQueryIDs: req.RequiredQueries,
+26 -4
View File
@@ -16,10 +16,14 @@ import (
"queryorchestration/internal/query"
"queryorchestration/internal/query/result"
querytest "queryorchestration/internal/query/test"
queryupdate "queryorchestration/internal/query/update"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/queryversionsync"
queuemock "queryorchestration/mocks/queue"
"strings"
"testing"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/go-playground/validator/v10"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
@@ -27,6 +31,7 @@ import (
"github.com/labstack/echo/v4"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
)
func TestCreateQuery(t *testing.T) {
@@ -159,12 +164,20 @@ func TestUpdateQuery(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &struct {
serviceconfig.BaseConfig
queryversionsync.QueryVersionSyncConfig
}{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.QueryVersionSyncURL = "here"
cons := queryservice.NewControllers(validator.New(), &queryservice.Services{
Query: query.New(cfg),
QueryUpdate: queryupdate.New(cfg, &queryupdate.Services{
Query: query.New(cfg),
}),
})
av := int32(2)
@@ -190,9 +203,18 @@ func TestUpdateQuery(t *testing.T) {
)
pool.ExpectBeginTx(pgx.TxOptions{})
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(av, int32(3), database.MustToDBUUID(id)).
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(av, int32(2), database.MustToDBUUID(id)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectCommit()
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.QueryVersionSyncURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", id.String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
err = cons.UpdateQuery(ctx, id)
assert.NoError(t, err)
@@ -268,7 +290,7 @@ func TestTestQuery(t *testing.T) {
pgxmock.NewRows([]string{"id", "jobId", "minCleanVersion", "minTextVersion", "activeVersion", "latestVersion", "fields"}).
AddRow(database.MustToDBUUID(coll.ID), database.MustToDBUUID(doc.JobID), coll.MinCleanVersion, coll.MinTextVersion, int32(1), int32(2), []byte("")),
)
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(params.QueryID), params.QueryVersion).WillReturnRows(
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(params.QueryID), &params.QueryVersion).WillReturnRows(
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
AddRow(database.MustToDBUUID(params.QueryID), repository.QuerytypeJsonExtractor, int32(1), params.QueryVersion+1, []byte("{\"path\":\"oldkey\"}"), []pgtype.UUID{reqID}),
)
+2 -10
View File
@@ -7,9 +7,7 @@ import (
querysyncrunner "queryorchestration/api/querySyncRunner"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/query"
"queryorchestration/internal/query/result"
resultset "queryorchestration/internal/query/result/set"
resultsync "queryorchestration/internal/query/result/sync"
querysync "queryorchestration/internal/query/sync"
"queryorchestration/internal/server/runner"
queryc "queryorchestration/internal/serviceconfig/queue/query"
@@ -45,14 +43,8 @@ func TestQueryRunner(t *testing.T) {
cfg.QueueClient = mockSQS
cfg.QueryURL = "/i/am/here"
que := query.New(cfg)
svc := querysync.New(cfg, &querysync.Services{
ResultSet: resultset.New(cfg, &resultset.Services{
Query: que,
Result: result.New(cfg, &result.Services{
Query: que,
}),
}),
ResultSync: resultsync.New(cfg),
})
runner := querysyncrunner.New(validator.New(), &querysyncrunner.Services{
+54
View File
@@ -0,0 +1,54 @@
package queryversionsyncrunner
import (
"context"
"encoding/json"
queryversionsync "queryorchestration/internal/query/versionsync"
"github.com/go-playground/validator/v10"
"github.com/google/uuid"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
)
const Name = "queryVersionSyncRunner"
type Services struct {
Sync *queryversionsync.Service
}
type Runner struct {
validator *validator.Validate
svc *Services
}
func New(validator *validator.Validate, svc *Services) Runner {
return Runner{
validator: validator,
svc: svc,
}
}
type Body struct {
ID uuid.UUID `json:"id" validate:"required,uuid"`
}
func (s *Runner) Process(ctx context.Context, req *types.Message) error {
var body Body
err := json.Unmarshal([]byte(*req.Body), &body)
if err != nil {
return err
}
err = s.validator.Struct(body)
if err != nil {
return err
}
err = s.svc.Sync.Sync(ctx, body.ID)
if err != nil {
return err
}
return nil
}
+95
View File
@@ -0,0 +1,95 @@
package queryversionsyncrunner_test
import (
"context"
"encoding/json"
"fmt"
queryversionsyncrunner "queryorchestration/api/queryVersionSyncRunner"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
queryversionsync "queryorchestration/internal/query/versionsync"
"queryorchestration/internal/server/runner"
"queryorchestration/internal/serviceconfig/queue/jobsync"
queuemock "queryorchestration/mocks/queue"
"testing"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
"github.com/go-playground/validator/v10"
"github.com/google/uuid"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
)
type QuerySyncConfig struct {
runner.BaseConfig
jobsync.JobSyncConfig
}
func TestQueryRunner(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &QuerySyncConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.JobSyncURL = "/i/am/here"
svc := queryversionsync.New(cfg)
runner := queryversionsyncrunner.New(validator.New(), &queryversionsyncrunner.Services{
Sync: svc,
})
assert.NotNil(t, runner)
doc := queryversionsyncrunner.Body{
ID: uuid.New(),
}
bodyBytes, err := json.Marshal(doc)
assert.NoError(t, err)
body := string(bodyBytes)
msg := &types.Message{
Body: &body,
}
jobIds := []uuid.UUID{
uuid.New(),
uuid.New(),
}
pool.ExpectQuery("name: ListQueryJobIDs :many").WithArgs(database.MustToDBUUID(doc.ID)).
WillReturnRows(
pgxmock.NewRows([]string{"jobId"}).
AddRow(database.MustToDBUUID(jobIds[0])).
AddRow(database.MustToDBUUID(jobIds[1])),
)
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.JobSyncURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", jobIds[0].String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.JobSyncURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", jobIds[1].String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
err = runner.Process(ctx, msg)
assert.NoError(t, err)
}
+3
View File
@@ -8,6 +8,7 @@ import (
"queryorchestration/internal/query"
"queryorchestration/internal/query/result"
resultset "queryorchestration/internal/query/result/set"
resultsync "queryorchestration/internal/query/result/sync"
"queryorchestration/internal/server/runner"
queryc "queryorchestration/internal/serviceconfig/queue/query"
@@ -29,9 +30,11 @@ func main() {
res := result.New(cfg, &result.Services{
Query: que,
})
sync := resultsync.New(cfg)
resset := resultset.New(cfg, &resultset.Services{
Result: res,
Query: que,
Sync: sync,
})
c := queryrunner.New(cfg.GetValidator(), &queryrunner.Services{
+8 -1
View File
@@ -17,8 +17,10 @@ import (
"queryorchestration/internal/query"
"queryorchestration/internal/query/result"
querytest "queryorchestration/internal/query/test"
queryupdate "queryorchestration/internal/query/update"
service "queryorchestration/internal/server/service"
"queryorchestration/internal/serviceconfig/queue/jobsync"
"queryorchestration/internal/serviceconfig/queue/queryversionsync"
"github.com/getkin/kin-openapi/openapi3"
_ "github.com/lib/pq"
@@ -27,6 +29,7 @@ import (
type QueryServiceConfig struct {
service.BaseConfig
jobsync.JobSyncConfig
queryversionsync.QueryVersionSyncConfig
}
func main() {
@@ -62,15 +65,19 @@ func main() {
Result: res,
Document: doc,
})
qupdate := queryupdate.New(cfg, &queryupdate.Services{
Query: que,
})
services := &queryservice.Services{
Export: exp,
Collector: col,
CollectorUpdate: colupdate,
Query: que,
QueryUpdate: qupdate,
QueryTest: quetest,
Client: cli,
Job: jbb,
QueryTest: quetest,
}
cons := queryservice.NewControllers(cfg.GetValidator(), services)
+3 -12
View File
@@ -5,9 +5,7 @@ import (
"log/slog"
"os"
querysyncrunner "queryorchestration/api/querySyncRunner"
"queryorchestration/internal/query"
"queryorchestration/internal/query/result"
resultset "queryorchestration/internal/query/result/set"
resultsync "queryorchestration/internal/query/result/sync"
querysync "queryorchestration/internal/query/sync"
"queryorchestration/internal/server/runner"
queryc "queryorchestration/internal/serviceconfig/queue/query"
@@ -26,16 +24,9 @@ func main() {
cfg := &QuerySyncConfig{}
cfg.ControllerFunc = func() runner.Controller {
que := query.New(cfg)
res := result.New(cfg, &result.Services{
Query: que,
})
resset := resultset.New(cfg, &resultset.Services{
Query: que,
Result: res,
})
sync := resultsync.New(cfg)
svc := querysync.New(cfg, &querysync.Services{
ResultSet: resset,
ResultSync: sync,
})
c := querysyncrunner.New(cfg.GetValidator(), &querysyncrunner.Services{
+42
View File
@@ -0,0 +1,42 @@
package main
import (
"context"
"log/slog"
"os"
queryversionsyncrunner "queryorchestration/api/queryVersionSyncRunner"
queryversionsync "queryorchestration/internal/query/versionsync"
"queryorchestration/internal/server/runner"
"queryorchestration/internal/serviceconfig/queue/jobsync"
_ "github.com/lib/pq"
)
type QueryVersionSyncConfig struct {
runner.BaseConfig
jobsync.JobSyncConfig
}
func main() {
ctx := context.Background()
cfg := &QueryVersionSyncConfig{}
cfg.ControllerFunc = func() runner.Controller {
svc := queryversionsync.New(cfg)
c := queryversionsyncrunner.New(cfg.GetValidator(), &queryversionsyncrunner.Services{
Sync: svc,
})
return &c
}
server, err := runner.New(ctx, cfg)
if err != nil {
slog.Error(err.Error())
os.Exit(1)
}
server.Listen(ctx)
}
@@ -1,20 +1,30 @@
CREATE VIEW fullActiveQueries AS
SELECT DISTINCT
q.id, q.type, q.activeVersion, q.latestVersion,
coalesce(c.config, null) as config,
coalesce(
ARRAY_AGG(r.requiredQueryId)
FILTER (WHERE r.requiredQueryId != '00000000-0000-0000-0000-000000000000')::uuid[],
array[]::uuid[]
)::uuid[] as requiredIds
WITH config as (
SELECT c.queryId, c.config
FROM queries AS q
LEFT JOIN queryConfigs AS c ON q.id = c.queryId
and q.activeVersion >= c.addedVersion
and q.activeVersion < COALESCE(c.removedVersion, q.activeVersion + 1)
and (c.removedVersion is null or q.activeVersion < c.removedVersion)
),
requiredIds as (
SELECT r.queryId,
ARRAY_AGG(DISTINCT r.requiredQueryId)
FILTER (WHERE r.requiredQueryId != '00000000-0000-0000-0000-000000000000')::uuid[]
as requiredIds
FROM queries AS q
LEFT JOIN requiredQueries AS r ON q.id = r.queryId
and q.activeVersion >= r.addedVersion
and q.activeVersion < COALESCE(r.removedVersion, q.activeVersion + 1)
GROUP BY q.id, q.type, q.activeversion, q.latestversion, c.config;
and (r.removedVersion is null or q.activeVersion < r.removedVersion)
GROUP BY r.queryId
)
SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, c.config,
coalesce(
r.requiredIds,
array[]::uuid[]
)::uuid[] as requiredIds
FROM queries AS q
LEFT JOIN config AS c ON q.id = c.queryId
LEFT JOIN requiredIds AS r ON q.id = r.queryId;
CREATE VIEW queryActiveDependencies AS
WITH RECURSIVE queryActiveDependencies AS (
+1 -1
View File
@@ -17,4 +17,4 @@ INSERT INTO jobs (clientId) VALUES ($1) RETURNING id;
INSERT INTO jobCanSync (canSync, jobId) VALUES ($1, $2);
-- name: ListJobDocumentIDsBatch :many
SELECT id, totalCount FROM listJobDocumentIDs(@jobId, @batchSize, @pageOffset);
SELECT id, totalCount FROM listJobDocumentIDs(@jobId, @batchSize, @pageOffset);
+31 -9
View File
@@ -5,16 +5,35 @@ SELECT id, config FROM queryConfigs where queryId = $1 and addedVersion >= $2 an
SELECT * FROM fullActiveQueries WHERE id = $1;
-- name: GetQueryWithVersion :one
SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, coalesce(c.config, null) as config, ARRAY_AGG(DISTINCT r.requiredQueryId)::uuid[] as requiredIds
FROM queries AS q
WITH query as (
SELECT id, type, activeVersion, latestVersion FROM queries WHERE id = @id
),
config as (
SELECT c.queryId, c.config
FROM query AS q
LEFT JOIN queryConfigs AS c ON q.id = c.queryId
and $2 >= c.addedVersion
and $2 < COALESCE(c.removedVersion, $2 + 1)
and @version >= c.addedVersion
and (c.removedVersion is null or @version < c.removedVersion)
),
requiredIds as (
SELECT r.queryId,
ARRAY_AGG(DISTINCT r.requiredQueryId)
FILTER (WHERE r.requiredQueryId != '00000000-0000-0000-0000-000000000000')::uuid[]
as requiredIds
FROM query AS q
LEFT JOIN requiredQueries AS r ON q.id = r.queryId
and $2 >= r.addedVersion
and $2 < COALESCE(r.removedVersion, $2 + 1)
WHERE q.id = $1
GROUP BY q.id, q.type, q.activeversion, q.latestversion, c.config;
and @version >= r.addedVersion
and (r.removedVersion is null or @version < r.removedVersion)
GROUP BY r.queryId
)
SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, c.config,
coalesce(
r.requiredIds,
array[]::uuid[]
)::uuid[] as requiredIds
FROM query AS q
LEFT JOIN config AS c ON q.id = c.queryId
LEFT JOIN requiredIds AS r ON q.id = r.queryId;
-- name: ListQueries :many
SELECT * FROM fullActiveQueries;
@@ -57,4 +76,7 @@ WITH doc AS (
SELECT dt.queryId
FROM collectorQueryDependencyTree as dt
JOIN doc as d on d.jobId = dt.jobId
where $1 = any(dt.requiredIds);
where $1 = any(dt.requiredIds);
-- name: ListQueryJobIDs :many
SELECT jobId FROM collectorQueryDependencyTree WHERE queryId = $1;
+8 -5
View File
@@ -7,17 +7,21 @@ WITH reqQueries as (
and @version >= rq.addedVersion
and (rq.removedVersion is null or @version < rq.removedVersion)
),
docs as (
SELECT id, jobId
FROM documents
WHERE id = @documentId
),
codeVersions as (
SELECT
d.id as documentId,
coalesce(ccv.minCleanVersion, 1) as minCleanVersion,
coalesce(ccv.minTextVersion, 1) as minTextVersion
FROM documents as d
FROM docs as d
LEFT JOIN collectors as c on c.jobId = d.jobId
LEFT JOIN collectorCodeVersions as ccv on c.id = ccv.collectorId
and c.activeVersion >= ccv.addedVersion
and c.activeVersion < COALESCE(ccv.removedVersion, c.activeVersion)
WHERE d.id = @documentId
LIMIT 1
),
latestVersions AS (
@@ -34,13 +38,12 @@ latestVersions AS (
FROM reqQueries as rq
LEFT JOIN results as r ON rq.queryId = r.queryId
and r.documentId = @documentId
JOIN codeVersions as ccv on ccv.documentId = r.documentId
WHERE r.documentId = @documentId
and r.queryVersion = rq.activeVersion
JOIN codeVersions as ccv on ccv.documentId = r.documentId
and r.cleanVersion >= ccv.minCleanVersion
and r.textVersion >= ccv.minTextVersion
)
SELECT lv.queryId, lv.type, r.value
SELECT DISTINCT lv.queryId, lv.type, r.value
FROM latestVersions as lv
JOIN results as r ON r.queryId = lv.queryId
and r.documentId = @documentId
+38 -1
View File
@@ -1,4 +1,3 @@
# compose.local.yaml
---
services:
doc_init_runner:
@@ -195,6 +194,7 @@ services:
environment:
LOG_LEVEL: DEBUG
JOB_SYNC_URL: ${JOB_SYNC_URL}
QUERY_VERSION_SYNC_URL: ${QUERY_VERSION_SYNC_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
@@ -241,12 +241,49 @@ services:
networks:
- server-network
query_version_sync_runner:
image: queryorchestration:latest
command: ["./queryVersionSyncRunner"]
depends_on:
localstack:
condition: service_healthy
db:
condition: service_healthy
ports:
- "8088:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
QUEUE_URL: ${QUERY_VERSION_SYNC_URL}
JOB_SYNC_URL: ${JOB_SYNC_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
DB_USER: ${DB_USER}
DB_PASS: ${DB_PASS}
DB_HOST: db
DB_PORT: 5432
DB_NAME: ${DB_NAME}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- server-network
prometheus:
image: prom/prometheus:latest
ports:
- "9091:9090"
volumes:
- ./prometheus.allmetrics.yaml:/etc/prometheus/prometheus.yml
healthcheck:
test: ["CMD", "wget", "--spider", "-q", "http://localhost:9090/-/healthy"]
interval: 30s
timeout: 10s
retries: 3
start_period: 30s
networks:
- server-network
+6 -2
View File
@@ -26,13 +26,17 @@ services:
- "9090:9090"
volumes:
- ./prometheus.localservice.yaml:/etc/prometheus/prometheus.yml
healthcheck:
test: ["CMD", "wget", "--spider", "-q", "http://localhost:9090/-/healthy"]
interval: 30s
timeout: 10s
retries: 3
start_period: 30s
networks:
- test-server-network
extra_hosts:
- "host.docker.internal:host-gateway"
volumes:
test-db-data:
networks:
test-server-network:
+12
View File
@@ -39,3 +39,15 @@ scrape_configs:
- targets:
- 'query_service:8080'
metrics_path: '/metrics'
- job_name: 'job_sync_runner'
static_configs:
- targets:
- 'job_sync_runner:8080'
metrics_path: '/metrics'
- job_name: 'query_version_sync_runner'
static_configs:
- targets:
- 'query_version_sync_runner:8080'
metrics_path: '/metrics'
+3 -1
View File
@@ -47,13 +47,15 @@
"QNAME_QUERY_SYNC": "query_sync",
"QNAME_QUERY_RUNNER": "query_runner",
"QNAME_JOB_SYNC": "job_sync",
"QNAME_QUERY_VERSION_SYNC": "query_version_sync",
"DOCUMENT_INIT_URL": "http://localstack:4566/queue/us-east-1/000000000000/document_init",
"DOCUMENT_SYNC_URL": "http://localstack:4566/queue/us-east-1/000000000000/document_sync",
"DOCUMENT_CLEAN_URL": "http://localstack:4566/queue/us-east-1/000000000000/document_clean",
"DOCUMENT_TEXT_URL": "http://localstack:4566/queue/us-east-1/000000000000/document_text",
"QUERY_SYNC_URL": "http://localstack:4566/queue/us-east-1/000000000000/query_sync",
"QUERY_URL": "http://localstack:4566/queue/us-east-1/000000000000/query_runner",
"QUERY_URL": "http://localstack:4566/queue/us-east-1/000000000000/query_runner",
"JOB_SYNC_URL": "http://localstack:4566/queue/us-east-1/000000000000/job_sync",
"QUERY_VERSION_SYNC_URL": "http://localstack:4566/queue/us-east-1/000000000000/query_version_sync",
"BUCKET_IN": "documentin"
},
"env_from": ".env"
+3 -3
View File
@@ -3,7 +3,7 @@ module queryorchestration
go 1.23
require (
github.com/aws/aws-sdk-go-v2 v1.36.0
github.com/aws/aws-sdk-go-v2 v1.36.1
github.com/aws/aws-sdk-go-v2/config v1.29.4
github.com/aws/aws-sdk-go-v2/service/s3 v1.75.2
github.com/aws/aws-sdk-go-v2/service/sqs v1.37.12
@@ -92,8 +92,8 @@ require (
github.com/Microsoft/go-winio v0.6.2 // indirect
github.com/aws/aws-sdk-go-v2/credentials v1.17.57 // indirect
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.27 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.31 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.31 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.32 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.32 // indirect
github.com/aws/aws-sdk-go-v2/internal/ini v1.8.2 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.2 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.12 // indirect
+6 -6
View File
@@ -11,8 +11,8 @@ github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA
github.com/RaveNoX/go-jsoncommentstrip v1.0.0/go.mod h1:78ihd09MekBnJnxpICcwzCMzGrKSKYe4AqU6PDYYpjk=
github.com/apapsch/go-jsonmerge/v2 v2.0.0 h1:axGnT1gRIfimI7gJifB699GoE/oq+F2MU7Dml6nw9rQ=
github.com/apapsch/go-jsonmerge/v2 v2.0.0/go.mod h1:lvDnEdqiQrp0O42VQGgmlKpxL1AP2+08jFMw88y4klk=
github.com/aws/aws-sdk-go-v2 v1.36.0 h1:b1wM5CcE65Ujwn565qcwgtOTT1aT4ADOHHgglKjG7fk=
github.com/aws/aws-sdk-go-v2 v1.36.0/go.mod h1:5PMILGVKiW32oDzjj6RU52yrNrDPUHcbZQYr1sM7qmM=
github.com/aws/aws-sdk-go-v2 v1.36.1 h1:iTDl5U6oAhkNPba0e1t1hrwAo02ZMqbrGq4k5JBWM5E=
github.com/aws/aws-sdk-go-v2 v1.36.1/go.mod h1:5PMILGVKiW32oDzjj6RU52yrNrDPUHcbZQYr1sM7qmM=
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.8 h1:zAxi9p3wsZMIaVCdoiQp2uZ9k1LsZvmAnoTBeZPXom0=
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.8/go.mod h1:3XkePX5dSaxveLAYY7nsbsZZrKxCyEuE5pM4ziFxyGg=
github.com/aws/aws-sdk-go-v2/config v1.29.4 h1:ObNqKsDYFGr2WxnoXKOhCvTlf3HhwtoGgc+KmZ4H5yg=
@@ -21,10 +21,10 @@ github.com/aws/aws-sdk-go-v2/credentials v1.17.57 h1:kFQDsbdBAR3GZsB8xA+51ptEnq9
github.com/aws/aws-sdk-go-v2/credentials v1.17.57/go.mod h1:2kerxPUUbTagAr/kkaHiqvj/bcYHzi2qiJS/ZinllU0=
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.27 h1:7lOW8NUwE9UZekS1DYoiPdVAqZ6A+LheHWb+mHbNOq8=
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.27/go.mod h1:w1BASFIPOPUae7AgaH4SbjNbfdkxuggLyGfNFTn8ITY=
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.31 h1:lWm9ucLSRFiI4dQQafLrEOmEDGry3Swrz0BIRdiHJqQ=
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.31/go.mod h1:Huu6GG0YTfbPphQkDSo4dEGmQRTKb9k9G7RdtyQWxuI=
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.31 h1:ACxDklUKKXb48+eg5ROZXi1vDgfMyfIA/WyvqHcHI0o=
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.31/go.mod h1:yadnfsDwqXeVaohbGc/RaD287PuyRw2wugkh5ZL2J6k=
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.32 h1:BjUcr3X3K0wZPGFg2bxOWW3VPN8rkE3/61zhP+IHviA=
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.32/go.mod h1:80+OGC/bgzzFFTUmcuwD0lb4YutwQeKLFpmt6hoWapU=
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.32 h1:m1GeXHVMJsRsUAqG6HjZWx9dj7F5TR+cF1bjyfYyBd4=
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.32/go.mod h1:IitoQxGfaKdVLNg0hD8/DXmAqNy0H4K2H2Sf91ti8sI=
github.com/aws/aws-sdk-go-v2/internal/ini v1.8.2 h1:Pg9URiobXy85kgFev3og2CuOZ8JZUBENF+dcgWBaYNk=
github.com/aws/aws-sdk-go-v2/internal/ini v1.8.2/go.mod h1:FbtygfRFze9usAadmnGJNc8KsP346kEe+y2/oyhGAGc=
github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.31 h1:8IwBjuLdqIO1dGB+dZ9zJEl8wzY3bVYxcs0Xyu/Lsc0=
+80 -15
View File
@@ -125,21 +125,40 @@ func (q *Queries) GetQueryConfig(ctx context.Context, arg *GetQueryConfigParams)
}
const getQueryWithVersion = `-- name: GetQueryWithVersion :one
SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, coalesce(c.config, null) as config, ARRAY_AGG(DISTINCT r.requiredQueryId)::uuid[] as requiredIds
FROM queries AS q
WITH query as (
SELECT id, type, activeVersion, latestVersion FROM queries WHERE id = $1
),
config as (
SELECT c.queryId, c.config
FROM query AS q
LEFT JOIN queryConfigs AS c ON q.id = c.queryId
and $2 >= c.addedVersion
and $2 < COALESCE(c.removedVersion, $2 + 1)
and (c.removedVersion is null or $2 < c.removedVersion)
),
requiredIds as (
SELECT r.queryId,
ARRAY_AGG(DISTINCT r.requiredQueryId)
FILTER (WHERE r.requiredQueryId != '00000000-0000-0000-0000-000000000000')::uuid[]
as requiredIds
FROM query AS q
LEFT JOIN requiredQueries AS r ON q.id = r.queryId
and $2 >= r.addedVersion
and $2 < COALESCE(r.removedVersion, $2 + 1)
WHERE q.id = $1
GROUP BY q.id, q.type, q.activeversion, q.latestversion, c.config
and (r.removedVersion is null or $2 < r.removedVersion)
GROUP BY r.queryId
)
SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, c.config,
coalesce(
r.requiredIds,
array[]::uuid[]
)::uuid[] as requiredIds
FROM query AS q
LEFT JOIN config AS c ON q.id = c.queryId
LEFT JOIN requiredIds AS r ON q.id = r.queryId
`
type GetQueryWithVersionParams struct {
ID pgtype.UUID `db:"id"`
Addedversion int32 `db:"addedversion"`
ID pgtype.UUID `db:"id"`
Version *int32 `db:"version"`
}
type GetQueryWithVersionRow struct {
@@ -153,18 +172,37 @@ type GetQueryWithVersionRow struct {
// GetQueryWithVersion
//
// SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, coalesce(c.config, null) as config, ARRAY_AGG(DISTINCT r.requiredQueryId)::uuid[] as requiredIds
// FROM queries AS q
// WITH query as (
// SELECT id, type, activeVersion, latestVersion FROM queries WHERE id = $1
// ),
// config as (
// SELECT c.queryId, c.config
// FROM query AS q
// LEFT JOIN queryConfigs AS c ON q.id = c.queryId
// and $2 >= c.addedVersion
// and $2 < COALESCE(c.removedVersion, $2 + 1)
// and (c.removedVersion is null or $2 < c.removedVersion)
// ),
// requiredIds as (
// SELECT r.queryId,
// ARRAY_AGG(DISTINCT r.requiredQueryId)
// FILTER (WHERE r.requiredQueryId != '00000000-0000-0000-0000-000000000000')::uuid[]
// as requiredIds
// FROM query AS q
// LEFT JOIN requiredQueries AS r ON q.id = r.queryId
// and $2 >= r.addedVersion
// and $2 < COALESCE(r.removedVersion, $2 + 1)
// WHERE q.id = $1
// GROUP BY q.id, q.type, q.activeversion, q.latestversion, c.config
// and (r.removedVersion is null or $2 < r.removedVersion)
// GROUP BY r.queryId
// )
// SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, c.config,
// coalesce(
// r.requiredIds,
// array[]::uuid[]
// )::uuid[] as requiredIds
// FROM query AS q
// LEFT JOIN config AS c ON q.id = c.queryId
// LEFT JOIN requiredIds AS r ON q.id = r.queryId
func (q *Queries) GetQueryWithVersion(ctx context.Context, arg *GetQueryWithVersionParams) (*GetQueryWithVersionRow, error) {
row := q.db.QueryRow(ctx, getQueryWithVersion, arg.ID, arg.Addedversion)
row := q.db.QueryRow(ctx, getQueryWithVersion, arg.ID, arg.Version)
var i GetQueryWithVersionRow
err := row.Scan(
&i.ID,
@@ -312,6 +350,33 @@ func (q *Queries) ListQueryDirectDependentsByDocumentID(ctx context.Context, arg
return items, nil
}
const listQueryJobIDs = `-- name: ListQueryJobIDs :many
SELECT jobId FROM collectorQueryDependencyTree WHERE queryId = $1
`
// ListQueryJobIDs
//
// SELECT jobId FROM collectorQueryDependencyTree WHERE queryId = $1
func (q *Queries) ListQueryJobIDs(ctx context.Context, queryid pgtype.UUID) ([]pgtype.UUID, error) {
rows, err := q.db.Query(ctx, listQueryJobIDs, queryid)
if err != nil {
return nil, err
}
defer rows.Close()
items := []pgtype.UUID{}
for rows.Next() {
var jobid pgtype.UUID
if err := rows.Scan(&jobid); err != nil {
return nil, err
}
items = append(items, jobid)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const removeQueryConfig = `-- name: RemoveQueryConfig :exec
UPDATE queryConfigs SET removedVersion = $1 WHERE queryId = $2 and removedVersion is null
`
+111 -2
View File
@@ -67,6 +67,17 @@ func TestQueries(t *testing.T) {
})
assert.NoError(t, err)
jsonQuery, err = queries.GetQuery(ctx, jsonQueryID)
assert.NoError(t, err)
assert.EqualExportedValues(t, &repository.Fullactivequery{
ID: jsonQueryID,
Type: repository.QuerytypeJsonExtractor,
Activeversion: 1,
Latestversion: 2,
Config: nil,
Requiredids: []pgtype.UUID{contextQueryID},
}, jsonQuery)
removeV := int32(2)
err = queries.RemoveRequiredQuery(ctx, &repository.RemoveRequiredQueryParams{
Queryid: jsonQueryID,
@@ -131,9 +142,10 @@ func TestQueries(t *testing.T) {
Requiredids: []pgtype.UUID{},
}, jsonQuery)
v := int32(1)
versionedQuery, err := queries.GetQueryWithVersion(ctx, &repository.GetQueryWithVersionParams{
ID: jsonQueryID,
Addedversion: 1,
ID: jsonQueryID,
Version: &v,
})
assert.NoError(t, err)
assert.EqualExportedValues(t, &repository.GetQueryWithVersionRow{
@@ -420,3 +432,100 @@ func TestQueriesList(t *testing.T) {
},
}, qs)
}
func TestListQueryJobs(t *testing.T) {
if testing.Short() {
t.Skip("Skipping long test in short mode")
}
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
test.SetCfgProvider(t, cfg)
cfg.SetBasePath(path.Join(os.Getenv("PWD"), "../../.."))
_, cleanup := test.CreateDB(t, ctx, &test.CreateDatabaseConfig{
Cfg: cfg,
RunMigrations: true,
})
defer cleanup()
queries := cfg.GetDBQueries()
contextID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeContextFull))
assert.NoError(t, err)
jobs, err := queries.ListQueryJobIDs(ctx, contextID)
assert.NoError(t, err)
assert.ElementsMatch(t, []pgtype.UUID{}, jobs)
clientOneID, err := queries.CreateClient(ctx, "example_client")
assert.NoError(t, err)
jobOneID, err := queries.CreateJob(ctx, clientOneID)
assert.NoError(t, err)
collOneID, err := queries.CreateCollector(ctx, jobOneID)
assert.NoError(t, err)
err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{
Collectorid: collOneID,
Queryid: contextID,
Addedversion: 1,
Name: "example_key",
})
assert.NoError(t, err)
jobs, err = queries.ListQueryJobIDs(ctx, contextID)
assert.NoError(t, err)
assert.ElementsMatch(t, []pgtype.UUID{jobOneID}, jobs)
jobTwoID, err := queries.CreateJob(ctx, clientOneID)
assert.NoError(t, err)
collTwoID, err := queries.CreateCollector(ctx, jobTwoID)
assert.NoError(t, err)
err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{
Collectorid: collTwoID,
Queryid: contextID,
Addedversion: 1,
Name: "example_key",
})
assert.NoError(t, err)
jobs, err = queries.ListQueryJobIDs(ctx, contextID)
assert.NoError(t, err)
assert.ElementsMatch(t, []pgtype.UUID{jobOneID, jobTwoID}, jobs)
clientTwoID, err := queries.CreateClient(ctx, "example_client_two")
assert.NoError(t, err)
jobThreeID, err := queries.CreateJob(ctx, clientTwoID)
assert.NoError(t, err)
collThreeID, err := queries.CreateCollector(ctx, jobThreeID)
assert.NoError(t, err)
err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{
Collectorid: collThreeID,
Queryid: contextID,
Addedversion: 1,
Name: "example_key",
})
assert.NoError(t, err)
jobs, err = queries.ListQueryJobIDs(ctx, contextID)
assert.NoError(t, err)
assert.ElementsMatch(t, []pgtype.UUID{jobOneID, jobTwoID, jobThreeID}, jobs)
jsonID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor))
assert.NoError(t, err)
err = queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{
Queryid: jsonID,
Requiredqueryid: contextID,
Addedversion: 1,
})
assert.NoError(t, err)
err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{
Collectorid: collOneID,
Queryid: jsonID,
Addedversion: 1,
Name: "example_key",
})
assert.NoError(t, err)
jobs, err = queries.ListQueryJobIDs(ctx, jsonID)
assert.NoError(t, err)
assert.ElementsMatch(t, []pgtype.UUID{jobOneID}, jobs)
}
+16 -10
View File
@@ -53,17 +53,21 @@ WITH reqQueries as (
and $3 >= rq.addedVersion
and (rq.removedVersion is null or $3 < rq.removedVersion)
),
docs as (
SELECT id, jobId
FROM documents
WHERE id = $1
),
codeVersions as (
SELECT
d.id as documentId,
coalesce(ccv.minCleanVersion, 1) as minCleanVersion,
coalesce(ccv.minTextVersion, 1) as minTextVersion
FROM documents as d
FROM docs as d
LEFT JOIN collectors as c on c.jobId = d.jobId
LEFT JOIN collectorCodeVersions as ccv on c.id = ccv.collectorId
and c.activeVersion >= ccv.addedVersion
and c.activeVersion < COALESCE(ccv.removedVersion, c.activeVersion)
WHERE d.id = $1
LIMIT 1
),
latestVersions AS (
@@ -80,13 +84,12 @@ latestVersions AS (
FROM reqQueries as rq
LEFT JOIN results as r ON rq.queryId = r.queryId
and r.documentId = $1
JOIN codeVersions as ccv on ccv.documentId = r.documentId
WHERE r.documentId = $1
and r.queryVersion = rq.activeVersion
JOIN codeVersions as ccv on ccv.documentId = r.documentId
and r.cleanVersion >= ccv.minCleanVersion
and r.textVersion >= ccv.minTextVersion
)
SELECT lv.queryId, lv.type, r.value
SELECT DISTINCT lv.queryId, lv.type, r.value
FROM latestVersions as lv
JOIN results as r ON r.queryId = lv.queryId
and r.documentId = $1
@@ -118,17 +121,21 @@ type ListQueryRequirementValuesRow struct {
// and $3 >= rq.addedVersion
// and (rq.removedVersion is null or $3 < rq.removedVersion)
// ),
// docs as (
// SELECT id, jobId
// FROM documents
// WHERE id = $1
// ),
// codeVersions as (
// SELECT
// d.id as documentId,
// coalesce(ccv.minCleanVersion, 1) as minCleanVersion,
// coalesce(ccv.minTextVersion, 1) as minTextVersion
// FROM documents as d
// FROM docs as d
// LEFT JOIN collectors as c on c.jobId = d.jobId
// LEFT JOIN collectorCodeVersions as ccv on c.id = ccv.collectorId
// and c.activeVersion >= ccv.addedVersion
// and c.activeVersion < COALESCE(ccv.removedVersion, c.activeVersion)
// WHERE d.id = $1
// LIMIT 1
// ),
// latestVersions AS (
@@ -145,13 +152,12 @@ type ListQueryRequirementValuesRow struct {
// FROM reqQueries as rq
// LEFT JOIN results as r ON rq.queryId = r.queryId
// and r.documentId = $1
// JOIN codeVersions as ccv on ccv.documentId = r.documentId
// WHERE r.documentId = $1
// and r.queryVersion = rq.activeVersion
// JOIN codeVersions as ccv on ccv.documentId = r.documentId
// and r.cleanVersion >= ccv.minCleanVersion
// and r.textVersion >= ccv.minTextVersion
// )
// SELECT lv.queryId, lv.type, r.value
// SELECT DISTINCT lv.queryId, lv.type, r.value
// FROM latestVersions as lv
// JOIN results as r ON r.queryId = lv.queryId
// and r.documentId = $1
@@ -398,5 +398,4 @@ func TestUnsyncedNoDepsQueries(t *testing.T) {
assert.NoError(t, err)
assert.Len(t, qs, 1)
assert.ElementsMatch(t, []pgtype.UUID{contextQueryID}, qs)
}
+33 -21
View File
@@ -209,28 +209,30 @@ func (s *Service) submitUpdate(ctx context.Context, current *collector.Collector
}
}
removeIDs := getRemoveFields(current.Fields, params.Fields)
for _, field := range removeIDs {
err := qtx.RemoveCollectorQuery(ctx, &repository.RemoveCollectorQueryParams{
Collectorid: id,
Queryid: field,
Removedversion: &latestVersion,
})
if err != nil {
return err
if params.Fields != nil {
removeIDs := getRemoveFields(current.Fields, params.Fields)
for _, field := range removeIDs {
err := qtx.RemoveCollectorQuery(ctx, &repository.RemoveCollectorQueryParams{
Collectorid: id,
Queryid: field,
Removedversion: &latestVersion,
})
if err != nil {
return err
}
}
}
addIDs := getAddFields(current.Fields, params.Fields)
for key, field := range addIDs {
err := qtx.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{
Collectorid: id,
Name: key,
Queryid: field,
Addedversion: latestVersion,
})
if err != nil {
return err
addIDs := getAddFields(current.Fields, params.Fields)
for key, field := range addIDs {
err := qtx.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{
Collectorid: id,
Name: key,
Queryid: field,
Addedversion: latestVersion,
})
if err != nil {
return err
}
}
}
@@ -239,7 +241,17 @@ func (s *Service) submitUpdate(ctx context.Context, current *collector.Collector
activeVersion = &current.ActiveVersion
}
err := qtx.UpdateCollector(ctx, &repository.UpdateCollectorParams{
activeName, err := validation.GetFieldName(params, params.ActiveVersion)
if err != nil {
return err
}
onlyactive := validation.AreAllPointersNilExcept(params, activeName)
if onlyactive {
latestVersion--
}
err = qtx.UpdateCollector(ctx, &repository.UpdateCollectorParams{
ID: id,
Latestversion: latestVersion,
Activeversion: *activeVersion,
@@ -172,70 +172,113 @@ func TestGetUpdateParams(t *testing.T) {
}
func TestSubmitUpdate(t *testing.T) {
ctx := context.Background()
t.Run("all fields", func(t *testing.T) {
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &CollectorUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
ctx := context.Background()
svc := Service{
cfg: cfg,
svc: &Services{},
}
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &CollectorUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
minCleanV := int32(2)
minTextV := int32(4)
current := collector.Collector{
ID: uuid.New(),
JobID: uuid.New(),
ActiveVersion: 1,
LatestVersion: 1,
MinCleanVersion: 1,
MinTextVersion: 2,
Fields: map[string]uuid.UUID{
"original_key": uuid.New(),
"og_key": uuid.New(),
},
}
aV := int32(2)
params := dbUpdateParams{
JobID: database.MustToDBUUID(current.JobID),
ActiveVersion: &aV,
MinCleanVersion: &minCleanV,
MinTextVersion: &minTextV,
Fields: &map[string]pgtype.UUID{
"example_key": database.MustToDBUUID(uuid.New()),
"second_key": database.MustToDBUUID(uuid.New()),
"changed_key": database.MustToDBUUID(current.Fields["original_key"]),
"og_key": database.MustToDBUUID(current.Fields["og_key"]),
},
}
svc := Service{
cfg: cfg,
svc: &Services{},
}
pool.ExpectBeginTx(pgx.TxOptions{})
rv := int32(2)
pool.ExpectExec("name: RemoveCollectorCodeVersion :exec").WithArgs(&rv, database.MustToDBUUID(current.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: AddCollectorCodeVersion :exec").WithArgs(database.MustToDBUUID(current.ID), int32(2), *params.MinCleanVersion, *params.MinTextVersion).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: RemoveCollectorQuery :exec").WithArgs(&rv, database.MustToDBUUID(current.Fields["original_key"]), database.MustToDBUUID(current.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.MatchExpectationsInOrder(false)
pool.ExpectExec("name: AddCollectorQuery :exec").WithArgs(database.MustToDBUUID(current.ID), "example_key", (*params.Fields)["example_key"], int32(2)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: AddCollectorQuery :exec").WithArgs(database.MustToDBUUID(current.ID), "second_key", (*params.Fields)["second_key"], int32(2)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: AddCollectorQuery :exec").WithArgs(database.MustToDBUUID(current.ID), "changed_key", (*params.Fields)["changed_key"], int32(2)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: UpdateCollector :exec").WithArgs(int32(2), int32(2), database.MustToDBUUID(current.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectCommit()
minCleanV := int32(2)
minTextV := int32(4)
current := collector.Collector{
ID: uuid.New(),
JobID: uuid.New(),
ActiveVersion: 1,
LatestVersion: 1,
MinCleanVersion: 1,
MinTextVersion: 2,
Fields: map[string]uuid.UUID{
"original_key": uuid.New(),
"og_key": uuid.New(),
},
}
aV := int32(2)
params := dbUpdateParams{
JobID: database.MustToDBUUID(current.JobID),
ActiveVersion: &aV,
MinCleanVersion: &minCleanV,
MinTextVersion: &minTextV,
Fields: &map[string]pgtype.UUID{
"example_key": database.MustToDBUUID(uuid.New()),
"second_key": database.MustToDBUUID(uuid.New()),
"changed_key": database.MustToDBUUID(current.Fields["original_key"]),
"og_key": database.MustToDBUUID(current.Fields["og_key"]),
},
}
err = svc.submitUpdate(ctx, &current, &params)
assert.NoError(t, err)
pool.ExpectBeginTx(pgx.TxOptions{})
rv := int32(2)
pool.ExpectExec("name: RemoveCollectorCodeVersion :exec").WithArgs(&rv, database.MustToDBUUID(current.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: AddCollectorCodeVersion :exec").WithArgs(database.MustToDBUUID(current.ID), int32(2), *params.MinCleanVersion, *params.MinTextVersion).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: RemoveCollectorQuery :exec").WithArgs(&rv, database.MustToDBUUID(current.Fields["original_key"]), database.MustToDBUUID(current.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.MatchExpectationsInOrder(false)
pool.ExpectExec("name: AddCollectorQuery :exec").WithArgs(database.MustToDBUUID(current.ID), "example_key", (*params.Fields)["example_key"], int32(2)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: AddCollectorQuery :exec").WithArgs(database.MustToDBUUID(current.ID), "second_key", (*params.Fields)["second_key"], int32(2)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: AddCollectorQuery :exec").WithArgs(database.MustToDBUUID(current.ID), "changed_key", (*params.Fields)["changed_key"], int32(2)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: UpdateCollector :exec").WithArgs(int32(2), int32(2), database.MustToDBUUID(current.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectCommit()
err = svc.submitUpdate(ctx, &current, &params)
assert.NoError(t, err)
})
t.Run("only active version", func(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &CollectorUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := Service{
cfg: cfg,
svc: &Services{},
}
current := collector.Collector{
ID: uuid.New(),
JobID: uuid.New(),
ActiveVersion: 1,
LatestVersion: 1,
MinCleanVersion: 1,
MinTextVersion: 2,
Fields: map[string]uuid.UUID{},
}
av := int32(2)
params := dbUpdateParams{
JobID: database.MustToDBUUID(current.JobID),
ActiveVersion: &av,
}
pool.ExpectBeginTx(pgx.TxOptions{})
pool.ExpectExec("name: UpdateCollector :exec").WithArgs(int32(1), av, database.MustToDBUUID(current.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectCommit()
err = svc.submitUpdate(ctx, &current, &params)
assert.NoError(t, err)
})
}
func TestNormalizeActiveVersion(t *testing.T) {
+4 -1
View File
@@ -4,6 +4,7 @@ import (
"context"
"database/sql"
"errors"
"log/slog"
docsyncrunner "queryorchestration/api/docSyncRunner"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
@@ -29,6 +30,8 @@ func (s *Service) Sync(ctx context.Context, id uuid.UUID) error {
return err
}
slog.Debug("batch processing", "id", id, "offset", offset, "length", len(ids))
for _, id := range ids {
err := s.cfg.SendToQueue(ctx, &queue.SendParams{
QueueURL: s.cfg.GetDocumentSyncURL(),
@@ -42,7 +45,7 @@ func (s *Service) Sync(ctx context.Context, id uuid.UUID) error {
}
offset += s.batchSize
if int64(offset) >= *ids[0].Totalcount {
if ids == nil || int64(offset) >= *ids[0].Totalcount {
hasMore = false
}
}
+1 -1
View File
@@ -33,7 +33,7 @@ func (s *Service) normalizeCreate(ctx context.Context, entity *resultprocessor.C
return err
}
err = s.normalizeConfig(entity)
err = s.NormalizeConfig(entity)
if err != nil {
return err
}
+2 -2
View File
@@ -20,8 +20,8 @@ type Query struct {
func (s *Service) GetWithVersion(ctx context.Context, id uuid.UUID, version int32) (*Query, error) {
query, err := s.cfg.GetDBQueries().GetQueryWithVersion(ctx, &repository.GetQueryWithVersionParams{
ID: database.MustToDBUUID(id),
Addedversion: version,
ID: database.MustToDBUUID(id),
Version: &version,
})
if err != nil {
return nil, err
+1 -1
View File
@@ -79,7 +79,7 @@ func TestGetWithVersion(t *testing.T) {
dbReqIDs := database.MustToDBUUIDArray(*query.RequiredQueryIDs)
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(query.ID), version).WillReturnRows(
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(query.ID), &version).WillReturnRows(
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
AddRow(database.MustToDBUUID(query.ID), repository.QuerytypeJsonExtractor, query.ActiveVersion, query.LatestVersion, []byte(config), dbReqIDs),
)
+2 -2
View File
@@ -18,7 +18,7 @@ type Config interface {
SetConfig(*string)
}
func (s *Service) normalizeConfig(config Config) error {
func (s *Service) NormalizeConfig(config Config) error {
if config == nil || config.GetConfig() == nil {
return nil
}
@@ -83,7 +83,7 @@ func (s *Service) NormalizeQueryIDs(ctx context.Context, ids RequiredQueryIDs) e
return nil
}
func (s *Service) normalizeActiveVersion(current *Query, entity *resultprocessor.Update) error {
func (s *Service) NormalizeActiveVersion(current *Query, entity *resultprocessor.Update) error {
if current == nil {
return errors.New("current query required")
} else if entity == nil {
+18 -18
View File
@@ -17,54 +17,54 @@ import (
func TestNormalizeConfig(t *testing.T) {
s := Service{}
err := s.normalizeConfig(nil)
err := s.NormalizeConfig(nil)
assert.NoError(t, err)
entity := resultprocessor.Create{}
entity.Config = nil
err = s.normalizeConfig(&entity)
err = s.NormalizeConfig(&entity)
assert.NoError(t, err)
assert.Nil(t, entity.Config)
cfg := ""
entity.Config = &cfg
err = s.normalizeConfig(&entity)
err = s.NormalizeConfig(&entity)
assert.NoError(t, err)
assert.Nil(t, entity.Config)
cfg = " "
entity.Config = &cfg
err = s.normalizeConfig(&entity)
err = s.NormalizeConfig(&entity)
assert.NoError(t, err)
assert.Nil(t, entity.Config)
cfg = "{}"
entity.Config = &cfg
err = s.normalizeConfig(&entity)
err = s.NormalizeConfig(&entity)
assert.NoError(t, err)
assert.Equal(t, "{}", *(entity.Config))
cfg = "{\"hello\":\"bye\"}"
entity.Config = &cfg
err = s.normalizeConfig(&entity)
err = s.NormalizeConfig(&entity)
assert.NoError(t, err)
assert.Equal(t, "{\"hello\":\"bye\"}", *(entity.Config))
cfg = " { \"hello\" : \"bye\" } "
entity.Config = &cfg
err = s.normalizeConfig(&entity)
err = s.NormalizeConfig(&entity)
assert.NoError(t, err)
assert.Equal(t, "{\"hello\":\"bye\"}", *(entity.Config))
cfg = "{'hello':'bye'}"
entity.Config = &cfg
err = s.normalizeConfig(&entity)
err = s.NormalizeConfig(&entity)
assert.Error(t, err)
cfg = "{\"hello\":\"}"
entity.Config = &cfg
err = s.normalizeConfig(&entity)
err = s.NormalizeConfig(&entity)
assert.Error(t, err)
}
@@ -143,14 +143,14 @@ func TestNormalizeActiveVersion(t *testing.T) {
s := Service{}
t.Run("all nil", func(t *testing.T) {
err := s.normalizeActiveVersion(nil, nil)
err := s.NormalizeActiveVersion(nil, nil)
assert.Error(t, err)
})
t.Run("nil current", func(t *testing.T) {
entity := resultprocessor.Update{}
err := s.normalizeActiveVersion(nil, &entity)
err := s.NormalizeActiveVersion(nil, &entity)
assert.Error(t, err)
})
@@ -159,7 +159,7 @@ func TestNormalizeActiveVersion(t *testing.T) {
ActiveVersion: 2,
LatestVersion: 4,
}
err := s.normalizeActiveVersion(&current, nil)
err := s.NormalizeActiveVersion(&current, nil)
assert.NoError(t, err)
})
@@ -169,7 +169,7 @@ func TestNormalizeActiveVersion(t *testing.T) {
LatestVersion: 4,
}
entity := resultprocessor.Update{}
err := s.normalizeActiveVersion(&current, &entity)
err := s.NormalizeActiveVersion(&current, &entity)
assert.NoError(t, err)
assert.Nil(t, entity.ActiveVersion)
})
@@ -183,7 +183,7 @@ func TestNormalizeActiveVersion(t *testing.T) {
entity := resultprocessor.Update{
ActiveVersion: &version,
}
err := s.normalizeActiveVersion(&current, &entity)
err := s.NormalizeActiveVersion(&current, &entity)
assert.NoError(t, err)
assert.Nil(t, entity.ActiveVersion)
})
@@ -197,7 +197,7 @@ func TestNormalizeActiveVersion(t *testing.T) {
entity := resultprocessor.Update{
ActiveVersion: &version,
}
err := s.normalizeActiveVersion(&current, &entity)
err := s.NormalizeActiveVersion(&current, &entity)
assert.NoError(t, err)
assert.Equal(t, version, *entity.ActiveVersion)
})
@@ -211,7 +211,7 @@ func TestNormalizeActiveVersion(t *testing.T) {
entity := resultprocessor.Update{
ActiveVersion: &version,
}
err := s.normalizeActiveVersion(&current, &entity)
err := s.NormalizeActiveVersion(&current, &entity)
assert.NoError(t, err)
assert.Equal(t, version, *entity.ActiveVersion)
})
@@ -225,7 +225,7 @@ func TestNormalizeActiveVersion(t *testing.T) {
entity := resultprocessor.Update{
ActiveVersion: &version,
}
err := s.normalizeActiveVersion(&current, &entity)
err := s.NormalizeActiveVersion(&current, &entity)
assert.Error(t, err)
})
@@ -238,7 +238,7 @@ func TestNormalizeActiveVersion(t *testing.T) {
entity := resultprocessor.Update{
ActiveVersion: &version,
}
err := s.normalizeActiveVersion(&current, &entity)
err := s.NormalizeActiveVersion(&current, &entity)
assert.Error(t, err)
})
}
+10
View File
@@ -19,6 +19,16 @@ func ParseQueryWithVersion(q *repository.GetQueryWithVersionRow) (*Query, error)
})
}
func ParseQuery(q *Query) *resultprocessor.Query {
return &resultprocessor.Query{
ID: q.ID,
Type: q.Type,
Version: q.ActiveVersion,
RequiredQueryIDs: q.RequiredQueryIDs,
Config: q.Config,
}
}
func ParseFullActiveQuery(q *repository.Fullactivequery) (*Query, error) {
var reqQueryIDs *[]uuid.UUID
if len(q.Requiredids) > 0 {
+13 -3
View File
@@ -42,6 +42,14 @@ func (s *Service) Process(ctx context.Context, p *Process) (resultprocessor.Valu
return nil, err
}
if s.cfg.GetLogLevel() == slog.LevelDebug {
elength := 0
if processQuery.RequiredQueryIDs != nil {
elength = len(*processQuery.RequiredQueryIDs)
}
slog.Debug("requirements", "length", len(values), "expected_length", elength)
}
processor, err := s.getProcessor(query.Type)
if err != nil {
return nil, err
@@ -55,7 +63,7 @@ func (s *Service) Process(ctx context.Context, p *Process) (resultprocessor.Valu
return getValueByType(query.Type, val)
}
func (s *Service) listRequiredValues(ctx context.Context, p *Process, query *resultprocessor.Query) (*[]resultprocessor.Value, error) {
func (s *Service) listRequiredValues(ctx context.Context, p *Process, query *resultprocessor.Query) ([]resultprocessor.Value, error) {
if p == nil || query == nil || query.RequiredQueryIDs == nil || len(*query.RequiredQueryIDs) == 0 {
return nil, nil
}
@@ -71,10 +79,12 @@ func (s *Service) listRequiredValues(ctx context.Context, p *Process, query *res
return nil, fmt.Errorf("required results not found for the query (%s) in document (%s)", p.QueryID, p.DocumentID)
}
slog.Debug("found results", "length", len(qResults))
return parseQueryRequirementValueArray(qResults)
}
func parseQueryRequirementValueArray(v []*repository.ListQueryRequirementValuesRow) (*[]resultprocessor.Value, error) {
func parseQueryRequirementValueArray(v []*repository.ListQueryRequirementValuesRow) ([]resultprocessor.Value, error) {
values := make([]resultprocessor.Value, len(v))
for index, r := range v {
qType, err := resultprocessor.ParseDBType(r.Type)
@@ -90,7 +100,7 @@ func parseQueryRequirementValueArray(v []*repository.ListQueryRequirementValuesR
values[index] = cleanValue
}
return &values, nil
return values, nil
}
func (s *Service) getProcessor(queryType resultprocessor.Type) (resultprocessor.Processor, error) {
+3 -3
View File
@@ -47,7 +47,7 @@ func TestProcess(t *testing.T) {
QueryVersion: query.Version,
}
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(query.ID), query.Version).WillReturnRows(
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(query.ID), &query.Version).WillReturnRows(
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
AddRow(database.MustToDBUUID(query.ID), repository.QuerytypeJsonExtractor, query.Version, query.Version, []byte(*query.Config), database.MustToDBUUIDArray(*query.RequiredQueryIDs)),
)
@@ -119,7 +119,7 @@ func TestListRequiredValue(t *testing.T) {
assert.NoError(t, err)
assert.ElementsMatch(t, []resultprocessor.Value{
jsonextractor.NewResult("example_value"),
}, *pr)
}, pr)
}
func TestGetProcessor(t *testing.T) {
@@ -150,5 +150,5 @@ func TestParseQueryRequirementValueArray(t *testing.T) {
assert.NoError(t, err)
assert.ElementsMatch(t, []resultprocessor.Value{
jsonextractor.NewResult("example_value"),
}, *out)
}, out)
}
+1 -1
View File
@@ -80,5 +80,5 @@ type Updator interface {
}
type Processor interface {
Process(ctx context.Context, query *Query, values *[]Value) (string, error)
Process(ctx context.Context, query *Query, values []Value) (string, error)
}
+4 -8
View File
@@ -3,26 +3,22 @@ package resultset
import (
"queryorchestration/internal/query"
"queryorchestration/internal/query/result"
resultsync "queryorchestration/internal/query/result/sync"
"queryorchestration/internal/serviceconfig"
queryc "queryorchestration/internal/serviceconfig/queue/query"
)
type ConfigProvider interface {
serviceconfig.ConfigProvider
queryc.ConfigProvider
}
type Services struct {
Result *result.Service
Query *query.Service
Sync *resultsync.Service
}
type Service struct {
cfg ConfigProvider
cfg serviceconfig.ConfigProvider
svc *Services
}
func New(cfg ConfigProvider, svc *Services) *Service {
func New(cfg serviceconfig.ConfigProvider, svc *Services) *Service {
return &Service{
cfg,
svc,
+1 -24
View File
@@ -6,7 +6,6 @@ import (
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/query/result"
"queryorchestration/internal/serviceconfig/queue"
"github.com/google/uuid"
)
@@ -68,27 +67,5 @@ func (s *Service) informQueryDependents(ctx context.Context, params *Set) error
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,
})
return s.svc.Sync.TriggerMultiSync(ctx, params.DocumentID, database.MustToUUIDArray(ids))
}
+4 -100
View File
@@ -8,6 +8,7 @@ import (
"queryorchestration/internal/query"
"queryorchestration/internal/query/result"
resultprocessor "queryorchestration/internal/query/result/processor"
resultsync "queryorchestration/internal/query/result/sync"
"queryorchestration/internal/serviceconfig"
queryc "queryorchestration/internal/serviceconfig/queue/query"
queuemock "queryorchestration/mocks/queue"
@@ -48,6 +49,7 @@ func TestSet(t *testing.T) {
Result: result.New(cfg, &result.Services{
Query: que,
}),
Sync: resultsync.New(cfg),
},
}
@@ -77,7 +79,7 @@ func TestSet(t *testing.T) {
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
AddRow(database.MustToDBUUID(query.ID), repository.QuerytypeJsonExtractor, query.Version, query.Version, []byte(*query.Config), database.MustToDBUUIDArray(*query.RequiredQueryIDs)),
)
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(query.ID), query.Version).WillReturnRows(
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(query.ID), &query.Version).WillReturnRows(
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
AddRow(database.MustToDBUUID(query.ID), repository.QuerytypeJsonExtractor, query.Version, query.Version, []byte(*query.Config), database.MustToDBUUIDArray(*query.RequiredQueryIDs)),
)
@@ -112,105 +114,6 @@ func TestSet(t *testing.T) {
assert.NoError(t, err)
}
func TestTriggerQuerySync(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &ResultSetConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.QueryURL = "/i/am/here"
que := query.New(cfg)
svc := Service{
cfg: cfg,
svc: &Services{
Query: que,
Result: result.New(cfg, &result.Services{
Query: que,
}),
},
}
params := &Set{
DocumentID: uuid.New(),
QueryID: uuid.New(),
}
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.QueryURL && *in.MessageBody == fmt.Sprintf("{\"document_id\":\"%s\",\"query_id\":\"%s\"}", params.DocumentID.String(), params.QueryID.String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
err = svc.TriggerQuerySync(ctx, params)
assert.NoError(t, err)
}
func TestTriggerQueriesSync(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &ResultSetConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.QueryURL = "/i/am/here"
que := query.New(cfg)
svc := Service{
cfg: cfg,
svc: &Services{
Query: que,
Result: result.New(cfg, &result.Services{
Query: que,
}),
},
}
docId := uuid.New()
qIds := []uuid.UUID{
uuid.New(),
uuid.New(),
}
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.QueryURL && *in.MessageBody == fmt.Sprintf("{\"document_id\":\"%s\",\"query_id\":\"%s\"}", docId.String(), qIds[0].String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.QueryURL && *in.MessageBody == fmt.Sprintf("{\"document_id\":\"%s\",\"query_id\":\"%s\"}", docId.String(), qIds[1].String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
err = svc.TriggerQueriesSync(ctx, docId, qIds)
assert.NoError(t, err)
}
func TestInformQueryDependents(t *testing.T) {
ctx := context.Background()
@@ -233,6 +136,7 @@ func TestInformQueryDependents(t *testing.T) {
Result: result.New(cfg, &result.Services{
Query: que,
}),
Sync: resultsync.New(cfg),
},
}
+21
View File
@@ -0,0 +1,21 @@
package resultsync
import (
"queryorchestration/internal/serviceconfig"
queryc "queryorchestration/internal/serviceconfig/queue/query"
)
type ConfigProvider interface {
serviceconfig.ConfigProvider
queryc.ConfigProvider
}
type Service struct {
cfg ConfigProvider
}
func New(cfg ConfigProvider) *Service {
return &Service{
cfg,
}
}
@@ -0,0 +1,21 @@
package resultsync_test
import (
resultset "queryorchestration/internal/query/result/set"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/query"
"testing"
"github.com/stretchr/testify/assert"
)
type QuerySyncConfig struct {
serviceconfig.BaseConfig
query.QueryConfig
}
func TestService(t *testing.T) {
cfg := &QuerySyncConfig{}
svc := resultset.New(cfg, &resultset.Services{})
assert.NotNil(t, svc)
}
+35
View File
@@ -0,0 +1,35 @@
package resultsync
import (
"context"
"queryorchestration/internal/serviceconfig/queue"
"github.com/google/uuid"
)
type Body struct {
DocumentID uuid.UUID `json:"document_id" validate:"required,uuid"`
QueryID uuid.UUID `json:"query_id" validate:"required,uuid"`
}
func (s *Service) TriggerMultiSync(ctx context.Context, documentID uuid.UUID, queryIDs []uuid.UUID) error {
for _, id := range queryIDs {
err := s.TriggerSync(ctx, &Body{
DocumentID: documentID,
QueryID: id,
})
if err != nil {
return err
}
}
return nil
}
func (s *Service) TriggerSync(ctx context.Context, params *Body) error {
return s.cfg.SendToQueue(ctx, &queue.SendParams{
QueueURL: s.cfg.GetQueryURL(),
Body: params,
})
}
+107
View File
@@ -0,0 +1,107 @@
package resultsync
import (
"context"
"fmt"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/serviceconfig"
queryc "queryorchestration/internal/serviceconfig/queue/query"
queuemock "queryorchestration/mocks/queue"
"testing"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/google/uuid"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
)
type ResultSyncConfig struct {
serviceconfig.BaseConfig
queryc.QueryConfig
}
func TestTriggerSync(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &ResultSyncConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.QueryURL = "/i/am/here"
svc := Service{
cfg: cfg,
}
params := &Body{
DocumentID: uuid.New(),
QueryID: uuid.New(),
}
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.QueryURL && *in.MessageBody == fmt.Sprintf("{\"document_id\":\"%s\",\"query_id\":\"%s\"}", params.DocumentID.String(), params.QueryID.String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
err = svc.TriggerSync(ctx, params)
assert.NoError(t, err)
}
func TestTriggerMultiSync(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &ResultSyncConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.QueryURL = "/i/am/here"
svc := Service{
cfg: cfg,
}
docId := uuid.New()
qIds := []uuid.UUID{
uuid.New(),
uuid.New(),
}
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.QueryURL && *in.MessageBody == fmt.Sprintf("{\"document_id\":\"%s\",\"query_id\":\"%s\"}", docId.String(), qIds[0].String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.QueryURL && *in.MessageBody == fmt.Sprintf("{\"document_id\":\"%s\",\"query_id\":\"%s\"}", docId.String(), qIds[1].String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
err = svc.TriggerMultiSync(ctx, docId, qIds)
assert.NoError(t, err)
}
+2 -2
View File
@@ -1,7 +1,7 @@
package querysync
import (
resultset "queryorchestration/internal/query/result/set"
resultsync "queryorchestration/internal/query/result/sync"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/query"
)
@@ -12,7 +12,7 @@ type ConfigProvider interface {
}
type Services struct {
ResultSet *resultset.Service
ResultSync *resultsync.Service
}
type Service struct {
+1 -1
View File
@@ -31,5 +31,5 @@ func (s *Service) Sync(ctx context.Context, id uuid.UUID) error {
ids[i] = database.MustToUUID(dbId)
}
return s.svc.ResultSet.TriggerQueriesSync(ctx, id, ids)
return s.svc.ResultSync.TriggerMultiSync(ctx, id, ids)
}
+2 -6
View File
@@ -5,8 +5,7 @@ import (
"fmt"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/query"
resultset "queryorchestration/internal/query/result/set"
resultsync "queryorchestration/internal/query/result/sync"
"queryorchestration/internal/serviceconfig"
queryc "queryorchestration/internal/serviceconfig/queue/query"
queuemock "queryorchestration/mocks/queue"
@@ -36,13 +35,10 @@ func TestSync(t *testing.T) {
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.QueryURL = "/i/am/here"
que := query.New(cfg)
svc := Service{
cfg: cfg,
svc: &Services{
ResultSet: resultset.New(cfg, &resultset.Services{
Query: que,
}),
ResultSync: resultsync.New(cfg),
},
}
+1 -1
View File
@@ -69,7 +69,7 @@ func TestTest(t *testing.T) {
pgxmock.NewRows([]string{"id", "jobId", "minCleanVersion", "minTextVersion", "activeVersion", "latestVersion", "fields"}).
AddRow(database.MustToDBUUID(coll.ID), database.MustToDBUUID(doc.JobID), coll.MinCleanVersion, coll.MinTextVersion, int32(1), int32(2), []byte("")),
)
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(params.QueryID), params.QueryVersion).WillReturnRows(
pool.ExpectQuery("name: GetQueryWithVersion :one").WithArgs(database.MustToDBUUID(params.QueryID), &params.QueryVersion).WillReturnRows(
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
AddRow(database.MustToDBUUID(params.QueryID), repository.QuerytypeJsonExtractor, int32(1), params.QueryVersion+1, []byte("{\"path\":\"oldkey\"}"), []pgtype.UUID{reqID}),
)
@@ -2,7 +2,6 @@ package contextfull_test
import (
"context"
"fmt"
resultprocessor "queryorchestration/internal/query/result/processor"
contextfull "queryorchestration/internal/query/types/contextFull"
"testing"
@@ -24,14 +23,14 @@ func TestContextFull(t *testing.T) {
values := []resultprocessor.Value{}
value, err := extractor.Process(ctx, query, &values)
value, err := extractor.Process(ctx, query, values)
assert.NoError(t, err)
assert.Equal(t, fmt.Sprintf("{\"id\":\"%s\"}", query.ID), value)
assert.Equal(t, "{\"id\":\"aaaa\",\"two\":\"bbbb\",\"three\":\"ccc\"}", value)
values = []resultprocessor.Value{
contextfull.NewResult("example_result"),
}
_, err = extractor.Process(ctx, query, &values)
_, err = extractor.Process(ctx, query, values)
assert.EqualError(t, err, "no requirements expected")
}
+3 -4
View File
@@ -3,7 +3,6 @@ package contextfull
import (
"context"
"errors"
"fmt"
resultprocessor "queryorchestration/internal/query/result/processor"
)
@@ -14,11 +13,11 @@ func NewExtractor() *Extractor {
return &Extractor{}
}
func (e *Extractor) Process(ctx context.Context, query *resultprocessor.Query, values *[]resultprocessor.Value) (string, error) {
if values != nil && len(*values) > 0 {
func (e *Extractor) Process(ctx context.Context, query *resultprocessor.Query, values []resultprocessor.Value) (string, error) {
if len(values) > 0 {
return "", errors.New("no requirements expected")
}
// TODO
return fmt.Sprintf(`{"id":"%s"}`, query.ID.String()), nil
return `{"id":"aaaa","two":"bbbb","three":"ccc"}`, nil
}
@@ -50,7 +50,7 @@ func TestJSONProcess(t *testing.T) {
AddRow(pgtype.UUID{}, []byte(config)),
)
value, err := extractor.Process(ctx, query, &values)
value, err := extractor.Process(ctx, query, values)
assert.NoError(t, err)
assert.Equal(t, entryValue, value)
@@ -65,7 +65,7 @@ func TestJSONProcess(t *testing.T) {
AddRow(pgtype.UUID{}, []byte(config)),
)
value, err = extractor.Process(ctx, query, &values)
value, err = extractor.Process(ctx, query, values)
assert.NoError(t, err)
assert.Equal(t, entryValue, value)
@@ -80,7 +80,7 @@ func TestJSONProcess(t *testing.T) {
AddRow(pgtype.UUID{}, []byte(config)),
)
value, err = extractor.Process(ctx, query, &values)
value, err = extractor.Process(ctx, query, values)
assert.NoError(t, err)
assert.Equal(t, entryValue, value)
}
@@ -118,7 +118,7 @@ func TestJSONProcessJSON(t *testing.T) {
AddRow(pgtype.UUID{}, []byte(config)),
)
value, err := extractor.Process(ctx, query, &values)
value, err := extractor.Process(ctx, query, values)
assert.EqualError(t, err, "JSON path does not exist: invalid_key")
assert.Empty(t, value)
@@ -130,7 +130,7 @@ func TestJSONProcessJSON(t *testing.T) {
AddRow(pgtype.UUID{}, []byte(config)),
)
value, err = extractor.Process(ctx, query, &values)
value, err = extractor.Process(ctx, query, values)
assert.EqualError(t, err, "unexpected end of JSON input")
assert.Empty(t, value)
@@ -142,7 +142,7 @@ func TestJSONProcessJSON(t *testing.T) {
AddRow(pgtype.UUID{}, []byte(config)),
)
value, err = extractor.Process(ctx, query, &values)
value, err = extractor.Process(ctx, query, values)
assert.EqualError(t, err, "JSON path does not exist: ")
assert.Empty(t, value)
@@ -154,7 +154,7 @@ func TestJSONProcessJSON(t *testing.T) {
AddRow(pgtype.UUID{}, []byte(config)),
)
value, err = extractor.Process(ctx, query, &values)
value, err = extractor.Process(ctx, query, values)
assert.EqualError(t, err, "invalid character '}' looking for beginning of value")
assert.Empty(t, value)
@@ -163,7 +163,7 @@ func TestJSONProcessJSON(t *testing.T) {
pgxmock.NewRows([]string{"id", "config"}),
)
value, err = extractor.Process(ctx, query, &values)
value, err = extractor.Process(ctx, query, values)
assert.EqualError(t, err, "no rows in result set")
assert.Empty(t, value)
}
@@ -188,7 +188,7 @@ func TestJSONProcessResults(t *testing.T) {
}
results := []resultprocessor.Value{}
value, err := extractor.Process(ctx, query, &results)
value, err := extractor.Process(ctx, query, results)
assert.EqualError(t, err, "JSON Extraction requires 1 result")
assert.Empty(t, value)
@@ -196,7 +196,7 @@ func TestJSONProcessResults(t *testing.T) {
contextfull.NewResult(""),
contextfull.NewResult(""),
}
value, err = extractor.Process(ctx, query, &results)
value, err = extractor.Process(ctx, query, results)
assert.EqualError(t, err, "JSON Extraction requires 1 result")
assert.Empty(t, value)
@@ -205,7 +205,7 @@ func TestJSONProcessResults(t *testing.T) {
contextfull.NewResult(""),
contextfull.NewResult(""),
}
value, err = extractor.Process(ctx, query, &results)
value, err = extractor.Process(ctx, query, results)
assert.EqualError(t, err, "JSON Extraction requires 1 result")
assert.Empty(t, value)
}
@@ -24,12 +24,12 @@ func NewExtractor(cfg serviceconfig.ConfigProvider) *Extractor {
return &Extractor{cfg}
}
func (e *Extractor) Process(ctx context.Context, query *resultprocessor.Query, values *[]resultprocessor.Value) (string, error) {
if values == nil || len(*values) != 1 {
func (e *Extractor) Process(ctx context.Context, query *resultprocessor.Query, values []resultprocessor.Value) (string, error) {
if len(values) != 1 {
return "", fmt.Errorf("JSON Extraction requires 1 result")
}
value, err := (*values)[0].GetValue(ctx)
value, err := values[0].GetValue(ctx)
if err != nil {
return "", err
}
+28
View File
@@ -0,0 +1,28 @@
package queryupdate
import (
"queryorchestration/internal/query"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/queryversionsync"
)
type ConfigProvider interface {
serviceconfig.ConfigProvider
queryversionsync.ConfigProvider
}
type Services struct {
Query *query.Service
}
type Service struct {
cfg ConfigProvider
svc *Services
}
func New(cfg ConfigProvider, svc *Services) *Service {
return &Service{
cfg,
svc,
}
}
+23
View File
@@ -0,0 +1,23 @@
package queryupdate_test
import (
"queryorchestration/internal/database/repository"
"queryorchestration/internal/query"
"queryorchestration/internal/serviceconfig"
"testing"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
)
func TestService(t *testing.T) {
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := query.New(cfg)
assert.NotNil(t, svc)
}
@@ -1,22 +1,25 @@
package query
package queryupdate
import (
"context"
"errors"
"fmt"
"log/slog"
queryversionsyncrunner "queryorchestration/api/queryVersionSyncRunner"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/query"
resultprocessor "queryorchestration/internal/query/result/processor"
contextfull "queryorchestration/internal/query/types/contextFull"
jsonextractor "queryorchestration/internal/query/types/jsonExtractor"
"queryorchestration/internal/serviceconfig/queue"
"queryorchestration/internal/validation"
"github.com/google/uuid"
)
func (s *Service) Update(ctx context.Context, entity *resultprocessor.Update) error {
current, err := s.Get(ctx, entity.ID)
current, err := s.svc.Query.Get(ctx, entity.ID)
if err != nil {
return err
}
@@ -31,11 +34,29 @@ func (s *Service) Update(ctx context.Context, entity *resultprocessor.Update) er
return err
}
err = s.informUpdate(ctx, entity)
if err != nil {
return err
}
return nil
}
func (s *Service) normalizeUpdateRequiredQueryIDs(ctx context.Context, current *Query, entity RequiredQueryIDs) error {
err := s.NormalizeQueryIDs(ctx, entity)
func (s *Service) informUpdate(ctx context.Context, entity *resultprocessor.Update) error {
if entity.ActiveVersion == nil {
return nil
}
return s.cfg.SendToQueue(ctx, &queue.SendParams{
QueueURL: s.cfg.GetQueryVersionSyncURL(),
Body: queryversionsyncrunner.Body{
ID: entity.ID,
},
})
}
func (s *Service) normalizeUpdateRequiredQueryIDs(ctx context.Context, current *query.Query, entity query.RequiredQueryIDs) error {
err := s.svc.Query.NormalizeQueryIDs(ctx, entity)
if err != nil {
return err
}
@@ -64,8 +85,8 @@ func (s *Service) normalizeUpdateRequiredQueryIDs(ctx context.Context, current *
return nil
}
func (s *Service) normalizeUpdate(ctx context.Context, current *Query, entity *resultprocessor.Update) error {
err := s.normalizeActiveVersion(current, entity)
func (s *Service) normalizeUpdate(ctx context.Context, current *query.Query, entity *resultprocessor.Update) error {
err := s.svc.Query.NormalizeActiveVersion(current, entity)
if err != nil {
return err
}
@@ -75,7 +96,7 @@ func (s *Service) normalizeUpdate(ctx context.Context, current *Query, entity *r
return err
}
err = s.normalizeConfig(entity)
err = s.svc.Query.NormalizeConfig(entity)
if err != nil {
return err
}
@@ -95,7 +116,7 @@ func (s *Service) normalizeUpdate(ctx context.Context, current *Query, entity *r
return err
}
err = validator.Validate(ctx, ParseQuery(current), entity)
err = validator.Validate(ctx, query.ParseQuery(current), entity)
if err != nil {
return err
}
@@ -103,32 +124,34 @@ func (s *Service) normalizeUpdate(ctx context.Context, current *Query, entity *r
return nil
}
func (s *Service) submitUpdate(ctx context.Context, current *Query, entity *resultprocessor.Update) error {
func (s *Service) submitUpdate(ctx context.Context, current *query.Query, entity *resultprocessor.Update) error {
err := s.cfg.ExecuteDBTransaction(ctx, func(ctx context.Context, q *repository.Queries) error {
latestVersion := current.LatestVersion + 1
id := database.MustToDBUUID(entity.ID)
addIDs := getSetDifference(entity.RequiredQueryIDs, current.RequiredQueryIDs)
for _, qID := range addIDs {
err := q.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{
Queryid: id,
Requiredqueryid: database.MustToDBUUID(qID),
Addedversion: latestVersion,
})
if err != nil {
return err
if entity.RequiredQueryIDs != nil {
addIDs := getSetDifference(entity.RequiredQueryIDs, current.RequiredQueryIDs)
for _, qID := range addIDs {
err := q.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{
Queryid: id,
Requiredqueryid: database.MustToDBUUID(qID),
Addedversion: latestVersion,
})
if err != nil {
return err
}
}
}
removeIDs := getSetDifference(current.RequiredQueryIDs, entity.RequiredQueryIDs)
for _, qID := range removeIDs {
err := q.RemoveRequiredQuery(ctx, &repository.RemoveRequiredQueryParams{
Queryid: id,
Requiredqueryid: database.MustToDBUUID(qID),
Removedversion: &latestVersion,
})
if err != nil {
return err
removeIDs := getSetDifference(current.RequiredQueryIDs, entity.RequiredQueryIDs)
for _, qID := range removeIDs {
err := q.RemoveRequiredQuery(ctx, &repository.RemoveRequiredQueryParams{
Queryid: id,
Requiredqueryid: database.MustToDBUUID(qID),
Removedversion: &latestVersion,
})
if err != nil {
return err
}
}
}
@@ -156,7 +179,17 @@ func (s *Service) submitUpdate(ctx context.Context, current *Query, entity *resu
activeVersion = &current.ActiveVersion
}
err := q.UpdateQuery(ctx, &repository.UpdateQueryParams{
activeName, err := validation.GetFieldName(entity, entity.ActiveVersion)
if err != nil {
return err
}
onlyactive := validation.AreAllPointersNilExcept(entity, activeName)
if onlyactive {
latestVersion--
}
err = q.UpdateQuery(ctx, &repository.UpdateQueryParams{
Latestversion: latestVersion,
Activeversion: *activeVersion,
ID: id,
@@ -210,13 +243,3 @@ func (s *Service) getUpdator(qType resultprocessor.Type) (resultprocessor.Updato
return nil, fmt.Errorf("attempting to process invalid query type")
}
}
func ParseQuery(q *Query) *resultprocessor.Query {
return &resultprocessor.Query{
ID: q.ID,
Type: q.Type,
Version: q.ActiveVersion,
RequiredQueryIDs: q.RequiredQueryIDs,
Config: q.Config,
}
}
+86
View File
@@ -0,0 +1,86 @@
package queryupdate_test
import (
"context"
"fmt"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/query"
resultprocessor "queryorchestration/internal/query/result/processor"
queryupdate "queryorchestration/internal/query/update"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/queryversionsync"
queuemock "queryorchestration/mocks/queue"
"testing"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgtype"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
)
type QueryUpdateConfig struct {
serviceconfig.BaseConfig
queryversionsync.QueryVersionSyncConfig
}
func TestUpdate(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.QueryVersionSyncURL = "/here"
svc := queryupdate.New(cfg, &queryupdate.Services{
Query: query.New(cfg),
})
config := "{\"path\":\"example_path\"}"
existing := query.Query{
ID: uuid.New(),
ActiveVersion: int32(1),
LatestVersion: int32(1),
}
av := int32(2)
jcfg := `{"path":"id"}`
update := &resultprocessor.Update{
ID: existing.ID,
ActiveVersion: &av,
Config: &jcfg,
}
pool.ExpectQuery("name: GetQuery :one").WithArgs(database.MustToDBUUID(update.ID)).WillReturnRows(
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
AddRow(database.MustToDBUUID(existing.ID), repository.QuerytypeJsonExtractor, existing.ActiveVersion, existing.LatestVersion, []byte(config), []pgtype.UUID{}),
)
pool.ExpectBeginTx(pgx.TxOptions{})
pool.ExpectExec("name: RemoveQueryConfig :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID(update.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: AddQueryConfig :exec").WithArgs(database.MustToDBUUID(update.ID), []byte(*update.Config), int32(2)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(av, int32(2), database.MustToDBUUID(update.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectCommit()
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.QueryVersionSyncURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", existing.ID.String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
err = svc.Update(ctx, update)
assert.NoError(t, err)
}
@@ -1,29 +1,42 @@
package query
package queryupdate
import (
"context"
"errors"
"fmt"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/query"
resultprocessor "queryorchestration/internal/query/result/processor"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/queryversionsync"
queuemock "queryorchestration/mocks/queue"
"testing"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
)
type QueryUpdateConfig struct {
serviceconfig.BaseConfig
queryversionsync.QueryVersionSyncConfig
}
func TestGetUpdator(t *testing.T) {
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
queryType := resultprocessor.Type(resultprocessor.TypeContextFull)
updator, err := svc.getUpdator(queryType)
@@ -47,13 +60,15 @@ func TestSubmitUpdate(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
config := "{\"path\":\"example_path\"}"
q := Query{
q := query.Query{
ID: uuid.New(),
RequiredQueryIDs: &[]uuid.UUID{
uuid.New(),
@@ -96,12 +111,14 @@ func TestSubmitUpdateRollback(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
q := Query{
q := query.Query{
ID: uuid.New(),
ActiveVersion: int32(1),
LatestVersion: int32(2),
@@ -114,7 +131,7 @@ func TestSubmitUpdateRollback(t *testing.T) {
pool.ExpectBeginTx(pgx.TxOptions{})
msg := "database failure"
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(aV, int32(3), database.MustToDBUUID(update.ID)).
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(aV, int32(2), database.MustToDBUUID(update.ID)).
WillReturnError(errors.New(msg))
pool.ExpectCommit()
@@ -123,53 +140,99 @@ func TestSubmitUpdateRollback(t *testing.T) {
}
func TestSubmitUpdateRequiredQueries(t *testing.T) {
ctx := context.Background()
t.Run("all fields", func(t *testing.T) {
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
ctx := context.Background()
config := "{\"path\":\"example_path\"}"
q := Query{
ID: uuid.New(),
RequiredQueryIDs: &[]uuid.UUID{
uuid.New(),
uuid.New(),
uuid.New(),
},
Config: &config,
ActiveVersion: int32(1),
LatestVersion: int32(2),
}
update := &resultprocessor.Update{
ID: q.ID,
RequiredQueryIDs: &[]uuid.UUID{
uuid.New(),
uuid.New(),
(*q.RequiredQueryIDs)[0],
},
}
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
pool.ExpectBeginTx(pgx.TxOptions{})
pool.ExpectExec("name: AddRequiredQuery :exec").WithArgs(database.MustToDBUUID(update.ID), database.MustToDBUUID((*update.RequiredQueryIDs)[0]), pgxmock.AnyArg()).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: AddRequiredQuery :exec").WithArgs(database.MustToDBUUID(update.ID), database.MustToDBUUID((*update.RequiredQueryIDs)[1]), pgxmock.AnyArg()).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: RemoveRequiredQuery :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID((*q.RequiredQueryIDs)[1]), database.MustToDBUUID(update.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: RemoveRequiredQuery :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID((*q.RequiredQueryIDs)[2]), database.MustToDBUUID(update.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(q.ActiveVersion, int32(3), database.MustToDBUUID(update.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectCommit()
config := "{\"path\":\"example_path\"}"
q := query.Query{
ID: uuid.New(),
RequiredQueryIDs: &[]uuid.UUID{
uuid.New(),
uuid.New(),
uuid.New(),
},
Config: &config,
ActiveVersion: int32(1),
LatestVersion: int32(2),
}
update := &resultprocessor.Update{
ID: q.ID,
RequiredQueryIDs: &[]uuid.UUID{
uuid.New(),
uuid.New(),
(*q.RequiredQueryIDs)[0],
},
}
err = svc.submitUpdate(ctx, &q, update)
assert.NoError(t, err)
pool.ExpectBeginTx(pgx.TxOptions{})
pool.ExpectExec("name: AddRequiredQuery :exec").WithArgs(database.MustToDBUUID(update.ID), database.MustToDBUUID((*update.RequiredQueryIDs)[0]), pgxmock.AnyArg()).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: AddRequiredQuery :exec").WithArgs(database.MustToDBUUID(update.ID), database.MustToDBUUID((*update.RequiredQueryIDs)[1]), pgxmock.AnyArg()).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: RemoveRequiredQuery :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID((*q.RequiredQueryIDs)[1]), database.MustToDBUUID(update.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: RemoveRequiredQuery :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID((*q.RequiredQueryIDs)[2]), database.MustToDBUUID(update.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(q.ActiveVersion, int32(3), database.MustToDBUUID(update.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectCommit()
err = svc.submitUpdate(ctx, &q, update)
assert.NoError(t, err)
})
t.Run("all nil", func(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
config := "{\"path\":\"example_path\"}"
q := query.Query{
ID: uuid.New(),
RequiredQueryIDs: &[]uuid.UUID{
uuid.New(),
uuid.New(),
uuid.New(),
},
Config: &config,
ActiveVersion: 1,
LatestVersion: 2,
}
av := int32(2)
update := &resultprocessor.Update{
ID: q.ID,
ActiveVersion: &av,
}
pool.ExpectBeginTx(pgx.TxOptions{})
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(av, int32(2), database.MustToDBUUID(update.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectCommit()
err = svc.submitUpdate(ctx, &q, update)
assert.NoError(t, err)
})
}
func TestSubmitUpdateActiveVersion(t *testing.T) {
@@ -179,12 +242,14 @@ func TestSubmitUpdateActiveVersion(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
q := Query{
q := query.Query{
ID: uuid.New(),
ActiveVersion: int32(1),
LatestVersion: int32(2),
@@ -196,7 +261,7 @@ func TestSubmitUpdateActiveVersion(t *testing.T) {
}
pool.ExpectBeginTx(pgx.TxOptions{})
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(aV, int32(3), database.MustToDBUUID(update.ID)).
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(aV, int32(2), database.MustToDBUUID(update.ID)).
WillReturnResult(pgxmock.NewResult("", 1))
pool.ExpectCommit()
@@ -225,12 +290,14 @@ func TestNormalizeUpdate(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
current := &Query{
current := &query.Query{
ID: uuid.New(),
ActiveVersion: int32(1),
LatestVersion: int32(2),
@@ -262,12 +329,14 @@ func TestNormalizeUpdate(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
current := &Query{
current := &query.Query{
ID: uuid.New(),
ActiveVersion: int32(1),
LatestVersion: int32(2),
@@ -287,13 +356,15 @@ func TestNormalizeUpdate(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
version := int32(1)
current := &Query{
current := &query.Query{
ID: uuid.New(),
ActiveVersion: version,
LatestVersion: int32(2),
@@ -314,12 +385,14 @@ func TestNormalizeUpdate(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
current := &Query{
current := &query.Query{
ID: uuid.New(),
ActiveVersion: 1,
LatestVersion: 2,
@@ -343,12 +416,14 @@ func TestNormalizeUpdateRequiredQueryIDs(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
current := &Query{
current := &query.Query{
ID: uuid.New(),
Type: resultprocessor.TypeJsonExtractor,
}
@@ -371,12 +446,14 @@ func TestNormalizeUpdateRequiredQueryIDs(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
current := &Query{
current := &query.Query{
ID: uuid.New(),
Type: resultprocessor.TypeJsonExtractor,
}
@@ -401,12 +478,14 @@ func TestNormalizeUpdateRequiredQueryIDs(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
current := &Query{
current := &query.Query{
ID: uuid.New(),
Type: resultprocessor.TypeJsonExtractor,
}
@@ -434,12 +513,14 @@ func TestNormalizeUpdateRequiredQueryIDs(t *testing.T) {
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg := &QueryUpdateConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Query: query.New(cfg),
})
current := &Query{
current := &query.Query{
ID: uuid.New(),
Type: resultprocessor.TypeJsonExtractor,
}
@@ -467,3 +548,55 @@ func TestNormalizeUpdateRequiredQueryIDs(t *testing.T) {
}, *update)
})
}
func TestInformUpdate(t *testing.T) {
t.Run("nil active version", func(t *testing.T) {
ctx := context.Background()
cfg := &QueryUpdateConfig{}
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.QueryVersionSyncURL = "/here"
svc := Service{
cfg: cfg,
}
entity := &resultprocessor.Update{
ID: uuid.New(),
ActiveVersion: nil,
}
err := svc.informUpdate(ctx, entity)
assert.NoError(t, err)
})
t.Run("update active version", func(t *testing.T) {
ctx := context.Background()
cfg := &QueryUpdateConfig{}
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.QueryVersionSyncURL = "/here"
svc := Service{
cfg: cfg,
}
av := int32(2)
entity := &resultprocessor.Update{
ID: uuid.New(),
ActiveVersion: &av,
}
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.QueryVersionSyncURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", entity.ID.String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
err := svc.informUpdate(ctx, entity)
assert.NoError(t, err)
})
}
-47
View File
@@ -1,47 +0,0 @@
package query_test
import (
"context"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/query"
resultprocessor "queryorchestration/internal/query/result/processor"
"queryorchestration/internal/serviceconfig"
"testing"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgtype"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
)
func TestUpdate(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &serviceconfig.BaseConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := query.New(cfg)
config := "{\"path\":\"example_path\"}"
existing := query.Query{
ID: uuid.New(),
ActiveVersion: int32(1),
LatestVersion: int32(1),
}
update := &resultprocessor.Update{
ID: existing.ID,
}
pool.ExpectQuery("name: GetQuery :one").WithArgs(database.MustToDBUUID(update.ID)).WillReturnRows(
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
AddRow(database.MustToDBUUID(existing.ID), repository.QuerytypeJsonExtractor, existing.ActiveVersion, existing.LatestVersion, []byte(config), []pgtype.UUID{}),
)
err = svc.Update(ctx, update)
assert.Error(t, err)
}
+21
View File
@@ -0,0 +1,21 @@
package queryversionsync
import (
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/jobsync"
)
type ConfigProvider interface {
serviceconfig.ConfigProvider
jobsync.ConfigProvider
}
type Service struct {
cfg ConfigProvider
}
func New(cfg ConfigProvider) *Service {
return &Service{
cfg,
}
}
@@ -0,0 +1,15 @@
package queryversionsync_test
import (
"queryorchestration/internal/query"
"queryorchestration/internal/serviceconfig"
"testing"
"github.com/stretchr/testify/assert"
)
func TestService(t *testing.T) {
cfg := &serviceconfig.BaseConfig{}
svc := query.New(cfg)
assert.NotNil(t, svc)
}
+35
View File
@@ -0,0 +1,35 @@
package queryversionsync
import (
"context"
"database/sql"
"errors"
jobsyncrunner "queryorchestration/api/jobSyncRunner"
"queryorchestration/internal/database"
"queryorchestration/internal/serviceconfig/queue"
"github.com/google/uuid"
)
func (s *Service) Sync(ctx context.Context, id uuid.UUID) error {
jobIds, err := s.cfg.GetDBQueries().ListQueryJobIDs(ctx, database.MustToDBUUID(id))
if err != nil && errors.Is(err, sql.ErrNoRows) {
return nil
} else if err != nil {
return err
}
for _, id := range jobIds {
err := s.cfg.SendToQueue(ctx, &queue.SendParams{
QueueURL: s.cfg.GetJobSyncURL(),
Body: jobsyncrunner.Body{
ID: database.MustToUUID(id),
},
})
if err != nil {
return err
}
}
return nil
}
+73
View File
@@ -0,0 +1,73 @@
package queryversionsync
import (
"context"
"fmt"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/jobsync"
queuemock "queryorchestration/mocks/queue"
"testing"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/google/uuid"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
)
type QueryVersionSyncConfig struct {
serviceconfig.BaseConfig
jobsync.JobSyncConfig
}
func TestSync(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
cfg := &QueryVersionSyncConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
cfg.JobSyncURL = "/i/am/here"
svc := Service{cfg}
queryId := uuid.New()
jobIds := []uuid.UUID{
uuid.New(),
uuid.New(),
}
pool.ExpectQuery("name: ListQueryJobIDs :many").WithArgs(database.MustToDBUUID(queryId)).
WillReturnRows(
pgxmock.NewRows([]string{"jobId"}).
AddRow(database.MustToDBUUID(jobIds[0])).
AddRow(database.MustToDBUUID(jobIds[1])),
)
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.JobSyncURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", jobIds[0].String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.JobSyncURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", jobIds[1].String())
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
assert.Nil(t, svc.Sync(ctx, queryId))
}
@@ -0,0 +1,13 @@
package queryversionsync
type QueryVersionSyncConfig struct {
QueryVersionSyncURL string `env:"QUERY_VERSION_SYNC_URL,required,notEmpty"`
}
func (c *QueryVersionSyncConfig) GetQueryVersionSyncURL() string {
return c.QueryVersionSyncURL
}
type ConfigProvider interface {
GetQueryVersionSyncURL() string
}
@@ -0,0 +1,20 @@
package queryversionsync_test
import (
"queryorchestration/internal/serviceconfig/queue/queryversionsync"
"testing"
"github.com/stretchr/testify/assert"
)
func TestGetQueryVersionSyncURL(t *testing.T) {
cfg := queryversionsync.QueryVersionSyncConfig{}
name := cfg.GetQueryVersionSyncURL()
assert.Equal(t, "", name)
cfg.QueryVersionSyncURL = "name"
name = cfg.GetQueryVersionSyncURL()
assert.Equal(t, "name", name)
assert.Equal(t, cfg.QueryVersionSyncURL, name)
}
+13 -10
View File
@@ -25,7 +25,7 @@ type CreateAWSConfig struct {
Cfg serviceconfig.ConfigProvider
}
func CreateAWSContainer(t *testing.T, ctx context.Context, cfg *CreateAWSConfig) (*AWSContainerConfig, func()) {
func CreateAWSContainer(t testing.TB, ctx context.Context, cfg *CreateAWSConfig) (*AWSContainerConfig, func()) {
alias := "localstack"
port, err := nat.NewPort("tcp", "4566")
@@ -36,15 +36,18 @@ func CreateAWSContainer(t *testing.T, ctx context.Context, cfg *CreateAWSConfig)
req := testcontainers.ContainerRequest{
Image: "localstack/localstack:4.1.0",
Env: map[string]string{
"AWS_ACCESS_KEY_ID": cfg.Cfg.GetAWSKeyID(),
"AWS_SECRET_ACCESS_KEY": cfg.Cfg.GetAWSSecretKey(),
"AWS_SESSION_TOKEN": cfg.Cfg.GetAWSSessionToken(),
"AWS_REGION": cfg.Cfg.GetAWSRegion(),
"SERVICES": "s3,sqs",
"SKIP_SSL_CERT_DOWNLOAD": "1",
"LOCALSTACK_HOST": alias,
"SQS_ENDPOINT_STRATEGY": "path",
"EAGER_SERVICE_LOADING": "1",
"AWS_ACCESS_KEY_ID": cfg.Cfg.GetAWSKeyID(),
"AWS_SECRET_ACCESS_KEY": cfg.Cfg.GetAWSSecretKey(),
"AWS_SESSION_TOKEN": cfg.Cfg.GetAWSSessionToken(),
"AWS_REGION": cfg.Cfg.GetAWSRegion(),
"SERVICES": "s3,sqs,cloudwatch,logs",
"SKIP_SSL_CERT_DOWNLOAD": "1",
"LOCALSTACK_HOST": alias,
"SQS_ENDPOINT_STRATEGY": "path",
"EAGER_SERVICE_LOADING": "1",
"DEBUG": "1",
"LS_LOG": "trace",
"SQS_CLOUDWATCH_METRICS_REPORT_INTERVAL": "5",
},
ExposedPorts: []string{port.Port()},
WaitingFor: wait.ForAll(
+1 -6
View File
@@ -28,7 +28,7 @@ type containerConfig struct {
ExposedPorts []nat.Port
}
func createContainer(t *testing.T, ctx context.Context, cfg *containerConfig) (testcontainers.Container, func()) {
func createContainer(t testing.TB, ctx context.Context, cfg *containerConfig) (testcontainers.Container, func()) {
env := map[string]string{
"DB_USER": cfg.Cfg.GetDBUser(),
"DB_PASS": cfg.Cfg.GetDBSecret(),
@@ -65,11 +65,6 @@ func createContainer(t *testing.T, ctx context.Context, cfg *containerConfig) (t
ports[index] = port.Port()
}
req.ExposedPorts = ports
// req.WaitingFor = wait.ForAll(
// // wait.ForExposedPort(),
// // wait.ForListeningPort(cfg.ExposedPorts[0]),
// wait.ForLog(cfg.WaitForMsg),
// )
}
container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{
+1 -1
View File
@@ -20,7 +20,7 @@ type CreateDatabaseConfig struct {
RunMigrations bool
}
func CreateDB(t *testing.T, ctx context.Context, cfg *CreateDatabaseConfig) (testcontainers.Container, func()) {
func CreateDB(t testing.TB, ctx context.Context, cfg *CreateDatabaseConfig) (testcontainers.Container, func()) {
alias := "postgres"
port, err := nat.NewPort("tcp", "5432")
+4 -4
View File
@@ -24,7 +24,7 @@ type EcosystemNetworkConfig struct {
AWSContainer *AWSContainerConfig
}
func CreateRunnersAndServicesNetwork(t *testing.T, ctx context.Context, ncfg *EcosystemNetworkConfig) (*EcosystemConfig, func()) {
func CreateRunnersAndServicesNetwork(t testing.TB, ctx context.Context, ncfg *EcosystemNetworkConfig) (*EcosystemConfig, func()) {
network := ncfg.Network
var ncleanup func()
if network == nil {
@@ -106,7 +106,7 @@ func CreateRunnersAndServicesNetwork(t *testing.T, ctx context.Context, ncfg *Ec
}
}
func SetCfgProvider(t *testing.T, cfg serviceconfig.ConfigProvider) {
func SetCfgProvider(t testing.TB, cfg serviceconfig.ConfigProvider) {
cfg.SetAWSConfig(&aws.AWSConfig{
AWSKeyID: "test",
AWSSecretKey: "test",
@@ -135,7 +135,7 @@ type ServiceNetworkConfig struct {
Env map[string]string
}
func CreateServiceNetwork(t *testing.T, ctx context.Context, scfg *ServiceNetworkConfig) (*Container, func()) {
func CreateServiceNetwork(t testing.TB, ctx context.Context, scfg *ServiceNetworkConfig) (*Container, func()) {
if scfg.Cfg == nil {
scfg.Cfg = &serviceconfig.BaseConfig{}
SetCfgProvider(t, scfg.Cfg)
@@ -177,7 +177,7 @@ type RunnerNetworkConfig struct {
Env map[string]string
}
func CreateRunnerNetwork(t *testing.T, ctx context.Context, rcfg *RunnerNetworkConfig) (*Container, func()) {
func CreateRunnerNetwork(t testing.TB, ctx context.Context, rcfg *RunnerNetworkConfig) (*Container, func()) {
network := rcfg.Network
var ncleanup func()
if network == nil {
+6 -3
View File
@@ -22,7 +22,8 @@ func TestCreateRunnerNetwork(t *testing.T) {
Cfg: cfg,
Name: QueryRunner,
Env: map[string]string{
"QUERY_URL": "/i/am/here",
"QUERY_URL": "/i/am/here",
"QUERY_VERSION_SYNC_URL": "/here/there/every/where",
},
})
@@ -41,7 +42,8 @@ func TestCreateServiceNetwork(t *testing.T) {
conn, cleanup := CreateServiceNetwork(t, ctx, &ServiceNetworkConfig{
Name: QueryService,
Env: map[string]string{
"JOB_SYNC_URL": "/i/am/here",
"JOB_SYNC_URL": "/i/am/here",
"QUERY_VERSION_SYNC_URL": "/here/there/every/where",
},
})
@@ -74,7 +76,8 @@ func TestCreateRunnersAndServicesNetwork(t *testing.T) {
{
Name: QueryService,
Env: map[string]string{
"JOB_SYNC_URL": "/i/am/here",
"JOB_SYNC_URL": "/i/am/here",
"QUERY_VERSION_SYNC_URL": "/here/there/every/where",
},
},
},
+1 -1
View File
@@ -8,7 +8,7 @@ import (
"github.com/testcontainers/testcontainers-go/network"
)
func CreateNetwork(t *testing.T, ctx context.Context) (*testcontainers.DockerNetwork, func()) {
func CreateNetwork(t testing.TB, ctx context.Context) (*testcontainers.DockerNetwork, func()) {
network, err := network.New(ctx, network.WithDriver("bridge"))
if err != nil {
t.Fatal(err)
+2 -2
View File
@@ -9,7 +9,7 @@ import (
"github.com/aws/aws-sdk-go-v2/service/s3"
)
func CreateBucket(t *testing.T, ctx context.Context, cfg objectstore.ConfigProvider, name string) {
func CreateBucket(t testing.TB, ctx context.Context, cfg objectstore.ConfigProvider, name string) {
_, err := cfg.GetStoreClient().CreateBucket(ctx, &s3.CreateBucketInput{
Bucket: aws.String(name),
})
@@ -23,7 +23,7 @@ func CreateBucket(t *testing.T, ctx context.Context, cfg objectstore.ConfigProvi
}
}
func CreateStoreClient(t *testing.T, ctx context.Context, cfg objectstore.ConfigProvider, endpoint string) {
func CreateStoreClient(t testing.TB, ctx context.Context, cfg objectstore.ConfigProvider, endpoint string) {
cfg.SetS3Endpoint(endpoint)
cfg.SetS3UsePathStyle(true)
err := cfg.SetStoreClient(ctx)
+4 -4
View File
@@ -13,7 +13,7 @@ import (
"github.com/stretchr/testify/assert"
)
func CreateQueue(t *testing.T, ctx context.Context, cfg serviceconfig.ConfigProvider, name string) string {
func CreateQueue(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, name string) string {
queueM, err := cfg.GetQueueClient().CreateQueue(ctx, &sqs.CreateQueueInput{
QueueName: aws.String(name),
})
@@ -29,7 +29,7 @@ func CreateQueue(t *testing.T, ctx context.Context, cfg serviceconfig.ConfigProv
return *queueM.QueueUrl
}
func AssertMessage(t *testing.T, ctx context.Context, cfg serviceconfig.ConfigProvider, params *queue.ReceiveParams) types.Message {
func AssertMessage(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, params *queue.ReceiveParams) types.Message {
result, err := cfg.ReceiveFromQueue(ctx, params)
assert.NoError(t, err)
assert.NotNil(t, result.Messages)
@@ -39,7 +39,7 @@ func AssertMessage(t *testing.T, ctx context.Context, cfg serviceconfig.ConfigPr
return result.Messages[0]
}
func AssertMessageBody(t *testing.T, ctx context.Context, cfg serviceconfig.ConfigProvider, url string, body *regexp.Regexp) {
func AssertMessageBody(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, url string, body *regexp.Regexp) {
message := AssertMessage(t, ctx, cfg, &queue.ReceiveParams{
QueueURL: url,
})
@@ -47,7 +47,7 @@ func AssertMessageBody(t *testing.T, ctx context.Context, cfg serviceconfig.Conf
assert.Regexp(t, body, *message.Body)
}
func AssertMessageAttr(t *testing.T, ctx context.Context, cfg serviceconfig.ConfigProvider, url string, name string, value *regexp.Regexp) {
func AssertMessageAttr(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, url string, name string, value *regexp.Regexp) {
message := AssertMessage(t, ctx, cfg, &queue.ReceiveParams{
QueueURL: url,
Attributes: []string{name},
+10 -8
View File
@@ -9,6 +9,7 @@ import (
jobsyncrunner "queryorchestration/api/jobSyncRunner"
queryrunner "queryorchestration/api/queryRunner"
querysyncrunner "queryorchestration/api/querySyncRunner"
queryversionsyncrunner "queryorchestration/api/queryVersionSyncRunner"
"testing"
"queryorchestration/internal/serviceconfig"
@@ -19,13 +20,14 @@ import (
type Runner = string
const (
DocInitRunner = docinitrunner.Name
DocSyncRunner = docsyncrunner.Name
DocCleanRunner = doccleanrunner.Name
DocTextRunner = doctextrunner.Name
QuerySyncRunner = querysyncrunner.Name
QueryRunner = queryrunner.Name
JobSyncRunner = jobsyncrunner.Name
DocInitRunner = docinitrunner.Name
DocSyncRunner = docsyncrunner.Name
DocCleanRunner = doccleanrunner.Name
DocTextRunner = doctextrunner.Name
QuerySyncRunner = querysyncrunner.Name
QueryRunner = queryrunner.Name
JobSyncRunner = jobsyncrunner.Name
QueryVersionSyncRunner = queryversionsyncrunner.Name
)
type RunnerConfig struct {
@@ -36,7 +38,7 @@ type RunnerConfig struct {
Env map[string]string
}
func CreateRunner(t *testing.T, ctx context.Context, config *RunnerConfig) (*Container, func()) {
func CreateRunner(t testing.TB, ctx context.Context, config *RunnerConfig) (*Container, func()) {
if config.Env == nil {
config.Env = map[string]string{}
}
+6 -7
View File
@@ -24,19 +24,18 @@ type ServiceConfig struct {
Env map[string]string
}
func CreateService(t *testing.T, ctx context.Context, config *ServiceConfig) (*Container, func()) {
func CreateService(t testing.TB, ctx context.Context, config *ServiceConfig) (*Container, func()) {
port, err := nat.NewPort("tcp", "8080")
if err != nil {
t.Fatalf("Failed to create port: %v", err)
}
container, cleanup := createContainer(t, ctx, &containerConfig{
Cfg: config.Cfg,
Name: config.Name,
Network: config.Network,
Env: config.Env,
ExposedPorts: []nat.Port{port},
WaitForMsg: "⇨ http server started on [::]:8080",
Cfg: config.Cfg,
Name: config.Name,
Network: config.Network,
Env: config.Env,
WaitForMsg: "⇨ http server started on [::]:8080",
})
host, err := container.Host(ctx)
+2 -1
View File
@@ -34,7 +34,8 @@ func TestCreateService(t *testing.T) {
Cfg: cfg,
Network: ncfg,
Env: map[string]string{
"JOB_SYNC_URL": "/i/am/here",
"JOB_SYNC_URL": "/i/am/here",
"QUERY_VERSION_SYNC_URL": "/here/there/every/where",
},
}
+4
View File
@@ -63,6 +63,7 @@ tasks:
- aws sqs create-queue --queue-name $QNAME_QUERY_SYNC
- aws sqs create-queue --queue-name $QNAME_QUERY_RUNNER
- aws sqs create-queue --queue-name $QNAME_JOB_SYNC
- aws sqs create-queue --queue-name $QNAME_QUERY_VERSION_SYNC
- |
aws s3api put-bucket-notification-configuration \
--bucket $BUCKET_IN \
@@ -105,3 +106,6 @@ tasks:
- |
docker compose -f {{.TEST_COMPOSE_FILE}} config -q \
--resolve-image-digests
- |
docker compose -f {{.GENERATE_COMPOSE_FILE}} config -q \
--resolve-image-digests
+10 -1
View File
@@ -52,6 +52,7 @@ func TestProcess(t *testing.T) {
querysyncurl := test.CreateQueue(t, ctx, cfg, test.QuerySyncRunner)
queryurl := test.CreateQueue(t, ctx, cfg, test.QueryRunner)
jobsyncurl := test.CreateQueue(t, ctx, cfg, test.JobSyncRunner)
queryversionsyncurl := test.CreateQueue(t, ctx, cfg, test.QueryVersionSyncRunner)
bucketName := "docinitbucket"
test.CreateBucket(t, ctx, cfg, bucketName)
@@ -125,12 +126,20 @@ func TestProcess(t *testing.T) {
"DOCUMENT_SYNC_URL": docsyncurl,
},
},
{
Name: test.QueryVersionSyncRunner,
QueueURL: &queryversionsyncurl,
Env: map[string]string{
"JOB_SYNC_URL": jobsyncurl,
},
},
},
Services: []*test.ServiceNetworkConfig{
{
Name: test.QueryService,
Env: map[string]string{
"JOB_SYNC_URL": jobsyncurl,
"JOB_SYNC_URL": jobsyncurl,
"QUERY_VERSION_SYNC_URL": queryversionsyncurl,
},
},
},
+2 -1
View File
@@ -15,7 +15,8 @@ func TestClient(t *testing.T) {
c, cleanup := test.CreateServiceNetwork(t, ctx, &test.ServiceNetworkConfig{
Name: test.QueryService,
Env: map[string]string{
"JOB_SYNC_URL": "/i/am/here",
"JOB_SYNC_URL": "/i/am/here",
"QUERY_VERSION_SYNC_URL": "iamthere",
},
})
defer cleanup()
+2 -1
View File
@@ -15,7 +15,8 @@ func TestExportService(t *testing.T) {
c, cleanup := test.CreateServiceNetwork(t, ctx, &test.ServiceNetworkConfig{
Name: test.QueryService,
Env: map[string]string{
"JOB_SYNC_URL": "/i/am/here",
"JOB_SYNC_URL": "/i/am/here",
"QUERY_VERSION_SYNC_URL": "iamthere",
},
})
defer cleanup()
+2 -1
View File
@@ -15,7 +15,8 @@ func TestJob(t *testing.T) {
c, cleanup := test.CreateServiceNetwork(t, ctx, &test.ServiceNetworkConfig{
Name: test.QueryService,
Env: map[string]string{
"JOB_SYNC_URL": "/i/am/here",
"JOB_SYNC_URL": "/i/am/here",
"QUERY_VERSION_SYNC_URL": "iamthere",
},
})
defer cleanup()
@@ -46,7 +46,8 @@ func TestJobCollectorService(t *testing.T) {
Network: network,
Name: test.QueryService,
Env: map[string]string{
"JOB_SYNC_URL": jobsyncurl,
"JOB_SYNC_URL": jobsyncurl,
"QUERY_VERSION_SYNC_URL": "iamthere",
},
})
defer cleanup()
+2 -1
View File
@@ -16,7 +16,8 @@ func TestQueryServiceOpenAPI(t *testing.T) {
c, cleanup := test.CreateServiceNetwork(t, ctx, &test.ServiceNetworkConfig{
Name: test.QueryService,
Env: map[string]string{
"JOB_SYNC_URL": "/i/am/here",
"JOB_SYNC_URL": "/i/am/here",
"QUERY_VERSION_SYNC_URL": "iamthere",
},
})
defer cleanup()
+34 -2
View File
@@ -2,21 +2,51 @@ package endtoend_test
import (
"context"
"os"
"path"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue"
"queryorchestration/internal/test"
queryservice "queryorchestration/pkg/queryService"
"regexp"
"testing"
"github.com/oapi-codegen/runtime/types"
"github.com/stretchr/testify/assert"
)
type QueryConfig struct {
serviceconfig.BaseConfig
queue.QueueConfig
}
func TestQueryService(t *testing.T) {
ctx := context.Background()
cfg := &QueryConfig{}
test.SetCfgProvider(t, cfg)
cfg.SetBasePath(path.Join(os.Getenv("PWD"), "../.."))
network, ncleanup := test.CreateNetwork(t, ctx)
defer ncleanup()
_, clean := test.CreateAWSContainer(t, ctx, &test.CreateAWSConfig{
Cfg: cfg,
Network: network,
})
defer clean()
err := cfg.SetQueueClient(ctx)
assert.NoError(t, err)
queryversionsyncurl := test.CreateQueue(t, ctx, cfg, test.QueryVersionSyncRunner)
c, cleanup := test.CreateServiceNetwork(t, ctx, &test.ServiceNetworkConfig{
Name: test.QueryService,
Cfg: cfg,
Network: network,
Name: test.QueryService,
Env: map[string]string{
"JOB_SYNC_URL": "/i/am/here",
"JOB_SYNC_URL": "/i/am/here",
"QUERY_VERSION_SYNC_URL": queryversionsyncurl,
},
})
defer cleanup()
@@ -64,6 +94,8 @@ func TestQueryService(t *testing.T) {
assert.NoError(t, err)
assert.NotNil(t, res)
test.AssertMessageBody(t, ctx, cfg, queryversionsyncurl, regexp.MustCompile(`{"id":".+"}`))
queryRes, err = client.GetQueryWithResponse(ctx, jsonID)
assert.NoError(t, err)
assert.Equal(t, jsonID, queryRes.JSON200.Id)
+1 -1
View File
@@ -3,4 +3,4 @@
package aws
// goModuleVersion is the tagged release for this module
const goModuleVersion = "1.36.0"
const goModuleVersion = "1.36.1"
+33 -22
View File
@@ -76,28 +76,39 @@ type UserAgentFeature string
// Enumerates UserAgentFeature.
const (
UserAgentFeatureResourceModel UserAgentFeature = "A" // n/a (we don't generate separate resource types)
UserAgentFeatureWaiter = "B"
UserAgentFeaturePaginator = "C"
UserAgentFeatureRetryModeLegacy = "D" // n/a (equivalent to standard)
UserAgentFeatureRetryModeStandard = "E"
UserAgentFeatureRetryModeAdaptive = "F"
UserAgentFeatureS3Transfer = "G"
UserAgentFeatureS3CryptoV1N = "H" // n/a (crypto client is external)
UserAgentFeatureS3CryptoV2 = "I" // n/a
UserAgentFeatureS3ExpressBucket = "J"
UserAgentFeatureS3AccessGrants = "K" // not yet implemented
UserAgentFeatureGZIPRequestCompression = "L"
UserAgentFeatureProtocolRPCV2CBOR = "M"
UserAgentFeatureRequestChecksumCRC32 = "U"
UserAgentFeatureRequestChecksumCRC32C = "V"
UserAgentFeatureRequestChecksumCRC64 = "W"
UserAgentFeatureRequestChecksumSHA1 = "X"
UserAgentFeatureRequestChecksumSHA256 = "Y"
UserAgentFeatureRequestChecksumWhenSupported = "Z"
UserAgentFeatureRequestChecksumWhenRequired = "a"
UserAgentFeatureResponseChecksumWhenSupported = "b"
UserAgentFeatureResponseChecksumWhenRequired = "c"
UserAgentFeatureResourceModel UserAgentFeature = "A" // n/a (we don't generate separate resource types)
UserAgentFeatureWaiter = "B"
UserAgentFeaturePaginator = "C"
UserAgentFeatureRetryModeLegacy = "D" // n/a (equivalent to standard)
UserAgentFeatureRetryModeStandard = "E"
UserAgentFeatureRetryModeAdaptive = "F"
UserAgentFeatureS3Transfer = "G"
UserAgentFeatureS3CryptoV1N = "H" // n/a (crypto client is external)
UserAgentFeatureS3CryptoV2 = "I" // n/a
UserAgentFeatureS3ExpressBucket = "J"
UserAgentFeatureS3AccessGrants = "K" // not yet implemented
UserAgentFeatureGZIPRequestCompression = "L"
UserAgentFeatureProtocolRPCV2CBOR = "M"
UserAgentFeatureAccountIDEndpoint = "O" // DO NOT IMPLEMENT: rules output is not currently defined. SDKs should not parse endpoints for feature information.
UserAgentFeatureAccountIDModePreferred = "P"
UserAgentFeatureAccountIDModeDisabled = "Q"
UserAgentFeatureAccountIDModeRequired = "R"
UserAgentFeatureRequestChecksumCRC32 = "U"
UserAgentFeatureRequestChecksumCRC32C = "V"
UserAgentFeatureRequestChecksumCRC64 = "W"
UserAgentFeatureRequestChecksumSHA1 = "X"
UserAgentFeatureRequestChecksumSHA256 = "Y"
UserAgentFeatureRequestChecksumWhenSupported = "Z"
UserAgentFeatureRequestChecksumWhenRequired = "a"
UserAgentFeatureResponseChecksumWhenSupported = "b"
UserAgentFeatureResponseChecksumWhenRequired = "c"
)
// RequestUserAgent is a build middleware that set the User-Agent for the request.
@@ -1,3 +1,7 @@
# v1.3.32 (2025-02-05)
* **Dependency Update**: Updated to the latest SDK module versions
# v1.3.31 (2025-01-31)
* **Dependency Update**: Updated to the latest SDK module versions
@@ -3,4 +3,4 @@
package configsources
// goModuleVersion is the tagged release for this module
const goModuleVersion = "1.3.31"
const goModuleVersion = "1.3.32"
@@ -1,3 +1,7 @@
# v2.6.32 (2025-02-05)
* **Dependency Update**: Updated to the latest SDK module versions
# v2.6.31 (2025-01-31)
* **Dependency Update**: Updated to the latest SDK module versions
@@ -3,4 +3,4 @@
package endpoints
// goModuleVersion is the tagged release for this module
const goModuleVersion = "2.6.31"
const goModuleVersion = "2.6.32"
+3 -3
View File
@@ -20,7 +20,7 @@ github.com/Microsoft/go-winio/pkg/guid
# github.com/apapsch/go-jsonmerge/v2 v2.0.0
## explicit; go 1.12
github.com/apapsch/go-jsonmerge/v2
# github.com/aws/aws-sdk-go-v2 v1.36.0
# github.com/aws/aws-sdk-go-v2 v1.36.1
## explicit; go 1.21
github.com/aws/aws-sdk-go-v2/aws
github.com/aws/aws-sdk-go-v2/aws/arn
@@ -67,10 +67,10 @@ github.com/aws/aws-sdk-go-v2/credentials/stscreds
## explicit; go 1.21
github.com/aws/aws-sdk-go-v2/feature/ec2/imds
github.com/aws/aws-sdk-go-v2/feature/ec2/imds/internal/config
# github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.31
# github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.32
## explicit; go 1.21
github.com/aws/aws-sdk-go-v2/internal/configsources
# github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.31
# github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.32
## explicit; go 1.21
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2
# github.com/aws/aws-sdk-go-v2/internal/ini v1.8.2