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.SkipNow() } ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() id := "EXAMPLE" err := queries.CreateClient(ctx, &repository.CreateClientParams{ Name: "example_client", Clientid: 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", BatchID: nil, }) 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", BatchID: nil, }) 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) } type docSyncStateParams struct { clientId string documentId uuid.UUID hash string fields string isSynced bool } func getDocumentSyncState(t testing.TB, ctx context.Context, queries *repository.Queries, params docSyncStateParams) { t.Helper() isSynced, err := queries.IsClientSynced(ctx, ¶ms.clientId) require.NoError(t, err) assert.Equal(t, params.isSynced, isSynced) doc, err := queries.GetDocumentExternal(ctx, params.documentId) require.NoError(t, err) assert.EqualExportedValues(t, &repository.GetDocumentExternalRow{ ID: params.documentId, Clientid: params.clientId, Hash: params.hash, Fields: []byte(params.fields), }, doc) } func createQuery(t testing.TB, queries *repository.Queries, queryType repository.Querytype) uuid.UUID { id, err := queries.CreateQuery(t.Context(), queryType) require.NoError(t, err) version, err := queries.AddLatestQueryVersion(t.Context(), id) require.NoError(t, err) err = queries.AddActiveQueryVersion(t.Context(), &repository.AddActiveQueryVersionParams{ Queryid: id, Versionid: version, }) require.NoError(t, err) return id } func createDependentQueries(t testing.TB, ctx context.Context, queries *repository.Queries) (uuid.UUID, uuid.UUID) { contextQueryID := createQuery(t, queries, repository.QuerytypeContextFull) jsonQueryID := createQuery(t, queries, repository.QuerytypeJsonExtractor) err := queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{ Queryid: jsonQueryID, Requiredqueryid: contextQueryID, Addedversion: 1, }) require.NoError(t, err) return contextQueryID, jsonQueryID } func createClientWithCollector(t testing.TB, queries *repository.Queries) string { clientId := "EXAMPLE" err := queries.CreateClient(t.Context(), &repository.CreateClientParams{ Name: "example_client", Clientid: clientId, }) require.NoError(t, err) collectorVersion, err := queries.AddLatestCollectorVersion(t.Context(), clientId) require.NoError(t, err) err = queries.SetActiveCollectorVersion(t.Context(), &repository.SetActiveCollectorVersionParams{ Versionid: collectorVersion, Clientid: clientId, }) require.NoError(t, err) return clientId } func TestClientSync(t *testing.T) { if testing.Short() { t.SkipNow() } t.Run("no collector", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() clientId := "EXAMPLE" err := queries.CreateClient(ctx, &repository.CreateClientParams{ Name: "example_client", Clientid: clientId, }) require.NoError(t, err) issynced, err := queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, issynced) }) t.Run("no documents", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() _, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, queries) isSynced, err := queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, isSynced) 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) }) t.Run("create document", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() queryId := createQuery(t, queries, repository.QuerytypeContextFull) clientId := createClientWithCollector(t, queries) err := queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Clientid: clientId, Queryid: queryId, Addedversion: 1, Name: "first_key", }) require.NoError(t, err) documentID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Clientid: clientId, Hash: "example_noclean", BatchID: nil, }) require.NoError(t, err) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_noclean", isSynced: false, fields: `{"first_key": null}`, }) }) t.Run("document fail clean", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() queryId := createQuery(t, queries, repository.QuerytypeContextFull) clientId := createClientWithCollector(t, queries) documentID := createDocumentWithCollector(t, queries, clientId, queryId) 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) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": null}`, }) }) t.Run("clean document", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() queryId := createQuery(t, queries, repository.QuerytypeContextFull) clientId := createClientWithCollector(t, queries) documentID := createDocumentWithCollector(t, queries, clientId, queryId) bucket := "buck" 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) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: false, fields: `{"first_key": null}`, }) }) t.Run("extracted text", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() queryId := createQuery(t, queries, repository.QuerytypeContextFull) clientId := createClientWithCollector(t, queries) documentID, cleanId := createCleanDocumentWithCollector(t, queries, clientId, queryId) 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) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: false, fields: `{"first_key": null}`, }) }) t.Run("single query", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() queryId := createQuery(t, queries, repository.QuerytypeContextFull) clientId := createClientWithCollector(t, queries) documentID, _, textId := createDocumentWithCollectorAndText(t, queries, clientId, queryId) _, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: queryId, Value: "json_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": "json_value"}`, }) }) t.Run("change super query", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) 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 := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, queries) documentID, _, textId := createDocumentWithCollectorAndText(t, queries, clientId, jsonQueryID) 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) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": "json_value"}`, }) }) t.Run("update upstream query", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, queries) documentID, _, textId := createDocumentWithCollectorAndResults(t, 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) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: false, fields: `{"first_key": null}`, }) contextResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "context_version_2", Textentryid: textId, Queryversion: contextLatestVersion, }) require.NoError(t, err) jsonResultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: "updated_context_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{ Resultid: jsonResultId, Requiredresultid: contextResultId, }) require.NoError(t, err) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": "updated_context_value"}`, }) }) t.Run("update text entry", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) 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) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: false, fields: `{"first_key": null}`, }) }) t.Run("update text entry result", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) 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, Value: "updated_text_json", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": "updated_text_json"}`, }) }) t.Run("update clean entry", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) 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) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: false, fields: `{"first_key": null}`, }) }) t.Run("update clean entry with result", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) 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", 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) _, err = queries.AddResult(ctx, &repository.AddResultParams{ Queryid: queryId, Value: "update_clean_json", Textentryid: textThreeId, Queryversion: 1, }) require.NoError(t, err) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": "update_clean_json"}`, }) }) t.Run("add new collector version", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() queryId := createQuery(t, queries, repository.QuerytypeContextFull) clientId := createClientWithCollector(t, queries) documentID, _, _ := createDocumentWithCollectorAndResult(t, queries, clientId, queryId) _, err := queries.AddLatestCollectorVersion(ctx, clientId) require.NoError(t, err) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": "json_value"}`, }) }) t.Run("add new active collector version", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) 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.SetActiveCollectorVersion(ctx, &repository.SetActiveCollectorVersionParams{ Versionid: latestCollectorVersion, Clientid: clientId, }) require.NoError(t, err) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": "json_value"}`, }) }) t.Run("remove query from new collector version", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) 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) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": "json_value"}`, }) }) t.Run("remove query and update collector version", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) 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, Clientid: clientId, }) require.NoError(t, err) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{}`, }) }) t.Run("add existing query", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries) clientId := createClientWithCollector(t, queries) documentID, _, _ := createDocumentWithCollectorAndResults(t, queries, clientId, contextQueryID, jsonQueryID) err := queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Clientid: clientId, Queryid: contextQueryID, Addedversion: 1, Name: "example_key", }) require.NoError(t, err) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": "json_value", "example_key": "example_context"}`, }) }) t.Run("add super query", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) queries := cfg.GetDBQueries() queryId := createQuery(t, queries, repository.QuerytypeContextFull) clientId := createClientWithCollector(t, queries) documentID, _, _ := 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) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: false, fields: `{"first_key": null}`, }) }) t.Run("add super query with result", func(t *testing.T) { t.Parallel() ctx := t.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(t, cfg) 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, Value: "super_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) resultId, err := queries.AddResult(ctx, &repository.AddResultParams{ Queryid: queryId, Value: "json_with_super", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{ Resultid: resultId, Requiredresultid: superResultId, }) require.NoError(t, err) getDocumentSyncState(t, ctx, queries, docSyncStateParams{ clientId: clientId, documentId: documentID, hash: "example_hash", isSynced: true, fields: `{"first_key": "json_with_super"}`, }) }) } func createDocumentWithCollector(t testing.TB, queries *repository.Queries, clientId string, jsonQueryID uuid.UUID) uuid.UUID { err := queries.AddCollectorQuery(t.Context(), &repository.AddCollectorQueryParams{ Clientid: clientId, Queryid: jsonQueryID, Addedversion: 1, Name: "first_key", }) require.NoError(t, err) documentID, err := queries.CreateDocument(t.Context(), &repository.CreateDocumentParams{ Clientid: clientId, Hash: "example_hash", BatchID: nil, }) require.NoError(t, err) return documentID } func createCleanDocumentWithCollector(t testing.TB, queries *repository.Queries, clientId string, jsonQueryID uuid.UUID) (uuid.UUID, uuid.UUID) { documentID := createDocumentWithCollector(t, queries, clientId, jsonQueryID) bucket := "example_bucket" key := "example_key" hash := "hahssh" cleanId, err := queries.AddDocumentClean(t.Context(), &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(t.Context(), &repository.AddDocumentCleanEntryParams{ Cleanid: cleanId, Version: 1, }) require.NoError(t, err) return documentID, cleanId } func createDocumentWithCollectorAndText(t testing.TB, queries *repository.Queries, clientId string, jsonQueryID uuid.UUID) (uuid.UUID, uuid.UUID, uuid.UUID) { documentID, cleanId := createCleanDocumentWithCollector(t, queries, clientId, jsonQueryID) textId, err := queries.AddDocumentText(t.Context(), &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(t.Context(), &repository.AddDocumentTextEntryParams{ Version: 1, Textid: textId, }) require.NoError(t, err) return documentID, cleanId, textId } func createDocumentWithCollectorAndResult(t testing.TB, queries *repository.Queries, clientId string, jsonQueryID uuid.UUID) (uuid.UUID, uuid.UUID, uuid.UUID) { documentID, cleanId, textId := createDocumentWithCollectorAndText(t, queries, clientId, jsonQueryID) _, err := queries.AddResult(t.Context(), &repository.AddResultParams{ Queryid: jsonQueryID, Value: "json_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) return documentID, cleanId, textId } func createDocumentWithCollectorAndResults(t testing.TB, queries *repository.Queries, clientId string, contextQueryID, jsonQueryID uuid.UUID) (uuid.UUID, uuid.UUID, uuid.UUID) { documentID, cleanId, textId := createDocumentWithCollectorAndText(t, queries, clientId, jsonQueryID) contextResultId, err := queries.AddResult(t.Context(), &repository.AddResultParams{ Queryid: contextQueryID, Value: "example_context", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) jsonResultId, err := queries.AddResult(t.Context(), &repository.AddResultParams{ Queryid: jsonQueryID, Value: "json_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) err = queries.AddResultDependency(t.Context(), &repository.AddResultDependencyParams{ Resultid: jsonResultId, Requiredresultid: contextResultId, }) require.NoError(t, err) return documentID, cleanId, textId } func BenchmarkIsClientSynced(b *testing.B) { ctx := b.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(b, cfg) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(b, ctx, queries) clientId := createClientWithCollector(b, queries) _, _, _ = createDocumentWithCollectorAndResults(b, queries, clientId, contextQueryID, jsonQueryID) b.ResetTimer() for b.Loop() { _, _ = queries.IsClientSynced(ctx, &clientId) } } func BenchmarkGetDocumentExternal(b *testing.B) { ctx := b.Context() cfg := &serviceconfig.BaseConfig{} test.CreateDB(b, cfg) queries := cfg.GetDBQueries() contextQueryID, jsonQueryID := createDependentQueries(b, ctx, queries) clientId := createClientWithCollector(b, queries) documentID, _, _ := createDocumentWithCollectorAndResults(b, queries, clientId, contextQueryID, jsonQueryID) b.ResetTimer() for b.Loop() { _, _ = queries.GetDocumentExternal(ctx, documentID) } }