From 0ebb8a21a1ba1d030536f7f147f17cf94eb33564 Mon Sep 17 00:00:00 2001 From: Michael McGuinness Date: Fri, 10 Jan 2025 11:12:03 +0000 Subject: [PATCH] Merged in bugfix/cleanup (pull request #10) Clean Up Testing and Linting * somescriptcleanup * mostgolangci * deplatest * imageversions * usingscratch * movedqueuefrtesting * finishedunittesting * linting * taskfilecontext --- .envrc | 1 + .mockery.yml | 5 +- .yamllint.yml | 2 +- Taskfile.yml | 1 + api/{controllers => api}/export.go | 0 api/{controllers => api}/jobcollector.go | 0 api/{controllers => api}/parse.go | 0 api/{controllers => api}/query.go | 0 api/queue/poll.go | 48 -------- api/queue/process.go | 18 --- api/queue/{document.go => queryrunner.go} | 25 ++-- build/Dockerfile | 69 +---------- cmd/queryRunner/main.go | 52 ++------- cmd/queryService/main.go | 48 +++----- deployments/compose.yaml | 5 +- internal/api/listener.go | 57 +++++++++ internal/api/listener_test.go | 106 +++++++++++++++++ internal/contextFull/service.go | 4 +- internal/database/connection_test.go | 27 +++-- internal/database/migrations_test.go | 4 +- internal/database/parseuuid.go | 8 +- internal/database/parseuuid_test.go | 4 +- .../database/repository/collector_test.go | 8 +- internal/database/repository/query_test.go | 8 +- internal/database/repository/result_test.go | 10 +- internal/database/repository/sqlquery_test.go | 19 ++- internal/document/sync_test.go | 18 +-- internal/jsonExtractor/process_test.go | 2 - internal/otel/service.go | 7 +- internal/query/get_test.go | 6 +- internal/query/list_test.go | 6 +- internal/queryQueue/create_test.go | 3 +- internal/queryQueue/service_test.go | 6 +- internal/queue/config.go | 9 +- internal/queue/delete.go | 2 +- internal/queue/delete_test.go | 32 ++++++ internal/queue/listener.go | 54 +++++++++ internal/queue/listener_test.go | 108 ++++++++++++++++++ internal/queue/poll.go | 51 +++++++++ internal/queue/poll_test.go | 47 ++++++++ internal/queue/queue_test.go | 11 +- internal/queue/receive.go | 17 +++ internal/queue/receive_test.go | 32 ++++++ internal/queue/send.go | 13 +-- internal/queue/send_test.go | 25 ++++ internal/result/store_test.go | 8 +- internal/server/server.go | 41 +++++++ internal/server/server_test.go | 92 +++++++++++++++ mocks/repository/mock_DBTX.go | 2 +- scripts/Taskfile.yml | 15 ++- scripts/compose.yml | 13 ++- scripts/database.yml | 11 +- scripts/dependencies.yml | 25 ++-- scripts/docker.yml | 6 +- scripts/proto.yml | 11 +- scripts/tests.yml | 23 ++-- sqlc.yml | 5 +- test/apiContainer_test.go | 34 ++++-- test/container_test.go | 2 +- test/queryrunner_test.go | 2 +- test/queryservice_test.go | 2 +- test/queueContainer_test.go | 27 ++++- 62 files changed, 940 insertions(+), 357 deletions(-) rename api/{controllers => api}/export.go (100%) rename api/{controllers => api}/jobcollector.go (100%) rename api/{controllers => api}/parse.go (100%) rename api/{controllers => api}/query.go (100%) delete mode 100644 api/queue/poll.go delete mode 100644 api/queue/process.go rename api/queue/{document.go => queryrunner.go} (59%) create mode 100644 internal/api/listener.go create mode 100644 internal/api/listener_test.go create mode 100644 internal/queue/delete_test.go create mode 100644 internal/queue/listener.go create mode 100644 internal/queue/listener_test.go create mode 100644 internal/queue/poll.go create mode 100644 internal/queue/poll_test.go create mode 100644 internal/queue/receive.go create mode 100644 internal/queue/receive_test.go create mode 100644 internal/queue/send_test.go create mode 100644 internal/server/server.go create mode 100644 internal/server/server_test.go diff --git a/.envrc b/.envrc index bc369071..7d00784b 100644 --- a/.envrc +++ b/.envrc @@ -4,6 +4,7 @@ export DB_PASS="${DB_PASS:-pass}" export DB_HOST="${DB_HOST:-localhost}" export DB_PORT="${DB_PORT:-5432}" export DB_NAME="${DB_NAME:-query_orchestration}" +export DB_URI="${DB_URI:-"postgres://${DB_USER}:${DB_PASS}@${DB_HOST}:${DB_PORT}/${DB_NAME}?sslmode=disable"}" dotenv_if_exists .env diff --git a/.mockery.yml b/.mockery.yml index 374293f6..8a8732c2 100644 --- a/.mockery.yml +++ b/.mockery.yml @@ -1,7 +1,8 @@ +--- with-expecter: true packages: - queryorchestration/internal: + queryorchestration/internal/database/repository: config: recursive: true all: true - dir: "mocks/{{.PackageName}}" \ No newline at end of file + dir: "mocks/{{.PackageName}}" diff --git a/.yamllint.yml b/.yamllint.yml index d7e61184..8815833d 100644 --- a/.yamllint.yml +++ b/.yamllint.yml @@ -1,4 +1,4 @@ --- extends: default ignore: | - vendor/ \ No newline at end of file + vendor/ diff --git a/Taskfile.yml b/Taskfile.yml index 8e74433f..24a2b550 100644 --- a/Taskfile.yml +++ b/Taskfile.yml @@ -1,3 +1,4 @@ +--- # https://taskfile.dev version: '3' diff --git a/api/controllers/export.go b/api/api/export.go similarity index 100% rename from api/controllers/export.go rename to api/api/export.go diff --git a/api/controllers/jobcollector.go b/api/api/jobcollector.go similarity index 100% rename from api/controllers/jobcollector.go rename to api/api/jobcollector.go diff --git a/api/controllers/parse.go b/api/api/parse.go similarity index 100% rename from api/controllers/parse.go rename to api/api/parse.go diff --git a/api/controllers/query.go b/api/api/query.go similarity index 100% rename from api/controllers/query.go rename to api/api/query.go diff --git a/api/queue/poll.go b/api/queue/poll.go deleted file mode 100644 index 7430c26f..00000000 --- a/api/queue/poll.go +++ /dev/null @@ -1,48 +0,0 @@ -package queue - -import ( - "context" - "log" - "queryorchestration/internal/queue" - - "github.com/aws/aws-sdk-go-v2/service/sqs" -) - -type Controllers struct { - Document DocumentController -} - -type Queue struct { - Config *queue.QueueConfig - Controllers *Controllers -} - -func PollMessages(ctx context.Context, queue *Queue) { - for { - result, err := queue.Config.Client.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{ - QueueUrl: &queue.Config.URL, - MaxNumberOfMessages: 1, - WaitTimeSeconds: 2, - VisibilityTimeout: 2, - MessageAttributeNames: []string{ - "type", - }, - }) - if err != nil { - log.Printf("Message Fetch Fail: %v", err) - continue - } - - for _, message := range result.Messages { - go func() { - toProcess, err := processMessage(ctx, queue, message) - if !toProcess { - return - } - if err != nil { - log.Printf("Message Process Fail: %v", err) - } - }() - } - } -} diff --git a/api/queue/process.go b/api/queue/process.go deleted file mode 100644 index 30e01a02..00000000 --- a/api/queue/process.go +++ /dev/null @@ -1,18 +0,0 @@ -package queue - -import ( - "context" - - "github.com/aws/aws-sdk-go-v2/service/sqs/types" -) - -func processMessage(ctx context.Context, queue *Queue, message types.Message) (bool, error) { - // Process the message here and return false if not message to process - switch *message.MessageAttributes["type"].StringValue { - case "DOCTEXT": - err := queue.Controllers.Document.Sync(ctx, queue.Config, &message) - return true, err - } - - return false, nil -} diff --git a/api/queue/document.go b/api/queue/queryrunner.go similarity index 59% rename from api/queue/document.go rename to api/queue/queryrunner.go index e4c60205..926ce331 100644 --- a/api/queue/document.go +++ b/api/queue/queryrunner.go @@ -1,4 +1,4 @@ -package queue +package controllers import ( "context" @@ -8,17 +8,18 @@ import ( "github.com/go-playground/validator/v10" + "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/service/sqs/types" "github.com/google/uuid" ) -type DocumentController struct { +type QueryRunner struct { validator *validator.Validate - document document.Service + document *document.Service } -func NewDocumentController(svc document.Service, validator *validator.Validate) *DocumentController { - return &DocumentController{ +func NewQueryRunner(svc *document.Service, validator *validator.Validate) QueryRunner { + return QueryRunner{ validator: validator, document: svc, } @@ -28,7 +29,7 @@ type DocumentQueryEvent struct { ID uuid.UUID `json:"id"` } -func (s *DocumentController) Sync(ctx context.Context, config *queue.QueueConfig, msg *types.Message) error { +func (s QueryRunner) Process(ctx context.Context, config *queue.Config, msg *types.Message) error { var body document.Document err := json.Unmarshal([]byte(*msg.Body), &body) if err != nil { @@ -49,12 +50,12 @@ func (s *DocumentController) Sync(ctx context.Context, config *queue.QueueConfig ID: body.ID, } - err = queue.Send(ctx, config, "DOCQUERY", queryEvent) - if err != nil { - return err - } - - err = queue.Delete(ctx, config, msg) + err = queue.Send(ctx, config, queryEvent, map[string]types.MessageAttributeValue{ + "type": { + DataType: aws.String("String"), + StringValue: aws.String("DOCQUERY"), + }, + }) if err != nil { return err } diff --git a/build/Dockerfile b/build/Dockerfile index 7347849d..4d158c37 100644 --- a/build/Dockerfile +++ b/build/Dockerfile @@ -1,85 +1,22 @@ # syntax=docker/dockerfile:1 -# Comments are provided throughout this file to help you get started. -# If you need more help, visit the Dockerfile reference guide at -# https://docs.docker.com/go/dockerfile-reference/ - -# Want to help us make this template better? Share your feedback here: https://forms.gle/ybq9Krt8jtBL3iCk7 - -################################################################################ -# Create a stage for building the application. ARG GO_VERSION=1.23 -#FROM --platform=$BUILDPLATFORM golang:${GO_VERSION} AS build FROM golang:${GO_VERSION} AS build WORKDIR /src -# Download dependencies as a separate step to take advantage of Docker's caching. -# Leverage a cache mount to /go/pkg/mod/ to speed up subsequent builds. -# Leverage bind mounts to go.sum and go.mod to avoid having to copy them into -# the container. RUN --mount=type=cache,target=/go/pkg/mod/ \ --mount=type=bind,source=go.sum,target=go.sum \ --mount=type=bind,source=go.mod,target=go.mod \ go mod download -x -# This is the architecture you're building for, which is passed in by the builder. -# Placing it here allows the previous steps to be cached across architectures. ARG TARGETARCH -ARG TARGETCMD - -RUN if [ -z "$TARGETCMD" ]; then \ - echo "Error: TARGETCMD is not set."; \ - exit 1; \ - fi - -# Build the application. -# Leverage a cache mount to /go/pkg/mod/ to speed up subsequent builds. -# Leverage a bind mount to the current directory to avoid having to copy the -# source code into the container. RUN --mount=type=cache,target=/go/pkg/mod/ \ --mount=type=bind,target=. \ - CGO_ENABLED=0 GOARCH=$TARGETARCH go build -o /bin/server ./cmd/$TARGETCMD/main.go + CGO_ENABLED=0 GOARCH=$TARGETARCH go build -o /bin ./cmd/... -################################################################################ -# Create a new stage for running the application that contains the minimal -# runtime dependencies for the application. This often uses a different base -# image from the build stage where the necessary files are copied from the build -# stage. -# -# The example below uses the alpine image as the foundation for running the app. -# By specifying the "latest" tag, it will also use whatever happens to be the -# most recent version of that image when you build your Dockerfile. If -# reproducability is important, consider using a versioned tag -# (e.g., alpine:3.17.2) or SHA (e.g., alpine@sha256:c41ab5c992deb4fe7e5da09f67a8804a46bd0592bfdf0b1847dde0e0889d2bff). -FROM alpine:latest AS final - -# Install any runtime dependencies that are needed to run your application. -# Leverage a cache mount to /var/cache/apk/ to speed up subsequent builds. -RUN --mount=type=cache,target=/var/cache/apk \ - apk --update add \ - ca-certificates \ - tzdata \ - && \ - update-ca-certificates - -# Create a non-privileged user that the app will run under. -# See https://docs.docker.com/go/dockerfile-user-best-practices/ -ARG UID=10001 -RUN adduser \ - --disabled-password \ - --gecos "" \ - --home "/nonexistent" \ - --shell "/sbin/nologin" \ - --no-create-home \ - --uid "${UID}" \ - appuser -USER appuser +FROM scratch AS final COPY database/migrations/ database/migrations/ -# Copy the executable from the "build" stage. -COPY --from=build /bin/server /bin/ - -# What the container should run when it is started. -ENTRYPOINT [ "/bin/server" ] +COPY --from=build /bin/ /bin/ diff --git a/cmd/queryRunner/main.go b/cmd/queryRunner/main.go index dc56aa34..5a192743 100644 --- a/cmd/queryRunner/main.go +++ b/cmd/queryRunner/main.go @@ -2,56 +2,24 @@ package main import ( "context" - "log" - "queryorchestration/api/queue" - "queryorchestration/internal/database" - "queryorchestration/internal/database/repository" + controllers "queryorchestration/api/queue" "queryorchestration/internal/document" - "queryorchestration/internal/env" - "queryorchestration/internal/otel" - queueSVC "queryorchestration/internal/queue" - - "github.com/aws/aws-sdk-go-v2/config" - "github.com/aws/aws-sdk-go-v2/service/sqs" - "github.com/go-playground/validator/v10" + "queryorchestration/internal/queue" + "queryorchestration/internal/server" ) func main() { ctx := context.Background() - closeTracer := otel.New(ctx) - defer closeTracer() + queryrunner := func(cfg *server.Config) queue.Controller { + svc := document.New(cfg.Database) - database.RunMigrations(&database.MigrationConfig{}) - - cfg, err := config.LoadDefaultConfig(ctx) - if err != nil { - log.Panicf("Unable to load SDK config: %v", err) + return controllers.NewQueryRunner(svc, cfg.Validator) } - queueURL := env.GetPanic("QUEUE_URL") + server := queue.NewServer(ctx, &queue.ListenerConfig{ + Controller: queryrunner, + }) - sqsClient := sqs.NewFromConfig(cfg) - - dbPool := database.GetDBPool(ctx) - dbQueries := repository.New(dbPool) - db := &database.Connection{ - Pool: dbPool, - Queries: dbQueries, - } - valid := validator.New() - - controllers := queue.Controllers{ - Document: *queue.NewDocumentController(*document.New(db), valid), - } - - config := queue.Queue{ - Controllers: &controllers, - Config: &queueSVC.QueueConfig{ - URL: queueURL, - Client: sqsClient, - }, - } - - queue.PollMessages(ctx, &config) + server.Listen(ctx) } diff --git a/cmd/queryService/main.go b/cmd/queryService/main.go index 7705a2ff..26067354 100644 --- a/cmd/queryService/main.go +++ b/cmd/queryService/main.go @@ -2,19 +2,14 @@ package main import ( "context" - "log" - "net" - "queryorchestration/api/controllers" + controllers "queryorchestration/api/api" serviceinterfaces "queryorchestration/api/serviceInterfaces" + "queryorchestration/internal/api" "queryorchestration/internal/collector" - "queryorchestration/internal/database" - "queryorchestration/internal/database/repository" "queryorchestration/internal/export" - "queryorchestration/internal/otel" "queryorchestration/internal/query" - "strconv" + "queryorchestration/internal/server" - "github.com/go-playground/validator/v10" _ "github.com/lib/pq" "google.golang.org/grpc" ) @@ -22,35 +17,18 @@ import ( func main() { ctx := context.Background() - closeTracer := otel.New(ctx) - defer closeTracer() + controllers := func(cfg *server.Config) *grpc.Server { + grpcServer := grpc.NewServer() + serviceinterfaces.RegisterQueryServiceServer(grpcServer, controllers.NewQueryController(*query.New(cfg.Database), cfg.Validator)) + serviceinterfaces.RegisterExportServiceServer(grpcServer, controllers.NewExportController(*export.New(cfg.Database), cfg.Validator)) + serviceinterfaces.RegisterJobCollectorServiceServer(grpcServer, controllers.NewJobCollectorController(*collector.New(cfg.Database), cfg.Validator)) - database.RunMigrations(&database.MigrationConfig{}) - - host := "0.0.0.0" - port := 8080 - address := net.JoinHostPort(host, strconv.Itoa(port)) - - lis, err := net.Listen("tcp", address) - if err != nil { - log.Panicf("failed to listen: %v", err) + return grpcServer } - dbPool := database.GetDBPool(ctx) - dbQueries := repository.New(dbPool) - db := &database.Connection{ - Pool: dbPool, - Queries: dbQueries, - } - valid := validator.New() + server := api.New(ctx, &api.Config{ + GetControllers: controllers, + }) - grpcServer := grpc.NewServer() - serviceinterfaces.RegisterQueryServiceServer(grpcServer, controllers.NewQueryController(*query.New(db), valid)) - serviceinterfaces.RegisterExportServiceServer(grpcServer, controllers.NewExportController(*export.New(db), valid)) - serviceinterfaces.RegisterJobCollectorServiceServer(grpcServer, controllers.NewJobCollectorController(*collector.New(db), valid)) - - log.Printf("Listening on port %d", port) - if err := grpcServer.Serve(lis); err != nil { - log.Panicf("Failed to serve: %v", err) - } + server.Listen() } diff --git a/deployments/compose.yaml b/deployments/compose.yaml index 4324eff0..6b48afa6 100644 --- a/deployments/compose.yaml +++ b/deployments/compose.yaml @@ -1,6 +1,7 @@ +--- services: db: - image: postgres:latest + image: postgres:17.2-alpine3.21 ports: - 5432:5432 environment: @@ -12,7 +13,7 @@ services: volumes: - db-data:/var/lib/postgresql/data healthcheck: - test: [ "CMD-SHELL", "pg_isready -U ${DB_USER}" ] + test: ["CMD-SHELL", "pg_isready -U ${DB_USER}"] interval: 10s timeout: 5s retries: 5 diff --git a/internal/api/listener.go b/internal/api/listener.go new file mode 100644 index 00000000..ba59fa69 --- /dev/null +++ b/internal/api/listener.go @@ -0,0 +1,57 @@ +package api + +import ( + "context" + "log" + "net" + "queryorchestration/internal/server" + "strconv" + + _ "github.com/lib/pq" + "google.golang.org/grpc" +) + +type Config struct { + GetControllers func(*server.Config) *grpc.Server + BasePath string +} + +type Server struct { + controllers *grpc.Server + port int + host string + address string +} + +func New(ctx context.Context, cfg *Config) *Server { + serverCfg := server.New(ctx, &server.NewConfig{ + BasePath: cfg.BasePath, + }) + + host := "0.0.0.0" + port := 8080 + address := net.JoinHostPort(host, strconv.Itoa(port)) + + controllers := cfg.GetControllers(serverCfg) + + return &Server{ + controllers: controllers, + port: port, + host: host, + address: address, + } +} + +func (s *Server) Listen() { + lis, err := net.Listen("tcp", s.address) + if err != nil { + log.Panicf("failed to listen: %v", err) + } + + log.Printf("Listening on port %d", s.port) + + err = s.controllers.Serve(lis) + if err != nil { + log.Panicf("Failed to serve: %v", err) + } +} diff --git a/internal/api/listener_test.go b/internal/api/listener_test.go new file mode 100644 index 00000000..1ba5cdd3 --- /dev/null +++ b/internal/api/listener_test.go @@ -0,0 +1,106 @@ +package api + +import ( + "context" + "fmt" + "queryorchestration/internal/database" + "queryorchestration/internal/server" + "testing" + + "github.com/docker/go-connections/nat" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/testcontainers/testcontainers-go" + "github.com/testcontainers/testcontainers-go/wait" + "google.golang.org/grpc" +) + +func TestNew(t *testing.T) { + ctx := context.Background() + + _, cleanup := createDB(t, ctx) + defer cleanup() + + getControllers := func(c *server.Config) *grpc.Server { + return nil + } + cfg := &Config{ + GetControllers: getControllers, + BasePath: "../..", + } + + server := New(ctx, cfg) + assert.Equal(t, 8080, server.port) + assert.Equal(t, "0.0.0.0", server.host) + assert.Equal(t, "0.0.0.0:8080", server.address) +} + +func TestListen(t *testing.T) { + t.Skip("Must test appropriately") + server := Server{} + + assert.Panics(t, func() { server.Listen() }) +} + +type db struct { + pool *pgxpool.Pool + container *testcontainers.Container +} + +func createDB(t *testing.T, ctx context.Context) (*db, func()) { + port, err := nat.NewPort("tcp", "5432") + if err != nil { + t.Fatalf("Failed to create port: %v", err) + } + + name := "queryorchestration" + pass := "pass" + user := "postgres" + + req := testcontainers.ContainerRequest{ + Image: "postgres:17.2-alpine3.21", + Env: map[string]string{ + "POSTGRES_DB": name, + "POSTGRES_USER": user, + "POSTGRES_PASSWORD": pass, + }, + ExposedPorts: []string{port.Port()}, + WaitingFor: wait.ForListeningPort(port), + } + + container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{ + ContainerRequest: req, + Started: true, + }) + if err != nil { + t.Fatalf("Failed to start container: %v", err) + } + + host, err := container.Host(ctx) + if err != nil { + t.Fatalf("Failed to extract host: %v", err) + } + mappedPort, err := container.MappedPort(ctx, port) + if err != nil { + t.Fatalf("Failed to extract port: %v", err) + } + + t.Setenv("DB_USER", user) + t.Setenv("DB_PASS", pass) + t.Setenv("DB_HOST", host) + t.Setenv("DB_PORT", fmt.Sprint(mappedPort.Int())) + t.Setenv("DB_NAME", name) + t.Setenv("DB_NOSSL", "1") + + pool := database.GetDBPool(ctx) + + return &db{ + pool: pool, + container: &container, + }, func() { + err := container.Terminate(ctx) + if err != nil { + t.Error(err) + } + } +} diff --git a/internal/contextFull/service.go b/internal/contextFull/service.go index 39242b71..bc23e1b9 100644 --- a/internal/contextFull/service.go +++ b/internal/contextFull/service.go @@ -2,7 +2,7 @@ package contextfull import ( "context" - "fmt" + "errors" queryprocessor "queryorchestration/internal/queryProcessor" "queryorchestration/internal/result" ) @@ -16,7 +16,7 @@ func NewExtractor() Extractor { func (e Extractor) Process(ctx context.Context, query *queryprocessor.Query, values []result.Value) (string, error) { if len(values) > 0 { - return "", fmt.Errorf("no requirements expected") + return "", errors.New("no requirements expected") } // TODO diff --git a/internal/database/connection_test.go b/internal/database/connection_test.go index 4043ad5a..b6df6c4b 100644 --- a/internal/database/connection_test.go +++ b/internal/database/connection_test.go @@ -14,8 +14,8 @@ import ( func TestDBConn(t *testing.T) { ctx := context.Background() - _, container := createDB(t, ctx) - defer container.Terminate(ctx) + _, cleanup := createDB(t, ctx) + defer cleanup() conn := database.GetDBConn(ctx) assert.NotNil(t, conn) @@ -34,8 +34,8 @@ func TestDBConnNoDB(t *testing.T) { func TestDBPool(t *testing.T) { ctx := context.Background() - _, container := createDB(t, ctx) - defer container.Terminate(ctx) + _, cleanup := createDB(t, ctx) + defer cleanup() pool := database.GetDBPool(ctx) assert.NotNil(t, pool) @@ -49,7 +49,12 @@ type dbConfig struct { Password string } -func createDB(t *testing.T, ctx context.Context) (*dbConfig, testcontainers.Container) { +type db struct { + config *dbConfig + container *testcontainers.Container +} + +func createDB(t *testing.T, ctx context.Context) (*db, func()) { alias := "postgres" port, err := nat.NewPort("tcp", "5432") if err != nil { @@ -65,7 +70,7 @@ func createDB(t *testing.T, ctx context.Context) (*dbConfig, testcontainers.Cont } req := testcontainers.ContainerRequest{ - Image: "postgres:latest", + Image: "postgres:17.2-alpine3.21", Env: map[string]string{ "POSTGRES_DB": config.Name, "POSTGRES_USER": config.User, @@ -102,5 +107,13 @@ func createDB(t *testing.T, ctx context.Context) (*dbConfig, testcontainers.Cont t.Setenv("DB_PORT", fmt.Sprint(config.Port)) t.Setenv("DB_NAME", config.Name) - return &config, container + return &db{ + config: &config, + container: &container, + }, func() { + err := container.Terminate(ctx) + if err != nil { + t.Error(err) + } + } } diff --git a/internal/database/migrations_test.go b/internal/database/migrations_test.go index 47cabf57..84ebbfb8 100644 --- a/internal/database/migrations_test.go +++ b/internal/database/migrations_test.go @@ -12,8 +12,8 @@ import ( func TestRunMigrations(t *testing.T) { ctx := context.Background() - _, container := createDB(t, ctx) - defer container.Terminate(ctx) + _, cleanup := createDB(t, ctx) + defer cleanup() t.Setenv("DB_NOSSL", "1") diff --git a/internal/database/parseuuid.go b/internal/database/parseuuid.go index 2a0ec948..4861ce44 100644 --- a/internal/database/parseuuid.go +++ b/internal/database/parseuuid.go @@ -1,6 +1,8 @@ package database import ( + "log" + "github.com/google/uuid" "github.com/jackc/pgx/v5/pgtype" ) @@ -16,7 +18,11 @@ func MustToDBUUIDArray(ids []uuid.UUID) []pgtype.UUID { func MustToDBUUID(id uuid.UUID) pgtype.UUID { var dbID pgtype.UUID - dbID.Scan(id.String()) + err := dbID.Scan(id.String()) + if err != nil { + log.Panic(err) + } + return dbID } diff --git a/internal/database/parseuuid_test.go b/internal/database/parseuuid_test.go index 14b098ef..ee1ac14e 100644 --- a/internal/database/parseuuid_test.go +++ b/internal/database/parseuuid_test.go @@ -23,7 +23,7 @@ func TestMustToDBUUIDArray(t *testing.T) { dbIDs := database.MustToDBUUIDArray(ids) - assert.Equal(t, len(ids), len(dbIDs)) + assert.Len(t, dbIDs, len(ids)) for index, id := range dbIDs { assert.Equal(t, database.MustToDBUUID(ids[index]), id) } @@ -42,7 +42,7 @@ func TestMustToUUIDArray(t *testing.T) { ids := database.MustToUUIDArray(dbIDs) - assert.Equal(t, len(ids), len(dbIDs)) + assert.Len(t, ids, len(dbIDs)) for index, id := range dbIDs { assert.Equal(t, database.MustToDBUUID(ids[index]), id) } diff --git a/internal/database/repository/collector_test.go b/internal/database/repository/collector_test.go index 7348914a..3f75944c 100644 --- a/internal/database/repository/collector_test.go +++ b/internal/database/repository/collector_test.go @@ -12,15 +12,15 @@ import ( func TestCollector(t *testing.T) { ctx := context.Background() - pool, container := createDB(t, ctx) - defer container.Terminate(ctx) - queries := repository.New(pool) + db, cleanup := createDB(t, ctx) + defer cleanup() + queries := repository.New(db.pool) collectorID := database.MustToDBUUID(uuid.New()) collectorQueries, err := queries.GetCollectorQueries(ctx, collectorID) assert.Nil(t, err) - assert.Equal(t, 0, len(collectorQueries)) + assert.Len(t, collectorQueries, 0) assert.ElementsMatch(t, []repository.GetCollectorQueriesRow{}, collectorQueries) jobID := database.MustToDBUUID(uuid.New()) diff --git a/internal/database/repository/query_test.go b/internal/database/repository/query_test.go index ea41fc9d..29ce3eaf 100644 --- a/internal/database/repository/query_test.go +++ b/internal/database/repository/query_test.go @@ -12,9 +12,9 @@ import ( func TestQueries(t *testing.T) { ctx := context.Background() - pool, container := createDB(t, ctx) - defer container.Terminate(ctx) - queries := repository.New(pool) + db, cleanup := createDB(t, ctx) + defer cleanup() + queries := repository.New(db.pool) contextQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeContextFull)) assert.Nil(t, err) @@ -60,7 +60,7 @@ func TestQueries(t *testing.T) { qs, err := queries.ListQueries(ctx) assert.Nil(t, err) - assert.Equal(t, 2, len(qs)) + assert.Len(t, qs, 2) log.Print(qs) assert.ElementsMatch(t, []repository.Fullactivequery{ {ID: jsonQueryID, diff --git a/internal/database/repository/result_test.go b/internal/database/repository/result_test.go index d960c5a4..b68c708c 100644 --- a/internal/database/repository/result_test.go +++ b/internal/database/repository/result_test.go @@ -13,9 +13,9 @@ import ( func TestResults(t *testing.T) { ctx := context.Background() - pool, container := createDB(t, ctx) - defer container.Terminate(ctx) - queries := repository.New(pool) + db, cleanup := createDB(t, ctx) + defer cleanup() + queries := repository.New(db.pool) jsonQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor)) assert.Nil(t, err) @@ -46,7 +46,7 @@ func TestResults(t *testing.T) { Textversion: textVersion, }) assert.Nil(t, err) - assert.Equal(t, 1, len(resultsByDoc)) + assert.Len(t, resultsByDoc, 1) assert.ElementsMatch(t, []repository.ListResultsByDocumentIDRow{ { ID: jsonResultID, @@ -57,7 +57,7 @@ func TestResults(t *testing.T) { results, err := queries.ListResultValuesByID(ctx, []pgtype.UUID{jsonResultID}) assert.Nil(t, err) - assert.Equal(t, 1, len(results)) + assert.Len(t, results, 1) assert.ElementsMatch(t, []repository.ListResultValuesByIDRow{ { ID: jsonResultID, diff --git a/internal/database/repository/sqlquery_test.go b/internal/database/repository/sqlquery_test.go index 8c6b84c1..7c64b0eb 100644 --- a/internal/database/repository/sqlquery_test.go +++ b/internal/database/repository/sqlquery_test.go @@ -14,7 +14,12 @@ import ( "github.com/testcontainers/testcontainers-go/wait" ) -func createDB(t *testing.T, ctx context.Context) (*pgxpool.Pool, testcontainers.Container) { +type db struct { + pool *pgxpool.Pool + container *testcontainers.Container +} + +func createDB(t *testing.T, ctx context.Context) (*db, func()) { port, err := nat.NewPort("tcp", "5432") if err != nil { t.Fatalf("Failed to create port: %v", err) @@ -25,7 +30,7 @@ func createDB(t *testing.T, ctx context.Context) (*pgxpool.Pool, testcontainers. user := "postgres" req := testcontainers.ContainerRequest{ - Image: "postgres:latest", + Image: "postgres:17.2-alpine3.21", Env: map[string]string{ "POSTGRES_DB": name, "POSTGRES_USER": user, @@ -65,5 +70,13 @@ func createDB(t *testing.T, ctx context.Context) (*pgxpool.Pool, testcontainers. pool := database.GetDBPool(ctx) - return pool, container + return &db{ + pool: pool, + container: &container, + }, func() { + err := container.Terminate(ctx) + if err != nil { + t.Error(err) + } + } } diff --git a/internal/document/sync_test.go b/internal/document/sync_test.go index 298f5ea2..f3dbb656 100644 --- a/internal/document/sync_test.go +++ b/internal/document/sync_test.go @@ -98,13 +98,13 @@ func TestSyncDBFail(t *testing.T) { pgxmock.NewRows([]string{"id", "jobId", "minCleanVersion", "minTextVersion"}). AddRow(dbCollectorId, dbJobID, minCleanVersion, minTextVersion), ) - errr := "database failure" + dbErr := "database failure" pool.ExpectQuery("name: ListResultsByDocumentID :many").WithArgs(database.MustToDBUUID(doc.ID), minCleanVersion, minTextVersion). - WillReturnError(errors.New(errr)) + WillReturnError(errors.New(dbErr)) docSvc = document.New(db) err = docSvc.Sync(ctx, &doc) - assert.EqualError(t, err, errr) + assert.EqualError(t, err, dbErr) pool.ExpectQuery("name: GetCollectorFromJobID :one").WithArgs(dbJobID). WillReturnRows( @@ -115,13 +115,13 @@ func TestSyncDBFail(t *testing.T) { WillReturnRows( pgxmock.NewRows([]string{"id", "queryId", "queryVersion"}), ) - errr = "database failure" + dbErr = "database failure" pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorId). - WillReturnError(errors.New(errr)) + WillReturnError(errors.New(dbErr)) docSvc = document.New(db) err = docSvc.Sync(ctx, &doc) - assert.EqualError(t, err, errr) + assert.EqualError(t, err, dbErr) pool.ExpectQuery("name: GetCollectorFromJobID :one").WithArgs(dbJobID). WillReturnRows( @@ -132,18 +132,18 @@ func TestSyncDBFail(t *testing.T) { WillReturnRows( pgxmock.NewRows([]string{"id", "queryId", "queryVersion"}), ) - errr = "database failure" + dbErr = "database failure" pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorId). WillReturnRows( pgxmock.NewRows([]string{"collectorId", "queryId", "type", "queryVersion", "requiredIds"}). AddRow(dbCollectorId, dbQueryID, repository.NullQuerytype{Querytype: repository.QuerytypeJsonExtractor, Valid: true}, pgtype.Int4{Int32: int32(1), Valid: true}, []pgtype.UUID{}), ) pool.ExpectQuery("name: ListResultValuesByID :many").WithArgs([]pgtype.UUID{}). - WillReturnError(errors.New(errr)) + WillReturnError(errors.New(dbErr)) docSvc = document.New(db) err = docSvc.Sync(ctx, &doc) - assert.EqualError(t, err, errr) + assert.EqualError(t, err, dbErr) } func TestSync(t *testing.T) { diff --git a/internal/jsonExtractor/process_test.go b/internal/jsonExtractor/process_test.go index f9eae4ca..177cdc94 100644 --- a/internal/jsonExtractor/process_test.go +++ b/internal/jsonExtractor/process_test.go @@ -164,8 +164,6 @@ func TestJSONProcessJSON(t *testing.T) { assert.EqualError(t, err, "invalid character '}' looking for beginning of value") assert.Empty(t, value) - config = "{\"path\":\"key\"}" - pool.ExpectQuery("name: GetQueryConfig :one").WithArgs(database.MustToDBUUID(query.ID), query.Version). WillReturnRows( pgxmock.NewRows([]string{"id", "config"}), diff --git a/internal/otel/service.go b/internal/otel/service.go index 1023ac71..845753a4 100644 --- a/internal/otel/service.go +++ b/internal/otel/service.go @@ -13,7 +13,7 @@ import ( func New(ctx context.Context) func() { if os.Getenv("ENABLE_OTEL") != "1" { log.Println("OpenTelemetry is disabled. Set ENABLE_OTEL to enable.") - return func() {} // No-op shutdown function + return func() {} } exporter, err := otlptracegrpc.New(ctx) @@ -28,6 +28,9 @@ func New(ctx context.Context) func() { otel.SetTracerProvider(tp) return func() { - _ = tp.Shutdown(ctx) + err := tp.Shutdown(ctx) + if err != nil { + log.Panic(err) + } } } diff --git a/internal/query/get_test.go b/internal/query/get_test.go index 8dd2e997..704a3779 100644 --- a/internal/query/get_test.go +++ b/internal/query/get_test.go @@ -9,7 +9,6 @@ import ( "testing" "github.com/google/uuid" - "github.com/jackc/pgx/v5/pgtype" "github.com/pashagolub/pgxmock/v3" "github.com/stretchr/testify/assert" ) @@ -40,10 +39,7 @@ func TestGet(t *testing.T) { Config: config, } - dbReqIDs := make([]pgtype.UUID, len(query.RequiredQueryIDs)) - for index, id := range query.RequiredQueryIDs { - dbReqIDs[index] = database.MustToDBUUID(id) - } + dbReqIDs := database.MustToDBUUIDArray(query.RequiredQueryIDs) pool.ExpectQuery("name: GetQuery :one").WithArgs(database.MustToDBUUID(query.ID)).WillReturnRows( pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}). diff --git a/internal/query/list_test.go b/internal/query/list_test.go index a51a3a0e..23315160 100644 --- a/internal/query/list_test.go +++ b/internal/query/list_test.go @@ -9,7 +9,6 @@ import ( "testing" "github.com/google/uuid" - "github.com/jackc/pgx/v5/pgtype" "github.com/pashagolub/pgxmock/v3" "github.com/stretchr/testify/assert" ) @@ -40,10 +39,7 @@ func TestList(t *testing.T) { Config: config, } - dbReqIDs := make([]pgtype.UUID, len(q.RequiredQueryIDs)) - for index, id := range q.RequiredQueryIDs { - dbReqIDs[index] = database.MustToDBUUID(id) - } + dbReqIDs := database.MustToDBUUIDArray(q.RequiredQueryIDs) filters := query.ListFilters{} diff --git a/internal/queryQueue/create_test.go b/internal/queryQueue/create_test.go index dff8e42e..2e54365d 100644 --- a/internal/queryQueue/create_test.go +++ b/internal/queryQueue/create_test.go @@ -76,7 +76,8 @@ func TestGetCollectorQueries(t *testing.T) { pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorID).WillReturnRows(rows) - svc.getCollectorQueries(ctx) + err = svc.getCollectorQueries(ctx) + assert.Nil(t, err) assert.EqualExportedValues(t, collectorQueries, svc.collectorQueries) } diff --git a/internal/queryQueue/service_test.go b/internal/queryQueue/service_test.go index e62aa49a..60516ed2 100644 --- a/internal/queryQueue/service_test.go +++ b/internal/queryQueue/service_test.go @@ -130,9 +130,9 @@ func TestQueueFail(t *testing.T) { coll, err := collector.NewByJobId(ctx, db, jobID) assert.Nil(t, err) - errr := "database failure" + dbErr := "database failure" pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorID). - WillReturnError(errors.New(errr)) + WillReturnError(errors.New(dbErr)) results := []*result.Result{} @@ -141,5 +141,5 @@ func TestQueueFail(t *testing.T) { textVersion := int32(1) _, err = queryqueue.New(ctx, db, coll, results, docID, cleanVersion, textVersion) - assert.EqualError(t, err, errr) + assert.EqualError(t, err, dbErr) } diff --git a/internal/queue/config.go b/internal/queue/config.go index 8a61050c..1bfa85cb 100644 --- a/internal/queue/config.go +++ b/internal/queue/config.go @@ -1,10 +1,17 @@ package queue import ( + "context" + "github.com/aws/aws-sdk-go-v2/service/sqs" + "github.com/aws/aws-sdk-go-v2/service/sqs/types" ) -type QueueConfig struct { +type Config struct { URL string Client *sqs.Client } + +type Controller interface { + Process(ctx context.Context, config *Config, msg *types.Message) error +} diff --git a/internal/queue/delete.go b/internal/queue/delete.go index b9641275..aa91287e 100644 --- a/internal/queue/delete.go +++ b/internal/queue/delete.go @@ -8,7 +8,7 @@ import ( "github.com/aws/aws-sdk-go-v2/service/sqs/types" ) -func Delete(ctx context.Context, config *QueueConfig, msg *types.Message) error { +func Delete(ctx context.Context, config *Config, msg *types.Message) error { _, err := config.Client.DeleteMessage(ctx, &sqs.DeleteMessageInput{ QueueUrl: aws.String(config.URL), ReceiptHandle: msg.ReceiptHandle, diff --git a/internal/queue/delete_test.go b/internal/queue/delete_test.go new file mode 100644 index 00000000..731e2432 --- /dev/null +++ b/internal/queue/delete_test.go @@ -0,0 +1,32 @@ +package queue_test + +import ( + "context" + "queryorchestration/internal/queue" + "testing" + + "github.com/aws/aws-sdk-go-v2/service/sqs/types" + "github.com/stretchr/testify/assert" +) + +func TestDelete(t *testing.T) { + ctx := context.Background() + queueConfig, cleanup := createQueue(t, ctx) + defer cleanup() + + err := queue.Send(ctx, queueConfig.Config, "{}", map[string]types.MessageAttributeValue{}) + assert.Nil(t, err) + + result, err := queue.Receive(ctx, queueConfig.Config, []string{}) + assert.Nil(t, err) + + assert.Len(t, result.Messages, 1) + message := result.Messages[0] + + err = queue.Delete(ctx, queueConfig.Config, &message) + assert.Nil(t, err) + + result, err = queue.Receive(ctx, queueConfig.Config, []string{}) + assert.Nil(t, err) + assert.Len(t, result.Messages, 0) +} diff --git a/internal/queue/listener.go b/internal/queue/listener.go new file mode 100644 index 00000000..7dc28d98 --- /dev/null +++ b/internal/queue/listener.go @@ -0,0 +1,54 @@ +package queue + +import ( + "context" + "log" + "queryorchestration/internal/env" + "queryorchestration/internal/server" + + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/service/sqs" +) + +type ListenerConfig struct { + Controller func(*server.Config) Controller + BasePath string +} + +type Server struct { + controller Controller + queueConnection *Config +} + +func NewServer(ctx context.Context, lConfig *ListenerConfig) *Server { + serverCfg := server.New(ctx, &server.NewConfig{ + BasePath: lConfig.BasePath, + }) + + cfg, err := config.LoadDefaultConfig(ctx) + if err != nil { + log.Panicf("Unable to load SDK config: %v", err) + } + + queueURL := env.GetPanic("QUEUE_URL") + + sqsClient := sqs.NewFromConfig(cfg) + + return &Server{ + controller: lConfig.Controller(serverCfg), + queueConnection: &Config{ + URL: queueURL, + Client: sqsClient, + }, + } +} + +func (s *Server) Listen(ctx context.Context) { + config := &PollConfig{ + Controller: s.controller, + Config: s.queueConnection, + } + + log.Print("Listening to queue") + PollMessages(ctx, config) +} diff --git a/internal/queue/listener_test.go b/internal/queue/listener_test.go new file mode 100644 index 00000000..83073b9e --- /dev/null +++ b/internal/queue/listener_test.go @@ -0,0 +1,108 @@ +package queue_test + +import ( + "context" + "fmt" + "queryorchestration/internal/database" + "queryorchestration/internal/queue" + "queryorchestration/internal/server" + "testing" + "time" + + "github.com/docker/go-connections/nat" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/testcontainers/testcontainers-go" + "github.com/testcontainers/testcontainers-go/wait" +) + +func TestNew(t *testing.T) { + ctx := context.Background() + _, cleanup := createQueue(t, ctx) + defer cleanup() + _, cleanup = createDB(t, ctx) + defer cleanup() + + ctx, cancel := context.WithTimeout(ctx, time.Second) + defer cancel() + + controller := func(cfg *server.Config) queue.Controller { + return MockController{} + } + + queue.NewServer(ctx, &queue.ListenerConfig{ + Controller: controller, + BasePath: "../..", + }) +} + +func TestListen(t *testing.T) { + ctx := context.Background() + + queue := queue.Server{} + + assert.Panics(t, func() { queue.Listen(ctx) }) +} + +type db struct { + pool *pgxpool.Pool + container *testcontainers.Container +} + +func createDB(t *testing.T, ctx context.Context) (*db, func()) { + port, err := nat.NewPort("tcp", "5432") + if err != nil { + t.Fatalf("Failed to create port: %v", err) + } + + name := "queryorchestration" + pass := "pass" + user := "postgres" + + req := testcontainers.ContainerRequest{ + Image: "postgres:17.2-alpine3.21", + Env: map[string]string{ + "POSTGRES_DB": name, + "POSTGRES_USER": user, + "POSTGRES_PASSWORD": pass, + }, + ExposedPorts: []string{port.Port()}, + WaitingFor: wait.ForListeningPort(port), + } + + container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{ + ContainerRequest: req, + Started: true, + }) + if err != nil { + t.Fatalf("Failed to start container: %v", err) + } + + host, err := container.Host(ctx) + if err != nil { + t.Fatalf("Failed to extract host: %v", err) + } + mappedPort, err := container.MappedPort(ctx, port) + if err != nil { + t.Fatalf("Failed to extract port: %v", err) + } + + t.Setenv("DB_USER", user) + t.Setenv("DB_PASS", pass) + t.Setenv("DB_HOST", host) + t.Setenv("DB_PORT", fmt.Sprint(mappedPort.Int())) + t.Setenv("DB_NAME", name) + t.Setenv("DB_NOSSL", "1") + + pool := database.GetDBPool(ctx) + + return &db{ + pool: pool, + container: &container, + }, func() { + err := container.Terminate(ctx) + if err != nil { + t.Error(err) + } + } +} diff --git a/internal/queue/poll.go b/internal/queue/poll.go new file mode 100644 index 00000000..4e6c2d5c --- /dev/null +++ b/internal/queue/poll.go @@ -0,0 +1,51 @@ +package queue + +import ( + "context" + "fmt" + "log" +) + +type PollConfig struct { + Config *Config + Controller Controller +} + +func PollMessages(ctx context.Context, queueConfig *PollConfig) { + for { + select { + case <-ctx.Done(): + return + default: + err := PollMessage(ctx, queueConfig) + if err != nil { + log.Print(err) + } + } + } +} + +func PollMessage(ctx context.Context, queueConfig *PollConfig) error { + result, err := Receive(ctx, queueConfig.Config, []string{ + "type", + }) + if err != nil { + return fmt.Errorf("message fetch fail: %v", err) + } + + for _, message := range result.Messages { + go func() { + err := queueConfig.Controller.Process(ctx, queueConfig.Config, &message) + if err != nil { + log.Printf("message process fail: %v", err) + } + + err = Delete(ctx, queueConfig.Config, &message) + if err != nil { + log.Printf("message delete fail: %v", err) + } + }() + } + + return nil +} diff --git a/internal/queue/poll_test.go b/internal/queue/poll_test.go new file mode 100644 index 00000000..d2ed0a12 --- /dev/null +++ b/internal/queue/poll_test.go @@ -0,0 +1,47 @@ +package queue_test + +import ( + "context" + "queryorchestration/internal/queue" + "testing" + "time" + + "github.com/aws/aws-sdk-go-v2/service/sqs/types" + "github.com/stretchr/testify/assert" +) + +type MockController struct{} + +func (s MockController) Process(ctx context.Context, config *queue.Config, msg *types.Message) error { + return nil +} + +func TestPollMessages(t *testing.T) { + ctx := context.Background() + queueConfig, cleanup := createQueue(t, ctx) + defer cleanup() + + controller := MockController{} + + ctx, cancel := context.WithTimeout(ctx, time.Second) + defer cancel() + + queue.PollMessages(ctx, &queue.PollConfig{ + Config: queueConfig.Config, + Controller: controller, + }) +} + +func TestPollMessage(t *testing.T) { + ctx := context.Background() + queueConfig, cleanup := createQueue(t, ctx) + defer cleanup() + + controller := MockController{} + + err := queue.PollMessage(ctx, &queue.PollConfig{ + Config: queueConfig.Config, + Controller: controller, + }) + assert.Nil(t, err) +} diff --git a/internal/queue/queue_test.go b/internal/queue/queue_test.go index a914377d..4e56414a 100644 --- a/internal/queue/queue_test.go +++ b/internal/queue/queue_test.go @@ -17,7 +17,7 @@ import ( type queueConfig struct { Container *testcontainers.Container - Config *queue.QueueConfig + Config *queue.Config } func createQueue(t *testing.T, ctx context.Context) (*queueConfig, func()) { @@ -32,7 +32,7 @@ func createQueue(t *testing.T, ctx context.Context) (*queueConfig, func()) { } req := testcontainers.ContainerRequest{ - Image: "localstack/localstack:latest", + Image: "localstack/localstack:4.0.3", Env: map[string]string{ "AWS_ACCESS_KEY_ID": provider.Value.AccessKeyID, "AWS_SECRET_ACCESS_KEY": provider.Value.SecretAccessKey, @@ -82,11 +82,14 @@ func createQueue(t *testing.T, ctx context.Context) (*queueConfig, func()) { return &queueConfig{ Container: &container, - Config: &queue.QueueConfig{ + Config: &queue.Config{ Client: client, URL: *queueM.QueueUrl, }, }, func() { - container.Terminate(ctx) + err := container.Terminate(ctx) + if err != nil { + t.Error(err) + } } } diff --git a/internal/queue/receive.go b/internal/queue/receive.go new file mode 100644 index 00000000..133cf01a --- /dev/null +++ b/internal/queue/receive.go @@ -0,0 +1,17 @@ +package queue + +import ( + "context" + + "github.com/aws/aws-sdk-go-v2/service/sqs" +) + +func Receive(ctx context.Context, config *Config, attributes []string) (*sqs.ReceiveMessageOutput, error) { + return config.Client.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{ + QueueUrl: &config.URL, + MaxNumberOfMessages: 1, + WaitTimeSeconds: 2, + VisibilityTimeout: 2, + MessageAttributeNames: attributes, + }) +} diff --git a/internal/queue/receive_test.go b/internal/queue/receive_test.go new file mode 100644 index 00000000..b9c93bb0 --- /dev/null +++ b/internal/queue/receive_test.go @@ -0,0 +1,32 @@ +package queue_test + +import ( + "context" + "log" + "queryorchestration/internal/queue" + "testing" + + "github.com/aws/aws-sdk-go-v2/service/sqs/types" + "github.com/stretchr/testify/assert" +) + +func TestReceive(t *testing.T) { + ctx := context.Background() + queueConfig, cleanup := createQueue(t, ctx) + defer cleanup() + + attributes := map[string]types.MessageAttributeValue{} + + err := queue.Send(ctx, queueConfig.Config, "example_body", attributes) + assert.Nil(t, err) + + result, err := queue.Receive(ctx, queueConfig.Config, []string{}) + assert.Nil(t, err) + + assert.Len(t, result.Messages, 1) + message := result.Messages[0] + + log.Print(message.Attributes) + + assert.Equal(t, "\"example_body\"", *message.Body) +} diff --git a/internal/queue/send.go b/internal/queue/send.go index 81611ef3..0a300238 100644 --- a/internal/queue/send.go +++ b/internal/queue/send.go @@ -9,7 +9,7 @@ import ( "github.com/aws/aws-sdk-go-v2/service/sqs/types" ) -func Send(ctx context.Context, config *QueueConfig, typeName string, body interface{}) error { +func Send(ctx context.Context, config *Config, body interface{}, attributes map[string]types.MessageAttributeValue) error { jsonBytes, err := json.Marshal(body) if err != nil { return err @@ -18,14 +18,9 @@ func Send(ctx context.Context, config *QueueConfig, typeName string, body interf strBody := string(jsonBytes) _, err = config.Client.SendMessage(ctx, &sqs.SendMessageInput{ - MessageAttributes: map[string]types.MessageAttributeValue{ - "type": { - DataType: aws.String("String"), - StringValue: aws.String(typeName), - }, - }, - QueueUrl: aws.String(config.URL), - MessageBody: aws.String(strBody), + MessageAttributes: attributes, + QueueUrl: aws.String(config.URL), + MessageBody: aws.String(strBody), }) if err != nil { return err diff --git a/internal/queue/send_test.go b/internal/queue/send_test.go new file mode 100644 index 00000000..66dd86e5 --- /dev/null +++ b/internal/queue/send_test.go @@ -0,0 +1,25 @@ +package queue_test + +import ( + "context" + "queryorchestration/internal/queue" + "testing" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/service/sqs/types" + "github.com/stretchr/testify/assert" +) + +func TestSend(t *testing.T) { + ctx := context.Background() + queueConfig, cleanup := createQueue(t, ctx) + defer cleanup() + + err := queue.Send(ctx, queueConfig.Config, "{}", map[string]types.MessageAttributeValue{ + "type": { + DataType: aws.String("String"), + StringValue: aws.String("EXAMPLE_TYPE"), + }, + }) + assert.Nil(t, err) +} diff --git a/internal/result/store_test.go b/internal/result/store_test.go index d5d286bc..34cc94ee 100644 --- a/internal/result/store_test.go +++ b/internal/result/store_test.go @@ -2,7 +2,7 @@ package result_test import ( "context" - "fmt" + "errors" "queryorchestration/internal/database" "queryorchestration/internal/database/repository" "queryorchestration/internal/result" @@ -40,11 +40,11 @@ func TestStore(t *testing.T) { assert.Nil(t, err) assert.NotNil(t, id) - errr := "database failing" + dbErr := "database failing" pool.ExpectQuery("name: SetResult :one").WithArgs(database.MustToDBUUID(resultStore.QueryID), database.MustToDBUUID(resultStore.DocumentID), resultStore.Value, resultStore.CleanVersion, resultStore.TextVersion, resultStore.QueryVersion). - WillReturnError(fmt.Errorf(errr)) + WillReturnError(errors.New(dbErr)) id, err = result.Store(ctx, queries, &resultStore) - assert.EqualError(t, err, errr) + assert.EqualError(t, err, dbErr) assert.Equal(t, uuid.Nil, id) } diff --git a/internal/server/server.go b/internal/server/server.go new file mode 100644 index 00000000..635c6b75 --- /dev/null +++ b/internal/server/server.go @@ -0,0 +1,41 @@ +package server + +import ( + "context" + "queryorchestration/internal/database" + "queryorchestration/internal/database/repository" + "queryorchestration/internal/otel" + + "github.com/go-playground/validator/v10" +) + +type Config struct { + Database *database.Connection + Validator *validator.Validate +} + +type NewConfig struct { + BasePath string +} + +func New(ctx context.Context, cfg *NewConfig) *Config { + closeTracer := otel.New(ctx) + defer closeTracer() + + database.RunMigrations(&database.MigrationConfig{ + BasePath: cfg.BasePath, + }) + + dbPool := database.GetDBPool(ctx) + dbQueries := repository.New(dbPool) + db := &database.Connection{ + Pool: dbPool, + Queries: dbQueries, + } + valid := validator.New() + + return &Config{ + Database: db, + Validator: valid, + } +} diff --git a/internal/server/server_test.go b/internal/server/server_test.go new file mode 100644 index 00000000..1a6437dc --- /dev/null +++ b/internal/server/server_test.go @@ -0,0 +1,92 @@ +package server_test + +import ( + "context" + "fmt" + "queryorchestration/internal/database" + "queryorchestration/internal/server" + "testing" + + "github.com/docker/go-connections/nat" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/testcontainers/testcontainers-go" + "github.com/testcontainers/testcontainers-go/wait" +) + +func TestNew(t *testing.T) { + ctx := context.Background() + + _, cleanup := createDB(t, ctx) + defer cleanup() + + newCfg := &server.NewConfig{ + BasePath: "../..", + } + + cfg := server.New(ctx, newCfg) + assert.NotNil(t, cfg) +} + +type db struct { + pool *pgxpool.Pool + container *testcontainers.Container +} + +func createDB(t *testing.T, ctx context.Context) (*db, func()) { + port, err := nat.NewPort("tcp", "5432") + if err != nil { + t.Fatalf("Failed to create port: %v", err) + } + + name := "queryorchestration" + pass := "pass" + user := "postgres" + + req := testcontainers.ContainerRequest{ + Image: "postgres:17.2-alpine3.21", + Env: map[string]string{ + "POSTGRES_DB": name, + "POSTGRES_USER": user, + "POSTGRES_PASSWORD": pass, + }, + ExposedPorts: []string{port.Port()}, + WaitingFor: wait.ForListeningPort(port), + } + + container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{ + ContainerRequest: req, + Started: true, + }) + if err != nil { + t.Fatalf("Failed to start container: %v", err) + } + + host, err := container.Host(ctx) + if err != nil { + t.Fatalf("Failed to extract host: %v", err) + } + mappedPort, err := container.MappedPort(ctx, port) + if err != nil { + t.Fatalf("Failed to extract port: %v", err) + } + + t.Setenv("DB_USER", user) + t.Setenv("DB_PASS", pass) + t.Setenv("DB_HOST", host) + t.Setenv("DB_PORT", fmt.Sprint(mappedPort.Int())) + t.Setenv("DB_NAME", name) + t.Setenv("DB_NOSSL", "1") + + pool := database.GetDBPool(ctx) + + return &db{ + pool: pool, + container: &container, + }, func() { + err := container.Terminate(ctx) + if err != nil { + t.Error(err) + } + } +} diff --git a/mocks/repository/mock_DBTX.go b/mocks/repository/mock_DBTX.go index ea1a7911..b94d7e42 100644 --- a/mocks/repository/mock_DBTX.go +++ b/mocks/repository/mock_DBTX.go @@ -1,4 +1,4 @@ -// Code generated by mockery v2.49.1. DO NOT EDIT. +// Code generated by mockery v2.50.0. DO NOT EDIT. package repository diff --git a/scripts/Taskfile.yml b/scripts/Taskfile.yml index 685a991d..70550f98 100644 --- a/scripts/Taskfile.yml +++ b/scripts/Taskfile.yml @@ -1,19 +1,26 @@ +--- # https://taskfile.dev version: '3' includes: deps: + dir: "{{.CONTEXT}}" taskfile: dependencies.yml test: + dir: "{{.CONTEXT}}" taskfile: tests.yml docker: + dir: "{{.CONTEXT}}" taskfile: docker.yml proto: + dir: "{{.CONTEXT}}" taskfile: proto.yml compose: + dir: "{{.CONTEXT}}" taskfile: compose.yml db: + dir: "{{.CONTEXT}}" taskfile: database.yml vars: @@ -39,12 +46,18 @@ tasks: cmds: - task code:lint:fix - task proto:lint:fix + - task db:lint + - task yaml:lint + - task docker:lint + - task compose:lint + - task shell:lint code:lint: cmds: - golangci-lint run code:lint:fix: cmds: - gofmt -w . + - golangci-lint run yaml:lint: cmds: - yamllint . -s @@ -53,4 +66,4 @@ tasks: - shellcheck --shell=sh scripts/install-deps.sh docs: cmds: - - godoc -http=:6060 \ No newline at end of file + - godoc -http=:6060 diff --git a/scripts/compose.yml b/scripts/compose.yml index 10287127..915c1985 100644 --- a/scripts/compose.yml +++ b/scripts/compose.yml @@ -1,9 +1,10 @@ +--- # https://taskfile.dev version: '3' vars: - COMPOSE_FILE: "{{.CONTEXT}}/deployments/compose.yaml" + COMPOSE_FILE: "deployments/compose.yaml" tasks: build: @@ -11,7 +12,13 @@ tasks: - docker compose -f {{.COMPOSE_FILE}} build up: cmds: - - docker compose -f {{.COMPOSE_FILE}} up + - docker compose -f {{.COMPOSE_FILE}} up --no-recreate + up:bg: + cmds: + - docker compose -f {{.COMPOSE_FILE}} up --no-recreate -d + down: + cmds: + - docker compose -f {{.COMPOSE_FILE}} down lint: cmds: - - docker compose -f {{.COMPOSE_FILE}} config \ No newline at end of file + - docker compose -f {{.COMPOSE_FILE}} config diff --git a/scripts/database.yml b/scripts/database.yml index 1353c9b0..548d8eb4 100644 --- a/scripts/database.yml +++ b/scripts/database.yml @@ -1,15 +1,18 @@ +--- # https://taskfile.dev version: '3' vars: - MIGRATIONS: "{{.CONTEXT}}/database/migrations" - DATABASE_URI: "postgres://${DB_USER}:${DB_PASS}@${DB_HOST}:${DB_PORT}/${DB_NAME}?sslmode=disable" + MIGRATIONS: "database/migrations" tasks: generate: cmds: - - sqlc generate --file ../sqlc.yml + - | + task compose:up:bg + task db:mig:run + sqlc generate --file ../sqlc.yml lint: cmds: - sqlc vet --file ../sqlc.yml @@ -21,4 +24,4 @@ tasks: migrate create -ext sql -dir {{.MIGRATIONS}} $name mig:run: cmds: - - migrate -path {{.MIGRATIONS}} -database {{.DATABASE_URI}} up \ No newline at end of file + - migrate -path {{.MIGRATIONS}} -database {{.DB_URI}} up diff --git a/scripts/dependencies.yml b/scripts/dependencies.yml index 34db6d5c..23105643 100644 --- a/scripts/dependencies.yml +++ b/scripts/dependencies.yml @@ -1,3 +1,4 @@ +--- # https://taskfile.dev version: '3' @@ -5,20 +6,22 @@ version: '3' tasks: install: cmds: - - go install golang.org/x/tools/cmd/godoc@latest - - go install github.com/golangci/golangci-lint/cmd/golangci-lint@latest - - go install github.com/yoheimuta/protolint/cmd/protolint@latest - - go install google.golang.org/protobuf/cmd/protoc-gen-go@latest - - go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@latest - - go install github.com/zabio3/godolint@latest - - go install github.com/wasilibs/go-yamllint/cmd/yamllint@latest - - go install github.com/sqlc-dev/sqlc/cmd/sqlc@latest - - go install -tags 'postgres' github.com/golang-migrate/migrate/v4/cmd/migrate@latest - - go install github.com/vektra/mockery/v2@latest + - go install golang.org/x/tools/cmd/godoc@v0.29.0 + - go install github.com/golangci/golangci-lint/cmd/golangci-lint@v1.62.2 + - go install github.com/yoheimuta/protolint/cmd/protolint@v0.52.0 + - go install google.golang.org/protobuf/cmd/protoc-gen-go@v1.36.0 + - go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@v1.5.1 + - go install github.com/zabio3/godolint@v1.0.3 + - go install github.com/wasilibs/go-yamllint/cmd/yamllint@v1.35.1 + - go install github.com/sqlc-dev/sqlc/cmd/sqlc@v1.27.0 + - | + go install -tags 'postgres' \ + github.com/golang-migrate/migrate/v4/cmd/migrate@v4.18.1 + - go install github.com/vektra/mockery/v2@v2.50.0 - curl -sfL https://direnv.net/install.sh | bash - direnv allow - go mod download tidy: cmds: - go mod tidy - - go mod vendor \ No newline at end of file + - go mod vendor diff --git a/scripts/docker.yml b/scripts/docker.yml index ec0caa80..c02e8481 100644 --- a/scripts/docker.yml +++ b/scripts/docker.yml @@ -1,9 +1,10 @@ +--- # https://taskfile.dev version: '3' vars: - DOCKERFILE: "{{.CONTEXT}}/build/Dockerfile" + DOCKERFILE: "build/Dockerfile" tasks: lint: @@ -11,5 +12,4 @@ tasks: - godolint {{.DOCKERFILE}} build: cmds: - - docker build --build-arg TARGETCMD=queryService -t {{.IMAGE_NAME}}_queryservice -f {{.DOCKERFILE}} {{.CONTEXT}} - - docker build --build-arg TARGETCMD=queryRunner -t {{.IMAGE_NAME}}_queryrunner -f {{.DOCKERFILE}} {{.CONTEXT}} \ No newline at end of file + - docker build -t {{.IMAGE_NAME}} -f {{.DOCKERFILE}} . diff --git a/scripts/proto.yml b/scripts/proto.yml index 9644a489..2b3f773e 100644 --- a/scripts/proto.yml +++ b/scripts/proto.yml @@ -1,10 +1,11 @@ +--- # https://taskfile.dev version: '3' vars: - PROTO_DIR: "{{.CONTEXT}}/serviceInterfaces" - API_DIR: "{{.CONTEXT}}/api" + PROTO_DIR: serviceInterfaces + API_DIR: api tasks: lint: @@ -15,4 +16,8 @@ tasks: - protolint lint -fix {{.PROTO_DIR}} generate: cmds: - - protoc --proto_path={{.PROTO_DIR}} --go_out={{.API_DIR}}/serviceInterfaces --go-grpc_out={{.API_DIR}} --go_opt=paths=source_relative --experimental_allow_proto3_optional main.proto + - | + protoc --proto_path={{.PROTO_DIR}} \ + --go_out={{.API_DIR}}/serviceInterfaces --go-grpc_out={{.API_DIR}} \ + --go_opt=paths=source_relative --experimental_allow_proto3_optional \ + main.proto diff --git a/scripts/tests.yml b/scripts/tests.yml index c5bf7e04..e8a31bea 100644 --- a/scripts/tests.yml +++ b/scripts/tests.yml @@ -1,3 +1,4 @@ +--- # https://taskfile.dev version: '3' @@ -8,25 +9,31 @@ includes: internal: true vars: - COVERAGE_FILE: "{{.CONTEXT}}/{{.OUT_DIR}}/coverage.out" + COVERAGE_FILE: "{{.OUT_DIR}}/coverage.out" tasks: + mocks: + cmds: + - mockery unit: vars: - INTERNAL: "{{.CONTEXT}}/internal/..." - API: "{{.CONTEXT}}/api/..." + INTERNAL: "./internal/..." + API: "./api/..." cmds: - - go test {{.INTERNAL}} {{.API}} -coverpkg={{.INTERNAL}},{{.API}} -coverprofile={{.COVERAGE_FILE}} + - | + go test {{.INTERNAL}} {{.API}} \ + -coverpkg={{.INTERNAL}},{{.API}} -coverprofile={{.COVERAGE_FILE}} unit:coverage: deps: - unit vars: - TMP_FILE: "{{.CONTEXT}}/{{.OUT_DIR}}/coverage.tmp" + TMP_FILE: "{{.OUT_DIR}}/coverage.tmp" cmds: - go tool cover -func={{.COVERAGE_FILE}} > {{.TMP_FILE}} - cat {{.TMP_FILE}} - | - COVERAGE=$(grep total: {{.TMP_FILE}} | awk '{print $3}' | sed 's/%//' | bc) + COVERAGE=$(grep total: {{.TMP_FILE}} \ + | awk '{print $3}' | sed 's/%//' | bc) echo "" echo "Coverage Threshold: {{.COVERAGE_THRESHOLD}}%" echo "Total Coverage: $COVERAGE%" @@ -44,7 +51,7 @@ tasks: deps: - docker:build cmds: - - go test -count=1 -v {{.CONTEXT}}/test/... + - go test -count=1 -v test/... integration:nobuild: cmds: - - go test -v {{.CONTEXT}}/test/... \ No newline at end of file + - go test -v test/... diff --git a/sqlc.yml b/sqlc.yml index da0fe787..c3e153e0 100644 --- a/sqlc.yml +++ b/sqlc.yml @@ -1,7 +1,8 @@ +--- version: "2" servers: - engine: postgresql - uri: "postgres://${DB_USER}:${DB_PASS}@${DB_HOST}:${DB_PORT}/${DB_NAME}?sslmode=disable" + uri: "${DB_URI}" sql: - name: "db" engine: "postgresql" @@ -13,4 +14,4 @@ sql: go: package: "repository" out: "internal/database/repository" - sql_package: "pgx/v5" \ No newline at end of file + sql_package: "pgx/v5" diff --git a/test/apiContainer_test.go b/test/apiContainer_test.go index 2b90b273..3da539e9 100644 --- a/test/apiContainer_test.go +++ b/test/apiContainer_test.go @@ -26,17 +26,19 @@ func createAPIContainer(t *testing.T, ctx context.Context, config *apiContainerC } req := testcontainers.ContainerRequest{ - Image: fmt.Sprintf("queryorchestration_%s:latest", config.ServiceName), + Image: "queryorchestration:latest", Env: map[string]string{ - "DB_USER": config.DB.User, - "DB_PASS": config.DB.Password, - "DB_HOST": config.DB.Host, - "DB_NAME": config.DB.Name, - "DB_PORT": strconv.Itoa(config.DB.Port), + "DB_USER": config.DB.User, + "DB_PASS": config.DB.Password, + "DB_HOST": config.DB.Host, + "DB_NAME": config.DB.Name, + "DB_PORT": strconv.Itoa(config.DB.Port), + "DB_NOSSL": "1", }, ExposedPorts: []string{port.Port()}, - WaitingFor: wait.ForListeningPort(port), + WaitingFor: wait.ForLog("Listening on port 8080"), Networks: []string{config.Network.Name}, + Entrypoint: []string{fmt.Sprintf("./bin/%s", config.ServiceName)}, } container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{ @@ -65,7 +67,10 @@ func createAPIContainer(t *testing.T, ctx context.Context, config *apiContainerC } return conn, func() { - container.Terminate(ctx) + err := container.Terminate(ctx) + if err != nil { + t.Error(err) + } } } @@ -82,8 +87,17 @@ func createAPIDependencies(t *testing.T, ctx context.Context, serviceName string return conn, func() { testcontainers.CleanupNetwork(t, network) - dbContainer.Terminate(ctx) + + err := dbContainer.Terminate(ctx) + if err != nil { + t.Error(err) + } + + err = conn.Close() + if err != nil { + t.Error(err) + } + containerCleanup() - conn.Close() } } diff --git a/test/container_test.go b/test/container_test.go index ec008b66..687aec98 100644 --- a/test/container_test.go +++ b/test/container_test.go @@ -43,7 +43,7 @@ func createDB(t *testing.T, ctx context.Context, network *testcontainers.DockerN } req := testcontainers.ContainerRequest{ - Image: "postgres:latest", + Image: "postgres:17.2-alpine3.21", Env: map[string]string{ "POSTGRES_DB": config.Name, "POSTGRES_USER": config.User, diff --git a/test/queryrunner_test.go b/test/queryrunner_test.go index 2337757f..17afad57 100644 --- a/test/queryrunner_test.go +++ b/test/queryrunner_test.go @@ -13,7 +13,7 @@ import ( func TestQueryRunner(t *testing.T) { ctx := context.Background() - queue, cleanup := createQueueDependencies(t, ctx, "queryrunner") + queue, cleanup := createQueueDependencies(t, ctx, "queryRunner") defer cleanup() document := document.Document{ diff --git a/test/queryservice_test.go b/test/queryservice_test.go index 4b2f0808..895d96ca 100644 --- a/test/queryservice_test.go +++ b/test/queryservice_test.go @@ -11,7 +11,7 @@ import ( func TestQueryService(t *testing.T) { ctx := context.Background() - conn, cleanup := createAPIDependencies(t, ctx, "queryservice") + conn, cleanup := createAPIDependencies(t, ctx, "queryService") defer cleanup() client := serviceinterfaces.NewQueryServiceClient(conn) diff --git a/test/queueContainer_test.go b/test/queueContainer_test.go index 9cf6222f..189d8c2f 100644 --- a/test/queueContainer_test.go +++ b/test/queueContainer_test.go @@ -32,7 +32,7 @@ func createQueueContainer(t *testing.T, ctx context.Context, config *queueContai } req := testcontainers.ContainerRequest{ - Image: fmt.Sprintf("queryorchestration_%s:latest", config.ServiceName), + Image: "queryorchestration:latest", Env: map[string]string{ "QUEUE_URL": config.Queue.URL, "AWS_DEFAULT_REGION": config.Queue.Region, @@ -46,8 +46,11 @@ func createQueueContainer(t *testing.T, ctx context.Context, config *queueContai "DB_HOST": config.DB.Host, "DB_NAME": config.DB.Name, "DB_PORT": strconv.Itoa(config.DB.Port), + "DB_NOSSL": "1", }, - Networks: []string{config.Network.Name}, + Networks: []string{config.Network.Name}, + WaitingFor: wait.ForLog("Listening to queue"), + Entrypoint: []string{fmt.Sprintf("./bin/%s", config.ServiceName)}, } container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{ @@ -86,7 +89,7 @@ func createQueue(t *testing.T, ctx context.Context, network *testcontainers.Dock } req := testcontainers.ContainerRequest{ - Image: "localstack/localstack:latest", + Image: "localstack/localstack:4.0.3", Env: map[string]string{ "AWS_ACCESS_KEY_ID": provider.Value.AccessKeyID, "AWS_SECRET_ACCESS_KEY": provider.Value.SecretAccessKey, @@ -210,8 +213,20 @@ func createQueueDependencies(t *testing.T, ctx context.Context, serviceName stri return queue, func() { testcontainers.CleanupNetwork(t, network) - dbContainer.Terminate(ctx) - (*queue.Container).Terminate(ctx) - container.Terminate(ctx) + + err := dbContainer.Terminate(ctx) + if err != nil { + t.Error(err) + } + + err = (*queue.Container).Terminate(ctx) + if err != nil { + t.Error(err) + } + + err = container.Terminate(ctx) + if err != nil { + t.Error(err) + } } }