package repository_test import ( "context" "testing" "time" "queryorchestration/internal/database/repository" "queryorchestration/internal/serviceconfig" "queryorchestration/internal/test" "github.com/google/uuid" "github.com/jackc/pgx/v5/pgtype" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) func TestListClientDocumentIDs(t *testing.T) { t.Parallel() if testing.Short() { t.Skip("Skipping long test in short mode") } ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() id := "EXAMPLE" err := queries.CreateClient(ctx, &repository.CreateClientParams{ Name: "example_client", ID: id, }) require.NoError(t, err) ids, err := queries.ListDocumentIDsBatch(ctx, &repository.ListDocumentIDsBatchParams{ Clientid: id, Batchsize: 1, Pageoffset: 0, }) require.NoError(t, err) assert.Len(t, ids, 0) docOne, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Clientid: id, Hash: "example_hash", }) require.NoError(t, err) ids, err = queries.ListDocumentIDsBatch(ctx, &repository.ListDocumentIDsBatchParams{ Clientid: id, Batchsize: 1, Pageoffset: 0, }) require.NoError(t, err) assert.Len(t, ids, 1) total := int64(1) assert.ElementsMatch(t, []*repository.ListDocumentIDsBatchRow{ { ID: &docOne, Totalcount: &total, }, }, ids) docTwo, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Clientid: id, Hash: "example_hash_two", }) require.NoError(t, err) ids, err = queries.ListDocumentIDsBatch(ctx, &repository.ListDocumentIDsBatchParams{ Clientid: id, Batchsize: 1, Pageoffset: 0, }) require.NoError(t, err) assert.Len(t, ids, 1) total = int64(2) assert.ElementsMatch(t, []*repository.ListDocumentIDsBatchRow{ { ID: &docOne, Totalcount: &total, }, }, ids) ids, err = queries.ListDocumentIDsBatch(ctx, &repository.ListDocumentIDsBatchParams{ Clientid: id, Batchsize: 1, Pageoffset: 1, }) require.NoError(t, err) assert.Len(t, ids, 1) assert.ElementsMatch(t, []*repository.ListDocumentIDsBatchRow{ { ID: &docTwo, Totalcount: &total, }, }, ids) ids, err = queries.ListDocumentIDsBatch(ctx, &repository.ListDocumentIDsBatchParams{ Clientid: id, Batchsize: 2, Pageoffset: 0, }) require.NoError(t, err) assert.Len(t, ids, 2) assert.ElementsMatch(t, []*repository.ListDocumentIDsBatchRow{ { ID: &docOne, Totalcount: &total, }, { ID: &docTwo, Totalcount: &total, }, }, ids) } func createDependentQueries(t testing.TB, ctx context.Context, queries *repository.Queries) (uuid.UUID, uuid.UUID) { contextQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeContextFull)) require.NoError(t, err) contextQueryVersion, err := queries.AddLatestQueryVersion(ctx, contextQueryID) require.NoError(t, err) err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{ Queryid: contextQueryID, Versionid: contextQueryVersion, }) require.NoError(t, err) jsonQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor)) require.NoError(t, err) jsonQueryVersion, err := queries.AddLatestQueryVersion(ctx, jsonQueryID) require.NoError(t, err) err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{ Queryid: jsonQueryID, Versionid: jsonQueryVersion, }) require.NoError(t, err) err = queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{ Queryid: jsonQueryID, Requiredqueryid: contextQueryID, Addedversion: jsonQueryVersion, }) require.NoError(t, err) return contextQueryID, jsonQueryID } func createClientWithCollector(t testing.TB, ctx context.Context, queries *repository.Queries) string { clientId := "EXAMPLE" err := queries.CreateClient(ctx, &repository.CreateClientParams{ Name: "example_client", ID: clientId, }) require.NoError(t, err) collectorVersion, err := queries.AddLatestCollectorVersion(ctx, clientId) require.NoError(t, err) err = queries.SetActiveCollectorVersion(ctx, &repository.SetActiveCollectorVersionParams{ Versionid: collectorVersion, Clientid: clientId, }) require.NoError(t, err) return clientId } func TestClientSync(t *testing.T) { if testing.Short() { t.Skip("Skipping long test in short mode") } t.Run("document fail clean", func(t *testing.T) { t.Parallel() ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() _, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, ctx, queries) err := queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Clientid: clientId, Queryid: jsonQueryID, Addedversion: 1, Name: "first_key", }) require.NoError(t, err) isSynced, err := queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) documentID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Clientid: clientId, Hash: "example_noclean", }) require.NoError(t, err) docExternal, err := queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_noclean", Fields: []byte(`{"first_key": null}`), }, docExternal) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) cleanid, err := queries.AddDocumentClean(ctx, &repository.AddDocumentCleanParams{ Documentid: documentID, Fail: repository.NullCleanfailtype{ Valid: true, Cleanfailtype: repository.CleanfailtypeInvalidMimetype, }, }) require.NoError(t, err) err = queries.AddDocumentCleanEntry(ctx, &repository.AddDocumentCleanEntryParams{ Cleanid: cleanid, Version: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_noclean", Fields: []byte(`{"first_key": null}`), }, docExternal) }) t.Run("standard extraction", func(t *testing.T) { t.Parallel() ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, ctx, queries) err := queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Clientid: clientId, Queryid: jsonQueryID, Addedversion: 1, Name: "first_key", }) require.NoError(t, err) isSynced, err := queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) bucket := "example_bucket" documentID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Clientid: clientId, Hash: "example_hash", }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err := queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) key := "example_key" hash := "hahssh" cleanId, 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: cleanId, Version: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) textId, err := queries.AddDocumentText(ctx, &repository.AddDocumentTextParams{ Cleanid: cleanId, Part: 0, Createdat: pgtype.Timestamp{ Time: time.Now().UTC(), Valid: true, }, Bucket: "hi", Key: "hello", Hash: "example", }) require.NoError(t, err) err = queries.AddDocumentTextEntry(ctx, &repository.AddDocumentTextEntryParams{ Version: 1, Textid: textId, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) contextResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "example_context", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) jsonResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: "json_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{ Resultid: jsonResultId, Requiredresultid: contextResultId, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": "json_value"}`), }, docExternal) }) t.Run("update upstream query", func(t *testing.T) { t.Parallel() ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, ctx, queries) documentID, _, textId := createDocumentWithCollectorAndResults(t, ctx, queries, clientId, contextQueryID, jsonQueryID) contextLatestVersion, err := queries.AddLatestQueryVersion(ctx, contextQueryID) require.NoError(t, err) err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{ Queryid: contextQueryID, Versionid: contextLatestVersion, }) require.NoError(t, err) isSynced, err := queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err := queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) contextResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "context_version_2", Textentryid: textId, Queryversion: contextLatestVersion, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) jsonResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: "updated_context_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{ Resultid: jsonResultId, Requiredresultid: contextResultId, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": "updated_context_value"}`), }, docExternal) }) t.Run("update text extry", func(t *testing.T) { t.Parallel() ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, ctx, queries) documentID, cleanId, _ := createDocumentWithCollectorAndResults(t, ctx, queries, clientId, contextQueryID, jsonQueryID) 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) isSynced, err := queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err := queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) contextResultID, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "updated_text_context", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) jsonResultID, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: "updated_text_json", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{ Resultid: jsonResultID, Requiredresultid: contextResultID, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": "updated_text_json"}`), }, docExternal) }) t.Run("update clean entry", func(t *testing.T) { t.Parallel() ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, ctx, queries) documentID, _, _ := createDocumentWithCollectorAndResults(t, ctx, queries, clientId, contextQueryID, jsonQueryID) 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) isSynced, err := queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err := queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) textThreeId, err := queries.AddDocumentText(ctx, &repository.AddDocumentTextParams{ Bucket: "hi", Key: "hello", Part: 0, Createdat: pgtype.Timestamp{ Time: time.Now().UTC(), Valid: true, }, Cleanid: cleanthreeid, }) require.NoError(t, err) err = queries.AddDocumentTextEntry(ctx, &repository.AddDocumentTextEntryParams{ Version: 1, Textid: textThreeId, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) contextResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "update_clean_context", Textentryid: textThreeId, Queryversion: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) jsonResultID, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: "update_clean_json", Textentryid: textThreeId, Queryversion: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{ Resultid: jsonResultID, Requiredresultid: contextResultId, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": "update_clean_json"}`), }, docExternal) }) t.Run("change collector name", func(t *testing.T) { t.Parallel() ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, ctx, queries) documentID, _, _ := createDocumentWithCollectorAndResults(t, ctx, queries, clientId, contextQueryID, jsonQueryID) latestCollectorVersion, err := queries.AddLatestCollectorVersion(ctx, clientId) require.NoError(t, err) isSynced, err := queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err := queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": "json_value"}`), }, docExternal) err = queries.SetActiveCollectorVersion(ctx, &repository.SetActiveCollectorVersionParams{ Versionid: latestCollectorVersion, Clientid: clientId, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": "json_value"}`), }, docExternal) err = queries.RemoveCollectorQuery(ctx, &repository.RemoveCollectorQueryParams{ Clientid: clientId, Queryid: jsonQueryID, Removedversion: &latestCollectorVersion, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{}`), }, docExternal) err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Clientid: clientId, Queryid: jsonQueryID, Addedversion: latestCollectorVersion, Name: "second_key", }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"second_key": "json_value"}`), }, docExternal) }) t.Run("add existing query", func(t *testing.T) { t.Parallel() ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, ctx, queries) documentID, _, _ := createDocumentWithCollectorAndResults(t, ctx, queries, clientId, contextQueryID, jsonQueryID) err := queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Clientid: clientId, Queryid: contextQueryID, Addedversion: 1, Name: "example_key", }) require.NoError(t, err) isSynced, err := queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err := queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": "json_value", "example_key": "example_context"}`), }, docExternal) }) t.Run("add super query", func(t *testing.T) { t.Parallel() ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, ctx, queries) documentID, _, textId := createDocumentWithCollectorAndResults(t, ctx, queries, clientId, contextQueryID, jsonQueryID) superQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor)) require.NoError(t, err) superQueryVersion, err := queries.AddLatestQueryVersion(ctx, superQueryID) require.NoError(t, err) err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{ Queryid: superQueryID, Versionid: superQueryVersion, }) require.NoError(t, err) err = queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{ Queryid: contextQueryID, Requiredqueryid: superQueryID, Addedversion: superQueryVersion, }) require.NoError(t, err) superResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: superQueryID, Value: "super_value", Textentryid: textId, Queryversion: superQueryVersion, }) require.NoError(t, err) isSynced, err := queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err := queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) contextResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "context_with_super_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) jsonResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: "json_with_super", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{ Resultid: contextResultId, Requiredresultid: superResultId, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": null}`), }, docExternal) err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{ Resultid: jsonResultId, Requiredresultid: contextResultId, }) require.NoError(t, err) isSynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) docExternal, err = queries.GetDocumentExternal(ctx, documentID) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: documentID, Clientid: clientId, Hash: "example_hash", Fields: []byte(`{"first_key": "json_with_super"}`), }, docExternal) }) } func createDocumentWithCollectorAndResults(t testing.TB, ctx context.Context, queries *repository.Queries, clientId string, contextQueryID, jsonQueryID uuid.UUID) (uuid.UUID, uuid.UUID, uuid.UUID) { err := queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Clientid: clientId, Queryid: jsonQueryID, Addedversion: 1, Name: "first_key", }) require.NoError(t, err) bucket := "example_bucket" documentID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Clientid: clientId, Hash: "example_hash", }) require.NoError(t, err) key := "example_key" hash := "hahssh" cleanId, 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: cleanId, Version: 1, }) require.NoError(t, err) textId, err := queries.AddDocumentText(ctx, &repository.AddDocumentTextParams{ Cleanid: cleanId, Part: 0, Createdat: pgtype.Timestamp{ Time: time.Now().UTC(), Valid: true, }, Bucket: "hi", Key: "hello", Hash: "example", }) require.NoError(t, err) err = queries.AddDocumentTextEntry(ctx, &repository.AddDocumentTextEntryParams{ Version: 1, Textid: textId, }) require.NoError(t, err) contextResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "example_context", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) jsonResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: "json_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{ Resultid: jsonResultId, Requiredresultid: contextResultId, }) require.NoError(t, err) return documentID, cleanId, textId } func BenchmarkIsClientSynced(b *testing.B) { ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(b, ctx) _ = test.CreateDB(b, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(b, ctx, queries) clientId := createClientWithCollector(b, ctx, queries) _, _, _ = createDocumentWithCollectorAndResults(b, ctx, queries, clientId, contextQueryID, jsonQueryID) b.ResetTimer() for b.Loop() { _, _ = queries.IsClientSynced(ctx, &clientId) } } func BenchmarkGetDocumentExternal(b *testing.B) { ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.DepNetwork.Get(b, ctx) _ = test.CreateDB(b, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(b, ctx, queries) clientId := createClientWithCollector(b, ctx, queries) documentID, _, _ := createDocumentWithCollectorAndResults(b, ctx, queries, clientId, contextQueryID, jsonQueryID) b.ResetTimer() for b.Loop() { _, _ = queries.GetDocumentExternal(ctx, documentID) } }