Merged in feature/postprocessing (pull request #114)

Feature/postprocessing

* tests

* passtest

* fixshorttests

* mosttests

* improvingbasedockerfile

* testspeeds

* testing

* host

* canparallel

* clean

* passfullsuite

* singlepagemax

* test

* findfeatures

* findstables

* tbls

* tablestoo

* tablestoo

* lateraltests

* tableloc

* cleanup

* inlinetable

* childids

* cleanup

* tests
This commit is contained in:
Michael McGuinness
2025-04-22 14:40:16 +00:00
parent edc50c7510
commit fee71e7740
404 changed files with 211048 additions and 2820 deletions
+2 -5
View File
@@ -10,7 +10,6 @@ import (
"github.com/docker/go-connections/nat"
"github.com/stretchr/testify/require"
"github.com/testcontainers/testcontainers-go"
)
type APIName string
@@ -26,19 +25,17 @@ type API struct {
type APIConfig struct {
API API
Network *testcontainers.DockerNetwork
MockHTTP string
}
func CreateAPI(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, config *APIConfig) (*Container, func()) {
func CreateAPI(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, network string, config *APIConfig) (*Container, func()) {
port, err := nat.NewPort("tcp", "8080")
require.NoError(t, err)
container, cleanup := createContainer(t, ctx, &containerConfig{
container, cleanup := createContainer(t, ctx, network, &containerConfig{
Cfg: cfg,
Name: string(config.API.Name),
DownstreamQueues: config.API.DownstreamQueues,
Network: config.Network,
MockHTTP: config.MockHTTP,
WaitForMsg: "⇨ http server started on [::]:8080",
})
+5 -12
View File
@@ -1,7 +1,6 @@
package test_test
import (
"context"
"testing"
"queryorchestration/internal/serviceconfig"
@@ -14,24 +13,18 @@ func TestCreateAPI(t *testing.T) {
if testing.Short() {
t.Skip("Skipping long test in short mode")
}
ctx := context.Background()
ncfg, ncleanup := test.CreateNetwork(t, ctx)
defer ncleanup()
ctx := t.Context()
cfg := &serviceconfig.BaseConfig{}
test.SetCfgProvider(t, cfg)
_, dbcleanup := test.CreateDB(t, ctx, cfg, &test.CreateDatabaseConfig{
Network: ncfg,
})
defer dbcleanup()
net := test.DepNetwork.Get(t, ctx)
_ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{})
acfg := &test.APIConfig{
API: test.QueryAPI,
Network: ncfg,
API: test.QueryAPI,
}
conn, cleanup := test.CreateAPI(t, ctx, cfg, acfg)
conn, cleanup := test.CreateAPI(t, ctx, cfg, net, acfg)
assert.NotNil(t, conn)
assert.NotNil(t, cleanup)
+36
View File
@@ -0,0 +1,36 @@
package test
import (
"encoding/json"
"io"
"os"
"testing"
"github.com/aws/aws-sdk-go-v2/service/textract"
"github.com/google/go-cmp/cmp"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func GetTextractFileResponse(t testing.TB, filename string) map[string][]textract.DetectDocumentTextOutput {
file, err := os.Open(filename)
require.NoError(t, err)
defer file.Close()
decoder := json.NewDecoder(file)
var out map[string][]textract.DetectDocumentTextOutput
err = decoder.Decode(&out)
require.NoError(t, err)
return out
}
func AssertStringToReader(t testing.TB, expected string, actual io.Reader) {
actualBytes, err := io.ReadAll(actual)
require.NoError(t, err)
actualStr := string(actualBytes)
if !assert.Equal(t, expected, actualStr) {
diff := cmp.Diff(expected, actualStr)
t.Logf("Detailed diff (-expected +actual):\n%s", diff)
}
}
+29
View File
@@ -0,0 +1,29 @@
package test_test
import (
"strings"
"testing"
"queryorchestration/internal/test"
"github.com/aws/aws-sdk-go-v2/service/textract"
"github.com/stretchr/testify/assert"
)
func TestAssertStringToReader(t *testing.T) {
t.Run("pass", func(t *testing.T) {
mockT := &testing.T{}
test.AssertStringToReader(t, "hi", strings.NewReader("hi"))
assert.False(t, mockT.Failed())
})
t.Run("fail", func(t *testing.T) {
mockT := &testing.T{}
test.AssertStringToReader(mockT, "bye", strings.NewReader("hi"))
assert.True(t, mockT.Failed())
})
}
func TestGetTextractFileResponse(t *testing.T) {
out := test.GetTextractFileResponse(t, "../../assets/sampleGeneration/generated/helloWorld.gen")
assert.IsType(t, map[string][]textract.DetectDocumentTextOutput{}, out)
}
+16 -27
View File
@@ -5,6 +5,7 @@ import (
"fmt"
"io"
"net/http"
"strconv"
"testing"
"time"
@@ -20,11 +21,6 @@ import (
type AWSContainerConfig struct {
Container testcontainers.Container
ExternalEndpoint string
NetworkEndpoint string
}
type CreateAWSConfig struct {
Network *testcontainers.DockerNetwork
}
type AWSConfigProvider interface {
@@ -32,10 +28,15 @@ type AWSConfigProvider interface {
objectstore.ConfigProvider
}
func CreateAWSContainer(t testing.TB, ctx context.Context, cfg AWSConfigProvider, acfg *CreateAWSConfig) (*AWSContainerConfig, func()) {
alias := "localstack"
const (
awsAlias = "localstack"
awsPort = 4566
)
port, err := nat.NewPort("tcp", "4566")
func CreateAWSContainer(t testing.TB, ctx context.Context, cfg AWSConfigProvider, network string) (*AWSContainerConfig, func()) {
alias := NormaliseAlias(fmt.Sprintf("%s_%s", awsAlias, t.Name()))
port, err := nat.NewPort("tcp", strconv.FormatInt(int64(awsPort), 10))
require.NoError(t, err)
req := testcontainers.ContainerRequest{
@@ -69,13 +70,10 @@ func CreateAWSContainer(t testing.TB, ctx context.Context, cfg AWSConfigProvider
return statusCode == http.StatusOK
}),
),
}
if acfg.Network != nil {
req.Networks = []string{acfg.Network.Name}
req.NetworkAliases = map[string][]string{
acfg.Network.Name: {alias},
}
Networks: []string{network},
NetworkAliases: map[string][]string{
network: {alias},
},
}
container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{
@@ -90,25 +88,16 @@ func CreateAWSContainer(t testing.TB, ctx context.Context, cfg AWSConfigProvider
require.NoError(t, err)
extEndpoint := fmt.Sprintf("http://%s:%s", host, mappedPort.Port())
t.Setenv("AWS_ENDPOINT_URL", extEndpoint)
cfg.SetAWSEndpoint(extEndpoint)
t.Setenv("AWS_ENDPOINT_URL_SQS", extEndpoint)
cfg.SetSQSEndpoint(extEndpoint)
t.Setenv("AWS_ENDPOINT_URL_S3", extEndpoint)
cfg.SetS3Endpoint(extEndpoint)
var endpoint string
if acfg.Network != nil {
endpoint = fmt.Sprintf("http://%s:%s", alias, port.Port())
cfg.SetAWSEndpoint(endpoint)
cfg.SetSQSEndpoint(endpoint)
cfg.SetS3Endpoint(endpoint)
}
t.Setenv("AWS_ENDPOINT_URL", extEndpoint)
t.Setenv("AWS_ENDPOINT_URL_SQS", extEndpoint)
t.Setenv("AWS_ENDPOINT_URL_S3", extEndpoint)
return &AWSContainerConfig{
Container: container,
ExternalEndpoint: extEndpoint,
NetworkEndpoint: endpoint,
}, func() {
err := container.Terminate(ctx)
if err != nil {
+3 -1
View File
@@ -23,7 +23,9 @@ func TestCreateQueueContainer(t *testing.T) {
cfg := &TestAWSConfig{}
SetCfgProvider(t, cfg)
qcfg, cleanup := CreateAWSContainer(t, ctx, cfg, &CreateAWSConfig{})
net := DepNetwork.Get(t, ctx)
qcfg, cleanup := CreateAWSContainer(t, ctx, cfg, net)
assert.NotNil(t, qcfg)
assert.NotNil(t, cleanup)
+9 -8
View File
@@ -24,7 +24,6 @@ type Container struct {
type containerConfig struct {
Name string
Cfg serviceconfig.ConfigProvider
Network *testcontainers.DockerNetwork
DownstreamQueues []RunnerName
Env map[string]string
WaitForMsg string
@@ -32,21 +31,23 @@ type containerConfig struct {
MockHTTP string
}
func createContainer(t testing.TB, ctx context.Context, cfg *containerConfig) (testcontainers.Container, func()) {
func createContainer(t testing.TB, ctx context.Context, network string, cfg *containerConfig) (testcontainers.Container, func()) {
awsAlias := NormaliseAlias(fmt.Sprintf("%s_%s", awsAlias, t.Name()))
awsEndpoint := fmt.Sprintf("http://%s:%d", awsAlias, awsPort)
env := map[string]string{
"PGUSER": cfg.Cfg.GetDBUser(),
"PGPASSWORD": cfg.Cfg.GetDBSecret(),
"PGHOST": cfg.Cfg.GetDBHost(),
"PGHOST": dbAlias,
"PGDATABASE": cfg.Cfg.GetDBName(),
"PGPORT": strconv.Itoa(cfg.Cfg.GetDBPort()),
"PGPORT": strconv.Itoa(dbPort),
"DB_NOSSL": strconv.FormatBool(cfg.Cfg.IsDBNoSSL()),
"AWS_ACCESS_KEY_ID": cfg.Cfg.GetAWSKeyID(),
"AWS_SECRET_ACCESS_KEY": cfg.Cfg.GetAWSSecretKey(),
"AWS_REGION": cfg.Cfg.GetAWSRegion(),
"AWS_DEFAULT_REGION": cfg.Cfg.GetAWSRegion(),
"AWS_ENDPOINT_URL": cfg.Cfg.GetAWSEndpoint(),
"AWS_ENDPOINT_URL_SQS": cfg.Cfg.GetAWSEndpoint(),
"AWS_ENDPOINT_URL_S3": cfg.Cfg.GetAWSEndpoint(),
"AWS_ENDPOINT_URL": awsEndpoint,
"AWS_ENDPOINT_URL_SQS": awsEndpoint,
"AWS_ENDPOINT_URL_S3": awsEndpoint,
"AWS_ENDPOINT_URL_TEXTRACT": cfg.MockHTTP,
"AWS_S3_USE_PATH_STYLE": strconv.FormatBool(true),
"LOG_LEVEL": "DEBUG",
@@ -64,9 +65,9 @@ func createContainer(t testing.TB, ctx context.Context, cfg *containerConfig) (t
req := testcontainers.ContainerRequest{
Image: "queryorchestration:latest",
Env: env,
Networks: []string{cfg.Network.Name},
WaitingFor: wait.ForLog(cfg.WaitForMsg),
Entrypoint: []string{fmt.Sprintf("./%s", cfg.Name)},
Networks: []string{network},
}
if len(cfg.ExposedPorts) > 0 {
+4 -9
View File
@@ -16,28 +16,23 @@ func TestCreateContainer(t *testing.T) {
}
ctx := context.Background()
ncfg, ncleanup := CreateNetwork(t, ctx)
defer ncleanup()
cfg := &serviceconfig.BaseConfig{}
SetCfgProvider(t, cfg)
_, dbcleanup := CreateDB(t, ctx, cfg, &CreateDatabaseConfig{
Network: ncfg,
})
defer dbcleanup()
net := DepNetwork.Get(t, ctx)
_ = CreateDB(t, ctx, cfg, net, &CreateDatabaseConfig{})
ccfg := &containerConfig{
Name: string(QueryAPIName),
DownstreamQueues: QueryAPI.DownstreamQueues,
Cfg: cfg,
Network: ncfg,
ExposedPorts: []nat.Port{
nat.Port("8080/tcp"),
},
}
container, cleanup := createContainer(t, ctx, ccfg)
container, cleanup := createContainer(t, ctx, net, ccfg)
assert.NotNil(t, container)
assert.NotNil(t, cleanup)
+23 -53
View File
@@ -3,12 +3,12 @@ package test
import (
"context"
"fmt"
"strconv"
"testing"
"time"
db "queryorchestration/internal/database"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/database"
"github.com/docker/go-connections/nat"
"github.com/stretchr/testify/require"
@@ -17,24 +17,26 @@ import (
)
type CreateDatabaseConfig struct {
Network *testcontainers.DockerNetwork
RunMigrations bool
}
func CreateDB(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, dcfg *CreateDatabaseConfig) (testcontainers.Container, func()) {
alias := "postgres"
const (
dbAlias = "postgres"
dbPort = 5432
)
port, err := nat.NewPort("tcp", "5432")
func CreateDB(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, network string, dcfg *CreateDatabaseConfig) testcontainers.Container {
port, err := nat.NewPort("tcp", strconv.Itoa(dbPort))
require.NoError(t, err)
name := "queryorchestration"
name := NormaliseAlias(fmt.Sprintf("queryorchestration_%s", t.Name()))
pass := "pass"
user := "postgres"
req := testcontainers.ContainerRequest{
Image: "postgres:17.2-alpine3.21",
Name: "postgres_test_queryorchestration",
Env: map[string]string{
"POSTGRES_DB": name,
"POSTGRES_USER": user,
"POSTGRES_PASSWORD": pass,
},
@@ -44,24 +46,22 @@ func CreateDB(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvide
wait.ForListeningPort(port),
wait.ForLog("database system is ready to accept connections"),
wait.ForSQL(port, "postgres", func(host string, port nat.Port) string {
return fmt.Sprintf("postgres://%s:%s@localhost:%s/%s?sslmode=disable",
user, pass, port.Port(), name)
return fmt.Sprintf("postgres://%s:%s@%s:%d/postgres?sslmode=disable",
user, pass, host, port.Int())
}).
WithStartupTimeout(25*time.Second).
WithPollInterval(10*time.Second),
),
}
if dcfg.Network != nil {
req.Networks = []string{dcfg.Network.Name}
req.NetworkAliases = map[string][]string{
dcfg.Network.Name: {alias},
}
Networks: []string{network},
NetworkAliases: map[string][]string{
network: {dbAlias},
},
}
container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{
ContainerRequest: req,
Started: true,
Reuse: true,
})
require.NoError(t, err)
@@ -70,28 +70,12 @@ func CreateDB(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvide
mappedPort, err := container.MappedPort(ctx, port)
require.NoError(t, err)
dbcfg := &database.DBConfig{
DBName: name,
DBSecret: pass,
DBUser: user,
DBPort: mappedPort.Int(),
DBHost: host,
DBNoSSL: true,
}
t.Setenv("PGUSER", dbcfg.DBUser)
t.Setenv("PGPASSWORD", dbcfg.DBSecret)
t.Setenv("PGDATABASE", dbcfg.DBName)
t.Setenv("DB_NOSSL", fmt.Sprintf("%v", dbcfg.DBNoSSL))
t.Setenv("PGHOST", dbcfg.DBHost)
t.Setenv("PGPORT", fmt.Sprint(dbcfg.DBPort))
cfg.SetDBUser(dbcfg.DBUser)
cfg.SetDBSecret(dbcfg.DBSecret)
cfg.SetDBName(dbcfg.DBName)
cfg.SetDBNoSSL(dbcfg.DBNoSSL)
cfg.SetDBHost(dbcfg.DBHost)
cfg.SetDBPort(dbcfg.DBPort)
cfg.SetDBUser(user)
cfg.SetDBSecret(pass)
cfg.SetDBName(name)
cfg.SetDBNoSSL(true)
cfg.SetDBHost(host)
cfg.SetDBPort(mappedPort.Int())
if dcfg.RunMigrations {
err := db.RunMigrations(ctx, cfg)
@@ -101,19 +85,5 @@ func CreateDB(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvide
require.NoError(t, err)
}
if dcfg.Network != nil {
dbcfg.DBHost = alias
dbcfg.DBPort = port.Int()
cfg.SetDBUser(dbcfg.DBUser)
cfg.SetDBSecret(dbcfg.DBSecret)
cfg.SetDBName(dbcfg.DBName)
cfg.SetDBNoSSL(dbcfg.DBNoSSL)
cfg.SetDBHost(dbcfg.DBHost)
cfg.SetDBPort(dbcfg.DBPort)
}
return container, func() {
err := container.Terminate(ctx)
require.NoError(t, err)
}
return container
}
+5 -9
View File
@@ -17,14 +17,12 @@ func TestCreateDB(t *testing.T) {
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
test.SetCfgProvider(t, cfg)
dbcfg, cleanup := test.CreateDB(t, ctx, cfg, &test.CreateDatabaseConfig{})
net := test.DepNetwork.Get(t, ctx)
dbcfg := test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{})
assert.NotNil(t, dbcfg)
assert.Nil(t, cfg.GetDBPool())
assert.NotNil(t, cleanup)
cleanup()
}
func TestCreateDBWithMigrations(t *testing.T) {
@@ -34,14 +32,12 @@ func TestCreateDBWithMigrations(t *testing.T) {
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
test.SetCfgProvider(t, cfg)
dbcfg, cleanup := test.CreateDB(t, ctx, cfg, &test.CreateDatabaseConfig{
net := test.DepNetwork.Get(t, ctx)
dbcfg := test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{
RunMigrations: true,
})
defer cleanup()
assert.NotNil(t, dbcfg)
assert.NotNil(t, cfg.GetDBPool())
assert.NotNil(t, cleanup)
}
+90 -113
View File
@@ -10,14 +10,17 @@ import (
"io"
"log/slog"
"net/http"
"strings"
"testing"
"time"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/database"
"queryorchestration/internal/serviceconfig/objectstore"
queryapi "queryorchestration/pkg/queryAPI"
"github.com/docker/go-connections/nat"
"github.com/google/uuid"
"github.com/stretchr/testify/require"
"github.com/testcontainers/testcontainers-go"
"github.com/testcontainers/testcontainers-go/wait"
@@ -36,9 +39,8 @@ func CreateFullNetwork(t testing.TB, ctx context.Context, cfg FullDependenciesCo
apiContainers := make(map[APIName]*Container, len(apis))
apiClean := make([]func(), len(apis))
for i, s := range apis {
c, ccleanup := CreateAPI(t, ctx, cfg, &APIConfig{
c, ccleanup := CreateAPI(t, ctx, cfg, deps.Network, &APIConfig{
API: s,
Network: deps.Network,
MockHTTP: string(deps.MockServer.Internal),
})
apiContainers[s.Name] = c
@@ -51,9 +53,8 @@ func CreateFullNetwork(t testing.TB, ctx context.Context, cfg FullDependenciesCo
runnerContainers := make(map[RunnerName]*Container, len(runners))
runnerClean := make([]func(), len(runners))
for i, r := range runners {
c, ccleanup := CreateRunner(t, ctx, cfg, &RunnerConfig{
c, ccleanup := CreateRunner(t, ctx, cfg, deps.Network, &RunnerConfig{
Runner: r,
Network: deps.Network,
MockHTTP: string(deps.MockServer.Internal),
})
@@ -85,52 +86,51 @@ type FullDependenciesConfig interface {
type Dependencies struct {
BucketName string
QueueURLs map[RunnerName]string
Network *testcontainers.DockerNetwork
AWSConfig *AWSContainerConfig
DBConfig testcontainers.Container
MockServer MockServer
Network string
}
func CreateFullDependencies(t testing.TB, ctx context.Context, cfg FullDependenciesConfig) (Dependencies, func()) {
network, ncleanup := CreateNetwork(t, ctx)
network := DepNetwork.Get(t, ctx)
acfg, awsclean := CreateAWSContainer(t, ctx, cfg, &CreateAWSConfig{
Network: network,
})
deps := Dependencies{
BucketName: BucketName,
Network: network,
QueueURLs: map[RunnerName]string{},
}
cleanups := []func(){}
dbcfg, dbcleanup := CreateDB(t, ctx, cfg, &CreateDatabaseConfig{
Network: network,
mockServer, cleanMock := CreateMockServer(t, ctx, network)
deps.MockServer = mockServer
cleanups = append(cleanups, cleanMock)
dbcfg := CreateDB(t, ctx, cfg, network, &CreateDatabaseConfig{
RunMigrations: true,
})
deps.DBConfig = dbcfg
mockServer, cleanMock := CreateMockServer(t, ctx, MockServerConfig{
Network: network.Name,
})
acfg, awsclean := CreateAWSContainer(t, ctx, cfg, network)
deps.AWSConfig = acfg
cleanups = append(cleanups, awsclean)
SetQueueClient(t, ctx, cfg)
SetStoreClient(t, ctx, cfg, acfg.ExternalEndpoint)
SetStoreClient(t, ctx, cfg, deps.AWSConfig.ExternalEndpoint)
CreateBucket(t, ctx, cfg, deps.BucketName)
urls := map[RunnerName]string{}
for _, runner := range runners {
urls[runner.Name] = CreateQueue(t, ctx, cfg, string(runner.Name))
deps.QueueURLs[runner.Name] = CreateQueue(t, ctx, cfg, string(runner.Name))
}
CreateBucket(t, ctx, cfg, BucketName)
SetBucketNotifs(t, ctx, cfg, BucketName)
SetBucketNotifs(t, ctx, cfg, deps.BucketName)
SetDBCfg(t, cfg)
return Dependencies{
BucketName: BucketName,
QueueURLs: urls,
Network: network,
AWSConfig: acfg,
DBConfig: dbcfg,
MockServer: mockServer,
}, func() {
awsclean()
dbcleanup()
cleanMock()
ncleanup()
return deps, func() {
for _, c := range cleanups {
c()
}
}
}
func SetCfgProvider(t testing.TB, cfg serviceconfig.ConfigProvider) {
@@ -154,24 +154,35 @@ func SetCfgProvider(t testing.TB, cfg serviceconfig.ConfigProvider) {
cfg.SetDBNoSSL(true)
}
type ServiceNetworkConfig struct {
Network *testcontainers.DockerNetwork
API API
func SetDBCfg(t testing.TB, cfg database.ConfigProvider) {
t.Setenv("PGUSER", cfg.GetDBUser())
t.Setenv("PGPASSWORD", cfg.GetDBSecret())
t.Setenv("PGDATABASE", cfg.GetDBName())
t.Setenv("DB_NOSSL", "true")
t.Setenv("PGHOST", cfg.GetDBHost())
t.Setenv("PGPORT", fmt.Sprintf("%d", cfg.GetDBPort()))
}
func CreateAPINetwork(t testing.TB, ctx context.Context, cfg FullDependenciesConfig, scfg *ServiceNetworkConfig) (*Container, func()) {
type APINetwork struct {
Dependencies Dependencies
API *Container
}
func CreateAPINetwork(t testing.TB, ctx context.Context, cfg FullDependenciesConfig, api API) (*APINetwork, func()) {
deps, depsclean := CreateFullDependencies(t, ctx, cfg)
c, ccleanup := CreateAPI(t, ctx, cfg, &APIConfig{
API: scfg.API,
Network: deps.Network,
c, ccleanup := CreateAPI(t, ctx, cfg, deps.Network, &APIConfig{
API: api,
MockHTTP: string(deps.MockServer.Internal),
})
return c, func() {
depsclean()
ccleanup()
}
return &APINetwork{
Dependencies: deps,
API: c,
}, func() {
depsclean()
ccleanup()
}
}
type Address string
@@ -206,22 +217,14 @@ type MockExpectation struct {
Response MockResponse `json:"httpResponse"`
}
type MockServerConfig struct {
Network string
}
func CreateMockServer(t testing.TB, ctx context.Context, cfg MockServerConfig) (MockServer, func()) {
name := "mockserver"
func CreateMockServer(t testing.TB, ctx context.Context, network string) (MockServer, func()) {
name := NormaliseAlias(fmt.Sprintf("mockserver_%s", t.Name()))
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",
},
@@ -229,6 +232,10 @@ func CreateMockServer(t testing.TB, ctx context.Context, cfg MockServerConfig) (
wait.ForExposedPort(),
wait.ForListeningPort(port),
),
Networks: []string{network},
NetworkAliases: map[string][]string{
network: {name},
},
}
container, err := testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{
@@ -258,6 +265,13 @@ func CreateMockServer(t testing.TB, ctx context.Context, cfg MockServerConfig) (
}
}
func NormaliseAlias(fullName string) string {
name := strings.ToLower(fullName)
name = strings.ReplaceAll(name, "/", "_")
name = strings.ReplaceAll(name, " ", "_")
return name
}
func CreateMockExpectation(t testing.TB, server MockServer, expectation MockExpectation) {
jsonData, err := json.Marshal(expectation)
require.NoError(t, err)
@@ -277,30 +291,18 @@ func CreateMockExpectation(t testing.TB, server MockServer, expectation MockExpe
}
}
type StartDocTextDetectionExpectationParams struct {
Bucket string
JobID string
Key objectstore.BucketKey
}
func CreateStartDocTextDetectionExpectation(t testing.TB, mockServer MockServer, params StartDocTextDetectionExpectationParams) MockExpectation {
func CreateDetectDocumentTextExpectation(t testing.TB, mockServer MockServer, body string) MockExpectation {
childId := uuid.NewString()
expectation := MockExpectation{
Request: MockRequest{
Method: "POST",
Path: "/",
Headers: MockHeaders{
"X-Amz-Target": []string{"Textract.StartDocumentTextDetection"},
"X-Amz-Target": []string{"Textract.AnalyzeDocument"},
},
Body: map[string]map[string]any{
"DocumentLocation": {
"S3Object": map[string]any{
"Bucket": params.Bucket,
"Name": params.Key.String(),
},
},
"OutputConfig": {
"S3Bucket": params.Bucket,
},
Body: map[string]interface{}{
"Document": map[string]interface{}{},
"FeatureTypes": []string{"LAYOUT", "SIGNATURES"},
},
Query: MockQueries{},
},
@@ -309,8 +311,24 @@ func CreateStartDocTextDetectionExpectation(t testing.TB, mockServer MockServer,
Headers: MockHeaders{
"Content-Type": {"application/json"},
},
Body: map[string]string{
"JobId": params.JobID,
Body: map[string]interface{}{
"Blocks": []map[string]interface{}{
{
"BlockType": "PAGE",
"Relationships": []map[string]interface{}{
{
"Type": "CHILD",
"Ids": []string{
childId,
},
},
},
},
{
"Id": childId,
"Text": body,
},
},
},
},
}
@@ -404,44 +422,3 @@ func retrieveMatchingRequest(t testing.TB, server MockServer, request MockReques
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,
})
}
+67 -54
View File
@@ -2,16 +2,16 @@ package test
import (
"context"
"fmt"
"io"
"net/http"
"net/http/httptest"
"os"
"strings"
"testing"
"time"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/aws"
"queryorchestration/internal/serviceconfig/database"
"queryorchestration/internal/serviceconfig/objectstore"
"github.com/stretchr/testify/assert"
@@ -27,9 +27,7 @@ func TestCreateAPINetwork(t *testing.T) {
cfg := &FullDepsConfig{}
SetCfgProvider(t, cfg)
conn, cleanup := CreateAPINetwork(t, ctx, cfg, &ServiceNetworkConfig{
API: QueryAPI,
})
conn, cleanup := CreateAPINetwork(t, ctx, cfg, QueryAPI)
assert.NotNil(t, conn)
assert.NotNil(t, cleanup)
@@ -52,6 +50,24 @@ func TestCreateBaseConfig(t *testing.T) {
assert.True(t, cfg.DBNoSSL)
}
func TestCreateDBConfig(t *testing.T) {
cfg := &database.DBConfig{
DBUser: "test user",
DBSecret: "test secret",
DBHost: "test host",
DBPort: 999,
DBName: "test name",
DBNoSSL: true,
}
SetDBCfg(t, cfg)
assert.Equal(t, "test user", os.Getenv("PGUSER"))
assert.Equal(t, "test secret", os.Getenv("PGPASSWORD"))
assert.Equal(t, "test host", os.Getenv("PGHOST"))
assert.Equal(t, "999", os.Getenv("PGPORT"))
assert.Equal(t, "test name", os.Getenv("PGDATABASE"))
assert.Equal(t, "true", os.Getenv("DB_NOSSL"))
}
type FullDepsConfig struct {
aws.AWSConfig
serviceconfig.BaseConfig
@@ -74,7 +90,7 @@ func TestCreateFullDependencies(t *testing.T) {
cleanup()
}
func TestCreateNetwork(t *testing.T) {
func TestDepNetworkGet(t *testing.T) {
if testing.Short() {
t.Skip("Skipping long test in short mode")
}
@@ -96,9 +112,9 @@ func TestCreateMockServer(t *testing.T) {
}
ctx := context.Background()
server, cleanup := CreateMockServer(t, ctx, MockServerConfig{
Network: "hello",
})
net := DepNetwork.Get(t, ctx)
server, cleanup := CreateMockServer(t, ctx, net)
assert.NotNil(t, server)
assert.NotNil(t, cleanup)
@@ -120,15 +136,54 @@ func TestCreateMockExpectation(t *testing.T) {
CreateMockExpectation(t, server, expectation)
}
func TestCreateDetectDocumentTextExpectation(t *testing.T) {
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
}))
defer ts.Close()
server := MockServer{
Client: ts.Client(),
External: Address(ts.URL),
}
expectation := CreateDetectDocumentTextExpectation(t, server, "HI")
blocks, ok := expectation.Response.Body.(map[string]interface{})["Blocks"].([]map[string]interface{})
require.True(t, ok)
assert.Equal(
t,
"PAGE",
blocks[0]["BlockType"],
)
relationship, ok := blocks[0]["Relationships"].([]map[string]interface{})
require.True(t, ok)
assert.Equal(
t,
"CHILD",
relationship[0]["Type"],
)
ids, ok := relationship[0]["Ids"].([]string)
require.True(t, ok)
assert.Equal(
t,
ids[0],
blocks[1]["Id"],
)
assert.Equal(
t,
"HI",
blocks[1]["Text"],
)
}
func TestWaitForMockEndpoint(t *testing.T) {
if testing.Short() {
t.Skip("Skipping long test in short mode")
}
ctx := context.Background()
server, cleanup := CreateMockServer(t, ctx, MockServerConfig{
Network: "hello",
})
net := DepNetwork.Get(t, ctx)
server, cleanup := CreateMockServer(t, ctx, net)
defer cleanup()
body := strings.NewReader(`{"team":"hello"}`)
@@ -163,45 +218,3 @@ func TestWaitForMockEndpoint(t *testing.T) {
request := WaitForMockEndpoint(t, server, expectation.Request)
assert.Equal(t, []string([]string{"here"}), request.Headers["Hidden"])
}
func TestWaitForTextractMockEndpoint(t *testing.T) {
if testing.Short() {
t.Skip("Skipping long test in short mode")
}
ctx := context.Background()
cfg := &FullDepsConfig{}
SetCfgProvider(t, cfg)
deps, cleanup := CreateFullDependencies(t, ctx, cfg)
defer cleanup()
clientId := "CLIENTID"
importKey := objectstore.BucketKey{
ClientID: clientId,
CreatedAt: time.Now().UTC(),
Location: objectstore.Import,
}
outputKey := objectstore.BucketKey{
ClientID: clientId,
CreatedAt: time.Now().UTC(),
Location: objectstore.TextTextract,
}
e := CreateStartDocTextDetectionExpectation(t, deps.MockServer, StartDocTextDetectionExpectationParams{
Bucket: deps.BucketName,
Key: importKey,
JobID: "hi",
})
body := strings.NewReader(fmt.Sprintf(`{"DocumentLocation":{"S3Object":{"Bucket":"%s","Name":"%s"}},"OutputConfig":{"S3Bucket":"%s","S3Prefix":"%s"}}`, deps.BucketName, importKey.String(), deps.BucketName, outputKey.String()))
req, err := http.NewRequest("POST", string(deps.MockServer.External), body)
require.NoError(t, err)
req.Header.Add("X-Amz-Target", "Textract.StartDocumentTextDetection")
_, err = deps.MockServer.Client.Do(req)
require.NoError(t, err)
WaitForMockTextractOutput(t, ctx, cfg, deps.MockServer, TextractOutputParams{
Request: e.Request,
File: strings.NewReader("hello"),
})
}
+52 -9
View File
@@ -2,18 +2,61 @@ package test
import (
"context"
"log/slog"
"sync"
"testing"
"github.com/docker/docker/api/types/network"
"github.com/docker/docker/client"
"github.com/stretchr/testify/require"
"github.com/testcontainers/testcontainers-go"
"github.com/testcontainers/testcontainers-go/network"
)
func CreateNetwork(t testing.TB, ctx context.Context) (*testcontainers.DockerNetwork, func()) {
network, err := network.New(ctx, network.WithDriver("bridge"))
require.NoError(t, err)
return network, func() {
testcontainers.CleanupNetwork(t, network)
}
type NetworkManager struct {
networkName string
mutex sync.Mutex
initialized bool
}
var (
DepNetwork = &NetworkManager{}
)
const (
networkName = "queryorchestration_test"
)
func (nm *NetworkManager) Get(t testing.TB, ctx context.Context) string {
nm.mutex.Lock()
defer nm.mutex.Unlock()
if !nm.initialized {
cli, err := client.NewClientWithOpts(client.FromEnv, client.WithAPIVersionNegotiation())
require.NoError(t, err)
nm.networkName = networkName
current, err := cli.NetworkList(ctx, network.ListOptions{})
require.NoError(t, err)
found := false
for _, c := range current {
if c.Name == nm.networkName {
found = true
continue
}
}
if !found {
_, err = cli.NetworkCreate(ctx, nm.networkName, network.CreateOptions{
Driver: "bridge",
})
require.NoError(t, err)
}
nm.initialized = true
}
slog.Info("get network", "name", nm.networkName)
return nm.networkName
}
+5 -5
View File
@@ -9,12 +9,12 @@ import (
"github.com/stretchr/testify/assert"
)
func TestCreateNetwork(t *testing.T) {
func TestDepNetworkGet(t *testing.T) {
ctx := context.Background()
ncfg, cleanup := test.CreateNetwork(t, ctx)
assert.NotNil(t, ncfg)
assert.NotNil(t, cleanup)
name := test.DepNetwork.Get(t, ctx)
assert.NotNil(t, name)
cleanup()
newName := test.DepNetwork.Get(t, ctx)
assert.Equal(t, name, newName)
}
+6 -2
View File
@@ -30,7 +30,9 @@ func TestCreateBucket(t *testing.T) {
cfg := &StoreConfig{}
SetCfgProvider(t, cfg)
acfg, cleanup := CreateAWSContainer(t, ctx, cfg, &CreateAWSConfig{})
net := DepNetwork.Get(t, ctx)
acfg, cleanup := CreateAWSContainer(t, ctx, cfg, net)
defer cleanup()
SetStoreClient(t, ctx, cfg, acfg.ExternalEndpoint)
@@ -47,7 +49,9 @@ func TestCreateStoreClient(t *testing.T) {
cfg := &StoreConfig{}
SetCfgProvider(t, cfg)
acfg, cleanup := CreateAWSContainer(t, ctx, cfg, &CreateAWSConfig{})
net := DepNetwork.Get(t, ctx)
acfg, cleanup := CreateAWSContainer(t, ctx, cfg, net)
defer cleanup()
SetStoreClient(t, ctx, cfg, acfg.ExternalEndpoint)
+23 -6
View File
@@ -6,6 +6,7 @@ import (
"log/slog"
"regexp"
"testing"
"time"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue"
@@ -46,13 +47,29 @@ func CreateQueue(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProv
}
func AssertMessage(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, params *queue.ReceiveParams) types.Message {
result, err := cfg.ReceiveFromQueue(ctx, params)
require.NoError(t, err)
assert.NotNil(t, result.Messages)
assert.Len(t, result.Messages, 1)
assert.NotNil(t, result.Messages[0])
timeout := time.After(30 * time.Second)
tick := time.NewTicker(2 * time.Second)
defer tick.Stop()
return result.Messages[0]
for {
select {
case <-timeout:
t.Fatal("assert timeout")
case <-tick.C:
slog.Info("receiving from queue")
result, err := cfg.ReceiveFromQueue(ctx, params)
if err != nil {
continue
} else if len(result.Messages) < 1 {
continue
}
return result.Messages[0]
case <-ctx.Done():
t.Fatal(ctx.Err())
}
}
}
func AssertMessageBody(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, url string, body *regexp.Regexp) {
+13 -5
View File
@@ -31,14 +31,16 @@ func TestCreateQueue(t *testing.T) {
cfg := &TestConfig{}
SetCfgProvider(t, cfg)
_, cleanup := CreateAWSContainer(t, ctx, cfg, &CreateAWSConfig{})
net := DepNetwork.Get(t, ctx)
_, cleanup := CreateAWSContainer(t, ctx, cfg, net)
defer cleanup()
err := cfg.SetQueueClient(ctx)
require.NoError(t, err)
url := CreateQueue(t, ctx, cfg, "myname")
assert.Equal(t, "http://localstack:4566/queue/us-east-1/000000000000/myname", url)
assert.Equal(t, "http://localstack_testcreatequeue:4566/queue/us-east-1/000000000000/myname", url)
}
func TestAssertMessageWait(t *testing.T) {
@@ -50,7 +52,9 @@ func TestAssertMessageWait(t *testing.T) {
cfg := &TestConfig{}
SetCfgProvider(t, cfg)
_, cleanup := CreateAWSContainer(t, ctx, cfg, &CreateAWSConfig{})
net := DepNetwork.Get(t, ctx)
_, cleanup := CreateAWSContainer(t, ctx, cfg, net)
defer cleanup()
err := cfg.SetQueueClient(ctx)
require.NoError(t, err)
@@ -77,7 +81,9 @@ func TestAssertMessageBodyWait(t *testing.T) {
cfg := &TestConfig{}
SetCfgProvider(t, cfg)
_, cleanup := CreateAWSContainer(t, ctx, cfg, &CreateAWSConfig{})
net := DepNetwork.Get(t, ctx)
_, cleanup := CreateAWSContainer(t, ctx, cfg, net)
defer cleanup()
err := cfg.SetQueueClient(ctx)
require.NoError(t, err)
@@ -101,7 +107,9 @@ func TestAssertMessageAttrWait(t *testing.T) {
cfg := &TestConfig{}
SetCfgProvider(t, cfg)
_, cleanup := CreateAWSContainer(t, ctx, cfg, &CreateAWSConfig{})
net := DepNetwork.Get(t, ctx)
_, cleanup := CreateAWSContainer(t, ctx, cfg, net)
defer cleanup()
err := cfg.SetQueueClient(ctx)
require.NoError(t, err)
+3 -18
View File
@@ -8,7 +8,6 @@ import (
doccleanrunner "queryorchestration/api/docCleanRunner"
docinitrunner "queryorchestration/api/docInitRunner"
docsyncrunner "queryorchestration/api/docSyncRunner"
doctextprocessrunner "queryorchestration/api/docTextProcessRunner"
doctextrunner "queryorchestration/api/docTextRunner"
queryrunner "queryorchestration/api/queryRunner"
querysyncrunner "queryorchestration/api/querySyncRunner"
@@ -20,14 +19,11 @@ import (
"queryorchestration/internal/serviceconfig/queue/documentclean"
"queryorchestration/internal/serviceconfig/queue/documentinit"
"queryorchestration/internal/serviceconfig/queue/documentsync"
"queryorchestration/internal/serviceconfig/queue/documenttextprocess"
"queryorchestration/internal/serviceconfig/queue/documenttexttrigger"
"queryorchestration/internal/serviceconfig/queue/query"
"queryorchestration/internal/serviceconfig/queue/querysync"
"queryorchestration/internal/serviceconfig/queue/queryversionsync"
"queryorchestration/internal/serviceconfig/queue/storeevent"
"github.com/testcontainers/testcontainers-go"
)
type RunnerName string
@@ -38,7 +34,6 @@ const (
DocSyncRunnerName RunnerName = docsyncrunner.Name
DocCleanRunnerName RunnerName = doccleanrunner.Name
DocTextRunnerName RunnerName = doctextrunner.Name
DocTextProcessRunnerName RunnerName = doctextprocessrunner.Name
QuerySyncRunnerName RunnerName = querysyncrunner.Name
QueryRunnerName RunnerName = queryrunner.Name
ClientSyncRunnerName RunnerName = clientsyncrunner.Name
@@ -53,7 +48,6 @@ const (
DocSyncRunnerEnv RunnerEnv = documentsync.EnvName
DocCleanRunnerEnv RunnerEnv = documentclean.EnvName
DocTextRunnerEnv RunnerEnv = documenttexttrigger.EnvName
DocTextProcessRunnerEnv RunnerEnv = documenttextprocess.EnvName
QuerySyncRunnerEnv RunnerEnv = querysync.EnvName
QueryRunnerEnv RunnerEnv = query.EnvName
ClientSyncRunnerEnv RunnerEnv = clientsync.EnvName
@@ -75,8 +69,6 @@ func GetRunnerEnvFromName(name RunnerName) RunnerEnv {
return DocCleanRunnerEnv
case DocTextRunnerName:
return DocTextRunnerEnv
case DocTextProcessRunnerName:
return DocTextProcessRunnerEnv
case QuerySyncRunnerName:
return QuerySyncRunnerEnv
case QueryRunnerName:
@@ -91,17 +83,15 @@ func GetRunnerEnvFromName(name RunnerName) RunnerEnv {
type RunnerConfig struct {
Runner Runner
Network *testcontainers.DockerNetwork
MockHTTP string
}
func CreateRunner(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, config *RunnerConfig) (*Container, func()) {
func CreateRunner(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, network string, config *RunnerConfig) (*Container, func()) {
env := map[string]string{}
url := GetQueueURL(cfg, config.Runner.Name)
env["QUEUE_URL"] = url
c, cleanup := createContainer(t, ctx, &containerConfig{
Network: config.Network,
c, cleanup := createContainer(t, ctx, network, &containerConfig{
Name: string(config.Runner.Name),
DownstreamQueues: config.Runner.DownstreamQueues,
Cfg: cfg,
@@ -118,7 +108,7 @@ func CreateRunner(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigPro
var StoreEventRunner = Runner{
Name: StoreEventRunnerName,
DownstreamQueues: []RunnerName{DocInitRunnerName, DocTextProcessRunnerName},
DownstreamQueues: []RunnerName{DocInitRunnerName},
}
var DocInitRunner = Runner{
Name: DocInitRunnerName,
@@ -136,10 +126,6 @@ var DocTextRunner = Runner{
Name: DocTextRunnerName,
DownstreamQueues: []RunnerName{QuerySyncRunnerName},
}
var DocTextProcessRunner = Runner{
Name: DocTextProcessRunnerName,
DownstreamQueues: []RunnerName{QuerySyncRunnerName},
}
var QuerySyncRunner = Runner{
Name: QuerySyncRunnerName,
DownstreamQueues: []RunnerName{QueryRunnerName},
@@ -163,7 +149,6 @@ var runners = []Runner{
DocSyncRunner,
DocCleanRunner,
DocTextRunner,
DocTextProcessRunner,
QuerySyncRunner,
QueryRunner,
ClientSyncRunner,
+6 -14
View File
@@ -14,31 +14,24 @@ func TestCreateRunner(t *testing.T) {
}
ctx := context.Background()
ncfg, ncleanup := CreateNetwork(t, ctx)
defer ncleanup()
cfg := &TestConfig{}
SetCfgProvider(t, cfg)
_, qcleanup := CreateAWSContainer(t, ctx, cfg, &CreateAWSConfig{
Network: ncfg,
})
net := DepNetwork.Get(t, ctx)
_, qcleanup := CreateAWSContainer(t, ctx, cfg, net)
defer qcleanup()
err := cfg.SetQueueClient(ctx)
require.NoError(t, err)
_ = CreateQueue(t, ctx, cfg, string(QueryRunnerName))
_, dbcleanup := CreateDB(t, ctx, cfg, &CreateDatabaseConfig{
Network: ncfg,
})
defer dbcleanup()
_ = CreateDB(t, ctx, cfg, net, &CreateDatabaseConfig{})
qccfg := &RunnerConfig{
Runner: QueryRunner,
Network: ncfg,
Runner: QueryRunner,
}
c, cleanup := CreateRunner(t, ctx, cfg, qccfg)
c, cleanup := CreateRunner(t, ctx, cfg, net, qccfg)
assert.NotNil(t, cleanup)
assert.NotNil(t, c)
@@ -51,7 +44,6 @@ func TestGetRunnerEnvFromName(t *testing.T) {
assert.Equal(t, DocSyncRunnerEnv, GetRunnerEnvFromName(DocSyncRunnerName))
assert.Equal(t, DocCleanRunnerEnv, GetRunnerEnvFromName(DocCleanRunnerName))
assert.Equal(t, DocTextRunnerEnv, GetRunnerEnvFromName(DocTextRunnerName))
assert.Equal(t, DocTextProcessRunnerEnv, GetRunnerEnvFromName(DocTextProcessRunnerName))
assert.Equal(t, QuerySyncRunnerEnv, GetRunnerEnvFromName(QuerySyncRunnerName))
assert.Equal(t, QueryRunnerEnv, GetRunnerEnvFromName(QueryRunnerName))
assert.Equal(t, ClientSyncRunnerEnv, GetRunnerEnvFromName(ClientSyncRunnerName))