Merged in feature/demo (pull request #116)

Demo prep + Fix client sync

* firstversion

* clientsync

* configlint

* fixtests
This commit is contained in:
Michael McGuinness
2025-04-22 19:57:35 +00:00
parent fee71e7740
commit 47fec079e5
21 changed files with 642 additions and 41 deletions
+3
View File
@@ -4,3 +4,6 @@ ignore: |
vendor/
.devbox/
out/
rules:
line-length:
max: 140
+2 -1
View File
@@ -5,6 +5,7 @@ import (
"net/http"
"queryorchestration/internal/client"
clientupdate "queryorchestration/internal/client/update"
"github.com/labstack/echo/v4"
)
@@ -47,7 +48,7 @@ func (s *Controllers) UpdateClient(ctx echo.Context, id ClientID) error {
return echo.NewHTTPError(http.StatusBadRequest, err)
}
err := s.svc.Client.Update(ctx.Request().Context(), id, &client.Update{
err := s.svc.ClientUpdate.Update(ctx.Request().Context(), id, &clientupdate.Update{
Name: req.Name,
CanSync: req.CanSync,
})
+11 -2
View File
@@ -10,8 +10,10 @@ import (
queryapi "queryorchestration/api/queryAPI"
"queryorchestration/internal/client"
clientupdate "queryorchestration/internal/client/update"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/clientsync"
"github.com/labstack/echo/v4"
"github.com/pashagolub/pgxmock/v3"
@@ -19,6 +21,11 @@ import (
"github.com/stretchr/testify/require"
)
type ClientConfig struct {
serviceconfig.BaseConfig
clientsync.ConfigProvider
}
func TestCreateClient(t *testing.T) {
pool, err := pgxmock.NewPool()
require.NoError(t, err)
@@ -97,12 +104,14 @@ func TestUpdateClient(t *testing.T) {
pool, err := pgxmock.NewPool()
require.NoError(t, err)
cfg := &serviceconfig.BaseConfig{}
cfg := &ClientConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
cons := queryapi.NewControllers(&queryapi.Services{
Client: client.New(cfg),
ClientUpdate: clientupdate.New(cfg, &clientupdate.Services{
Client: client.New(cfg),
}),
})
cs := false
+2
View File
@@ -2,6 +2,7 @@ package queryapi
import (
"queryorchestration/internal/client"
clientupdate "queryorchestration/internal/client/update"
"queryorchestration/internal/collector"
collectorset "queryorchestration/internal/collector/set"
"queryorchestration/internal/document"
@@ -21,6 +22,7 @@ type Services struct {
QueryUpdate *queryupdate.Service
QueryTest *querytest.Service
Client *client.Service
ClientUpdate *clientupdate.Service
Document *document.Service
}
+2 -1
View File
@@ -56,7 +56,8 @@ type Region string
type EventS3 string
const (
EventS3ObjectCreatedPut = "ObjectCreated:Put"
EventS3ObjectCreatedPut = "ObjectCreated:Put"
EventS3ObjectCreatedCompleteMultipartUpload = "ObjectCreated:CompleteMultipartUpload"
)
func (s Runner) Process(ctx context.Context, body S3EventNotification) bool {
+2 -3
View File
@@ -21,9 +21,6 @@ ENV CGO_ENABLED=1
ENV GOARCH=${TARGETARCH}
ENV CC=musl-gcc
# Must be same a github.com/gen2brain/go-fitz in go.mod
# ENV FZ_VERSION=1.24.14
RUN --mount=type=cache,target=/go/pkg/mod/ \
go build \
-mod vendor \
@@ -37,6 +34,8 @@ ENV PWD=/app
WORKDIR ${PWD}
COPY --from=build /etc/ssl/certs/ca-certificates.crt /etc/ssl/certs/
COPY --from=build /bin/ .
EXPOSE 8080
+5
View File
@@ -21,6 +21,7 @@ import (
queryapi "queryorchestration/api/queryAPI"
"queryorchestration/internal/client"
clientupdate "queryorchestration/internal/client/update"
"queryorchestration/internal/collector"
collectorset "queryorchestration/internal/collector/set"
"queryorchestration/internal/document"
@@ -59,6 +60,9 @@ func main() {
Collector: col,
})
cli := client.New(cfg)
cliUpdate := clientupdate.New(cfg, &clientupdate.Services{
Client: cli,
})
doc := document.New(cfg)
quetest := querytest.New(cfg, &querytest.Services{
Collector: col,
@@ -77,6 +81,7 @@ func main() {
QueryUpdate: qupdate,
QueryTest: quetest,
Client: cli,
ClientUpdate: cliUpdate,
Document: doc,
}
+381
View File
@@ -0,0 +1,381 @@
---
services:
store_event_runner:
image: queryorchestration:latest
command: ["./storeEventRunner"]
depends_on:
localstack:
condition: service_healthy
aws-db:
condition: service_healthy
ports:
- "8081:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
QUEUE_URL: ${STORE_EVENT_URL}
DOCUMENT_INIT_URL: ${DOCUMENT_INIT_URL}
DOCUMENT_TEXT_PROCESS_URL: ${DOCUMENT_TEXT_PROCESS_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
PGUSER: ${PGUSER}
PGPASSWORD: ${PGPASSWORD}
PGHOST: aws-db
PGPORT: 5432
PGDATABASE: ${PGDATABASE}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- aws-server-network
doc_init_runner:
image: queryorchestration:latest
command: ["./docInitRunner"]
depends_on:
localstack:
condition: service_healthy
aws-db:
condition: service_healthy
ports:
- "8082:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
QUEUE_URL: ${DOCUMENT_INIT_URL}
DOCUMENT_SYNC_URL: ${DOCUMENT_SYNC_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
PGUSER: ${PGUSER}
PGPASSWORD: ${PGPASSWORD}
PGHOST: aws-db
PGPORT: 5432
PGDATABASE: ${PGDATABASE}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- aws-server-network
doc_sync_runner:
image: queryorchestration:latest
command: ["./docSyncRunner"]
depends_on:
localstack:
condition: service_healthy
aws-db:
condition: service_healthy
ports:
- "8083:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
QUEUE_URL: ${DOCUMENT_SYNC_URL}
DOCUMENT_CLEAN_URL: ${DOCUMENT_CLEAN_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
PGUSER: ${PGUSER}
PGPASSWORD: ${PGPASSWORD}
PGHOST: aws-db
PGPORT: 5432
PGDATABASE: ${PGDATABASE}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- aws-server-network
doc_clean_runner:
image: queryorchestration:latest
command: ["./docCleanRunner"]
depends_on:
localstack:
condition: service_healthy
aws-db:
condition: service_healthy
ports:
- "8084:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
QUEUE_URL: ${DOCUMENT_CLEAN_URL}
DOCUMENT_TEXT_TRIGGER_URL: ${DOCUMENT_TEXT_TRIGGER_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
PGUSER: ${PGUSER}
PGPASSWORD: ${PGPASSWORD}
PGHOST: aws-db
PGPORT: 5432
PGDATABASE: ${PGDATABASE}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- aws-server-network
doc_text_runner:
image: queryorchestration:latest
command: ["./docTextRunner"]
depends_on:
localstack:
condition: service_healthy
aws-db:
condition: service_healthy
ports:
- "8085:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
QUEUE_URL: ${DOCUMENT_TEXT_TRIGGER_URL}
QUERY_SYNC_URL: ${QUERY_SYNC_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
AWS_ACCESS_KEY_ID_TEXTRACT: ${AWS_ACCESS_KEY_ID_TEXTRACT}
AWS_SECRET_ACCESS_KEY_TEXTRACT: ${AWS_SECRET_ACCESS_KEY_TEXTRACT}
AWS_SESSION_TOKEN_TEXTRACT: ${AWS_SESSION_TOKEN_TEXTRACT}
AWS_REGION_TEXTRACT: ${AWS_REGION_TEXTRACT}
PGUSER: ${PGUSER}
PGPASSWORD: ${PGPASSWORD}
PGHOST: aws-db
PGPORT: 5432
PGDATABASE: ${PGDATABASE}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL_S3: "http://localstack:4566"
AWS_ENDPOINT_URL_SQS: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- aws-server-network
query_sync_runner:
image: queryorchestration:latest
command: ["./querySyncRunner"]
depends_on:
localstack:
condition: service_healthy
aws-db:
condition: service_healthy
ports:
- "8087:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
QUEUE_URL: ${QUERY_SYNC_URL}
QUERY_URL: ${QUERY_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
PGUSER: ${PGUSER}
PGPASSWORD: ${PGPASSWORD}
PGHOST: aws-db
PGPORT: 5432
PGDATABASE: ${PGDATABASE}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- aws-server-network
query_runner:
image: queryorchestration:latest
command: ["./queryRunner"]
depends_on:
localstack:
condition: service_healthy
aws-db:
condition: service_healthy
ports:
- "8088:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
QUEUE_URL: ${QUERY_URL}
QUERY_URL: ${QUERY_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
PGUSER: ${PGUSER}
PGPASSWORD: ${PGPASSWORD}
PGHOST: aws-db
PGPORT: 5432
PGDATABASE: ${PGDATABASE}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- aws-server-network
query_api:
image: queryorchestration:latest
command: ["./queryAPI"]
depends_on:
localstack:
condition: service_healthy
aws-db:
condition: service_healthy
ports:
- "8080:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
CLIENT_SYNC_URL: ${CLIENT_SYNC_URL}
QUERY_VERSION_SYNC_URL: ${QUERY_VERSION_SYNC_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
PGUSER: ${PGUSER}
PGPASSWORD: ${PGPASSWORD}
PGHOST: aws-db
PGPORT: 5432
PGDATABASE: ${PGDATABASE}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- aws-server-network
client_sync_runner:
image: queryorchestration:latest
command: ["./clientSyncRunner"]
depends_on:
localstack:
condition: service_healthy
aws-db:
condition: service_healthy
ports:
- "8089:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
QUEUE_URL: ${CLIENT_SYNC_URL}
DOCUMENT_SYNC_URL: ${DOCUMENT_SYNC_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
PGUSER: ${PGUSER}
PGPASSWORD: ${PGPASSWORD}
PGHOST: aws-db
PGPORT: 5432
PGDATABASE: ${PGDATABASE}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- aws-server-network
query_version_sync_runner:
image: queryorchestration:latest
command: ["./queryVersionSyncRunner"]
depends_on:
localstack:
condition: service_healthy
aws-db:
condition: service_healthy
ports:
- "8090:8080"
expose:
- 8080
environment:
LOG_LEVEL: DEBUG
QUEUE_URL: ${QUERY_VERSION_SYNC_URL}
CLIENT_SYNC_URL: ${CLIENT_SYNC_URL}
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
PGUSER: ${PGUSER}
PGPASSWORD: ${PGPASSWORD}
PGHOST: aws-db
PGPORT: 5432
PGDATABASE: ${PGDATABASE}
DB_NOSSL: ${DB_NOSSL}
AWS_ENDPOINT_URL: "http://localstack:4566"
AWS_S3_USE_PATH_STYLE: true
networks:
- aws-server-network
prometheus:
image: prom/prometheus:latest
ports:
- "9091:9090"
volumes:
- ./prometheus.allmetrics.yaml:/etc/prometheus/prometheus.yml
healthcheck:
test: ["CMD", "wget", "--spider", "-q", "http://localhost:9090/-/healthy"]
interval: 30s
timeout: 10s
retries: 3
start_period: 30s
networks:
- aws-server-network
aws-db:
image: postgres:17.2-alpine3.21
ports:
- 5432:5432
environment:
POSTGRES_B: ${PGDATABASE}
POSTGRES_PASSWORD: ${PGPASSWORD}
POSTGRES_USER: ${PGUSER}
expose:
- 5432
volumes:
- aws-db-data:/var/lib/postgresql/data
healthcheck:
test: ["CMD-SHELL", "pg_isready -U ${PGUSER}"]
interval: 10s
timeout: 5s
retries: 5
networks:
- aws-server-network
localstack:
image: localstack/localstack:4.1.0
ports:
- 4566:4566
environment:
AWS_ACCESS_KEY_ID: ${AWS_ACCESS_KEY_ID}
AWS_SECRET_ACCESS_KEY: ${AWS_SECRET_ACCESS_KEY}
AWS_SESSION_TOKEN: ${AWS_SESSION_TOKEN}
AWS_REGION: ${AWS_REGION}
SERVICES: "s3,sqs"
SKIP_SSL_CERT_DOWNLOAD: "1"
LOCALSTACK_HOST: localstack
SQS_ENDPOINT_STRATEGY: "path"
EAGER_SERVICE_LOADING: "1"
expose:
- 4566
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:4566/_localstack/health"]
interval: 5s
timeout: 5s
retries: 10
networks:
- aws-server-network
volumes:
aws-db-data:
networks:
aws-server-network:
+1 -1
View File
@@ -129,7 +129,7 @@ require (
dario.cat/mergo v1.0.1 // indirect
github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c // indirect
github.com/Microsoft/go-winio v0.6.2 // indirect
github.com/aws/aws-sdk-go-v2/credentials v1.17.62 // indirect
github.com/aws/aws-sdk-go-v2/credentials v1.17.62
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.30 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.34 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.34 // indirect
+2 -2
View File
@@ -14,7 +14,7 @@ type Client struct {
CanSync bool
}
func (c *Client) normalizeCanSyncUpdate(n **bool) {
func (c *Client) NormalizeCanSyncUpdate(n **bool) {
if n == nil || *n == nil {
return
}
@@ -24,7 +24,7 @@ func (c *Client) normalizeCanSyncUpdate(n **bool) {
}
}
func (c *Client) normalizeNameUpdate(n **string) error {
func (c *Client) NormalizeNameUpdate(n **string) error {
if n == nil || *n == nil {
return nil
}
+8 -8
View File
@@ -39,19 +39,19 @@ func TestNormalizeCanSyncUpdate(t *testing.T) {
CanSync: true,
}
c.normalizeCanSyncUpdate(nil)
c.NormalizeCanSyncUpdate(nil)
var val *bool
c.normalizeCanSyncUpdate(&val)
c.NormalizeCanSyncUpdate(&val)
assert.Nil(t, val)
v := true
val = &v
c.normalizeCanSyncUpdate(&val)
c.NormalizeCanSyncUpdate(&val)
assert.Nil(t, val)
v = false
val = &v
c.normalizeCanSyncUpdate(&val)
c.NormalizeCanSyncUpdate(&val)
assert.NotNil(t, *val)
}
@@ -61,24 +61,24 @@ func TestNormalizeNameUpdate(t *testing.T) {
Name: "example_name",
}
err := c.normalizeNameUpdate(nil)
err := c.NormalizeNameUpdate(nil)
require.NoError(t, err)
v := "update_name"
val := &v
err = c.normalizeNameUpdate(&val)
err = c.NormalizeNameUpdate(&val)
require.NoError(t, err)
assert.NotNil(t, val)
assert.Equal(t, "update_name", *val)
v = c.Name
val = &v
err = c.normalizeNameUpdate(&val)
err = c.NormalizeNameUpdate(&val)
require.NoError(t, err)
assert.Nil(t, val)
v = "###"
val = &v
err = c.normalizeNameUpdate(&val)
err = c.NormalizeNameUpdate(&val)
assert.Error(t, err)
}
+28
View File
@@ -0,0 +1,28 @@
package clientupdate
import (
"queryorchestration/internal/client"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/clientsync"
)
type ConfigProvider interface {
serviceconfig.ConfigProvider
clientsync.ConfigProvider
}
type Services struct {
Client *client.Service
}
type Service struct {
cfg ConfigProvider
svc *Services
}
func New(cfg ConfigProvider, svc *Services) *Service {
return &Service{
cfg,
svc,
}
}
+24
View File
@@ -0,0 +1,24 @@
package clientupdate_test
import (
"testing"
"queryorchestration/internal/client"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/serviceconfig"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestService(t *testing.T) {
pool, err := pgxmock.NewPool()
require.NoError(t, err)
cfg := &serviceconfig.BaseConfig{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := client.New(cfg)
assert.NotNil(t, svc)
}
@@ -1,11 +1,14 @@
package client
package clientupdate
import (
"context"
"errors"
"log/slog"
clientsyncrunner "queryorchestration/api/clientSyncRunner"
"queryorchestration/internal/client"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/serviceconfig/queue"
"queryorchestration/internal/validation"
)
@@ -15,7 +18,7 @@ type Update struct {
}
func (s *Service) Update(ctx context.Context, id string, entity *Update) error {
current, err := s.Get(ctx, id)
current, err := s.svc.Client.Get(ctx, id)
if err != nil {
return err
}
@@ -30,9 +33,27 @@ func (s *Service) Update(ctx context.Context, id string, entity *Update) error {
return err
}
err = s.informUpdate(ctx, id, entity)
if err != nil {
return err
}
return nil
}
func (s *Service) informUpdate(ctx context.Context, id string, entity *Update) error {
if entity.CanSync == nil || !*entity.CanSync {
return nil
}
return s.cfg.SendToQueue(ctx, &queue.SendParams{
QueueURL: s.cfg.GetClientSyncURL(),
Body: clientsyncrunner.Body{
ClientID: id,
},
})
}
func (s *Service) submitUpdate(ctx context.Context, id string, entity *Update) error {
return s.cfg.ExecuteDBTransaction(ctx, func(ctx context.Context, q *repository.Queries) error {
id := id
@@ -63,17 +84,17 @@ func (s *Service) submitUpdate(ctx context.Context, id string, entity *Update) e
})
}
func (s *Service) normalizeUpdateParams(id string, current *Client, entity *Update) error {
func (s *Service) normalizeUpdateParams(id string, current *client.Client, entity *Update) error {
if entity == nil || current == nil || current.ID != id {
return errors.New("no updates presented")
}
err := current.normalizeNameUpdate(&entity.Name)
err := current.NormalizeNameUpdate(&entity.Name)
if err != nil {
return err
}
current.normalizeCanSyncUpdate(&entity.CanSync)
current.NormalizeCanSyncUpdate(&entity.CanSync)
if validation.AreAllPointersNilExcept(entity) {
return errors.New("no updates presented")
@@ -1,29 +1,42 @@
package client
package clientupdate
import (
"context"
"fmt"
"testing"
"queryorchestration/internal/client"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/queue/clientsync"
queuemock "queryorchestration/mocks/queue"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
"github.com/stretchr/testify/require"
)
type Config struct {
serviceconfig.BaseConfig
clientsync.ClientSyncConfig
}
func TestUpdate(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
require.NoError(t, err)
cfg := &serviceconfig.BaseConfig{}
cfg := &Config{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{
Client: client.New(cfg),
})
c := Client{
c := client.Client{
ID: "huhu",
Name: "example_name",
CanSync: false,
@@ -72,7 +85,7 @@ func TestNormalizeUpdateParams(t *testing.T) {
err := svc.normalizeUpdateParams("huhu", nil, nil)
assert.Error(t, err)
current := &Client{
current := &client.Client{
ID: "hi",
Name: "client",
}
@@ -109,13 +122,13 @@ func TestSubmitUpdate(t *testing.T) {
pool, err := pgxmock.NewPool()
require.NoError(t, err)
cfg := &serviceconfig.BaseConfig{}
cfg := &Config{}
cfg.DBPool = pool
cfg.DBQueries = repository.New(pool)
svc := New(cfg)
svc := New(cfg, &Services{})
c := Client{
c := client.Client{
ID: "hello",
Name: "example_name",
CanSync: false,
@@ -151,3 +164,49 @@ func TestSubmitUpdate(t *testing.T) {
err = svc.submitUpdate(ctx, c.ID, &update)
require.NoError(t, err)
}
func TestInformUpdate(t *testing.T) {
ctx := context.Background()
cfg := &Config{}
mockSQS := queuemock.NewMockSQSClient(t)
cfg.QueueClient = mockSQS
svc := New(cfg, &Services{})
id := "HIELO"
t.Run("no cansync", func(t *testing.T) {
update := Update{}
err := svc.informUpdate(ctx, id, &update)
require.NoError(t, err)
})
t.Run("false cansync", func(t *testing.T) {
cansync := false
update := Update{
CanSync: &cansync,
}
err := svc.informUpdate(ctx, id, &update)
require.NoError(t, err)
})
t.Run("true cansync", func(t *testing.T) {
cansync := true
update := Update{
CanSync: &cansync,
}
mockSQS.EXPECT().
SendMessage(
mock.Anything,
mock.MatchedBy(func(in *sqs.SendMessageInput) bool {
return *in.QueueUrl == cfg.ClientSyncURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", id)
}),
mock.Anything,
).
Return(&sqs.SendMessageOutput{}, nil)
err := svc.informUpdate(ctx, id, &update)
require.NoError(t, err)
})
}
+4 -2
View File
@@ -19,7 +19,8 @@ type Params struct {
type EventS3 string
const (
EventS3ObjectCreatedPut = "ObjectCreated:Put"
EventS3ObjectCreatedPut = "ObjectCreated:Put"
EventS3ObjectCreatedCompleteMultipartUpload = "ObjectCreated:CompleteMultipartUpload"
)
func (s *Service) Process(ctx context.Context, params Params) error {
@@ -52,5 +53,6 @@ func (s *Service) Process(ctx context.Context, params Params) error {
}
func (s *Service) isSupportedEvent(name EventS3) bool {
return name == EventS3ObjectCreatedPut
return name == EventS3ObjectCreatedPut ||
name == EventS3ObjectCreatedCompleteMultipartUpload
}
+17 -2
View File
@@ -2,10 +2,12 @@ package textract
import (
"context"
"log/slog"
awsc "queryorchestration/internal/serviceconfig/aws"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/textract"
)
@@ -20,8 +22,12 @@ type ConfigProvider interface {
}
type TextractConfig struct {
AWSEndpointUrlTextract string `env:"AWS_ENDPOINT_URL_TEXTRACT"`
TextractClient TextractClient
AWSEndpointUrlTextract string `env:"AWS_ENDPOINT_URL_TEXTRACT"`
AWSTextractAccessKeyId string `env:"AWS_ACCESS_KEY_ID_TEXTRACT"`
AWSTextractSecretAccessKey string `env:"AWS_SECRET_ACCESS_KEY_TEXTRACT"`
AWSTextractSessionToken string `env:"AWS_SESSION_TOKEN_TEXTRACT"`
AWSTextractRegion string `env:"AWS_REGION_TEXTRACT"`
TextractClient TextractClient
}
func (c *TextractConfig) SetTextractClientWithCfg(ctx context.Context, cfg aws.Config) {
@@ -45,6 +51,15 @@ func (c *TextractConfig) SetTextractClient(ctx context.Context) error {
return err
}
if c.AWSTextractAccessKeyId != "" && c.AWSTextractSecretAccessKey != "" {
cfg.Credentials = credentials.NewStaticCredentialsProvider(c.AWSTextractAccessKeyId, c.AWSTextractSecretAccessKey, c.AWSTextractSessionToken)
if c.AWSTextractRegion != "" {
cfg.Region = c.AWSTextractRegion
}
slog.Debug("textract specific credentials", "key_id", c.AWSTextractAccessKeyId, "region", c.AWSTextractRegion)
}
c.SetTextractClientWithCfg(ctx, cfg)
return nil
+19 -5
View File
@@ -21,11 +21,25 @@ func TestGetTextractClient(t *testing.T) {
func TestTextractClient(t *testing.T) {
os.Clearenv()
ctx := context.Background()
c := textract.TextractConfig{}
err := c.SetTextractClient(ctx)
require.NoError(t, err)
assert.NotNil(t, c.TextractClient)
t.Run("stndard", func(t *testing.T) {
ctx := context.Background()
c := textract.TextractConfig{}
err := c.SetTextractClient(ctx)
require.NoError(t, err)
assert.NotNil(t, c.TextractClient)
})
t.Run("textract specific", func(t *testing.T) {
ctx := context.Background()
c := textract.TextractConfig{
AWSTextractAccessKeyId: "not empty",
AWSTextractSecretAccessKey: "not empty",
AWSTextractSessionToken: "not empty",
AWSTextractRegion: "not empty",
}
err := c.SetTextractClient(ctx)
require.NoError(t, err)
assert.NotNil(t, c.TextractClient)
})
}
func TestTextractClientWithProfile(t *testing.T) {
+5
View File
@@ -21,6 +21,7 @@ import (
)
func SetQueueClient(t testing.TB, ctx context.Context, cfg queue.ConfigProvider) {
t.Helper()
err := cfg.SetQueueClient(ctx)
require.NoError(t, err)
}
@@ -34,6 +35,7 @@ func GetQueueArn(cfg awsc.ConfigProvider, name RunnerName) string {
}
func CreateQueue(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, name string) string {
t.Helper()
queueM, err := cfg.GetQueueClient().CreateQueue(ctx, &sqs.CreateQueueInput{
QueueName: aws.String(name),
})
@@ -47,6 +49,7 @@ 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 {
t.Helper()
timeout := time.After(30 * time.Second)
tick := time.NewTicker(2 * time.Second)
defer tick.Stop()
@@ -73,6 +76,7 @@ func AssertMessage(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigPr
}
func AssertMessageBody(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, url string, body *regexp.Regexp) {
t.Helper()
message := AssertMessage(t, ctx, cfg, &queue.ReceiveParams{
QueueURL: url,
})
@@ -81,6 +85,7 @@ func AssertMessageBody(t testing.TB, ctx context.Context, cfg serviceconfig.Conf
}
func AssertMessageAttr(t testing.TB, ctx context.Context, cfg serviceconfig.ConfigProvider, url string, name string, value *regexp.Regexp) {
t.Helper()
message := AssertMessage(t, ctx, cfg, &queue.ReceiveParams{
QueueURL: url,
Attributes: []string{name},
+8
View File
@@ -85,3 +85,11 @@ tasks:
cmds:
- gomarkdoc --template-file file=docs/templates/root.gotxt ./cmd/...
aws:login: aws sso login
aws:textract:credentials:
- |
sed -i '/^AWS_REGION_TEXTRACT\|^AWS_ACCESS_KEY_ID_TEXTRACT\|^AWS_SECRET_ACCESS_KEY_TEXTRACT\|^AWS_SESSION_TOKEN_TEXTRACT/d' .env
echo "AWS_REGION_TEXTRACT=$(aws configure get region --profile $AWS_PROFILE)" >> .env
CREDENTIALS_JSON=$(aws configure export-credentials --profile $AWS_PROFILE)
echo "AWS_ACCESS_KEY_ID_TEXTRACT=$(echo $CREDENTIALS_JSON | jq -r '.AccessKeyId')" >> .env
echo "AWS_SECRET_ACCESS_KEY_TEXTRACT=$(echo $CREDENTIALS_JSON | jq -r '.SecretAccessKey')" >> .env
echo "AWS_SESSION_TOKEN_TEXTRACT=$(echo $CREDENTIALS_JSON | jq -r '.SessionToken')" >> .env
+25 -1
View File
@@ -7,6 +7,7 @@ vars:
LOCAL_COMPOSE_FILE: "deployments/compose.local.yaml"
TEST_COMPOSE_FILE: "deployments/compose.test.yaml"
GENERATE_COMPOSE_FILE: "deployments/compose.generate.yaml"
AWS_COMPOSE_FILE: "deployments/compose.aws.yaml"
tasks:
build:test:
@@ -41,14 +42,25 @@ tasks:
- task: up:cmd
vars:
COMPOSE_FILE: "{{.GENERATE_COMPOSE_FILE}}"
up:aws:
cmds:
- task deps:tidy
- task docker:build
- task: build
- task: up:cmd
vars:
COMPOSE_FILE: "{{.AWS_COMPOSE_FILE}}"
- task: init
up:
cmds:
- task deps:tidy
- task docker:build
- task: build
- task: up:cmd
- task: init
refresh:
cmds:
- task deps:tidy
- task docker:build
- task: build
- task: down
@@ -73,7 +85,7 @@ tasks:
--notification-configuration '{
"QueueConfigurations": [{
"QueueArn": "arn:aws:sqs:us-east-1:000000000000:store_event",
"Events": ["s3:ObjectCreated:Put"]
"Events": ["s3:ObjectCreated:*"]
}]
}'
down:test:
@@ -86,6 +98,11 @@ tasks:
- task: down
vars:
COMPOSE_FILE: "{{.GENERATE_COMPOSE_FILE}}"
down:aws:
cmds:
- task: down
vars:
COMPOSE_FILE: "{{.AWS_COMPOSE_FILE}}"
down:
vars:
COMPOSE_FILE: "{{.COMPOSE_FILE | default .LOCAL_COMPOSE_FILE}}"
@@ -100,6 +117,11 @@ tasks:
cmds:
- task: down:generate
- docker volume rm -f deployments_generate-db-data
clean:aws:
run: once
cmds:
- task: down:aws
- docker volume rm -f deployments_aws-db-data
clean:
cmds:
- task: down
@@ -113,8 +135,10 @@ tasks:
- |
docker compose -f {{.GENERATE_COMPOSE_FILE}} config -q \
--resolve-image-digests
- docker compose -f {{.AWS_COMPOSE_FILE}} config -q
lint:short:
cmds:
- docker compose -f {{.LOCAL_COMPOSE_FILE}} config -q
- docker compose -f {{.TEST_COMPOSE_FILE}} config -q
- docker compose -f {{.GENERATE_COMPOSE_FILE}} config -q
- docker compose -f {{.AWS_COMPOSE_FILE}} config -q