diff --git a/internal/database/repository/document_test.go b/internal/database/repository/document_test.go index 1ee1d70d..fcf17ca1 100644 --- a/internal/database/repository/document_test.go +++ b/internal/database/repository/document_test.go @@ -13,153 +13,334 @@ import ( ) func TestDocument(t *testing.T) { - t.Parallel() if testing.Short() { t.SkipNow() } - ctx := context.Background() + t.Run("List docs for a client", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - cfg := &serviceconfig.BaseConfig{} - net := test.GetNetwork(t) - test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) - queries := cfg.GetDBQueries() + queries := cfg.GetDBQueries() - clientId := "EXAMPLE" - err := queries.CreateClient(ctx, &repository.CreateClientParams{ - Name: "example_client", - Clientid: clientId, + clientId := "EXAMPLE" + err := queries.CreateClient(ctx, &repository.CreateClientParams{ + Name: "example_client", + Clientid: clientId, + }) + require.NoError(t, err) + + docs, err := queries.ListDocumentsByClient(ctx, clientId) + require.NoError(t, err) + assert.Len(t, docs, 0) }) - require.NoError(t, err) + t.Run("Create document", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - docs, err := queries.ListDocumentsByClient(ctx, clientId) - require.NoError(t, err) - assert.Len(t, docs, 0) + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) - hash := "example_hash" - id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ - Clientid: clientId, - Hash: hash, + queries := cfg.GetDBQueries() + + clientId := "EXAMPLE" + err := queries.CreateClient(ctx, &repository.CreateClientParams{ + Name: "example_client", + Clientid: clientId, + }) + require.NoError(t, err) + + hash := "example_hash" + id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Clientid: clientId, + Hash: hash, + }) + require.NoError(t, err) + assert.NotEmpty(t, id) + + docs, err := queries.ListDocumentsByClient(ctx, clientId) + require.NoError(t, err) + assert.Len(t, docs, 1) + assert.Equal(t, hash, docs[0].Hash) }) - require.NoError(t, err) - assert.NotEmpty(t, id) - bucketone := "bucketone" - keyone := "keyone" - err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ - Documentid: id, - Bucket: bucketone, - Key: keyone, + t.Run("Create two documents", func(t *testing.T) { + t.Parallel() + ctx := context.Background() + + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) + + queries := cfg.GetDBQueries() + + clientId := "EXAMPLE" + err := queries.CreateClient(ctx, &repository.CreateClientParams{ + Name: "example_client", + Clientid: clientId, + }) + require.NoError(t, err) + + hash := "example_hash" + id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Clientid: clientId, + Hash: hash, + }) + require.NoError(t, err) + assert.NotEmpty(t, id) + + documentTwoID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Clientid: clientId, + Hash: "example_hash_two", + }) + require.NoError(t, err) + assert.NotEmpty(t, documentTwoID) + + docs, err := queries.ListDocumentsByClient(ctx, clientId) + require.NoError(t, err) + assert.Len(t, docs, 2) }) - require.NoError(t, err) + t.Run("Two clients", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - docs, err = queries.ListDocumentsByClient(ctx, clientId) - require.NoError(t, err) - assert.Len(t, docs, 1) - assert.Equal(t, hash, docs[0].Hash) + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) - documentTwoID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ - Clientid: clientId, - Hash: "example_hash_two", + queries := cfg.GetDBQueries() + + clientId := "EXAMPLE" + err := queries.CreateClient(ctx, &repository.CreateClientParams{ + Name: "example_client", + Clientid: clientId, + }) + require.NoError(t, err) + + hash := "example_hash" + id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Clientid: clientId, + Hash: hash, + }) + require.NoError(t, err) + assert.NotEmpty(t, id) + + clientTwoId := "EXAMPLE TWO" + err = queries.CreateClient(ctx, &repository.CreateClientParams{ + Name: "example_name_two", + Clientid: clientTwoId, + }) + require.NoError(t, err) + + hashTwo := "example_hash_two" + idTwo, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Clientid: clientTwoId, + Hash: hashTwo, + }) + require.NoError(t, err) + assert.NotEmpty(t, idTwo) + + docs, err := queries.ListDocumentsByClient(ctx, clientId) + require.NoError(t, err) + assert.Len(t, docs, 1) + docs, err = queries.ListDocumentsByClient(ctx, clientTwoId) + require.NoError(t, err) + assert.Len(t, docs, 1) }) - require.NoError(t, err) - assert.NotEmpty(t, documentTwoID) + t.Run("document summary", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - buckettwo := "buckettwo" - keytwo := "keytwo" - err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ - Documentid: documentTwoID, - Bucket: buckettwo, - Key: keytwo, + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) + + queries := cfg.GetDBQueries() + + clientId := "EXAMPLE" + err := queries.CreateClient(ctx, &repository.CreateClientParams{ + Name: "example_client", + Clientid: clientId, + }) + require.NoError(t, err) + + hash := "example_hash" + id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Clientid: clientId, + Hash: hash, + }) + require.NoError(t, err) + assert.NotEmpty(t, id) + + doc, err := queries.GetDocumentSummary(ctx, id) + require.NoError(t, err) + assert.EqualExportedValues(t, &repository.Document{ + ID: id, + Clientid: clientId, + Hash: hash, + }, doc) }) - require.NoError(t, err) + t.Run("document external", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - docs, err = queries.ListDocumentsByClient(ctx, clientId) - require.NoError(t, err) - assert.Len(t, docs, 2) + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) - clientTwoId := "EXAMPLE TWO" - err = queries.CreateClient(ctx, &repository.CreateClientParams{ - Name: "example_name_two", - Clientid: clientTwoId, + queries := cfg.GetDBQueries() + + clientId := "EXAMPLE" + err := queries.CreateClient(ctx, &repository.CreateClientParams{ + Name: "example_client", + Clientid: clientId, + }) + require.NoError(t, err) + + hash := "example_hash" + id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Clientid: clientId, + Hash: hash, + }) + require.NoError(t, err) + assert.NotEmpty(t, id) + + docext, err := queries.GetDocumentExternal(ctx, id) + require.NoError(t, err) + assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ + ID: id, + Clientid: clientId, + Hash: hash, + Fields: []byte("{}"), + }, docext) }) - require.NoError(t, err) + t.Run("doc id by hash", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - docs, err = queries.ListDocumentsByClient(ctx, clientId) - require.NoError(t, err) - assert.Len(t, docs, 2) - docs, err = queries.ListDocumentsByClient(ctx, clientTwoId) - require.NoError(t, err) - assert.Len(t, docs, 0) + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) - documentThreeId, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ - Clientid: clientTwoId, - Hash: hash, + queries := cfg.GetDBQueries() + + clientId := "EXAMPLE" + err := queries.CreateClient(ctx, &repository.CreateClientParams{ + Name: "example_client", + Clientid: clientId, + }) + require.NoError(t, err) + + hash := "example_hash" + id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Clientid: clientId, + Hash: hash, + }) + require.NoError(t, err) + assert.NotEmpty(t, id) + + docid, err := queries.GetDocumentIDByHash(ctx, &repository.GetDocumentIDByHashParams{ + Hash: hash, + Clientid: clientId, + }) + require.NoError(t, err) + assert.EqualExportedValues(t, id, docid) }) - require.NoError(t, err) + t.Run("doc entry", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - bucketthree := "buckettwo" - keythree := "keytwo" - err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ - Documentid: documentThreeId, - Bucket: bucketthree, - Key: keythree, + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) + + queries := cfg.GetDBQueries() + + clientId := "EXAMPLE" + err := queries.CreateClient(ctx, &repository.CreateClientParams{ + Name: "example_client", + Clientid: clientId, + }) + require.NoError(t, err) + + hash := "example_hash" + id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Clientid: clientId, + Hash: hash, + }) + require.NoError(t, err) + assert.NotEmpty(t, id) + + bucketone := "bucketone" + keyone := "keyone" + err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ + Documentid: id, + Bucket: bucketone, + Key: keyone, + }) + require.NoError(t, err) + + entry, err := queries.GetDocumentEntry(ctx, id) + require.NoError(t, err) + assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ + Documentid: id, + Bucket: bucketone, + Key: keyone, + }, entry) }) - require.NoError(t, err) + t.Run("doc latest entry", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - docs, err = queries.ListDocumentsByClient(ctx, clientId) - require.NoError(t, err) - assert.Len(t, docs, 2) - docs, err = queries.ListDocumentsByClient(ctx, clientTwoId) - require.NoError(t, err) - assert.Len(t, docs, 1) - assert.Equal(t, hash, docs[0].Hash) + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) - doc, err := queries.GetDocumentSummary(ctx, id) - require.NoError(t, err) - assert.EqualExportedValues(t, &repository.Document{ - ID: id, - Clientid: clientId, - Hash: hash, - }, doc) + queries := cfg.GetDBQueries() - docext, err := queries.GetDocumentExternal(ctx, id) - require.NoError(t, err) - assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ - ID: id, - Clientid: clientId, - Hash: hash, - Fields: []byte("{}"), - }, docext) + clientId := "EXAMPLE" + err := queries.CreateClient(ctx, &repository.CreateClientParams{ + Name: "example_client", + Clientid: clientId, + }) + require.NoError(t, err) - docid, err := queries.GetDocumentIDByHash(ctx, &repository.GetDocumentIDByHashParams{ - Hash: hash, - Clientid: clientId, + hash := "example_hash" + id, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ + Clientid: clientId, + Hash: hash, + }) + require.NoError(t, err) + assert.NotEmpty(t, id) + + bucketone := "bucketone" + keyone := "keyone" + err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ + Documentid: id, + Bucket: bucketone, + Key: keyone, + }) + require.NoError(t, err) + + buckettwo := "buckettwo" + keytwo := "keytwo" + err = queries.AddDocumentEntry(ctx, &repository.AddDocumentEntryParams{ + Documentid: id, + Bucket: buckettwo, + Key: keytwo, + }) + require.NoError(t, err) + + entry, err := queries.GetDocumentEntry(ctx, id) + require.NoError(t, err) + assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ + Documentid: id, + Bucket: buckettwo, + Key: keytwo, + }, entry) }) - require.NoError(t, err) - assert.EqualExportedValues(t, id, docid) - - entry, err := queries.GetDocumentEntry(ctx, id) - require.NoError(t, err) - assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ - Documentid: id, - Bucket: bucketone, - Key: keyone, - }, entry) - - entry, err = queries.GetDocumentEntry(ctx, documentTwoID) - require.NoError(t, err) - assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ - Documentid: documentTwoID, - Bucket: buckettwo, - Key: keytwo, - }, entry) - - entry, err = queries.GetDocumentEntry(ctx, documentThreeId) - require.NoError(t, err) - assert.EqualExportedValues(t, &repository.GetDocumentEntryRow{ - Documentid: documentThreeId, - Bucket: bucketthree, - Key: keythree, - }, entry) } diff --git a/internal/database/repository/sync_test.go b/internal/database/repository/sync_test.go index dbe51bd3..fe0c90fe 100644 --- a/internal/database/repository/sync_test.go +++ b/internal/database/repository/sync_test.go @@ -436,7 +436,38 @@ func TestClientSync(t *testing.T) { }) }) - t.Run("multiple queries", func(t *testing.T) { + t.Run("change super query", func(t *testing.T) { + t.Parallel() + ctx := context.Background() + + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) + + queries := cfg.GetDBQueries() + + contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) + clientId := createClientWithCollector(t, queries) + documentID, _, textId := createDocumentWithCollectorAndText(t, queries, clientId, jsonQueryID) + + _, err := queries.AddResult(ctx, &repository.AddResultParams{ + Queryid: contextQueryID, + Value: "example_context", + Textentryid: textId, + Queryversion: 1, + }) + require.NoError(t, err) + + getDocumentSyncState(t, ctx, queries, docSyncStateParams{ + clientId: clientId, + documentId: documentID, + hash: "example_hash", + isSynced: false, + fields: `{"first_key": null}`, + }) + + }) + t.Run("upstream query update result", func(t *testing.T) { t.Parallel() ctx := context.Background() @@ -466,14 +497,6 @@ func TestClientSync(t *testing.T) { }) require.NoError(t, err) - getDocumentSyncState(t, ctx, queries, docSyncStateParams{ - clientId: clientId, - documentId: documentID, - hash: "example_hash", - isSynced: false, - fields: `{"first_key": null}`, - }) - err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{ Resultid: jsonResultId, Requiredresultid: contextResultId, @@ -550,7 +573,7 @@ func TestClientSync(t *testing.T) { }) }) - t.Run("update text extry", func(t *testing.T) { + t.Run("update text entry", func(t *testing.T) { t.Parallel() ctx := context.Background() @@ -588,6 +611,37 @@ func TestClientSync(t *testing.T) { isSynced: false, fields: `{"first_key": null}`, }) + }) + t.Run("update text entry result", func(t *testing.T) { + t.Parallel() + ctx := context.Background() + + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) + + queries := cfg.GetDBQueries() + + queryId := createQuery(t, queries, repository.QuerytypeContextFull) + clientId := createClientWithCollector(t, queries) + documentID, cleanId, _ := createDocumentWithCollectorAndResult(t, queries, clientId, queryId) + + textId, err := queries.AddDocumentText(ctx, &repository.AddDocumentTextParams{ + Bucket: "hi", + Key: "hello", + Part: 0, + Createdat: pgtype.Timestamp{ + Time: time.Now().UTC(), + Valid: true, + }, + Cleanid: cleanId, + }) + require.NoError(t, err) + err = queries.AddDocumentTextEntry(ctx, &repository.AddDocumentTextEntryParams{ + Version: 1, + Textid: textId, + }) + require.NoError(t, err) _, err = queries.AddResult(ctx, &repository.AddResultParams{ Queryid: queryId, @@ -605,7 +659,6 @@ func TestClientSync(t *testing.T) { fields: `{"first_key": "updated_text_json"}`, }) }) - t.Run("update clean entry", func(t *testing.T) { t.Parallel() ctx := context.Background() @@ -647,6 +700,40 @@ func TestClientSync(t *testing.T) { isSynced: false, fields: `{"first_key": null}`, }) + }) + t.Run("update clean entry with result", func(t *testing.T) { + t.Parallel() + ctx := context.Background() + + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) + + queries := cfg.GetDBQueries() + + queryId := createQuery(t, queries, repository.QuerytypeContextFull) + clientId := createClientWithCollector(t, queries) + documentID, _, _ := createDocumentWithCollectorAndResult(t, queries, clientId, queryId) + + bucket := "cu" + key := "kk" + hash := "hash" + cleanthreeid, err := queries.AddDocumentClean(ctx, &repository.AddDocumentCleanParams{ + Documentid: documentID, + Bucket: &bucket, + Key: &key, + Hash: &hash, + Mimetype: repository.NullCleanmimetype{ + Valid: true, + Cleanmimetype: repository.CleanmimetypeApplicationPdf, + }, + }) + require.NoError(t, err) + err = queries.AddDocumentCleanEntry(ctx, &repository.AddDocumentCleanEntryParams{ + Cleanid: cleanthreeid, + Version: 1, + }) + require.NoError(t, err) textThreeId, err := queries.AddDocumentText(ctx, &repository.AddDocumentTextParams{ Bucket: "hi", @@ -771,6 +858,30 @@ func TestClientSync(t *testing.T) { isSynced: true, fields: `{"first_key": "json_value"}`, }) + }) + t.Run("remove query and update collector version", func(t *testing.T) { + t.Parallel() + ctx := context.Background() + + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) + + queries := cfg.GetDBQueries() + + queryId := createQuery(t, queries, repository.QuerytypeContextFull) + clientId := createClientWithCollector(t, queries) + documentID, _, _ := createDocumentWithCollectorAndResult(t, queries, clientId, queryId) + + latestCollectorVersion, err := queries.AddLatestCollectorVersion(ctx, clientId) + require.NoError(t, err) + + err = queries.RemoveCollectorQuery(ctx, &repository.RemoveCollectorQueryParams{ + Clientid: clientId, + Queryid: queryId, + Removedversion: &latestCollectorVersion, + }) + require.NoError(t, err) err = queries.SetActiveCollectorVersion(ctx, &repository.SetActiveCollectorVersionParams{ Versionid: latestCollectorVersion, @@ -830,7 +941,7 @@ func TestClientSync(t *testing.T) { queryId := createQuery(t, queries, repository.QuerytypeContextFull) clientId := createClientWithCollector(t, queries) - documentID, _, textId := createDocumentWithCollectorAndResult(t, queries, clientId, queryId) + documentID, _, _ := createDocumentWithCollectorAndResult(t, queries, clientId, queryId) superQueryID := createQuery(t, queries, repository.QuerytypeJsonExtractor) err := queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{ @@ -847,6 +958,28 @@ func TestClientSync(t *testing.T) { isSynced: false, fields: `{"first_key": null}`, }) + }) + t.Run("add super query with result", func(t *testing.T) { + t.Parallel() + ctx := context.Background() + + cfg := &serviceconfig.BaseConfig{} + net := test.GetNetwork(t) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{}) + + queries := cfg.GetDBQueries() + + queryId := createQuery(t, queries, repository.QuerytypeContextFull) + clientId := createClientWithCollector(t, queries) + documentID, _, textId := createDocumentWithCollectorAndResult(t, queries, clientId, queryId) + + superQueryID := createQuery(t, queries, repository.QuerytypeJsonExtractor) + err := queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{ + Queryid: queryId, + Requiredqueryid: superQueryID, + Addedversion: 1, + }) + require.NoError(t, err) superResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: superQueryID, diff --git a/internal/server/runner/listener_test.go b/internal/server/runner/listener_test.go index 979001cf..b67fa009 100644 --- a/internal/server/runner/listener_test.go +++ b/internal/server/runner/listener_test.go @@ -31,25 +31,27 @@ func TestNewRunner(t *testing.T) { net := test.GetNetwork(t) var wg sync.WaitGroup + wg.Add(1) go func() { - a := test.CreateAWSContainer(t, cfg, net) - test.SetQueueClient(t, t.Context(), cfg, a.ExternalEndpoint) - cfg.QueueURL = test.CreateQueue(t, t.Context(), cfg, "queueName") - cfg.SetSQSEndpoint(a.ExternalEndpoint) + test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{ + NoMigrations: true, + }) + + cfg.ControllerFunc = func() Controller[interface{}] { + return runnermock.NewMockController[interface{}](t) + } wg.Done() }() + a := test.CreateAWSContainer(t, cfg, net) + test.SetQueueClient(t, t.Context(), cfg, a.ExternalEndpoint) + + cfg.QueueURL = test.CreateQueue(t, t.Context(), cfg, "queueName") + cfg.SetSQSEndpoint(a.ExternalEndpoint) + wg.Wait() - test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{ - NoMigrations: true, - }) - - cfg.ControllerFunc = func() Controller[interface{}] { - return runnermock.NewMockController[interface{}](t) - } - srvPtr, err := New(t.Context(), cfg) require.NoError(t, err) assert.NotNil(t, srvPtr) diff --git a/internal/test/ecosystem.go b/internal/test/ecosystem.go index b8dae807..17ab8b9d 100644 --- a/internal/test/ecosystem.go +++ b/internal/test/ecosystem.go @@ -113,7 +113,12 @@ type Dependencies struct { Network string } -func CreateFullDependencies(t testing.TB, ctx context.Context, cfg FullDependenciesConfig) Dependencies { +type FullDependenciesParams struct { + NoObjectStore bool + Runners *[]RunnerName +} + +func CreateFullDependenciesWithParams(t testing.TB, cfg FullDependenciesConfig, params *FullDependenciesParams) Dependencies { network := GetNetwork(t) deps := Dependencies{ @@ -127,29 +132,37 @@ func CreateFullDependencies(t testing.TB, ctx context.Context, cfg FullDependenc go func() { deps.AWSConfig = CreateAWSContainer(t, cfg, network) - SetQueueClient(t, ctx, cfg, deps.AWSConfig.ExternalEndpoint) + SetQueueClient(t, t.Context(), cfg, deps.AWSConfig.ExternalEndpoint) var lock sync.Mutex - for _, runner := range runners { + runs := params.Runners + if runs == nil { + allRuns := []RunnerName{} + for _, run := range runners { + allRuns = append(allRuns, run.Name) + } + runs = &allRuns + } + + for _, runner := range *runs { wg.Add(1) go func() { - - queue := CreateQueue(t, ctx, cfg, runner.Name) + queue := CreateQueue(t, t.Context(), cfg, runner) lock.Lock() - - deps.QueueURLs[runner.Name] = queue - + deps.QueueURLs[runner] = queue lock.Unlock() wg.Done() }() } - SetStoreClient(t, ctx, cfg, deps.AWSConfig.ExternalEndpoint) - CreateBucket(t, ctx, cfg) - SetBucketNotifs(t, ctx, cfg) + if !params.NoObjectStore { + SetStoreClient(t, t.Context(), cfg, deps.AWSConfig.ExternalEndpoint) + CreateBucket(t, t.Context(), cfg) + SetBucketNotifs(t, t.Context(), cfg) + } wg.Done() }() @@ -166,15 +179,19 @@ func CreateFullDependencies(t testing.TB, ctx context.Context, cfg FullDependenc return deps } +func CreateFullDependencies(t testing.TB, ctx context.Context, cfg FullDependenciesConfig) Dependencies { + return CreateFullDependenciesWithParams(t, cfg, &FullDependenciesParams{}) +} + type APINetwork struct { Dependencies Dependencies API *Container } -func CreateAPINetwork(t testing.TB, ctx context.Context, cfg FullDependenciesConfig, api API) (*APINetwork, func()) { - deps := CreateFullDependencies(t, ctx, cfg) +func CreateAPINetworkWithParams(t testing.TB, cfg FullDependenciesConfig, api API, params *FullDependenciesParams) (*APINetwork, func()) { + deps := CreateFullDependenciesWithParams(t, cfg, params) - c, ccleanup := CreateAPI(t, ctx, cfg, deps.Network, &APIConfig{ + c, ccleanup := CreateAPI(t, t.Context(), cfg, deps.Network, &APIConfig{ API: api, MockHTTP: string(deps.MockServer.Internal), }) @@ -185,6 +202,10 @@ func CreateAPINetwork(t testing.TB, ctx context.Context, cfg FullDependenciesCon }, ccleanup } +func CreateAPINetwork(t testing.TB, ctx context.Context, cfg FullDependenciesConfig, api API) (*APINetwork, func()) { + return CreateAPINetworkWithParams(t, cfg, api, &FullDependenciesParams{}) +} + func GetAlias(t testing.TB, baseName string) string { name := fmt.Sprintf("%s_%s", baseName, t.Name()) name = strings.ToLower(name) diff --git a/scripts/tests.yml b/scripts/tests.yml index a7d17d61..7f1f432c 100644 --- a/scripts/tests.yml +++ b/scripts/tests.yml @@ -50,7 +50,7 @@ tasks: NUM_RESULTS: 5 cmds: - | - GOMAXPROCS={{.TEST_PARALLEL}} go test -parallel {{.CPU_COUNT}} -count=1 -json ./... \ + GOMAXPROCS={{.TEST_PARALLEL}} go test -count=1 -parallel {{.TEST_PARALLEL}} -json ./... \ > {{.TEST_FILE}} - | echo -e "🐢 TOP {{.NUM_RESULTS}} SLOWEST INDIVIDUAL PACKAGE RUNS:" @@ -61,11 +61,6 @@ tasks: jq -r 'select(.Action == "pass" and .Test != null) | "\(.Elapsed)s \(.Test) (\(.Package))"' "{{.TEST_FILE}}" | \ sort -rn | head -n {{.NUM_RESULTS}} - echo -e "\nšŸ“¦ TOP {{.NUM_RESULTS}} SLOWEST PACKAGES (TOTAL TIME):" - jq -r 'select(.Action == "pass") | [.Package, .Elapsed] | @tsv' "{{.TEST_FILE}}" | - awk '{pkg[$1] += $2} END {for (p in pkg) print pkg[p] "s " p}' | - sort -rn | head -n {{.NUM_RESULTS}} - SLOWEST_PKG=$(jq -r 'select(.Action == "pass") | [.Package, .Elapsed] | @tsv' "{{.TEST_FILE}}" | awk '{pkg[$1] += $2} END {for (p in pkg) print pkg[p] "s " p}' | sort -rn | head -n 1 | awk '{print $2}') diff --git a/test/queryAPI/queryservice_test.go b/test/queryAPI/queryservice_test.go index 552f5021..0cac82ed 100644 --- a/test/queryAPI/queryservice_test.go +++ b/test/queryAPI/queryservice_test.go @@ -15,67 +15,134 @@ import ( ) func TestQueryAPI(t *testing.T) { - t.Parallel() - ctx := context.Background() + t.Run("create and get query", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - cfg := &Config{} + cfg := &Config{} - c, cleanup := test.CreateAPINetwork(t, ctx, cfg, test.QueryAPI) - defer cleanup() + c, cleanup := test.CreateAPINetworkWithParams(t, cfg, test.QueryAPI, &test.FullDependenciesParams{ + NoObjectStore: true, + Runners: &[]test.RunnerName{}, + }) + defer cleanup() - client, err := queryapi.NewClientWithResponses(c.API.URI) - require.NoError(t, err) + client, err := queryapi.NewClientWithResponses(c.API.URI) + require.NoError(t, err) - idRes, err := client.CreateQueryWithResponse(ctx, queryapi.QueryCreate{ - Type: queryapi.CONTEXTFULL, + jcfg := `{"path": "createkey"}` + idRes, err := client.CreateQueryWithResponse(ctx, queryapi.QueryCreate{ + Type: queryapi.JSONEXTRACTOR, + Config: &jcfg, + }) + require.NoError(t, err) + assert.NotNil(t, idRes) + assert.NotNil(t, idRes.JSON201) + jsonID := idRes.JSON201.Id + + queryRes, err := client.GetQueryWithResponse(ctx, jsonID) + require.NoError(t, err) + assert.Equal(t, jsonID, queryRes.JSON200.Id) + assert.Equal(t, queryapi.JSONEXTRACTOR, queryRes.JSON200.Type) + assert.Equal(t, int32(1), queryRes.JSON200.ActiveVersion) + assert.Equal(t, int32(1), queryRes.JSON200.LatestVersion) + assert.Equal(t, jcfg, *queryRes.JSON200.Config) + assert.Nil(t, queryRes.JSON200.RequiredQueries) + + queriesRes, err := client.ListQueriesWithResponse(ctx) + require.NoError(t, err) + assert.Len(t, queriesRes.JSON200.Queries, 1) + assert.Equal(t, jcfg, *queriesRes.JSON200.Queries[0].Config) }) - require.NoError(t, err) - assert.NotNil(t, idRes) - assert.NotNil(t, idRes.JSON201) - contextID := idRes.JSON201.Id - assert.NotEmpty(t, contextID) + t.Run("update query", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - jcfg := "{}" - idRes, err = client.CreateQueryWithResponse(ctx, queryapi.QueryCreate{ - Type: queryapi.JSONEXTRACTOR, - Config: &jcfg, + cfg := &Config{} + + c, cleanup := test.CreateAPINetworkWithParams(t, cfg, test.QueryAPI, &test.FullDependenciesParams{ + NoObjectStore: true, + Runners: &[]test.RunnerName{test.QueryVersionSyncRunnerName}, + }) + defer cleanup() + + client, err := queryapi.NewClientWithResponses(c.API.URI) + require.NoError(t, err) + + jcfg := "{}" + idRes, err := client.CreateQueryWithResponse(ctx, queryapi.QueryCreate{ + Type: queryapi.JSONEXTRACTOR, + Config: &jcfg, + }) + require.NoError(t, err) + assert.NotNil(t, idRes) + assert.NotNil(t, idRes.JSON201) + jsonID := idRes.JSON201.Id + + aV := int32(2) + newJcfg := `{"path": "keyone"}` + res, err := client.UpdateQueryWithResponse(ctx, jsonID, queryapi.QueryUpdate{ + ActiveVersion: &aV, + Config: &newJcfg, + }) + require.NoError(t, err) + assert.NotNil(t, res) + + test.AssertMessageBody(t, ctx, cfg, c.Dependencies.QueueURLs[test.QueryVersionSyncRunnerName], regexp.MustCompile(`{"id":".+"}`)) + + queryRes, err := client.GetQueryWithResponse(ctx, jsonID) + require.NoError(t, err) + assert.Equal(t, jsonID, queryRes.JSON200.Id) + assert.Equal(t, queryapi.JSONEXTRACTOR, queryRes.JSON200.Type) + assert.Equal(t, int32(2), queryRes.JSON200.ActiveVersion) + assert.Equal(t, int32(2), queryRes.JSON200.LatestVersion) + assert.Equal(t, newJcfg, *queryRes.JSON200.Config) + assert.Nil(t, queryRes.JSON200.RequiredQueries) }) - require.NoError(t, err) - assert.NotNil(t, idRes) - assert.NotNil(t, idRes.JSON201) - jsonID := idRes.JSON201.Id + t.Run("multiple dependent queries", func(t *testing.T) { + t.Parallel() + ctx := context.Background() - queryRes, err := client.GetQueryWithResponse(ctx, jsonID) - require.NoError(t, err) - assert.Equal(t, jsonID, queryRes.JSON200.Id) - assert.Equal(t, queryapi.JSONEXTRACTOR, queryRes.JSON200.Type) - assert.Equal(t, int32(1), queryRes.JSON200.ActiveVersion) - assert.Equal(t, int32(1), queryRes.JSON200.LatestVersion) - assert.Equal(t, jcfg, *queryRes.JSON200.Config) - assert.Nil(t, queryRes.JSON200.RequiredQueries) + cfg := &Config{} - queriesRes, err := client.ListQueriesWithResponse(ctx) - require.NoError(t, err) - assert.Len(t, queriesRes.JSON200.Queries, 2) + c, cleanup := test.CreateAPINetworkWithParams(t, cfg, test.QueryAPI, &test.FullDependenciesParams{ + NoObjectStore: true, + Runners: &[]test.RunnerName{}, + }) + defer cleanup() - aV := int32(2) - res, err := client.UpdateQueryWithResponse(ctx, jsonID, queryapi.QueryUpdate{ - ActiveVersion: &aV, - RequiredQueries: &[]types.UUID{ - contextID, - }, + client, err := queryapi.NewClientWithResponses(c.API.URI) + require.NoError(t, err) + + idRes, err := client.CreateQueryWithResponse(ctx, queryapi.QueryCreate{ + Type: queryapi.CONTEXTFULL, + }) + require.NoError(t, err) + assert.NotNil(t, idRes) + assert.NotNil(t, idRes.JSON201) + contextID := idRes.JSON201.Id + assert.NotEmpty(t, contextID) + + jcfg := "{}" + idRes, err = client.CreateQueryWithResponse(ctx, queryapi.QueryCreate{ + Type: queryapi.JSONEXTRACTOR, + Config: &jcfg, + RequiredQueries: &[]types.UUID{ + contextID, + }, + }) + require.NoError(t, err) + assert.NotNil(t, idRes) + assert.NotNil(t, idRes.JSON201) + jsonID := idRes.JSON201.Id + + queryRes, err := client.GetQueryWithResponse(ctx, jsonID) + require.NoError(t, err) + assert.Equal(t, jsonID, queryRes.JSON200.Id) + assert.Equal(t, queryapi.JSONEXTRACTOR, queryRes.JSON200.Type) + assert.Equal(t, int32(1), queryRes.JSON200.ActiveVersion) + assert.Equal(t, int32(1), queryRes.JSON200.LatestVersion) + assert.Equal(t, jcfg, *queryRes.JSON200.Config) + assert.ElementsMatch(t, []types.UUID{contextID}, *queryRes.JSON200.RequiredQueries) }) - require.NoError(t, err) - assert.NotNil(t, res) - - test.AssertMessageBody(t, ctx, cfg, c.Dependencies.QueueURLs[test.QueryVersionSyncRunnerName], regexp.MustCompile(`{"id":".+"}`)) - - queryRes, err = client.GetQueryWithResponse(ctx, jsonID) - require.NoError(t, err) - assert.Equal(t, jsonID, queryRes.JSON200.Id) - assert.Equal(t, queryapi.JSONEXTRACTOR, queryRes.JSON200.Type) - assert.Equal(t, int32(2), queryRes.JSON200.ActiveVersion) - assert.Equal(t, int32(2), queryRes.JSON200.LatestVersion) - assert.Equal(t, jcfg, *queryRes.JSON200.Config) - assert.ElementsMatch(t, []types.UUID{contextID}, *queryRes.JSON200.RequiredQueries) }