Files
query-orchestration/internal/test/ecosystem.go
T

375 lines
8.6 KiB
Go
Raw Normal View History

package test
import (
"context"
"os"
"path"
2025-03-05 12:05:46 +00:00
"testing"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/aws"
"queryorchestration/internal/serviceconfig/database"
"queryorchestration/internal/serviceconfig/objectstore"
queryservice "queryorchestration/pkg/queryService"
"github.com/stretchr/testify/require"
"github.com/testcontainers/testcontainers-go"
)
type EcosystemConfig struct {
Services map[Service]*Container
Runners map[Runner]*Container
Cfg serviceconfig.ConfigProvider
}
type EcosystemNetworkConfig struct {
Runners []*RunnerNetworkConfig
Services []*ServiceNetworkConfig
Cfg serviceconfig.ConfigProvider
Dependencies *Dependencies
}
func CreateRunnersAndServicesNetwork(t testing.TB, ctx context.Context, ncfg *EcosystemNetworkConfig) (*EcosystemConfig, func()) {
if ncfg.Dependencies == nil {
ncfg.Dependencies = &Dependencies{}
}
network := ncfg.Dependencies.Network
var ncleanup func()
if network == nil {
network, ncleanup = CreateNetwork(t, ctx)
}
_, dbcleanup := CreateDB(t, ctx, &CreateDatabaseConfig{
Network: network,
Cfg: ncfg.Cfg,
RunMigrations: true,
})
ss := make(map[string]*Container, len(ncfg.Services))
sclean := make([]func(), len(ncfg.Services))
for i, s := range ncfg.Services {
c, ccleanup := CreateService(t, ctx, &ServiceConfig{
Name: s.Name,
Env: s.Env,
Cfg: ncfg.Cfg,
Network: network,
})
ss[s.Name] = c
sclean[i] = ccleanup
}
var qcleanup func()
if ncfg.Dependencies.AWSConfig == nil {
_, qcleanup = CreateAWSContainer(t, ctx, &CreateAWSConfig{
Network: network,
Cfg: ncfg.Cfg,
})
}
if ncfg.Cfg.GetQueueClient() == nil {
SetQueueClient(t, ctx, ncfg.Cfg)
}
qs := make(map[string]*Container, len(ncfg.Runners))
qclean := make([]func(), len(ncfg.Runners))
for i, r := range ncfg.Runners {
url := ncfg.Dependencies.QueueURLs[r.Name]
if url == "" {
url = CreateQueue(t, ctx, ncfg.Cfg, r.Name)
}
c, ccleanup := CreateRunner(t, ctx, &RunnerConfig{
Name: r.Name,
Env: r.Env,
QueueURL: url,
Cfg: ncfg.Cfg,
Network: network,
})
qs[r.Name] = c
qclean[i] = ccleanup
}
return &EcosystemConfig{
Services: ss,
Runners: qs,
}, func() {
for _, c := range sclean {
c()
}
for _, c := range qclean {
c()
}
dbcleanup()
if qcleanup != nil {
qcleanup()
}
if ncleanup != nil {
ncleanup()
}
}
}
func SetCfgProviderWithBasePath(t testing.TB, cfg serviceconfig.ConfigProvider, basePath string) {
SetCfgProvider(t, cfg)
cfg.SetBasePath(path.Join(os.Getenv("PWD"), basePath))
}
func SetCfgProvider(t testing.TB, cfg serviceconfig.ConfigProvider) {
cfg.SetAWSConfig(&aws.AWSConfig{
AWSKeyID: "test",
AWSSecretKey: "test",
AWSRegion: "us-east-1",
})
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())
cfg.SetDBConfig(&database.DBConfig{
DBUser: "invalid_user",
DBSecret: "invalid_pass",
DBHost: "invalid_host",
DBPort: 5432,
DBName: "invalid_name",
DBNoSSL: true,
})
}
type ServiceNetworkConfig struct {
Cfg serviceconfig.ConfigProvider
Network *testcontainers.DockerNetwork
Name Runner
Env map[string]string
}
func CreateServiceNetwork(t testing.TB, ctx context.Context, scfg *ServiceNetworkConfig) (*Container, func()) {
if scfg.Cfg == nil {
scfg.Cfg = &serviceconfig.BaseConfig{}
SetCfgProvider(t, scfg.Cfg)
}
network := scfg.Network
var ncleanup func()
if network == nil {
network, ncleanup = CreateNetwork(t, ctx)
}
_, dbcleanup := CreateDB(t, ctx, &CreateDatabaseConfig{
Network: network,
Cfg: scfg.Cfg,
})
c, ccleanup := CreateService(t, ctx, &ServiceConfig{
Name: scfg.Name,
Env: scfg.Env,
Cfg: scfg.Cfg,
Network: network,
})
return c, func() {
dbcleanup()
ccleanup()
if ncleanup != nil {
ncleanup()
}
}
}
type RunnerNetworkConfig struct {
Network *testcontainers.DockerNetwork
AWSContainer *AWSContainerConfig
QueueURL *string
Cfg serviceconfig.ConfigProvider
Name Runner
Env map[string]string
}
func CreateRunnerNetwork(t testing.TB, ctx context.Context, rcfg *RunnerNetworkConfig) (*Container, func()) {
network := rcfg.Network
var ncleanup func()
if network == nil {
network, ncleanup = CreateNetwork(t, ctx)
}
var qcleanup func()
if rcfg.AWSContainer == nil {
_, qcleanup = CreateAWSContainer(t, ctx, &CreateAWSConfig{
Network: network,
Cfg: rcfg.Cfg,
})
}
if rcfg.Cfg.GetQueueClient() == nil {
err := rcfg.Cfg.SetQueueClient(ctx)
require.NoError(t, err)
}
url := rcfg.QueueURL
if url == nil {
u := CreateQueue(t, ctx, rcfg.Cfg, rcfg.Name)
url = &u
}
_, dbcleanup := CreateDB(t, ctx, &CreateDatabaseConfig{
Network: network,
Cfg: rcfg.Cfg,
})
c, ccleanup := CreateRunner(t, ctx, &RunnerConfig{
Name: rcfg.Name,
Env: rcfg.Env,
QueueURL: *url,
Cfg: rcfg.Cfg,
Network: network,
})
return c, func() {
ccleanup()
dbcleanup()
if qcleanup != nil {
qcleanup()
}
if ncleanup != nil {
ncleanup()
}
}
}
type FullDependenciesConfig interface {
serviceconfig.ConfigProvider
objectstore.ConfigProvider
}
type Dependencies struct {
BucketName string
QueueURLs map[string]string
Network *testcontainers.DockerNetwork
AWSConfig *AWSContainerConfig
}
func CreateFullDependencies(t testing.TB, ctx context.Context, cfg FullDependenciesConfig) (Dependencies, func()) {
network, ncleanup := CreateNetwork(t, ctx)
acfg, clean := CreateAWSContainer(t, ctx, &CreateAWSConfig{
Cfg: cfg,
Network: network,
})
SetQueueClient(t, ctx, cfg)
SetStoreClient(t, ctx, cfg, acfg.ExternalEndpoint)
urls := map[string]string{}
urls[DocInitRunner] = CreateQueue(t, ctx, cfg, DocInitRunner)
urls[DocSyncRunner] = CreateQueue(t, ctx, cfg, DocSyncRunner)
urls[DocCleanRunner] = CreateQueue(t, ctx, cfg, DocCleanRunner)
urls[DocTextRunner] = CreateQueue(t, ctx, cfg, DocTextRunner)
urls[QuerySyncRunner] = CreateQueue(t, ctx, cfg, QuerySyncRunner)
urls[QueryRunner] = CreateQueue(t, ctx, cfg, QueryRunner)
urls[ClientSyncRunner] = CreateQueue(t, ctx, cfg, ClientSyncRunner)
urls[QueryVersionSyncRunner] = CreateQueue(t, ctx, cfg, QueryVersionSyncRunner)
bucketName := "docinitbucket"
CreateBucket(t, ctx, cfg, bucketName)
SetBucketNotifs(t, ctx, cfg, bucketName)
return Dependencies{
BucketName: bucketName,
QueueURLs: urls,
Network: network,
AWSConfig: acfg,
}, func() {
ncleanup()
clean()
}
}
type Network struct {
Dependencies Dependencies
Network EcosystemConfig
Client *queryservice.ClientWithResponses
}
func CreateFullNetwork(t testing.TB, ctx context.Context, cfg FullDependenciesConfig) (Network, func()) {
deps, clean := CreateFullDependencies(t, ctx, cfg)
net, cleanup := CreateRunnersAndServicesNetwork(t, ctx, &EcosystemNetworkConfig{
Cfg: cfg,
Dependencies: &deps,
Runners: []*RunnerNetworkConfig{
{
Name: DocInitRunner,
Env: map[string]string{
"DOCUMENT_SYNC_URL": deps.QueueURLs[DocSyncRunner],
},
},
{
Name: DocSyncRunner,
Env: map[string]string{
"DOCUMENT_CLEAN_URL": deps.QueueURLs[DocCleanRunner],
},
},
{
Name: DocCleanRunner,
Env: map[string]string{
"DOCUMENT_TEXT_URL": deps.QueueURLs[DocTextRunner],
},
},
{
Name: DocTextRunner,
Env: map[string]string{
"QUERY_SYNC_URL": deps.QueueURLs[QuerySyncRunner],
},
},
{
Name: QuerySyncRunner,
Env: map[string]string{
"QUERY_URL": deps.QueueURLs[QueryRunner],
},
},
{
Name: QueryRunner,
Env: map[string]string{
"QUERY_URL": deps.QueueURLs[QueryRunner],
},
},
{
Name: ClientSyncRunner,
Env: map[string]string{
"DOCUMENT_SYNC_URL": deps.QueueURLs[DocSyncRunner],
},
},
{
Name: QueryVersionSyncRunner,
Env: map[string]string{
"CLIENT_SYNC_URL": deps.QueueURLs[ClientSyncRunner],
},
},
},
Services: []*ServiceNetworkConfig{
{
Name: QueryService,
Env: map[string]string{
"CLIENT_SYNC_URL": deps.QueueURLs[ClientSyncRunner],
"QUERY_VERSION_SYNC_URL": deps.QueueURLs[QueryVersionSyncRunner],
},
},
},
})
qService, err := queryservice.NewClientWithResponses(net.Services[QueryService].URI)
require.NoError(t, err)
return Network{
Dependencies: deps,
Network: *net,
Client: qService,
}, func() {
cleanup()
clean()
}
}