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) }