diff --git a/.yamllint.yml b/.yamllint.yml index d9d1b77d..5ff6f7c7 100644 --- a/.yamllint.yml +++ b/.yamllint.yml @@ -4,3 +4,6 @@ ignore: | vendor/ .devbox/ out/ +rules: + line-length: + max: 140 diff --git a/api/queryAPI/client.go b/api/queryAPI/client.go index 83526222..32397c31 100644 --- a/api/queryAPI/client.go +++ b/api/queryAPI/client.go @@ -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, }) diff --git a/api/queryAPI/client_test.go b/api/queryAPI/client_test.go index 1bd6c83e..d4b4d163 100644 --- a/api/queryAPI/client_test.go +++ b/api/queryAPI/client_test.go @@ -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 diff --git a/api/queryAPI/controllers.go b/api/queryAPI/controllers.go index dd750bcf..ea1c8b5e 100644 --- a/api/queryAPI/controllers.go +++ b/api/queryAPI/controllers.go @@ -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 } diff --git a/api/storeEventRunner/runner.go b/api/storeEventRunner/runner.go index fbc9e48d..d5132cc7 100644 --- a/api/storeEventRunner/runner.go +++ b/api/storeEventRunner/runner.go @@ -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 { diff --git a/build/Dockerfile b/build/Dockerfile index a004e417..7e9a1fb2 100644 --- a/build/Dockerfile +++ b/build/Dockerfile @@ -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 diff --git a/cmd/queryAPI/main.go b/cmd/queryAPI/main.go index d4348426..c00aa7e1 100644 --- a/cmd/queryAPI/main.go +++ b/cmd/queryAPI/main.go @@ -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, } diff --git a/deployments/compose.aws.yaml b/deployments/compose.aws.yaml new file mode 100644 index 00000000..765f4961 --- /dev/null +++ b/deployments/compose.aws.yaml @@ -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: diff --git a/go.mod b/go.mod index 266bfc78..f9f3ba9c 100644 --- a/go.mod +++ b/go.mod @@ -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 diff --git a/internal/client/service.go b/internal/client/service.go index 055d2483..6840b9b6 100644 --- a/internal/client/service.go +++ b/internal/client/service.go @@ -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 } diff --git a/internal/client/serviceprivate_test.go b/internal/client/serviceprivate_test.go index d03da022..8910e488 100644 --- a/internal/client/serviceprivate_test.go +++ b/internal/client/serviceprivate_test.go @@ -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) } diff --git a/internal/client/update/service.go b/internal/client/update/service.go new file mode 100644 index 00000000..aa43610f --- /dev/null +++ b/internal/client/update/service.go @@ -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, + } +} diff --git a/internal/client/update/service_test.go b/internal/client/update/service_test.go new file mode 100644 index 00000000..e836bc94 --- /dev/null +++ b/internal/client/update/service_test.go @@ -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) +} diff --git a/internal/client/update.go b/internal/client/update/update.go similarity index 63% rename from internal/client/update.go rename to internal/client/update/update.go index 90f495e6..b9c4be87 100644 --- a/internal/client/update.go +++ b/internal/client/update/update.go @@ -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") diff --git a/internal/client/update_test.go b/internal/client/update/update_test.go similarity index 67% rename from internal/client/update_test.go rename to internal/client/update/update_test.go index 7b6af969..9a74889e 100644 --- a/internal/client/update_test.go +++ b/internal/client/update/update_test.go @@ -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) + }) +} diff --git a/internal/document/store/process.go b/internal/document/store/process.go index 30ac61af..ba6b6f56 100644 --- a/internal/document/store/process.go +++ b/internal/document/store/process.go @@ -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 } diff --git a/internal/serviceconfig/textract/config.go b/internal/serviceconfig/textract/config.go index 1bc88d8e..987dbc6c 100644 --- a/internal/serviceconfig/textract/config.go +++ b/internal/serviceconfig/textract/config.go @@ -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 diff --git a/internal/serviceconfig/textract/config_test.go b/internal/serviceconfig/textract/config_test.go index 2b9acf10..733e8988 100644 --- a/internal/serviceconfig/textract/config_test.go +++ b/internal/serviceconfig/textract/config_test.go @@ -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) { diff --git a/internal/test/queue.go b/internal/test/queue.go index 88a53d91..89963e40 100644 --- a/internal/test/queue.go +++ b/internal/test/queue.go @@ -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}, diff --git a/scripts/Taskfile.yml b/scripts/Taskfile.yml index d8820adc..681d6cf3 100644 --- a/scripts/Taskfile.yml +++ b/scripts/Taskfile.yml @@ -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 diff --git a/scripts/local-deployments.yml b/scripts/local-deployments.yml index 50f46b13..cfaae72f 100644 --- a/scripts/local-deployments.yml +++ b/scripts/local-deployments.yml @@ -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