diff --git a/api/docCleanRunner/runner.go b/api/docCleanRunner/runner.go new file mode 100644 index 00000000..fc8a1b02 --- /dev/null +++ b/api/docCleanRunner/runner.go @@ -0,0 +1,54 @@ +package doccleanrunner + +import ( + "context" + "encoding/json" + documentclean "queryorchestration/internal/document/clean" + + "github.com/go-playground/validator/v10" + "github.com/google/uuid" + + "github.com/aws/aws-sdk-go-v2/service/sqs/types" +) + +const Name = "docCleanRunner" + +type Services struct { + Clean *documentclean.Service +} + +type Runner struct { + validator *validator.Validate + svc *Services +} + +func New(validator *validator.Validate, svc *Services) Runner { + return Runner{ + validator: validator, + svc: svc, + } +} + +type Create struct { + ID uuid.UUID `json:"id" validate:"required,uuid"` +} + +func (s Runner) Process(ctx context.Context, req *types.Message) error { + var body Create + err := json.Unmarshal([]byte(*req.Body), &body) + if err != nil { + return err + } + + err = s.validator.Struct(body) + if err != nil { + return err + } + + err = s.svc.Clean.Create(ctx, body.ID) + if err != nil { + return err + } + + return nil +} diff --git a/api/docCleanRunner/runner_test.go b/api/docCleanRunner/runner_test.go new file mode 100644 index 00000000..b559aed1 --- /dev/null +++ b/api/docCleanRunner/runner_test.go @@ -0,0 +1,96 @@ +package doccleanrunner_test + +import ( + "context" + "encoding/json" + "fmt" + doccleanrunner "queryorchestration/api/docCleanRunner" + "queryorchestration/internal/database" + "queryorchestration/internal/database/repository" + "queryorchestration/internal/document" + documentclean "queryorchestration/internal/document/clean" + "queryorchestration/internal/serviceconfig" + "queryorchestration/internal/serviceconfig/objectstore" + "queryorchestration/internal/serviceconfig/queue/documenttext" + objectstoremock "queryorchestration/mocks/objectstore" + queuemock "queryorchestration/mocks/queue" + "testing" + + "github.com/aws/aws-sdk-go-v2/service/sqs" + "github.com/aws/aws-sdk-go-v2/service/sqs/types" + "github.com/go-playground/validator/v10" + "github.com/google/uuid" + "github.com/pashagolub/pgxmock/v3" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" +) + +type DocTextConfig struct { + serviceconfig.BaseConfig + documenttext.DocTextConfig + objectstore.ObjectStoreConfig +} + +func TestDocCleanRunner(t *testing.T) { + ctx := context.Background() + + pool, err := pgxmock.NewPool() + if err != nil { + t.Fatalf("failed to open pgxmock database: %v", err) + } + + cfg := &DocTextConfig{} + cfg.DBPool = pool + cfg.DBQueries = repository.New(pool) + mockStore := objectstoremock.NewMockS3Client(t) + cfg.StoreClient = mockStore + mockSQS := queuemock.NewMockSQSClient(t) + cfg.QueueClient = mockSQS + cfg.DocumentTextURL = "/i/am/here" + + runner := doccleanrunner.New(validator.New(), &doccleanrunner.Services{ + Clean: documentclean.New(cfg, &documentclean.Services{ + Document: document.New(cfg), + }), + }) + assert.NotNil(t, runner) + + doc := doccleanrunner.Create{ + ID: uuid.New(), + } + bodyBytes, err := json.Marshal(doc) + assert.NoError(t, err) + body := string(bodyBytes) + msg := &types.Message{ + Body: &body, + } + + inloc := document.Location{ + Bucket: "bucket_name", + Key: "/i/am/here", + } + + pool.ExpectQuery("name: IsDocumentClean :one").WithArgs(database.MustToDBUUID(doc.ID)).WillReturnRows( + pgxmock.NewRows([]string{"isclean"}). + AddRow(false), + ) + pool.ExpectQuery("name: GetDocumentEntry :one").WithArgs(database.MustToDBUUID(doc.ID)).WillReturnRows( + pgxmock.NewRows([]string{"documentId", "bucket", "key"}). + AddRow(database.MustToDBUUID(doc.ID), inloc.Bucket, inloc.Key), + ) + pool.ExpectExec("name: AddDocumentCleanEntry :exec").WithArgs(database.MustToDBUUID(doc.ID), int32(1), inloc.Bucket, inloc.Key). + WillReturnResult(pgxmock.NewResult("", 1)) + + mockSQS.EXPECT(). + SendMessage( + mock.Anything, + mock.MatchedBy(func(in *sqs.SendMessageInput) bool { + return *in.QueueUrl == cfg.DocumentTextURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", doc.ID.String()) + }), + mock.Anything, + ). + Return(&sqs.SendMessageOutput{}, nil) + + err = runner.Process(ctx, msg) + assert.NoError(t, err) +} diff --git a/api/docInitRunner/runner.go b/api/docInitRunner/runner.go index 5c6a79d2..7bad2eca 100644 --- a/api/docInitRunner/runner.go +++ b/api/docInitRunner/runner.go @@ -6,6 +6,7 @@ import ( "errors" "fmt" "log/slog" + "queryorchestration/internal/document" documentinit "queryorchestration/internal/document/init" "github.com/go-playground/validator/v10" @@ -31,26 +32,33 @@ func New(validator *validator.Validate, svc *Services) Runner { } } +type S3Bucket struct { + Name string `json:"name"` + Arn string `json:"arn"` +} + +type S3Object struct { + Key string `json:"key"` + Size int64 `json:"size"` + ETag string `json:"eTag"` + VersionId string `json:"versionId"` +} + +type S3EventRecordDetails struct { + Bucket S3Bucket `json:"bucket"` + Object S3Object `json:"object"` +} + +type S3EventRecord struct { + EventVersion string `json:"eventVersion"` + EventSource string `json:"eventSource"` + AwsRegion string `json:"awsRegion"` + EventTime string `json:"eventTime"` + EventName string `json:"eventName"` + S3 S3EventRecordDetails `json:"s3"` +} type S3EventNotification struct { - Records []struct { - EventVersion string `json:"eventVersion"` - EventSource string `json:"eventSource"` - AwsRegion string `json:"awsRegion"` - EventTime string `json:"eventTime"` - EventName string `json:"eventName"` - S3 struct { - Bucket struct { - Name string `json:"name"` - Arn string `json:"arn"` - } `json:"bucket"` - Object struct { - Key string `json:"key"` - Size int64 `json:"size"` - ETag string `json:"eTag"` - VersionId string `json:"versionId"` - } `json:"object"` - } `json:"s3"` - } `json:"Records"` + Records []S3EventRecord `json:"Records"` } func (s Runner) Process(ctx context.Context, req *types.Message) error { @@ -79,10 +87,12 @@ func (s Runner) Process(ctx context.Context, req *types.Message) error { } _, err = s.svc.Document.Create(ctx, &documentinit.Create{ - JobID: jobID, - Location: record.S3.Object.Key, - Bucket: record.S3.Bucket.Name, - Hash: record.S3.Object.ETag, + JobID: jobID, + Location: document.Location{ + Bucket: record.S3.Bucket.Name, + Key: record.S3.Object.Key, + }, + Hash: record.S3.Object.ETag, }) if err != nil { return err diff --git a/api/docInitRunner/runner_test.go b/api/docInitRunner/runner_test.go index 1b1ea896..53fde512 100644 --- a/api/docInitRunner/runner_test.go +++ b/api/docInitRunner/runner_test.go @@ -3,6 +3,7 @@ package docinitrunner_test import ( "context" "encoding/json" + "fmt" docinitrunner "queryorchestration/api/docInitRunner" "queryorchestration/internal/client" "queryorchestration/internal/database" @@ -58,10 +59,27 @@ func TestDocInitRunner(t *testing.T) { ID: uuid.New(), ClientID: uuid.New(), } - doc := documentinit.Create{ - JobID: j.ID, - Bucket: "bucket_name", - Location: "/I/am/here", + bucketName := "bucketName" + location := fmt.Sprintf("aaa/%s/aaa", j.ID.String()) + docinfo := document.Document{ + ID: uuid.New(), + Hash: "example_hash", + } + doc := docinitrunner.S3EventNotification{ + Records: []docinitrunner.S3EventRecord{ + { + EventName: "ObjectCreated:Put", + S3: docinitrunner.S3EventRecordDetails{ + Bucket: docinitrunner.S3Bucket{ + Name: bucketName, + }, + Object: docinitrunner.S3Object{ + Key: location, + ETag: docinfo.Hash, + }, + }, + }, + }, } bodyBytes, err := json.Marshal(doc) assert.NoError(t, err) @@ -70,29 +88,24 @@ func TestDocInitRunner(t *testing.T) { Body: &body, } - docinfo := document.Document{ - ID: uuid.New(), - Hash: "example_hash", - } - pool.ExpectQuery("name: GetJob :one").WithArgs(database.MustToDBUUID(j.ID)).WillReturnRows( pgxmock.NewRows([]string{"id", "clientId", "canSync"}). AddRow(database.MustToDBUUID(j.ID), database.MustToDBUUID(j.ClientID), j.CanSync), ) - pool.ExpectQuery("-- name: GetClient :one").WithArgs(database.MustToDBUUID(j.ClientID)).WillReturnRows( + pool.ExpectQuery("name: GetClient :one").WithArgs(database.MustToDBUUID(j.ClientID)).WillReturnRows( pgxmock.NewRows([]string{"id", "name", "canSync"}). AddRow(database.MustToDBUUID(j.ClientID), "client_name", true), ) - pool.ExpectQuery("-- name: GetDocumentIDByHash :one").WithArgs(docinfo.Hash, database.MustToDBUUID(doc.JobID)).WillReturnRows( + pool.ExpectQuery("name: GetDocumentIDByHash :one").WithArgs(docinfo.Hash, database.MustToDBUUID(j.ID)).WillReturnRows( pgxmock.NewRows([]string{"id"}), ) pool.ExpectBegin() - pool.ExpectQuery("name: CreateDocument :one").WithArgs(database.MustToDBUUID(doc.JobID), docinfo.Hash). + pool.ExpectQuery("name: CreateDocument :one").WithArgs(database.MustToDBUUID(j.ID), docinfo.Hash). WillReturnRows( pgxmock.NewRows([]string{"id"}). AddRow(database.MustToDBUUID(docinfo.ID)), ) - pool.ExpectExec("name: AddDocumentEntry :exec").WithArgs(database.MustToDBUUID(docinfo.ID), doc.Bucket, doc.Location). + pool.ExpectExec("name: AddDocumentEntry :exec").WithArgs(database.MustToDBUUID(docinfo.ID), bucketName, location). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectCommit() diff --git a/api/queryService/client_test.go b/api/queryService/client_test.go index 1edc933b..cea3f26c 100644 --- a/api/queryService/client_test.go +++ b/api/queryService/client_test.go @@ -48,7 +48,7 @@ func TestCreateClient(t *testing.T) { id := uuid.New() - pool.ExpectQuery("-- name: CreateClient :one").WithArgs(body.Name).WillReturnRows( + pool.ExpectQuery("name: CreateClient :one").WithArgs(body.Name).WillReturnRows( pgxmock.NewRows([]string{"id"}). AddRow(database.MustToDBUUID(id)), ) @@ -82,7 +82,7 @@ func TestGetClient(t *testing.T) { ctx.Set("id", id) - pool.ExpectQuery("-- name: GetClient :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( + pool.ExpectQuery("name: GetClient :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( pgxmock.NewRows([]string{"id", "name", "canSync"}). AddRow(database.MustToDBUUID(id), "client_name", true), ) @@ -132,7 +132,7 @@ func TestUpdateClient(t *testing.T) { ctx.Set("id", id) - pool.ExpectQuery("-- name: GetClient :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( + pool.ExpectQuery("name: GetClient :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( pgxmock.NewRows([]string{"id", "name", "canSync"}). AddRow(database.MustToDBUUID(id), "client_name", true), ) diff --git a/api/queryService/job_test.go b/api/queryService/job_test.go index 5dffb944..5b320b76 100644 --- a/api/queryService/job_test.go +++ b/api/queryService/job_test.go @@ -9,8 +9,7 @@ import ( "queryorchestration/internal/client" "queryorchestration/internal/database" "queryorchestration/internal/database/repository" - documentclean "queryorchestration/internal/document/clean" - documenttext "queryorchestration/internal/document/text" + "queryorchestration/internal/document" "queryorchestration/internal/job" "queryorchestration/internal/job/collector" "queryorchestration/internal/serviceconfig" @@ -35,12 +34,10 @@ func TestCreateJob(t *testing.T) { cfg.DBPool = pool cfg.DBQueries = repository.New(pool) - extract := documenttext.New() cons := queryservice.NewControllers(validator.New(), &queryservice.Services{ Job: job.New(cfg, &job.Services{ Collector: collector.New(cfg, &collector.Services{ - Clean: documentclean.New(), - Text: extract, + Document: document.New(cfg), }), }), }) @@ -59,7 +56,7 @@ func TestCreateJob(t *testing.T) { id := uuid.New() - pool.ExpectQuery("-- name: CreateJob :one").WithArgs(database.MustToDBUUID(body.ClientId)).WillReturnRows( + pool.ExpectQuery("name: CreateJob :one").WithArgs(database.MustToDBUUID(body.ClientId)).WillReturnRows( pgxmock.NewRows([]string{"id"}). AddRow(database.MustToDBUUID(id)), ) diff --git a/api/queryService/jobcollector_test.go b/api/queryService/jobcollector_test.go index d99c0dc0..eb2f40fb 100644 --- a/api/queryService/jobcollector_test.go +++ b/api/queryService/jobcollector_test.go @@ -7,8 +7,7 @@ import ( queryservice "queryorchestration/api/queryService" "queryorchestration/internal/database" "queryorchestration/internal/database/repository" - documentclean "queryorchestration/internal/document/clean" - documenttext "queryorchestration/internal/document/text" + "queryorchestration/internal/document" "queryorchestration/internal/job/collector" "queryorchestration/internal/serviceconfig" "strings" @@ -44,11 +43,9 @@ func TestUpdateJobCollector(t *testing.T) { rec := httptest.NewRecorder() ctx := e.NewContext(req, rec) - extract := documenttext.New() cons := queryservice.NewControllers(validator.New(), &queryservice.Services{ Collector: collector.New(cfg, &collector.Services{ - Clean: documentclean.New(), - Text: extract, + Document: document.New(cfg), }), }) @@ -91,10 +88,8 @@ func TestGetJobCollectorByJobId(t *testing.T) { rec := httptest.NewRecorder() ctx := e.NewContext(req, rec) - extract := documenttext.New() svc := collector.New(cfg, &collector.Services{ - Clean: documentclean.New(), - Text: extract, + Document: document.New(cfg), }) cons := queryservice.NewControllers(validator.New(), &queryservice.Services{ Collector: svc, diff --git a/api/queryService/query_test.go b/api/queryService/query_test.go index 7af4580d..f2a06e30 100644 --- a/api/queryService/query_test.go +++ b/api/queryService/query_test.go @@ -9,8 +9,6 @@ import ( "queryorchestration/internal/database" "queryorchestration/internal/database/repository" "queryorchestration/internal/document" - documentclean "queryorchestration/internal/document/clean" - documenttext "queryorchestration/internal/document/text" "queryorchestration/internal/job/collector" "queryorchestration/internal/query" "queryorchestration/internal/query/result" @@ -55,7 +53,7 @@ func TestCreateQuery(t *testing.T) { id := uuid.New() pool.ExpectBeginTx(pgx.TxOptions{}) - pool.ExpectQuery("-- name: CreateQuery :one").WithArgs(repository.QuerytypeContextFull).WillReturnRows( + pool.ExpectQuery("name: CreateQuery :one").WithArgs(repository.QuerytypeContextFull).WillReturnRows( pgxmock.NewRows([]string{"id"}). AddRow(database.MustToDBUUID(id)), ) @@ -132,7 +130,7 @@ func TestGetQuery(t *testing.T) { ctx.Set("id", id) - pool.ExpectQuery("-- name: GetQuery :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( + pool.ExpectQuery("name: GetQuery :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}). AddRow(database.MustToDBUUID(id), repository.QuerytypeContextFull, int32(1), int32(2), nil, []pgtype.UUID{}), ) @@ -182,7 +180,7 @@ func TestUpdateQuery(t *testing.T) { ctx.Set("id", id) - pool.ExpectQuery("-- name: GetQuery :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( + pool.ExpectQuery("name: GetQuery :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}). AddRow(database.MustToDBUUID(id), repository.QuerytypeContextFull, int32(1), int32(2), []byte(""), []pgtype.UUID{}), ) @@ -209,8 +207,7 @@ func TestTestQuery(t *testing.T) { docsvc := document.New(cfg) col := collector.New(cfg, &collector.Services{ - Clean: documentclean.New(), - Text: documenttext.New(), + Document: docsvc, }) cons := queryservice.NewControllers(validator.New(), &queryservice.Services{ Collector: col, diff --git a/cmd/docCleanRunner/main.go b/cmd/docCleanRunner/main.go new file mode 100644 index 00000000..1fc0da3b --- /dev/null +++ b/cmd/docCleanRunner/main.go @@ -0,0 +1,46 @@ +package main + +import ( + "context" + "log/slog" + "os" + doccleanrunner "queryorchestration/api/docCleanRunner" + "queryorchestration/internal/document" + documentclean "queryorchestration/internal/document/clean" + "queryorchestration/internal/server/runner" + "queryorchestration/internal/serviceconfig/objectstore" + "queryorchestration/internal/serviceconfig/queue/documenttext" + + _ "github.com/lib/pq" +) + +type DocCleanConfig struct { + runner.BaseConfig + objectstore.ObjectStoreConfig + documenttext.DocTextConfig +} + +func main() { + ctx := context.Background() + + cfg := &DocCleanConfig{} + + cfg.ControllerFunc = func() runner.Controller { + doc := document.New(cfg) + clean := documentclean.New(cfg, &documentclean.Services{ + Document: doc, + }) + + return doccleanrunner.New(cfg.GetValidator(), &doccleanrunner.Services{ + Clean: clean, + }) + } + + server, err := runner.New(ctx, cfg) + if err != nil { + slog.Error(err.Error()) + os.Exit(1) + } + + server.Listen(ctx) +} diff --git a/cmd/docInitRunner/main.go b/cmd/docInitRunner/main.go index 9d5d7480..a702de2a 100644 --- a/cmd/docInitRunner/main.go +++ b/cmd/docInitRunner/main.go @@ -6,13 +6,11 @@ import ( "os" docinitrunner "queryorchestration/api/docInitRunner" "queryorchestration/internal/client" - documentclean "queryorchestration/internal/document/clean" + "queryorchestration/internal/document" documentinit "queryorchestration/internal/document/init" - documenttext "queryorchestration/internal/document/text" "queryorchestration/internal/job" "queryorchestration/internal/job/collector" "queryorchestration/internal/server/runner" - "queryorchestration/internal/serviceconfig" documentcleanc "queryorchestration/internal/serviceconfig/queue/documentclean" _ "github.com/lib/pq" @@ -28,29 +26,22 @@ func main() { cfg := &DocInitConfig{} - if err := serviceconfig.InitializeConfig(cfg); err != nil { - slog.Error("Error initializing config", "err", err) - os.Exit(1) - } - cfg.ControllerFunc = func() runner.Controller { - text := documenttext.New() - clean := documentclean.New() cli := client.New(cfg) + doc := document.New(cfg) col := collector.New(cfg, &collector.Services{ - Clean: clean, - Text: text, + Document: doc, }) j := job.New(cfg, &job.Services{ Collector: col, Client: cli, }) - doc := documentinit.New(cfg, &documentinit.Services{ + docinit := documentinit.New(cfg, &documentinit.Services{ Job: j, }) return docinitrunner.New(cfg.GetValidator(), &docinitrunner.Services{ - Document: doc, + Document: docinit, }) } diff --git a/cmd/queryRunner/main.go b/cmd/queryRunner/main.go index 4ed57fca..763e295a 100644 --- a/cmd/queryRunner/main.go +++ b/cmd/queryRunner/main.go @@ -6,13 +6,11 @@ import ( "os" controllers "queryorchestration/api/queryRunner" "queryorchestration/internal/document" - documentclean "queryorchestration/internal/document/clean" documenttext "queryorchestration/internal/document/text" "queryorchestration/internal/job/collector" "queryorchestration/internal/query" "queryorchestration/internal/query/result" "queryorchestration/internal/server/runner" - "queryorchestration/internal/serviceconfig" _ "github.com/lib/pq" ) @@ -22,24 +20,17 @@ func main() { cfg := &runner.BaseConfig{} - if err := serviceconfig.InitializeConfig(cfg); err != nil { - slog.Error("Error initializing config", "err", err) - os.Exit(1) - } - cfg.ControllerFunc = func() runner.Controller { text := documenttext.New() - clean := documentclean.New() - coll := collector.New(cfg, &collector.Services{ - Text: text, - Clean: clean, + doc := document.New(cfg) + col := collector.New(cfg, &collector.Services{ + Document: doc, }) res := result.New(cfg) - doc := document.New(cfg) svc := query.New(cfg, &query.Services{ Result: res, Text: text, - Collector: coll, + Collector: col, Document: doc, }) diff --git a/cmd/queryService/main.go b/cmd/queryService/main.go index ceecc278..9bb79ea5 100644 --- a/cmd/queryService/main.go +++ b/cmd/queryService/main.go @@ -8,7 +8,6 @@ import ( queryservice "queryorchestration/api/queryService" "queryorchestration/internal/client" "queryorchestration/internal/document" - documentclean "queryorchestration/internal/document/clean" documenttext "queryorchestration/internal/document/text" "queryorchestration/internal/export" "queryorchestration/internal/job" @@ -16,7 +15,6 @@ import ( "queryorchestration/internal/query" "queryorchestration/internal/query/result" service "queryorchestration/internal/server/service" - "queryorchestration/internal/serviceconfig" "github.com/getkin/kin-openapi/openapi3" _ "github.com/lib/pq" @@ -27,21 +25,14 @@ func main() { cfg := &service.BaseConfig{} - if err := serviceconfig.InitializeConfig(cfg); err != nil { - fmt.Printf("Error initializing config: %s", err) - os.Exit(1) - } - cfg.RegisterHandlersFunc = func() (*openapi3.T, error) { exp := export.New() extract := documenttext.New() - clean := documentclean.New() res := result.New(cfg) - col := collector.New(cfg, &collector.Services{ - Text: extract, - Clean: clean, - }) doc := document.New(cfg) + col := collector.New(cfg, &collector.Services{ + Document: doc, + }) que := query.New(cfg, &query.Services{ Text: extract, Result: res, diff --git a/database/migrations/00000000000001_extensions.up.sql b/database/migrations/00000000000001_extensions.up.sql index 7f823b17..5d6891ab 100644 --- a/database/migrations/00000000000001_extensions.up.sql +++ b/database/migrations/00000000000001_extensions.up.sql @@ -1 +1,20 @@ -CREATE EXTENSION IF NOT EXISTS "pgcrypto"; \ No newline at end of file +CREATE EXTENSION IF NOT EXISTS "pgcrypto"; + +create or replace function uuid_generate_v7() +returns uuid +as $$ +select encode( + set_bit( + set_bit( + overlay(uuid_send(gen_random_uuid()) + placing substring(int8send(floor(extract(epoch from clock_timestamp()) * 1000)::bigint) from 3) + from 1 for 6 + ), + 52, 1 + ), + 53, 1 + ), + 'hex')::uuid; +$$ +language SQL +volatile; \ No newline at end of file diff --git a/database/migrations/00000000000002_queries.up.sql b/database/migrations/00000000000002_queries.up.sql index 78caec6f..52232372 100644 --- a/database/migrations/00000000000002_queries.up.sql +++ b/database/migrations/00000000000002_queries.up.sql @@ -1,14 +1,14 @@ CREATE TYPE queryType AS ENUM ('context_full', 'json_extractor'); CREATE TABLE queries ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), latestVersion int not null default 1, activeVersion int not null default 1, type queryType not null ); CREATE TABLE requiredQueries ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), queryId uuid not null, requiredQueryId uuid not null, addedVersion int not null, @@ -19,7 +19,7 @@ CREATE TABLE requiredQueries ( ); CREATE TABLE queryConfigs ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), queryId uuid not null, config jsonb not null, addedVersion int not null, diff --git a/database/migrations/00000000000003_clients.up.sql b/database/migrations/00000000000003_clients.up.sql index f0e1a9ee..75d6d0a3 100644 --- a/database/migrations/00000000000003_clients.up.sql +++ b/database/migrations/00000000000003_clients.up.sql @@ -1,11 +1,11 @@ CREATE TABLE clients ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), name TEXT not null, UNIQUE (name) ); CREATE TABLE clientCanSync ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), clientId uuid not null, canSync boolean not null, addedAt timestamp not null default current_timestamp, diff --git a/database/migrations/00000000000004_jobs.up.sql b/database/migrations/00000000000004_jobs.up.sql index cc4a299e..352ed48a 100644 --- a/database/migrations/00000000000004_jobs.up.sql +++ b/database/migrations/00000000000004_jobs.up.sql @@ -1,11 +1,11 @@ CREATE TABLE jobs ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), clientId uuid not null, foreign key (clientId) references clients(id) ); CREATE TABLE jobCanSync ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), jobId uuid not null, canSync boolean not null, addedAt timestamp not null default current_timestamp, diff --git a/database/migrations/00000000000005_collectors.up.sql b/database/migrations/00000000000005_collectors.up.sql index ec8c5dbd..7b5c6004 100644 --- a/database/migrations/00000000000005_collectors.up.sql +++ b/database/migrations/00000000000005_collectors.up.sql @@ -1,5 +1,5 @@ CREATE TABLE collectors ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), jobId uuid not null, latestVersion int not null default 1, activeVersion int not null default 1, @@ -8,7 +8,7 @@ CREATE TABLE collectors ( ); CREATE TABLE collectorCodeVersions ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), collectorId uuid not null, minCleanVersion int not null, minTextVersion int not null, @@ -18,7 +18,7 @@ CREATE TABLE collectorCodeVersions ( ); CREATE TABLE collectorQueries ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), collectorId uuid not null, name varchar(255) not null, queryId uuid not null, diff --git a/database/migrations/00000000000006_documents.up.sql b/database/migrations/00000000000006_documents.up.sql index b64719fa..68ca6d3e 100644 --- a/database/migrations/00000000000006_documents.up.sql +++ b/database/migrations/00000000000006_documents.up.sql @@ -1,5 +1,5 @@ CREATE TABLE documents ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), jobId uuid not null, hash text not null, foreign key (jobId) references jobs(id), @@ -7,10 +7,18 @@ CREATE TABLE documents ( ); CREATE TABLE documentEntries ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), documentId uuid not null, bucket text not null, - location text not null, - createdAt timestamp default now(), + key text not null, foreign key (documentId) references documents(id) -); \ No newline at end of file +); + +CREATE TABLE documentCleans ( + id uuid primary key DEFAULT uuid_generate_v7(), + documentId uuid not null, + version int not null, + bucket text not null, + key text not null, + foreign key (documentId) references documents(id) +); diff --git a/database/migrations/00000000000007_results.up.sql b/database/migrations/00000000000007_results.up.sql index d41879f1..eee0748d 100644 --- a/database/migrations/00000000000007_results.up.sql +++ b/database/migrations/00000000000007_results.up.sql @@ -1,5 +1,5 @@ CREATE TABLE results ( - id uuid primary key DEFAULT gen_random_uuid(), + id uuid primary key DEFAULT uuid_generate_v7(), queryId uuid not null, documentId uuid not null, value TEXT not null, diff --git a/database/queries/clean.sql b/database/queries/clean.sql new file mode 100644 index 00000000..690826c0 --- /dev/null +++ b/database/queries/clean.sql @@ -0,0 +1,12 @@ +-- name: IsDocumentClean :one +SELECT EXISTS( + SELECT 1 + FROM documentCleans AS dc + JOIN documents AS d ON d.id = dc.documentId + JOIN collectors as c ON d.jobId = c.jobId + JOIN collectorCodeVersions as cv ON c.id = cv.collectorId + WHERE dc.documentId = $1 and dc.version >= cv.minCleanVersion +); + +-- name: AddDocumentCleanEntry :exec +INSERT INTO documentCleans (documentId, version, bucket, key) VALUES ($1, $2, $3, $4); \ No newline at end of file diff --git a/database/queries/document.sql b/database/queries/document.sql index 0338821c..2db398f6 100644 --- a/database/queries/document.sql +++ b/database/queries/document.sql @@ -5,7 +5,10 @@ SELECT id, jobId, hash FROM documents WHERE id = $1 LIMIT 1; INSERT INTO documents (jobId, hash) VALUES ($1, $2) RETURNING id; -- name: AddDocumentEntry :exec -INSERT INTO documentEntries (documentId, bucket, location) VALUES ($1, $2, $3); +INSERT INTO documentEntries (documentId, bucket, key) VALUES ($1, $2, $3); + +-- name: GetDocumentEntry :one +SELECT documentId, bucket, key FROM documentEntries WHERE documentId = $1 ORDER BY id DESC LIMIT 1; -- name: GetDocumentIDByHash :one -SELECT id FROM documents WHERE hash = $1 and jobId = $2; \ No newline at end of file +SELECT id FROM documents WHERE hash = $1 and jobId = $2; diff --git a/deployments/compose.local.yaml b/deployments/compose.local.yaml index 944daef2..589575e6 100644 --- a/deployments/compose.local.yaml +++ b/deployments/compose.local.yaml @@ -23,6 +23,29 @@ services: AWS_S3_USE_PATH_STYLE: true networks: - server-network + doc_clean_runner: + image: queryorchestration:latest + command: ["./docCleanRunner"] + depends_on: + - db + - localstack + environment: + QUEUE_URL: ${DOCUMENT_CLEAN_URL} + DOCUMENT_TEXT_URL: ${DOCUMENT_TEXT_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} + DB_USER: ${DB_USER} + DB_PASS: ${DB_PASS} + DB_HOST: db + DB_PORT: 5432 + DB_NAME: ${DB_NAME} + DB_NOSSL: ${DB_NOSSL} + AWS_ENDPOINT_URL: "http://localstack:4566" + AWS_S3_USE_PATH_STYLE: true + networks: + - server-network query_runner: image: queryorchestration:latest command: ["./queryRunner"] diff --git a/devbox.json b/devbox.json index b134e999..a360a08f 100644 --- a/devbox.json +++ b/devbox.json @@ -40,11 +40,13 @@ "AWS_SESSION_TOKEN": "", "AWS_REGION": "us-east-1", "AWS_ENDPOINT_URL": "http://localhost:4566", - "QNAME_DOCUMENT_CLEAN": "document_clean", "QNAME_DOCUMENT_INIT": "document_init", + "QNAME_DOCUMENT_CLEAN": "document_clean", + "QNAME_DOCUMENT_TEXT": "document_text", "QNAME_QUERY_RUNNER": "query_runner", - "DOCUMENT_CLEAN_URL": "http://localstack:4566/queue/us-east-1/000000000000/document_clean", "DOCUMENT_INIT_URL": "http://localstack:4566/queue/us-east-1/000000000000/document_init", + "DOCUMENT_CLEAN_URL": "http://localstack:4566/queue/us-east-1/000000000000/document_clean", + "DOCUMENT_TEXT_URL": "http://localstack:4566/queue/us-east-1/000000000000/document_text", "QUERY_RUNNER_URL": "http://localstack:4566/queue/us-east-1/000000000000/query_runner", "BUCKET_IN": "documentin" }, diff --git a/internal/database/repository/clean.sql.go b/internal/database/repository/clean.sql.go new file mode 100644 index 00000000..a92f92a8 --- /dev/null +++ b/internal/database/repository/clean.sql.go @@ -0,0 +1,64 @@ +// Code generated by sqlc. DO NOT EDIT. +// versions: +// sqlc v1.27.0 +// source: clean.sql + +package repository + +import ( + "context" + + "github.com/jackc/pgx/v5/pgtype" +) + +const addDocumentCleanEntry = `-- name: AddDocumentCleanEntry :exec +INSERT INTO documentCleans (documentId, version, bucket, key) VALUES ($1, $2, $3, $4) +` + +type AddDocumentCleanEntryParams struct { + Documentid pgtype.UUID `db:"documentid"` + Version int32 `db:"version"` + Bucket string `db:"bucket"` + Key string `db:"key"` +} + +// AddDocumentCleanEntry +// +// INSERT INTO documentCleans (documentId, version, bucket, key) VALUES ($1, $2, $3, $4) +func (q *Queries) AddDocumentCleanEntry(ctx context.Context, arg *AddDocumentCleanEntryParams) error { + _, err := q.db.Exec(ctx, addDocumentCleanEntry, + arg.Documentid, + arg.Version, + arg.Bucket, + arg.Key, + ) + return err +} + +const isDocumentClean = `-- name: IsDocumentClean :one +SELECT EXISTS( + SELECT 1 + FROM documentCleans AS dc + JOIN documents AS d ON d.id = dc.documentId + JOIN collectors as c ON d.jobId = c.jobId + JOIN collectorCodeVersions as cv ON c.id = cv.collectorId + WHERE dc.documentId = $1 and dc.version >= cv.minCleanVersion +) +` + +// IsDocumentClean +// +// SELECT EXISTS( +// SELECT 1 +// FROM documentCleans AS dc +// JOIN documents AS d ON d.id = dc.documentId +// JOIN collectors as c ON d.jobId = c.jobId +// JOIN collectorCodeVersions as cv ON c.id = cv.collectorId +// WHERE dc.documentId = $1 and dc.version >= cv.minCleanVersion +// ) +func (q *Queries) IsDocumentClean(ctx context.Context, documentid pgtype.UUID) (bool, error) { + row := q.db.QueryRow(ctx, isDocumentClean, documentid) + var exists bool + err := row.Scan(&exists) + return exists, err +} diff --git a/internal/database/repository/clean_test.go b/internal/database/repository/clean_test.go new file mode 100644 index 00000000..9b929c43 --- /dev/null +++ b/internal/database/repository/clean_test.go @@ -0,0 +1,77 @@ +package repository_test + +import ( + "context" + "os" + "path" + "queryorchestration/internal/database/repository" + "queryorchestration/internal/serviceconfig" + "queryorchestration/internal/test" + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestClean(t *testing.T) { + ctx := context.Background() + + cfg := &serviceconfig.BaseConfig{} + test.SetCfgProvider(t, cfg) + cfg.SetBasePath(path.Join(os.Getenv("PWD"), "../../..")) + _, cleanup := test.CreateDB(t, ctx, &test.CreateDatabaseConfig{ + Cfg: cfg, + RunMigrations: true, + }) + defer cleanup() + + queries := cfg.GetDBQueries() + + clientId, err := queries.CreateClient(ctx, "example_client") + assert.NoError(t, err) + jobId, err := queries.CreateJob(ctx, clientId) + assert.NoError(t, err) + + hash := "example_hash" + id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Jobid: jobId, + Hash: hash, + }) + assert.NoError(t, err) + assert.NotEmpty(t, id) + + collId, err := queries.CreateCollector(ctx, jobId) + assert.NoError(t, err) + err = queries.AddCollectorCodeVersion(ctx, &repository.AddCollectorCodeVersionParams{ + Collectorid: collId, + Addedversion: 1, + Mincleanversion: 2, + Mintextversion: 1, + }) + assert.NoError(t, err) + + isclean, err := queries.IsDocumentClean(ctx, id) + assert.NoError(t, err) + assert.False(t, isclean) + + err = queries.AddDocumentCleanEntry(ctx, &repository.AddDocumentCleanEntryParams{ + Documentid: id, + Version: 1, + }) + assert.NoError(t, err) + + isclean, err = queries.IsDocumentClean(ctx, id) + assert.NoError(t, err) + assert.False(t, isclean) + + err = queries.AddCollectorCodeVersion(ctx, &repository.AddCollectorCodeVersionParams{ + Collectorid: collId, + Addedversion: 1, + Mincleanversion: 1, + Mintextversion: 1, + }) + assert.NoError(t, err) + + isclean, err = queries.IsDocumentClean(ctx, id) + assert.NoError(t, err) + assert.True(t, isclean) +} diff --git a/internal/database/repository/document.sql.go b/internal/database/repository/document.sql.go index 0adceae5..32410885 100644 --- a/internal/database/repository/document.sql.go +++ b/internal/database/repository/document.sql.go @@ -12,20 +12,20 @@ import ( ) const addDocumentEntry = `-- name: AddDocumentEntry :exec -INSERT INTO documentEntries (documentId, bucket, location) VALUES ($1, $2, $3) +INSERT INTO documentEntries (documentId, bucket, key) VALUES ($1, $2, $3) ` type AddDocumentEntryParams struct { Documentid pgtype.UUID `db:"documentid"` Bucket string `db:"bucket"` - Location string `db:"location"` + Key string `db:"key"` } // AddDocumentEntry // -// INSERT INTO documentEntries (documentId, bucket, location) VALUES ($1, $2, $3) +// INSERT INTO documentEntries (documentId, bucket, key) VALUES ($1, $2, $3) func (q *Queries) AddDocumentEntry(ctx context.Context, arg *AddDocumentEntryParams) error { - _, err := q.db.Exec(ctx, addDocumentEntry, arg.Documentid, arg.Bucket, arg.Location) + _, err := q.db.Exec(ctx, addDocumentEntry, arg.Documentid, arg.Bucket, arg.Key) return err } @@ -62,6 +62,26 @@ func (q *Queries) GetDocument(ctx context.Context, id pgtype.UUID) (*Document, e return &i, err } +const getDocumentEntry = `-- name: GetDocumentEntry :one +SELECT documentId, bucket, key FROM documentEntries WHERE documentId = $1 ORDER BY id DESC LIMIT 1 +` + +type GetDocumentEntryRow struct { + Documentid pgtype.UUID `db:"documentid"` + Bucket string `db:"bucket"` + Key string `db:"key"` +} + +// GetDocumentEntry +// +// SELECT documentId, bucket, key FROM documentEntries WHERE documentId = $1 ORDER BY id DESC LIMIT 1 +func (q *Queries) GetDocumentEntry(ctx context.Context, documentid pgtype.UUID) (*GetDocumentEntryRow, error) { + row := q.db.QueryRow(ctx, getDocumentEntry, documentid) + var i GetDocumentEntryRow + err := row.Scan(&i.Documentid, &i.Bucket, &i.Key) + return &i, err +} + const getDocumentIDByHash = `-- name: GetDocumentIDByHash :one SELECT id FROM documents WHERE hash = $1 and jobId = $2 ` diff --git a/internal/database/repository/document_test.go b/internal/database/repository/document_test.go index 5f264028..20feb7f7 100644 --- a/internal/database/repository/document_test.go +++ b/internal/database/repository/document_test.go @@ -39,15 +39,6 @@ func TestDocument(t *testing.T) { assert.NoError(t, err) assert.NotEmpty(t, id) - bucket := "example_bucket" - location := "/i/am/here" - err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ - Documentid: id, - Bucket: bucket, - Location: location, - }) - assert.NoError(t, err) - doc, err := queries.GetDocument(ctx, id) assert.NoError(t, err) assert.EqualExportedValues(t, &repository.Document{ @@ -62,4 +53,49 @@ func TestDocument(t *testing.T) { }) assert.NoError(t, err) assert.EqualExportedValues(t, id, docid) + + err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ + Documentid: id, + Bucket: "bucket_one", + Key: "/i/am/here", + }) + assert.NoError(t, err) + + entry, err := queries.GetDocumentEntry(ctx, id) + assert.NoError(t, err) + assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ + Documentid: id, + Bucket: "bucket_one", + Key: "/i/am/here", + }, entry) + + err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ + Documentid: id, + Bucket: "bucket_two", + Key: "/you/is/there", + }) + assert.NoError(t, err) + + entry, err = queries.GetDocumentEntry(ctx, id) + assert.NoError(t, err) + assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ + Documentid: id, + Bucket: "bucket_two", + Key: "/you/is/there", + }, entry) + + err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ + Documentid: id, + Bucket: "bucket_three", + Key: "/who/is/where", + }) + assert.NoError(t, err) + + entry, err = queries.GetDocumentEntry(ctx, id) + assert.NoError(t, err) + assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ + Documentid: id, + Bucket: "bucket_three", + Key: "/who/is/where", + }, entry) } diff --git a/internal/database/repository/models.go b/internal/database/repository/models.go index 95f75243..e235c748 100644 --- a/internal/database/repository/models.go +++ b/internal/database/repository/models.go @@ -113,12 +113,19 @@ type Document struct { Hash string `db:"hash"` } +type Documentclean struct { + ID pgtype.UUID `db:"id"` + Documentid pgtype.UUID `db:"documentid"` + Version int32 `db:"version"` + Bucket string `db:"bucket"` + Key string `db:"key"` +} + type Documententry struct { - ID pgtype.UUID `db:"id"` - Documentid pgtype.UUID `db:"documentid"` - Bucket string `db:"bucket"` - Location string `db:"location"` - Createdat pgtype.Timestamp `db:"createdat"` + ID pgtype.UUID `db:"id"` + Documentid pgtype.UUID `db:"documentid"` + Bucket string `db:"bucket"` + Key string `db:"key"` } type Fullactivecollector struct { diff --git a/internal/document/clean/clean.go b/internal/document/clean/clean.go new file mode 100644 index 00000000..9c65842b --- /dev/null +++ b/internal/document/clean/clean.go @@ -0,0 +1,54 @@ +package documentclean + +import ( + "context" + "queryorchestration/internal/database" + "queryorchestration/internal/database/repository" + "queryorchestration/internal/document" + + "github.com/google/uuid" +) + +type CleanParams struct { + ID uuid.UUID + Location document.Location +} + +func (s *Service) executeCleanTasks(params *CleanParams) (*document.Location, error) { + // TODO - various cleaning tasks + return ¶ms.Location, nil +} + +func (s *Service) clean(ctx context.Context, id uuid.UUID) error { + docId := database.MustToDBUUID(id) + + entry, err := s.cfg.GetDBQueries().GetDocumentEntry(ctx, docId) + if err != nil { + return err + } + + outLocation, err := s.executeCleanTasks(&CleanParams{ + ID: id, + Location: document.Location{ + Bucket: entry.Bucket, + Key: entry.Key, + }, + }) + if err != nil { + return err + } + + version := s.svc.Document.GetCleanVersion() + + err = s.cfg.GetDBQueries().AddDocumentCleanEntry(ctx, &repository.AddDocumentCleanEntryParams{ + Documentid: docId, + Version: version, + Bucket: outLocation.Bucket, + Key: outLocation.Key, + }) + if err != nil { + return err + } + + return nil +} diff --git a/internal/document/clean/clean_test.go b/internal/document/clean/clean_test.go new file mode 100644 index 00000000..f77c627c --- /dev/null +++ b/internal/document/clean/clean_test.go @@ -0,0 +1,62 @@ +package documentclean + +import ( + "context" + "queryorchestration/internal/database" + "queryorchestration/internal/database/repository" + "queryorchestration/internal/document" + "testing" + + "github.com/google/uuid" + "github.com/pashagolub/pgxmock/v3" + "github.com/stretchr/testify/assert" +) + +func TestClean(t *testing.T) { + ctx := context.Background() + pool, err := pgxmock.NewPool() + if err != nil { + t.Fatalf("failed to open pgxmock database: %v", err) + } + + cfg := &DocCleanConfig{} + cfg.DBPool = pool + cfg.DBQueries = repository.New(pool) + + svc := Service{ + cfg: cfg, + svc: &Services{ + Document: document.New(cfg), + }, + } + + id := uuid.New() + inloc := document.Location{ + Bucket: "bucket_name", + Key: "/i/am/here", + } + + pool.ExpectQuery("name: GetDocumentEntry :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( + pgxmock.NewRows([]string{"documentId", "bucket", "key"}). + AddRow(database.MustToDBUUID(id), inloc.Bucket, inloc.Key), + ) + pool.ExpectExec("name: AddDocumentCleanEntry :exec").WithArgs(database.MustToDBUUID(id), int32(1), inloc.Bucket, inloc.Key). + WillReturnResult(pgxmock.NewResult("", 1)) + + err = svc.clean(ctx, id) + assert.NoError(t, err) +} + +func TestExecuteCleanTasks(t *testing.T) { + svc := Service{} + + id := uuid.New() + location := document.Location{} + + outloc, err := svc.executeCleanTasks(&CleanParams{ + ID: id, + Location: location, + }) + assert.NoError(t, err) + assert.EqualExportedValues(t, location, *outloc) +} diff --git a/internal/document/clean/create.go b/internal/document/clean/create.go new file mode 100644 index 00000000..6e2fa8d1 --- /dev/null +++ b/internal/document/clean/create.go @@ -0,0 +1,45 @@ +package documentclean + +import ( + "context" + "queryorchestration/internal/database" + documenttext "queryorchestration/internal/document/text" + "queryorchestration/internal/serviceconfig/queue" + + "github.com/google/uuid" +) + +func (s *Service) Create(ctx context.Context, id uuid.UUID) error { + isclean, err := s.cfg.GetDBQueries().IsDocumentClean(ctx, database.MustToDBUUID(id)) + if err != nil { + return err + } + + if !isclean { + err = s.clean(ctx, id) + if err != nil { + return err + } + } + + err = s.informClean(ctx, id) + if err != nil { + return err + } + + return nil +} + +func (s *Service) informClean(ctx context.Context, id uuid.UUID) error { + err := s.cfg.SendToQueue(ctx, &queue.SendParams{ + QueueURL: s.cfg.GetDocumentTextURL(), + Body: documenttext.Create{ + ID: id, + }, + }) + if err != nil { + return err + } + + return nil +} diff --git a/internal/document/clean/create_test.go b/internal/document/clean/create_test.go new file mode 100644 index 00000000..63b939ae --- /dev/null +++ b/internal/document/clean/create_test.go @@ -0,0 +1,105 @@ +package documentclean + +import ( + "context" + "fmt" + "queryorchestration/internal/database" + "queryorchestration/internal/database/repository" + "queryorchestration/internal/document" + "queryorchestration/internal/serviceconfig" + "queryorchestration/internal/serviceconfig/objectstore" + "queryorchestration/internal/serviceconfig/queue/documenttext" + queuemock "queryorchestration/mocks/queue" + "testing" + + "github.com/aws/aws-sdk-go-v2/service/sqs" + "github.com/google/uuid" + "github.com/pashagolub/pgxmock/v3" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" +) + +type DocCleanConfig struct { + serviceconfig.BaseConfig + documenttext.DocTextConfig + objectstore.ObjectStoreConfig +} + +func TestCreate(t *testing.T) { + ctx := context.Background() + pool, err := pgxmock.NewPool() + if err != nil { + t.Fatalf("failed to open pgxmock database: %v", err) + } + + mockSQS := queuemock.NewMockSQSClient(t) + + cfg := &DocCleanConfig{} + cfg.QueueClient = mockSQS + cfg.DocumentTextURL = "/i/am/here" + cfg.DBPool = pool + cfg.DBQueries = repository.New(pool) + + svc := Service{ + cfg: cfg, + svc: &Services{ + Document: document.New(cfg), + }, + } + + id := uuid.New() + inloc := document.Location{ + Bucket: "bucket_name", + Key: "/i/am/here", + } + + pool.ExpectQuery("name: IsDocumentClean :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( + pgxmock.NewRows([]string{"isclean"}). + AddRow(false), + ) + pool.ExpectQuery("name: GetDocumentEntry :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( + pgxmock.NewRows([]string{"documentId", "bucket", "key"}). + AddRow(database.MustToDBUUID(id), inloc.Bucket, inloc.Key), + ) + pool.ExpectExec("name: AddDocumentCleanEntry :exec").WithArgs(database.MustToDBUUID(id), int32(1), inloc.Bucket, inloc.Key). + WillReturnResult(pgxmock.NewResult("", 1)) + + mockSQS.EXPECT(). + SendMessage( + mock.Anything, + mock.MatchedBy(func(in *sqs.SendMessageInput) bool { + return *in.QueueUrl == cfg.DocumentTextURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", id) + }), + mock.Anything, + ). + Return(&sqs.SendMessageOutput{}, nil) + + err = svc.Create(ctx, id) + assert.NoError(t, err) +} + +func TestInformClean(t *testing.T) { + ctx := context.Background() + mockSQS := queuemock.NewMockSQSClient(t) + + cfg := &DocCleanConfig{} + cfg.QueueClient = mockSQS + cfg.DocumentTextURL = "/i/am/here" + svc := Service{ + cfg: cfg, + } + id := uuid.New() + + mockSQS.EXPECT(). + SendMessage( + mock.Anything, + mock.MatchedBy(func(in *sqs.SendMessageInput) bool { + return *in.QueueUrl == cfg.DocumentTextURL && *in.MessageBody == fmt.Sprintf("{\"id\":\"%s\"}", id) + }), + mock.Anything, + ). + Return(&sqs.SendMessageOutput{}, nil) + + err := svc.informClean(ctx, id) + assert.NoError(t, err) +} diff --git a/internal/document/clean/service.go b/internal/document/clean/service.go index aa5cd110..4e843406 100644 --- a/internal/document/clean/service.go +++ b/internal/document/clean/service.go @@ -1,26 +1,30 @@ package documentclean import ( - "errors" - - "github.com/google/uuid" + "queryorchestration/internal/document" + "queryorchestration/internal/serviceconfig" + "queryorchestration/internal/serviceconfig/objectstore" + "queryorchestration/internal/serviceconfig/queue/documenttext" ) -type Create struct { - ID uuid.UUID `json:"id" validate:"required,uuid"` +type ConfigProvider interface { + serviceconfig.ConfigProvider + objectstore.ConfigProvider + documenttext.ConfigProvider +} + +type Services struct { + Document *document.Service } type Service struct { + cfg ConfigProvider + svc *Services } -func New() *Service { - return &Service{} -} - -func (s *Service) IsValidVersion(v int32) error { - if v <= 0 { - return errors.New("document clean code version must be > 0") +func New(cfg ConfigProvider, svc *Services) *Service { + return &Service{ + cfg: cfg, + svc: svc, } - - return nil } diff --git a/internal/document/clean/service_test.go b/internal/document/clean/service_test.go index cfe9982b..b0dd5e45 100644 --- a/internal/document/clean/service_test.go +++ b/internal/document/clean/service_test.go @@ -2,19 +2,22 @@ package documentclean_test import ( documentclean "queryorchestration/internal/document/clean" + "queryorchestration/internal/serviceconfig" + "queryorchestration/internal/serviceconfig/objectstore" + "queryorchestration/internal/serviceconfig/queue/documenttext" "testing" "github.com/stretchr/testify/assert" ) +type DocCleanConfig struct { + serviceconfig.BaseConfig + documenttext.DocTextConfig + objectstore.ObjectStoreConfig +} + func TestService(t *testing.T) { - svc := documentclean.New() + cfg := &DocCleanConfig{} + svc := documentclean.New(cfg, nil) assert.NotNil(t, svc) } - -func TestIsValidVersion(t *testing.T) { - svc := documentclean.New() - - assert.Nil(t, svc.IsValidVersion(2)) - assert.Error(t, svc.IsValidVersion(-1)) -} diff --git a/internal/document/init/create.go b/internal/document/init/create.go index 9dab210c..fba9bb3b 100644 --- a/internal/document/init/create.go +++ b/internal/document/init/create.go @@ -5,15 +5,14 @@ import ( "database/sql" "errors" "fmt" + doccleanrunner "queryorchestration/api/docCleanRunner" "queryorchestration/internal/database" "queryorchestration/internal/database/repository" "queryorchestration/internal/document" - documentclean "queryorchestration/internal/document/clean" "queryorchestration/internal/job" "queryorchestration/internal/serviceconfig/queue" "regexp" - "github.com/aws/aws-sdk-go-v2/service/sqs/types" "github.com/google/uuid" "github.com/jackc/pgx/v5/pgtype" ) @@ -55,8 +54,7 @@ type createDocumentParams struct { ID *pgtype.UUID JobID pgtype.UUID Hash string - Bucket string - Location string + Location document.Location } func (s *Service) getCreateParams(ctx context.Context, doc *Create) (*createDocumentParams, error) { @@ -75,7 +73,6 @@ func (s *Service) getCreateParams(ctx context.Context, doc *Create) (*createDocu ID: docID, JobID: database.MustToDBUUID(doc.JobID), Hash: doc.Hash, - Bucket: doc.Bucket, Location: doc.Location, }, nil } @@ -101,8 +98,8 @@ func (s *Service) submitCreate(ctx context.Context, params *createDocumentParams err := s.cfg.GetDBQueries().AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ Documentid: dbid, - Bucket: params.Bucket, - Location: params.Location, + Bucket: params.Location.Bucket, + Key: params.Location.Key, }) if err != nil { return err @@ -126,10 +123,9 @@ func (s *Service) informCreate(ctx context.Context, id uuid.UUID, j *job.Job) er err := s.cfg.SendToQueue(ctx, &queue.SendParams{ QueueURL: s.cfg.GetDocumentCleanURL(), - Body: documentclean.Create{ + Body: doccleanrunner.Create{ ID: id, }, - Attributes: map[string]types.MessageAttributeValue{}, }) if err != nil { return err diff --git a/internal/document/init/create_test.go b/internal/document/init/create_test.go index f296e4f4..4d0fbd99 100644 --- a/internal/document/init/create_test.go +++ b/internal/document/init/create_test.go @@ -57,18 +57,20 @@ func TestCreate(t *testing.T) { JobID: j.ID, Hash: "example_hash", } - bucket := "eample_bucket" - location := "/i/am/here" + location := document.Location{ + Bucket: "example_bucket", + Key: "/i/am/here", + } pool.ExpectQuery("name: GetJob :one").WithArgs(database.MustToDBUUID(j.ID)).WillReturnRows( pgxmock.NewRows([]string{"id", "clientId", "canSync"}). AddRow(database.MustToDBUUID(j.ID), database.MustToDBUUID(j.ClientID), j.CanSync), ) - pool.ExpectQuery("-- name: GetClient :one").WithArgs(database.MustToDBUUID(j.ClientID)).WillReturnRows( + pool.ExpectQuery("name: GetClient :one").WithArgs(database.MustToDBUUID(j.ClientID)).WillReturnRows( pgxmock.NewRows([]string{"id", "name", "canSync"}). AddRow(database.MustToDBUUID(j.ClientID), "client_name", true), ) - pool.ExpectQuery("-- name: GetDocumentIDByHash :one").WithArgs(doc.Hash, database.MustToDBUUID(doc.JobID)).WillReturnRows( + pool.ExpectQuery("name: GetDocumentIDByHash :one").WithArgs(doc.Hash, database.MustToDBUUID(doc.JobID)).WillReturnRows( pgxmock.NewRows([]string{"id"}), ) pool.ExpectBegin() @@ -77,14 +79,13 @@ func TestCreate(t *testing.T) { pgxmock.NewRows([]string{"id"}). AddRow(database.MustToDBUUID(doc.ID)), ) - pool.ExpectExec("name: AddDocumentEntry :exec").WithArgs(database.MustToDBUUID(doc.ID), bucket, location). + pool.ExpectExec("name: AddDocumentEntry :exec").WithArgs(database.MustToDBUUID(doc.ID), location.Bucket, location.Key). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectCommit() id, err := svc.Create(ctx, &Create{ JobID: doc.JobID, Location: location, - Bucket: bucket, Hash: doc.Hash, }) assert.NoError(t, err) @@ -145,24 +146,24 @@ func TestGetCreateParams(t *testing.T) { JobID: uuid.New(), Hash: "example_hash", } - bucket := "eample_bucket" - location := "/i/am/here" + location := document.Location{ + Bucket: "example_bucket", + Key: "/i/am/here", + } - pool.ExpectQuery("-- name: GetDocumentIDByHash :one").WithArgs(doc.Hash, database.MustToDBUUID(doc.JobID)).WillReturnRows( + pool.ExpectQuery("name: GetDocumentIDByHash :one").WithArgs(doc.Hash, database.MustToDBUUID(doc.JobID)).WillReturnRows( pgxmock.NewRows([]string{"id"}), ) params, err := svc.getCreateParams(ctx, &Create{ JobID: doc.JobID, Location: location, - Bucket: bucket, Hash: doc.Hash, }) assert.NoError(t, err) assert.Equal(t, &createDocumentParams{ JobID: database.MustToDBUUID(doc.JobID), Hash: doc.Hash, - Bucket: bucket, Location: location, }, params) } @@ -191,10 +192,12 @@ func TestGetCreateParamsExisting(t *testing.T) { JobID: uuid.New(), Hash: "example_hash", } - bucket := "eample_bucket" - location := "/i/am/here" + location := document.Location{ + Bucket: "example_bucket", + Key: "/i/am/here", + } - pool.ExpectQuery("-- name: GetDocumentIDByHash :one").WithArgs(doc.Hash, database.MustToDBUUID(doc.JobID)).WillReturnRows( + pool.ExpectQuery("name: GetDocumentIDByHash :one").WithArgs(doc.Hash, database.MustToDBUUID(doc.JobID)).WillReturnRows( pgxmock.NewRows([]string{"id"}). AddRow(database.MustToDBUUID(doc.ID)), ) @@ -202,7 +205,6 @@ func TestGetCreateParamsExisting(t *testing.T) { params, err := svc.getCreateParams(ctx, &Create{ JobID: doc.JobID, Location: location, - Bucket: bucket, Hash: doc.Hash, }) assert.NoError(t, err) @@ -211,7 +213,6 @@ func TestGetCreateParamsExisting(t *testing.T) { ID: &dbid, JobID: database.MustToDBUUID(doc.JobID), Hash: doc.Hash, - Bucket: bucket, Location: location, }, params) } @@ -239,16 +240,17 @@ func TestGetCreateParamsCurrentDocErr(t *testing.T) { JobID: uuid.New(), Hash: "example_hash", } - bucket := "eample_bucket" - location := "/i/am/here" + location := document.Location{ + Bucket: "example_bucket", + Key: "/i/am/here", + } - pool.ExpectQuery("-- name: GetDocumentIDByHash :one").WithArgs(doc.Hash, database.MustToDBUUID(doc.JobID)). + pool.ExpectQuery("name: GetDocumentIDByHash :one").WithArgs(doc.Hash, database.MustToDBUUID(doc.JobID)). WillReturnError(errors.New("db err")) _, err = svc.getCreateParams(ctx, &Create{ JobID: doc.JobID, Location: location, - Bucket: bucket, }) assert.Error(t, err) } @@ -275,8 +277,10 @@ func TestSubmitCreate(t *testing.T) { JobID: uuid.New(), Hash: "example_hash", } - bucket := "eample_bucket" - location := "/i/am/here" + location := document.Location{ + Bucket: "example_bucket", + Key: "/i/am/here", + } pool.ExpectBegin() pool.ExpectQuery("name: CreateDocument :one").WithArgs(database.MustToDBUUID(doc.JobID), doc.Hash). @@ -284,14 +288,13 @@ func TestSubmitCreate(t *testing.T) { pgxmock.NewRows([]string{"id"}). AddRow(database.MustToDBUUID(doc.ID)), ) - pool.ExpectExec("name: AddDocumentEntry :exec").WithArgs(database.MustToDBUUID(doc.ID), bucket, location). + pool.ExpectExec("name: AddDocumentEntry :exec").WithArgs(database.MustToDBUUID(doc.ID), location.Bucket, location.Key). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectCommit() id, err := svc.submitCreate(ctx, &createDocumentParams{ JobID: database.MustToDBUUID(doc.JobID), Location: location, - Bucket: bucket, Hash: doc.Hash, }) assert.NoError(t, err) @@ -320,11 +323,13 @@ func TestSubmitCreateExists(t *testing.T) { JobID: uuid.New(), Hash: "example_hash", } - bucket := "eample_bucket" - location := "/i/am/here" + location := document.Location{ + Bucket: "example_bucket", + Key: "/i/am/here", + } pool.ExpectBegin() - pool.ExpectExec("name: AddDocumentEntry :exec").WithArgs(database.MustToDBUUID(doc.ID), bucket, location). + pool.ExpectExec("name: AddDocumentEntry :exec").WithArgs(database.MustToDBUUID(doc.ID), location.Bucket, location.Key). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectCommit() @@ -333,7 +338,6 @@ func TestSubmitCreateExists(t *testing.T) { ID: &docid, JobID: database.MustToDBUUID(doc.JobID), Location: location, - Bucket: bucket, Hash: doc.Hash, }) assert.NoError(t, err) diff --git a/internal/document/service.go b/internal/document/service.go index 1c20b776..1db1c8c3 100644 --- a/internal/document/service.go +++ b/internal/document/service.go @@ -6,7 +6,10 @@ import ( "github.com/google/uuid" ) -type Location = string +type Location struct { + Bucket string + Key string +} type Document struct { ID uuid.UUID diff --git a/internal/document/text/service.go b/internal/document/text/service.go index 6b4f5776..fed7e67d 100644 --- a/internal/document/text/service.go +++ b/internal/document/text/service.go @@ -1,8 +1,6 @@ package documenttext import ( - "errors" - "github.com/google/uuid" ) @@ -13,12 +11,9 @@ func New() *Service { return &Service{} } -func (s *Service) IsValidVersion(v int32) error { - if v <= 0 { - return errors.New("document clean code version must be > 0") - } - - return nil +// TODO - move to api/ +type Create struct { + ID uuid.UUID `json:"id" validate:"required,uuid"` } type IsExtractedParams struct { diff --git a/internal/document/text/service_test.go b/internal/document/text/service_test.go index d631b928..965ee2a8 100644 --- a/internal/document/text/service_test.go +++ b/internal/document/text/service_test.go @@ -13,13 +13,6 @@ func TestService(t *testing.T) { assert.NotNil(t, svc) } -func TestIsValidVersion(t *testing.T) { - svc := documenttext.New() - - assert.Nil(t, svc.IsValidVersion(2)) - assert.Error(t, svc.IsValidVersion(-1)) -} - func TestIsExtracted(t *testing.T) { svc := documenttext.New() diff --git a/internal/document/version.go b/internal/document/version.go new file mode 100644 index 00000000..222ffb1e --- /dev/null +++ b/internal/document/version.go @@ -0,0 +1,35 @@ +package document + +import ( + "fmt" +) + +func (s *Service) GetCleanVersion() int32 { + // TODO - actual version + return 1 +} + +func (s *Service) GetTextVersion() int32 { + // TODO - actual version + return 1 +} + +func (s *Service) IsValidCleanVersion(v int32) error { + currentVersion := s.GetCleanVersion() + + if v < 1 || v > currentVersion { + return fmt.Errorf("document clean code version must be in the range 0 < version <= %d", currentVersion) + } + + return nil +} + +func (s *Service) IsValidTextVersion(v int32) error { + currentVersion := s.GetTextVersion() + + if v < 1 || v > currentVersion { + return fmt.Errorf("document text code version must be in the range 0 < version <= %d", currentVersion) + } + + return nil +} diff --git a/internal/document/version_test.go b/internal/document/version_test.go new file mode 100644 index 00000000..96fad2ee --- /dev/null +++ b/internal/document/version_test.go @@ -0,0 +1,43 @@ +package document_test + +import ( + "queryorchestration/internal/document" + "queryorchestration/internal/serviceconfig" + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestIsValidCleanVersion(t *testing.T) { + cfg := &serviceconfig.BaseConfig{} + svc := document.New(cfg) + + assert.Nil(t, svc.IsValidCleanVersion(1)) + assert.Error(t, svc.IsValidCleanVersion(2)) + assert.Error(t, svc.IsValidCleanVersion(0)) + assert.Error(t, svc.IsValidCleanVersion(-1)) +} + +func TestIsValidTextVersion(t *testing.T) { + cfg := &serviceconfig.BaseConfig{} + svc := document.New(cfg) + + assert.Nil(t, svc.IsValidTextVersion(1)) + assert.Error(t, svc.IsValidTextVersion(2)) + assert.Error(t, svc.IsValidTextVersion(0)) + assert.Error(t, svc.IsValidTextVersion(-1)) +} + +func TestGetCleanVersion(t *testing.T) { + cfg := &serviceconfig.BaseConfig{} + svc := document.New(cfg) + + assert.Equal(t, int32(1), svc.GetCleanVersion()) +} + +func TestGetTextVersion(t *testing.T) { + cfg := &serviceconfig.BaseConfig{} + svc := document.New(cfg) + + assert.Equal(t, int32(1), svc.GetTextVersion()) +} diff --git a/internal/job/collector/create.go b/internal/job/collector/create.go index 984607e3..f09df4a3 100644 --- a/internal/job/collector/create.go +++ b/internal/job/collector/create.go @@ -41,7 +41,7 @@ type dbCreateParams struct { func (s *Service) getCreateParams(ctx context.Context, params *CreateParams) (*dbCreateParams, error) { minClean := params.MinCleanVersion if minClean != nil { - err := s.svc.Clean.IsValidVersion(*minClean) + err := s.svc.Document.IsValidCleanVersion(*minClean) if err != nil { return nil, err } @@ -49,7 +49,7 @@ func (s *Service) getCreateParams(ctx context.Context, params *CreateParams) (*d minText := params.MinTextVersion if minText != nil { - err := s.svc.Text.IsValidVersion(*minText) + err := s.svc.Document.IsValidTextVersion(*minText) if err != nil { return nil, err } diff --git a/internal/job/collector/create_test.go b/internal/job/collector/create_test.go index c20cd63f..be4aad76 100644 --- a/internal/job/collector/create_test.go +++ b/internal/job/collector/create_test.go @@ -28,8 +28,8 @@ func TestCreate(t *testing.T) { svc := collector.New(cfg, &collector.Services{}) id := uuid.New() - minCleanV := int32(2) - minTextV := int32(4) + minCleanV := int32(1) + minTextV := int32(1) create := collector.CreateParams{ JobID: uuid.New(), MinCleanVersion: &minCleanV, diff --git a/internal/job/collector/createprivate_test.go b/internal/job/collector/createprivate_test.go index 09994cf6..57ebb10f 100644 --- a/internal/job/collector/createprivate_test.go +++ b/internal/job/collector/createprivate_test.go @@ -30,8 +30,8 @@ func TestGetCreateParams(t *testing.T) { svc: &Services{}, } - minCleanV := int32(2) - minTextV := int32(4) + minCleanV := int32(1) + minTextV := int32(1) params := CreateParams{ JobID: uuid.New(), MinCleanVersion: &minCleanV, diff --git a/internal/job/collector/service.go b/internal/job/collector/service.go index 6024b356..cae0a169 100644 --- a/internal/job/collector/service.go +++ b/internal/job/collector/service.go @@ -1,8 +1,7 @@ package collector import ( - documentclean "queryorchestration/internal/document/clean" - documenttext "queryorchestration/internal/document/text" + "queryorchestration/internal/document" "queryorchestration/internal/serviceconfig" "github.com/google/uuid" @@ -19,8 +18,7 @@ type Collector struct { } type Services struct { - Clean *documentclean.Service - Text *documenttext.Service + Document *document.Service } type Service struct { diff --git a/internal/job/collector/update.go b/internal/job/collector/update.go index 3e236234..8379eec9 100644 --- a/internal/job/collector/update.go +++ b/internal/job/collector/update.go @@ -97,7 +97,7 @@ func (s *Service) normalizeCodeVersions(current *Collector, params *UpdateParams if params.MinCleanVersion == nil { params.MinCleanVersion = ¤t.MinCleanVersion } else { - err := s.svc.Clean.IsValidVersion(*params.MinCleanVersion) + err := s.svc.Document.IsValidCleanVersion(*params.MinCleanVersion) if err != nil { return err } @@ -106,7 +106,7 @@ func (s *Service) normalizeCodeVersions(current *Collector, params *UpdateParams if params.MinTextVersion == nil { params.MinTextVersion = ¤t.MinTextVersion } else { - err := s.svc.Text.IsValidVersion(*params.MinTextVersion) + err := s.svc.Document.IsValidTextVersion(*params.MinTextVersion) if err != nil { return err } diff --git a/internal/job/collector/update_test.go b/internal/job/collector/update_test.go index b4ef5b13..534759b4 100644 --- a/internal/job/collector/update_test.go +++ b/internal/job/collector/update_test.go @@ -4,8 +4,7 @@ import ( "context" "queryorchestration/internal/database" "queryorchestration/internal/database/repository" - documentclean "queryorchestration/internal/document/clean" - documenttext "queryorchestration/internal/document/text" + "queryorchestration/internal/document" "queryorchestration/internal/job/collector" "queryorchestration/internal/serviceconfig" "testing" @@ -28,8 +27,7 @@ func TestUpdate(t *testing.T) { cfg.DBQueries = repository.New(pool) svc := collector.New(cfg, &collector.Services{ - Clean: documentclean.New(), - Text: documenttext.New(), + Document: document.New(cfg), }) current := collector.Collector{ diff --git a/internal/job/collector/updateprivate_test.go b/internal/job/collector/updateprivate_test.go index 5b799455..0394adbd 100644 --- a/internal/job/collector/updateprivate_test.go +++ b/internal/job/collector/updateprivate_test.go @@ -4,8 +4,7 @@ import ( "context" "queryorchestration/internal/database" "queryorchestration/internal/database/repository" - documentclean "queryorchestration/internal/document/clean" - documenttext "queryorchestration/internal/document/text" + "queryorchestration/internal/document" "queryorchestration/internal/serviceconfig" "testing" @@ -32,8 +31,8 @@ func TestGetUpdateParams(t *testing.T) { svc: &Services{}, } - minCleanV := int32(2) - minTextV := int32(4) + minCleanV := int32(1) + minTextV := int32(1) aV := int32(3) current := Collector{ ActiveVersion: 1, @@ -225,10 +224,10 @@ func TestNormalizeActiveVersion(t *testing.T) { } func TestNormalizeCodeVersions(t *testing.T) { + cfg := &serviceconfig.BaseConfig{} svc := Service{ svc: &Services{ - Clean: documentclean.New(), - Text: documenttext.New(), + Document: document.New(cfg), }, } @@ -247,7 +246,7 @@ func TestNormalizeCodeVersions(t *testing.T) { assert.Nil(t, update.MinCleanVersion) assert.Nil(t, update.MinTextVersion) - cv := int32(2) + cv := int32(1) update.MinCleanVersion = &cv update.MinTextVersion = nil err = svc.normalizeCodeVersions(¤t, &update) @@ -256,15 +255,15 @@ func TestNormalizeCodeVersions(t *testing.T) { assert.Equal(t, int32(0), *update.MinTextVersion) update.MinCleanVersion = nil - tv := int32(2) + tv := int32(1) update.MinTextVersion = &tv err = svc.normalizeCodeVersions(¤t, &update) assert.NoError(t, err) assert.Equal(t, int32(0), *update.MinCleanVersion) assert.Equal(t, tv, *update.MinTextVersion) - current.MinCleanVersion = 2 - current.MinTextVersion = 2 + current.MinCleanVersion = 1 + current.MinTextVersion = 1 err = svc.normalizeCodeVersions(¤t, &update) assert.NoError(t, err) assert.Nil(t, update.MinCleanVersion) @@ -290,23 +289,6 @@ func TestNormalizeCodeVersions(t *testing.T) { assert.NoError(t, err) assert.Nil(t, update.MinCleanVersion) assert.Nil(t, update.MinTextVersion) - - cv = current.MinCleanVersion + 1 - update.MinCleanVersion = &cv - update.MinTextVersion = ¤t.MinTextVersion - err = svc.normalizeCodeVersions(¤t, &update) - assert.NoError(t, err) - assert.Equal(t, cv, *update.MinCleanVersion) - assert.Equal(t, current.MinTextVersion, *update.MinTextVersion) - - cv = current.MinCleanVersion + 1 - update.MinCleanVersion = &cv - tv = current.MinTextVersion + 1 - update.MinTextVersion = &tv - err = svc.normalizeCodeVersions(¤t, &update) - assert.NoError(t, err) - assert.Equal(t, cv, *update.MinCleanVersion) - assert.Equal(t, tv, *update.MinTextVersion) } func TestGetRemoveFields(t *testing.T) { diff --git a/internal/query/test_test.go b/internal/query/test_test.go index 4f253e9c..aae8d3bd 100644 --- a/internal/query/test_test.go +++ b/internal/query/test_test.go @@ -5,7 +5,6 @@ import ( "queryorchestration/internal/database" "queryorchestration/internal/database/repository" "queryorchestration/internal/document" - documentclean "queryorchestration/internal/document/clean" documenttext "queryorchestration/internal/document/text" "queryorchestration/internal/job/collector" "queryorchestration/internal/query" @@ -29,13 +28,11 @@ func TestTest(t *testing.T) { cfg := &serviceconfig.BaseConfig{} cfg.DBPool = pool cfg.DBQueries = repository.New(pool) - text := documenttext.New() - clean := documentclean.New() - col := collector.New(cfg, &collector.Services{ - Text: text, - Clean: clean, - }) docsvc := document.New(cfg) + col := collector.New(cfg, &collector.Services{ + Document: docsvc, + }) + text := documenttext.New() svc := query.New(cfg, &query.Services{ Text: text, Document: docsvc, diff --git a/internal/server/server_test.go b/internal/server/server_test.go index b7ad089d..046253a4 100644 --- a/internal/server/server_test.go +++ b/internal/server/server_test.go @@ -9,6 +9,7 @@ import ( "queryorchestration/internal/test" "testing" + "github.com/go-playground/validator/v10" "github.com/stretchr/testify/assert" ) @@ -36,3 +37,21 @@ func TestNew(t *testing.T) { } }() } + +func TestSetValidator(t *testing.T) { + cfg := server.BaseConfig{} + + assert.Nil(t, cfg.Validator) + + cfg.SetValidator() + assert.NotNil(t, cfg.Validator) +} + +func TestGetValidator(t *testing.T) { + cfg := server.BaseConfig{} + + assert.Nil(t, cfg.GetValidator()) + cfg.Validator = validator.New() + assert.NotNil(t, cfg.GetValidator()) + assert.EqualExportedValues(t, validator.New(), cfg.GetValidator()) +} diff --git a/internal/serviceconfig/queue/documenttext/config.go b/internal/serviceconfig/queue/documenttext/config.go new file mode 100644 index 00000000..9f4d7f6c --- /dev/null +++ b/internal/serviceconfig/queue/documenttext/config.go @@ -0,0 +1,13 @@ +package documenttext + +type DocTextConfig struct { + DocumentTextURL string `env:"DOCUMENT_TEXT_URL,required,notEmpty"` +} + +func (c *DocTextConfig) GetDocumentTextURL() string { + return c.DocumentTextURL +} + +type ConfigProvider interface { + GetDocumentTextURL() string +} diff --git a/internal/serviceconfig/queue/documenttext/config_test.go b/internal/serviceconfig/queue/documenttext/config_test.go new file mode 100644 index 00000000..dfeb9d42 --- /dev/null +++ b/internal/serviceconfig/queue/documenttext/config_test.go @@ -0,0 +1,20 @@ +package documenttext_test + +import ( + "queryorchestration/internal/serviceconfig/queue/documenttext" + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestGetDocumentTextURL(t *testing.T) { + cfg := documenttext.DocTextConfig{} + + name := cfg.GetDocumentTextURL() + assert.Equal(t, "", name) + + cfg.DocumentTextURL = "name" + name = cfg.GetDocumentTextURL() + assert.Equal(t, "name", name) + assert.Equal(t, cfg.DocumentTextURL, name) +} diff --git a/internal/test/runner.go b/internal/test/runner.go index cad1a1b0..b8f47b10 100644 --- a/internal/test/runner.go +++ b/internal/test/runner.go @@ -14,8 +14,9 @@ import ( type Runner = string const ( - QueryRunner = queryrunner.Name - DocInitRunner = docinitrunner.Name + DocInitRunner = docinitrunner.Name + DocCleanRunner = docinitrunner.Name + QueryRunner = queryrunner.Name ) type RunnerConfig struct { diff --git a/scripts/compose.yml b/scripts/compose.yml index 9a31594a..123d4f58 100644 --- a/scripts/compose.yml +++ b/scripts/compose.yml @@ -34,6 +34,9 @@ tasks: cmds: - aws s3 mb s3://$BUCKET_IN - aws sqs create-queue --queue-name $QNAME_DOCUMENT_INIT + - aws sqs create-queue --queue-name $QNAME_DOCUMENT_CLEAN + - aws sqs create-queue --queue-name $QNAME_DOCUMENT_TEXT + - aws sqs create-queue --queue-name $QNAME_QUERY_RUNNER - | aws s3api put-bucket-notification-configuration \ --bucket $BUCKET_IN \ @@ -43,8 +46,6 @@ tasks: "Events": ["s3:ObjectCreated:*"] }] }' - - aws sqs create-queue --queue-name $QNAME_DOCUMENT_CLEAN - - aws sqs create-queue --queue-name $QNAME_QUERY_RUNNER down:test: cmds: - task: down diff --git a/test/docInitRunner/docinitrunner_test.go b/test/process_test.go similarity index 85% rename from test/docInitRunner/docinitrunner_test.go rename to test/process_test.go index b931aa13..64110e9f 100644 --- a/test/docInitRunner/docinitrunner_test.go +++ b/test/process_test.go @@ -5,7 +5,7 @@ import ( "fmt" "os" "path" - "queryorchestration/internal/server/runner" + "queryorchestration/internal/serviceconfig" "queryorchestration/internal/serviceconfig/objectstore" documentcleanc "queryorchestration/internal/serviceconfig/queue/documentclean" "queryorchestration/internal/test" @@ -19,19 +19,19 @@ import ( "github.com/stretchr/testify/assert" ) -type DocInitConfig struct { - runner.BaseConfig +type ProcessConfig struct { + serviceconfig.BaseConfig documentcleanc.DocCleanConfig objectstore.ObjectStoreConfig } -func TestDocInitRunner(t *testing.T) { +func TestProcess(t *testing.T) { t.SkipNow() ctx := context.Background() - cfg := &DocInitConfig{} + cfg := &ProcessConfig{} test.SetCfgProvider(t, cfg) - cfg.SetBasePath(path.Join(os.Getenv("PWD"), "../..")) + cfg.SetBasePath(path.Join(os.Getenv("PWD"), "..")) network, ncleanup := test.CreateNetwork(t, ctx) defer ncleanup() @@ -46,7 +46,6 @@ func TestDocInitRunner(t *testing.T) { assert.NoError(t, err) test.CreateStoreClient(t, ctx, cfg, acfg.ExternalEndpoint) - doccleanurl := test.CreateQueue(t, ctx, cfg, "docclean") bucketName := "docinitbucket" test.CreateBucket(t, ctx, cfg, bucketName) arn := fmt.Sprintf("arn:aws:sqs:%s:000000000000:%s", cfg.AWSRegion, test.DocInitRunner) @@ -65,6 +64,9 @@ func TestDocInitRunner(t *testing.T) { }) assert.NoError(t, err) + doccleanurl := test.CreateQueue(t, ctx, cfg, test.DocCleanRunner) + doctexturl := test.CreateQueue(t, ctx, cfg, "doctextrunner") + net, cleanup := test.CreateRunnersAndServicesNetwork(t, ctx, &test.EcosystemNetworkConfig{ Cfg: cfg, Network: network, @@ -75,6 +77,12 @@ func TestDocInitRunner(t *testing.T) { "DOCUMENT_CLEAN_URL": doccleanurl, }, }, + { + Name: test.DocCleanRunner, + Env: map[string]string{ + "DOCUMENT_TEXT_URL": doctexturl, + }, + }, }, Services: []*test.ServiceNetworkConfig{ { @@ -116,5 +124,5 @@ func TestDocInitRunner(t *testing.T) { assert.NoError(t, err) resRegex := regexp.MustCompile(`{"id": "[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}"`) - test.AssertMessageBody(t, ctx, cfg, doccleanurl, resRegex) + test.AssertMessageBody(t, ctx, cfg, doctexturl, resRegex) } diff --git a/test/queryRunner/queryrunner_test.go b/test/queryRunner/queryrunner_test.go deleted file mode 100644 index d3f3a5ce..00000000 --- a/test/queryRunner/queryrunner_test.go +++ /dev/null @@ -1,41 +0,0 @@ -package integration_test - -import ( - "context" - "queryorchestration/internal/query" - "queryorchestration/internal/serviceconfig" - "queryorchestration/internal/serviceconfig/queue" - "queryorchestration/internal/test" - "testing" - - "github.com/google/uuid" - "github.com/stretchr/testify/assert" -) - -func TestQueryRunner(t *testing.T) { - ctx := context.Background() - - cfg := &serviceconfig.BaseConfig{} - test.SetCfgProvider(t, cfg) - - c, cleanup := test.CreateRunnerNetwork(t, ctx, &test.RunnerNetworkConfig{ - Cfg: cfg, - Name: test.QueryRunner, - }) - defer cleanup() - - document := query.Document{ - ID: uuid.New(), - JobID: uuid.New(), - CleanVersion: int32(1), - TextVersion: int32(1), - } - - err := cfg.SendToQueue(ctx, &queue.SendParams{ - QueueURL: c.URI, - Body: document, - }) - assert.NoError(t, err) - - // TODO - check document output -}