Files
query-orchestration/internal/test/ecosystem.go
T
Michael McGuinness 1d49313a9f Merged in feature/standardisefilepath (pull request #111)
Feature/standardisefilepath

* baseprocessing

* generalstructure
2025-04-03 12:13:16 +00:00

448 lines
11 KiB
Go

package test
import (
"bufio"
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"log/slog"
"net/http"
"testing"
"time"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/objectstore"
queryapi "queryorchestration/pkg/queryAPI"
"github.com/docker/go-connections/nat"
"github.com/stretchr/testify/require"
"github.com/testcontainers/testcontainers-go"
"github.com/testcontainers/testcontainers-go/wait"
)
type Network struct {
Dependencies Dependencies
APIs map[APIName]*Container
Runners map[RunnerName]*Container
Client *queryapi.ClientWithResponses
}
func CreateFullNetwork(t testing.TB, ctx context.Context, cfg FullDependenciesConfig) (Network, func()) {
deps, clean := CreateFullDependencies(t, ctx, cfg)
apiContainers := make(map[APIName]*Container, len(apis))
apiClean := make([]func(), len(apis))
for i, s := range apis {
c, ccleanup := CreateAPI(t, ctx, cfg, &APIConfig{
API: s,
Network: deps.Network,
MockHTTP: string(deps.MockServer.Internal),
})
apiContainers[s.Name] = c
apiClean[i] = ccleanup
}
qService, err := queryapi.NewClientWithResponses(apiContainers[QueryAPIName].URI)
require.NoError(t, err)
runnerContainers := make(map[RunnerName]*Container, len(runners))
runnerClean := make([]func(), len(runners))
for i, r := range runners {
c, ccleanup := CreateRunner(t, ctx, cfg, &RunnerConfig{
Runner: r,
Network: deps.Network,
MockHTTP: string(deps.MockServer.Internal),
})
runnerContainers[r.Name] = c
runnerClean[i] = ccleanup
}
return Network{
Dependencies: deps,
APIs: apiContainers,
Runners: runnerContainers,
Client: qService,
}, func() {
for _, c := range apiClean {
c()
}
for _, c := range runnerClean {
c()
}
clean()
}
}
type FullDependenciesConfig interface {
serviceconfig.ConfigProvider
objectstore.ConfigProvider
}
type Dependencies struct {
BucketName string
QueueURLs map[RunnerName]string
Network *testcontainers.DockerNetwork
AWSConfig *AWSContainerConfig
DBConfig testcontainers.Container
MockServer MockServer
}
func CreateFullDependencies(t testing.TB, ctx context.Context, cfg FullDependenciesConfig) (Dependencies, func()) {
network, ncleanup := CreateNetwork(t, ctx)
acfg, awsclean := CreateAWSContainer(t, ctx, cfg, &CreateAWSConfig{
Network: network,
})
dbcfg, dbcleanup := CreateDB(t, ctx, cfg, &CreateDatabaseConfig{
Network: network,
RunMigrations: true,
})
mockServer, cleanMock := CreateMockServer(t, ctx, MockServerConfig{
Network: network.Name,
})
SetQueueClient(t, ctx, cfg)
SetStoreClient(t, ctx, cfg, acfg.ExternalEndpoint)
urls := map[RunnerName]string{}
for _, runner := range runners {
urls[runner.Name] = CreateQueue(t, ctx, cfg, string(runner.Name))
}
CreateBucket(t, ctx, cfg, BucketName)
SetBucketNotifs(t, ctx, cfg, BucketName)
return Dependencies{
BucketName: BucketName,
QueueURLs: urls,
Network: network,
AWSConfig: acfg,
DBConfig: dbcfg,
MockServer: mockServer,
}, func() {
awsclean()
dbcleanup()
cleanMock()
ncleanup()
}
}
func SetCfgProvider(t testing.TB, cfg serviceconfig.ConfigProvider) {
cfg.SetAWSKeyID("test")
cfg.SetAWSSecretKey("test")
cfg.SetAWSRegion("us-east-1")
cfg.SetAWSProfile("")
t.Setenv("AWS_ACCESS_KEY_ID", cfg.GetAWSKeyID())
t.Setenv("AWS_SECRET_ACCESS_KEY", cfg.GetAWSSecretKey())
t.Setenv("AWS_SESSION_TOKEN", cfg.GetAWSSessionToken())
t.Setenv("AWS_REGION", cfg.GetAWSRegion())
t.Setenv("AWS_DEFAULT_REGION", cfg.GetAWSRegion())
t.Setenv("AWS_PROFILE", string(cfg.GetAWSProfile()))
cfg.SetDBUser("invalid_user")
cfg.SetDBSecret("invalid_pass")
cfg.SetDBHost("invalid_host")
cfg.SetDBPort(5432)
cfg.SetDBName("invalid_name")
cfg.SetDBNoSSL(true)
}
type ServiceNetworkConfig struct {
Network *testcontainers.DockerNetwork
API API
}
func CreateAPINetwork(t testing.TB, ctx context.Context, cfg FullDependenciesConfig, scfg *ServiceNetworkConfig) (*Container, func()) {
deps, depsclean := CreateFullDependencies(t, ctx, cfg)
c, ccleanup := CreateAPI(t, ctx, cfg, &APIConfig{
API: scfg.API,
Network: deps.Network,
MockHTTP: string(deps.MockServer.Internal),
})
return c, func() {
depsclean()
ccleanup()
}
}
type Address string
type MockServer struct {
Internal Address
External Address
Container testcontainers.Container
Client *http.Client
}
type MockBody any
type MockQueries map[string][]string
type MockHeaders map[string][]string
type MockRequest struct {
Method string `json:"method"`
Path string `json:"path"`
Headers MockHeaders `json:"headers"`
Body MockBody `json:"body"`
Query MockQueries `json:"queryStringParameters"`
}
type MockResponse struct {
Code int `json:"statusCode"`
Headers MockHeaders `json:"headers"`
Body MockBody `json:"body"`
}
type MockExpectation struct {
Request MockRequest `json:"httpRequest"`
Response MockResponse `json:"httpResponse"`
}
type MockServerConfig struct {
Network string
}
func CreateMockServer(t testing.TB, ctx context.Context, cfg MockServerConfig) (MockServer, func()) {
name := "mockserver"
port, err := nat.NewPort("tcp", "1080")
require.NoError(t, err)
req := testcontainers.ContainerRequest{
Image: "mockserver/mockserver:latest",
ExposedPorts: []string{port.Port()},
Networks: []string{cfg.Network},
NetworkAliases: map[string][]string{
cfg.Network: {name},
},
Env: map[string]string{
"MOCKSERVER_LOG_LEVEL": "INFO",
},
WaitingFor: wait.ForAll(
wait.ForExposedPort(),
wait.ForListeningPort(port),
),
}
container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{
ContainerRequest: req,
Started: true,
})
require.NoError(t, err)
time.Sleep(2 * time.Second)
host, err := container.Host(ctx)
require.NoError(t, err)
externalPort, err := container.MappedPort(ctx, port)
require.NoError(t, err)
server := MockServer{
Client: &http.Client{},
Internal: Address(fmt.Sprintf("http://%s:%d", name, port.Int())),
External: Address(fmt.Sprintf("http://%s:%d", host, externalPort.Int())),
Container: container,
}
return server, func() {
err := container.Terminate(ctx)
require.NoError(t, err)
}
}
func CreateMockExpectation(t testing.TB, server MockServer, expectation MockExpectation) {
jsonData, err := json.Marshal(expectation)
require.NoError(t, err)
url := fmt.Sprintf("%s/mockserver/expectation", server.External)
req, err := http.NewRequest("PUT", url, bytes.NewBuffer(jsonData))
require.NoError(t, err)
req.Header.Set("Content-Type", "application/json")
resp, err := server.Client.Do(req)
require.NoError(t, err)
defer resp.Body.Close()
if resp.StatusCode != http.StatusCreated && resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
t.Fatalf("failed to configure MockServer, status: %d, response: %s", resp.StatusCode, body)
}
}
type StartDocTextDetectionExpectationParams struct {
Bucket string
JobID string
Key objectstore.BucketKey
}
func CreateStartDocTextDetectionExpectation(t testing.TB, mockServer MockServer, params StartDocTextDetectionExpectationParams) MockExpectation {
expectation := MockExpectation{
Request: MockRequest{
Method: "POST",
Path: "/",
Headers: MockHeaders{
"X-Amz-Target": []string{"Textract.StartDocumentTextDetection"},
},
Body: map[string]map[string]any{
"DocumentLocation": {
"S3Object": map[string]any{
"Bucket": params.Bucket,
"Name": params.Key.String(),
},
},
"OutputConfig": {
"S3Bucket": params.Bucket,
},
},
Query: MockQueries{},
},
Response: MockResponse{
Code: 200,
Headers: MockHeaders{
"Content-Type": {"application/json"},
},
Body: map[string]string{
"JobId": params.JobID,
},
},
}
CreateMockExpectation(t, mockServer, expectation)
return expectation
}
func WaitForMockEndpoint(t testing.TB, server MockServer, request MockRequest) MockRequest {
t.Helper()
verificationRequest := map[string]any{
"httpRequest": request,
"times": map[string]any{
"atLeast": 1,
},
}
jsonData, err := json.Marshal(verificationRequest)
require.NoError(t, err)
verifyURL := fmt.Sprintf("%s/mockserver/verify", server.External)
client := &http.Client{
Timeout: 5 * time.Second,
}
timeout := time.After(60 * time.Second)
ticker := time.NewTicker(500 * time.Millisecond)
defer ticker.Stop()
slog.Info("Attempting to process request", "body", jsonData)
for {
select {
case <-timeout:
require.NoError(t, errors.New("Timeout waiting for mock http request to be fulfilled"))
case <-ticker.C:
req, err := http.NewRequest("PUT", verifyURL, bytes.NewBuffer(jsonData))
require.NoError(t, err)
req.Header.Set("Content-Type", "application/json")
resp, err := client.Do(req)
require.NoError(t, err)
scanner := bufio.NewScanner(resp.Body)
for scanner.Scan() {
line := scanner.Text()
slog.Info(line)
}
if resp.StatusCode == http.StatusAccepted {
slog.Info("request found")
resp.Body.Close()
return retrieveMatchingRequest(t, server, request)
}
resp.Body.Close()
slog.Error("no request found")
}
}
}
func retrieveMatchingRequest(t testing.TB, server MockServer, request MockRequest) MockRequest {
retrieveURL := fmt.Sprintf("%s/mockserver/retrieve?type=REQUESTS", server.External)
req, err := http.NewRequest("PUT", retrieveURL, nil)
require.NoError(t, err)
client := &http.Client{}
resp, err := client.Do(req)
require.NoError(t, err)
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
require.NoError(t, err)
var requests []MockRequest
err = json.Unmarshal(body, &requests)
require.NoError(t, err)
for i := len(requests) - 1; i >= 0; i-- {
attemptRequest := requests[i]
if attemptRequest.Method == request.Method && attemptRequest.Path == request.Path {
return attemptRequest
}
}
require.Fail(t, "no request found")
return MockRequest{}
}
type TextractOutputParams struct {
Request MockRequest
File io.Reader
}
func WaitForMockTextractOutput(t testing.TB, ctx context.Context, cfg objectstore.ConfigProvider, mockServer MockServer, params TextractOutputParams) {
response := WaitForMockEndpoint(t, mockServer, params.Request)
body, ok := response.Body.(map[string]any)
if !ok {
bodyStr, ok := response.Body.(string)
require.True(t, ok)
err := json.Unmarshal([]byte(bodyStr), &body)
require.NoError(t, err)
}
json, ok := body["json"].(map[string]any)
if !ok {
json = body
}
outputConfig, ok := json["OutputConfig"].(map[string]any)
require.True(t, ok)
bucket, ok := outputConfig["S3Bucket"].(string)
require.True(t, ok)
keyStr, ok := outputConfig["S3Prefix"].(string)
require.True(t, ok)
key, err := objectstore.ParseBucketKey(keyStr)
require.NoError(t, err)
require.Equal(t, objectstore.TextTextract, key.Location)
PutObject(t, ctx, cfg, PutObjectParams{
Key: key,
Bucket: bucket,
File: params.File,
})
}