0ddae4f91e
remove query from codebase part 1 * remove query * fix localstack run
114 lines
3.1 KiB
Go
114 lines
3.1 KiB
Go
package test
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log/slog"
|
|
"regexp"
|
|
"testing"
|
|
"time"
|
|
|
|
"queryorchestration/internal/serviceconfig"
|
|
"queryorchestration/internal/serviceconfig/queue"
|
|
|
|
awsc "queryorchestration/internal/serviceconfig/aws"
|
|
|
|
"github.com/aws/aws-sdk-go-v2/aws"
|
|
"github.com/aws/aws-sdk-go-v2/config"
|
|
"github.com/aws/aws-sdk-go-v2/credentials"
|
|
"github.com/aws/aws-sdk-go-v2/service/sqs"
|
|
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
)
|
|
|
|
func SetQueueClient(t testing.TB, ctx context.Context, cfg queue.ConfigProvider, endpoint string) {
|
|
t.Helper()
|
|
acfg, err := awsc.GetAWSConfigWithOpts(ctx, func(lo *config.LoadOptions) error {
|
|
lo.BaseEndpoint = endpoint
|
|
lo.Credentials = credentials.NewStaticCredentialsProvider(cfg.GetAWSKeyID(), cfg.GetAWSSecretKey(), cfg.GetAWSSessionToken())
|
|
lo.Region = cfg.GetAWSRegion()
|
|
|
|
return nil
|
|
})
|
|
require.NoError(t, err)
|
|
|
|
cfg.SetQueueClientFromSQSCfg(sqs.NewFromConfig(acfg))
|
|
}
|
|
|
|
func GetQueueURL(t testing.TB, cfg awsc.ConfigProvider, name RunnerName) string {
|
|
queueName := GetQueueName(t, name)
|
|
return fmt.Sprintf("http://localstack:4566/queue/%s/000000000000/%s", cfg.GetAWSRegion(), queueName)
|
|
}
|
|
|
|
func GetQueueArn(t testing.TB, cfg awsc.ConfigProvider, name RunnerName) string {
|
|
queueName := GetQueueName(t, name)
|
|
return fmt.Sprintf("arn:aws:sqs:%s:000000000000:%s", cfg.GetAWSRegion(), queueName)
|
|
}
|
|
|
|
func GetQueueName(t testing.TB, name RunnerName) string {
|
|
return GetAlias(t, string(name))
|
|
}
|
|
|
|
func CreateQueue(t testing.TB, cfg serviceconfig.ConfigProvider, name RunnerName) string {
|
|
t.Helper()
|
|
queueName := GetQueueName(t, name)
|
|
queueM, err := cfg.GetQueueClient().CreateQueue(t.Context(), &sqs.CreateQueueInput{
|
|
QueueName: aws.String(queueName),
|
|
})
|
|
require.NoError(t, err)
|
|
slog.Info("create queue", "name", queueName, "url", *queueM.QueueUrl)
|
|
|
|
err = cfg.PingQueueByURL(t.Context(), *queueM.QueueUrl)
|
|
require.NoError(t, err)
|
|
|
|
return *queueM.QueueUrl
|
|
}
|
|
|
|
func AssertMessage(t testing.TB, cfg serviceconfig.ConfigProvider, params *queue.ReceiveParams) types.Message {
|
|
t.Helper()
|
|
timeout := time.After(30 * time.Second)
|
|
tick := time.NewTicker(500 * time.Millisecond)
|
|
defer tick.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-timeout:
|
|
t.Fatal("assert timeout")
|
|
case <-tick.C:
|
|
slog.Info("receiving from queue")
|
|
|
|
result, err := cfg.ReceiveFromQueue(t.Context(), params)
|
|
if err != nil {
|
|
continue
|
|
} else if len(result.Messages) < 1 {
|
|
continue
|
|
}
|
|
|
|
return result.Messages[0]
|
|
case <-t.Context().Done():
|
|
t.Fatal(t.Context().Err())
|
|
}
|
|
}
|
|
}
|
|
|
|
func AssertMessageBody(t testing.TB, cfg serviceconfig.ConfigProvider, url string, body *regexp.Regexp) {
|
|
t.Helper()
|
|
message := AssertMessage(t, cfg, &queue.ReceiveParams{
|
|
QueueURL: url,
|
|
})
|
|
|
|
assert.Regexp(t, body, *message.Body)
|
|
}
|
|
|
|
func AssertMessageAttr(t testing.TB, cfg serviceconfig.ConfigProvider, url string, name string, value *regexp.Regexp) {
|
|
t.Helper()
|
|
message := AssertMessage(t, cfg, &queue.ReceiveParams{
|
|
QueueURL: url,
|
|
Attributes: []string{name},
|
|
})
|
|
|
|
assert.NotNil(t, message.MessageAttributes[name])
|
|
assert.Regexp(t, value, *(message.MessageAttributes[name]).StringValue)
|
|
}
|