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

114 lines
3.2 KiB
Go
Raw Normal View History

package test
import (
"context"
"fmt"
2025-03-19 11:54:14 +00:00
"log/slog"
"regexp"
"testing"
"time"
2025-03-05 12:05:46 +00:00
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue"
awsc "queryorchestration/internal/serviceconfig/aws"
"github.com/aws/aws-sdk-go-v2/aws"
2025-04-25 17:02:50 +00:00
"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"
)
2025-04-25 17:02:50 +00:00
func SetQueueClient(t testing.TB, ctx context.Context, cfg queue.ConfigProvider, endpoint string) {
2025-04-22 19:57:35 +00:00
t.Helper()
2025-04-25 17:02:50 +00:00
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)
2025-04-25 17:02:50 +00:00
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)
}
2025-04-25 17:02:50 +00:00
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)
}
2025-04-25 17:02:50 +00:00
func GetQueueName(t testing.TB, name RunnerName) string {
return NormaliseAlias(fmt.Sprintf("%s_%s", name, t.Name()))
}
2025-04-25 17:02:50 +00:00
func CreateQueue(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, name RunnerName) string {
2025-04-22 19:57:35 +00:00
t.Helper()
2025-04-25 17:02:50 +00:00
queueName := GetQueueName(t, name)
queueM, err := cfg.GetQueueClient().CreateQueue(ctx, &sqs.CreateQueueInput{
2025-04-25 17:02:50 +00:00
QueueName: aws.String(queueName),
})
require.NoError(t, err)
2025-04-25 17:02:50 +00:00
slog.Info("create queue", "name", queueName, "url", *queueM.QueueUrl)
err = cfg.PingQueueByURL(ctx, *queueM.QueueUrl)
require.NoError(t, err)
return *queueM.QueueUrl
}
func AssertMessage(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, params *queue.ReceiveParams) types.Message {
2025-04-22 19:57:35 +00:00
t.Helper()
timeout := time.After(30 * time.Second)
tick := time.NewTicker(2 * time.Second)
defer tick.Stop()
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) {
2025-04-22 19:57:35 +00:00
t.Helper()
message := AssertMessage(t, ctx, cfg, &queue.ReceiveParams{
QueueURL: url,
})
assert.Regexp(t, body, *message.Body)
}
func AssertMessageAttr(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, url string, name string, value *regexp.Regexp) {
2025-04-22 19:57:35 +00:00
t.Helper()
message := AssertMessage(t, ctx, cfg, &queue.ReceiveParams{
QueueURL: url,
Attributes: []string{name},
})
assert.NotNil(t, message.MessageAttributes[name])
assert.Regexp(t, value, *(message.MessageAttributes[name]).StringValue)
}