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 TestResults(t *testing.T) { t.Parallel() if testing.Short() { t.SkipNow() } ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.GetNetwork(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() jsonQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor)) require.NoError(t, err) _, err = queries.AddLatestQueryVersion(ctx, jsonQueryID) require.NoError(t, err) err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{ Queryid: jsonQueryID, Versionid: 1, }) require.NoError(t, err) clientId := "EXAMPLE" err = queries.CreateClient(ctx, &repository.CreateClientParams{ Name: "example_client", ID: clientId, }) 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_hash", }) require.NoError(t, err) issynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, issynced) version, err := queries.AddLatestCollectorVersion(ctx, clientId) require.NoError(t, err) err = queries.SetActiveCollectorVersion(ctx, &repository.SetActiveCollectorVersionParams{ Versionid: version, Clientid: clientId, }) require.NoError(t, err) issynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, issynced) bucket := "example_bucket" key := "example_key" hash := "hash" 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) textId, err := queries.AddDocumentText(ctx, &repository.AddDocumentTextParams{ Cleanid: cleanid, Bucket: "hi", Key: "hello", Hash: "example", Part: 0, Createdat: pgtype.Timestamp{ Time: time.Now().UTC(), Valid: true, }, }) 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.True(t, issynced) err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Clientid: clientId, Queryid: jsonQueryID, Addedversion: 1, Name: "example_key", }) require.NoError(t, err) issynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.False(t, issynced) jsonResultValue := "example_value" _, err = queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: jsonResultValue, Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) issynced, err = queries.IsClientSynced(ctx, &clientId) require.NoError(t, err) assert.True(t, issynced) qv := int32(1) res, err := queries.GetResultValueWithVersion(ctx, &repository.GetResultValueWithVersionParams{ Queryid: &jsonQueryID, Queryversion: &qv, Documentid: &documentID, }) require.NoError(t, err) assert.NotNil(t, res.Value) assert.Equal(t, jsonResultValue, *res.Value) } func TestResultValues(t *testing.T) { t.Parallel() if testing.Short() { t.SkipNow() } ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.GetNetwork(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() jsonQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor)) require.NoError(t, err) clientId := "EXAMPLE" err = queries.CreateClient(ctx, &repository.CreateClientParams{ Name: "example_client", ID: clientId, }) require.NoError(t, err) documentID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Clientid: clientId, Hash: "example_hash", }) require.NoError(t, err) documentTwoID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Clientid: clientId, Hash: "example_hash_two", }) require.NoError(t, err) contextQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeContextFull)) require.NoError(t, err) _, err = queries.AddLatestQueryVersion(ctx, jsonQueryID) require.NoError(t, err) err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{ Queryid: jsonQueryID, Versionid: 1, }) require.NoError(t, err) _, err = queries.AddLatestQueryVersion(ctx, contextQueryID) require.NoError(t, err) _, err = queries.AddLatestQueryVersion(ctx, contextQueryID) require.NoError(t, err) err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{ Queryid: contextQueryID, Versionid: 2, }) require.NoError(t, err) jsonVersion := int32(1) err = queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{ Queryid: jsonQueryID, Requiredqueryid: contextQueryID, Addedversion: jsonVersion, }) require.NoError(t, err) contextQuery, err := queries.GetQuery(ctx, contextQueryID) require.NoError(t, err) qResults, err := queries.ListQueryRequirementValues(ctx, &repository.ListQueryRequirementValuesParams{ Queryid: &jsonQueryID, Documentid: &documentID, Version: &jsonVersion, }) require.NoError(t, err) assert.Len(t, qResults, 0) qResults, err = queries.ListQueryRequirementValues(ctx, &repository.ListQueryRequirementValuesParams{ Queryid: &jsonQueryID, Documentid: &documentTwoID, Version: &jsonVersion, }) require.NoError(t, err) assert.Len(t, qResults, 0) version, err := queries.AddLatestCollectorVersion(ctx, clientId) require.NoError(t, err) err = queries.SetActiveCollectorVersion(ctx, &repository.SetActiveCollectorVersionParams{ Versionid: version, Clientid: clientId, }) require.NoError(t, err) bucket := "example_bucket" key := "example_key" hash := "hash" 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, Bucket: "hi", Key: "hello", Hash: "example", Part: 0, Createdat: pgtype.Timestamp{ Time: time.Now().UTC(), Valid: true, }, }) require.NoError(t, err) err = queries.AddDocumentTextEntry(ctx, &repository.AddDocumentTextEntryParams{ Version: 1, Textid: textId, }) require.NoError(t, err) result := repository.AddResultParams{ Queryid: contextQueryID, Value: "context_value_1", Textentryid: textId, Queryversion: contextQuery.Activeversion, } _, err = queries.AddResult(ctx, &result) require.NoError(t, err) qResults, err = queries.ListQueryRequirementValues(ctx, &repository.ListQueryRequirementValuesParams{ Queryid: &jsonQueryID, Documentid: &documentID, Version: &jsonVersion, }) require.NoError(t, err) assert.Len(t, qResults, 1) assert.Equal(t, contextQueryID, qResults[0].Queryid) assert.Equal(t, repository.QuerytypeContextFull, qResults[0].Type) assert.Equal(t, "context_value_1", *qResults[0].Value) assert.NotEqual(t, uuid.UUID{}, qResults[0].ID) qResults, err = queries.ListQueryRequirementValues(ctx, &repository.ListQueryRequirementValuesParams{ Queryid: &jsonQueryID, Documentid: &documentTwoID, Version: &jsonVersion, }) require.NoError(t, err) assert.Len(t, qResults, 0) _, err = queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "context_value_2", Textentryid: textId, Queryversion: contextQuery.Activeversion - 1, }) require.NoError(t, err) qResults, err = queries.ListQueryRequirementValues(ctx, &repository.ListQueryRequirementValuesParams{ Queryid: &jsonQueryID, Documentid: &documentID, Version: &jsonVersion, }) require.NoError(t, err) assert.Len(t, qResults, 1) assert.Equal(t, contextQueryID, qResults[0].Queryid) assert.Equal(t, repository.QuerytypeContextFull, qResults[0].Type) assert.Equal(t, "context_value_1", *qResults[0].Value) assert.NotEqual(t, uuid.UUID{}, qResults[0].ID) _, err = queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "context_value_3", Textentryid: textId, Queryversion: contextQuery.Activeversion, }) require.NoError(t, err) qResults, err = queries.ListQueryRequirementValues(ctx, &repository.ListQueryRequirementValuesParams{ Queryid: &jsonQueryID, Documentid: &documentID, Version: &jsonVersion, }) require.NoError(t, err) assert.Len(t, qResults, 1) assert.Equal(t, contextQueryID, qResults[0].Queryid) assert.Equal(t, repository.QuerytypeContextFull, qResults[0].Type) assert.Equal(t, "context_value_3", *qResults[0].Value) assert.NotEqual(t, uuid.UUID{}, qResults[0].ID) _, err = queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: "json_value_1", Textentryid: textId, Queryversion: jsonVersion, }) require.NoError(t, err) qResults, err = queries.ListQueryRequirementValues(ctx, &repository.ListQueryRequirementValuesParams{ Queryid: &jsonQueryID, Documentid: &documentID, Version: &jsonVersion, }) require.NoError(t, err) assert.Len(t, qResults, 1) assert.Equal(t, contextQueryID, qResults[0].Queryid) assert.Equal(t, repository.QuerytypeContextFull, qResults[0].Type) assert.Equal(t, "context_value_3", *qResults[0].Value) assert.NotEqual(t, uuid.UUID{}, qResults[0].ID) } func TestUnsyncedNoDepsQueries(t *testing.T) { t.Parallel() if testing.Short() { t.SkipNow() } ctx := context.Background() cfg := &serviceconfig.BaseConfig{} net := test.GetNetwork(t, ctx) _ = test.CreateDB(t, ctx, cfg, net, &test.CreateDatabaseConfig{}) queries := cfg.GetDBQueries() clientId := "EXAMPLE" err := queries.CreateClient(ctx, &repository.CreateClientParams{ Name: "example_client", ID: clientId, }) require.NoError(t, err) version, err := queries.AddLatestCollectorVersion(ctx, clientId) require.NoError(t, err) err = queries.SetActiveCollectorVersion(ctx, &repository.SetActiveCollectorVersionParams{ Versionid: version, Clientid: clientId, }) require.NoError(t, err) documentID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Clientid: clientId, Hash: "example_hash", }) require.NoError(t, err) documentTwoID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Clientid: clientId, Hash: "example_hash_two", }) require.NoError(t, err) contextQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeContextFull)) require.NoError(t, err) jsonQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor)) require.NoError(t, err) _, err = queries.AddLatestQueryVersion(ctx, jsonQueryID) require.NoError(t, err) err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{ Queryid: jsonQueryID, Versionid: 1, }) require.NoError(t, err) _, err = queries.AddLatestQueryVersion(ctx, contextQueryID) require.NoError(t, err) err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{ Queryid: contextQueryID, Versionid: 1, }) require.NoError(t, err) err = queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{ Queryid: jsonQueryID, Requiredqueryid: contextQueryID, Addedversion: 1, }) require.NoError(t, err) err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Clientid: clientId, Name: "example_name", Queryid: jsonQueryID, Addedversion: 1, }) require.NoError(t, err) qs, err := queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentID) require.NoError(t, err) assert.Len(t, qs, 0) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentTwoID) require.NoError(t, err) assert.Len(t, qs, 0) bucket := "example_bucket" key := "example_key" hash := "hahs" 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) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentID) require.NoError(t, err) assert.Len(t, qs, 0) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentTwoID) require.NoError(t, err) assert.Len(t, qs, 0) textId, err := queries.AddDocumentText(ctx, &repository.AddDocumentTextParams{ Cleanid: cleanid, Bucket: "hi", Key: "hello", Part: 0, Createdat: pgtype.Timestamp{ Time: time.Now().UTC(), Valid: true, }, }) require.NoError(t, err) err = queries.AddDocumentTextEntry(ctx, &repository.AddDocumentTextEntryParams{ Version: 1, Textid: textId, }) require.NoError(t, err) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentID) require.NoError(t, err) assert.Len(t, qs, 1) assert.ElementsMatch(t, []*uuid.UUID{&contextQueryID}, qs) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentTwoID) require.NoError(t, err) assert.Len(t, qs, 0) _, err = queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "context_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentID) require.NoError(t, err) assert.Len(t, qs, 1) assert.ElementsMatch(t, []*uuid.UUID{&jsonQueryID}, qs) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentTwoID) require.NoError(t, err) assert.Len(t, qs, 0) _, err = queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: "context_value", Textentryid: textId, Queryversion: 1, }) require.NoError(t, err) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentID) require.NoError(t, err) assert.Len(t, qs, 0) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentTwoID) require.NoError(t, err) assert.Len(t, qs, 0) cleantwoid, err := queries.AddDocumentClean(ctx, &repository.AddDocumentCleanParams{ Documentid: documentTwoID, 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: cleantwoid, Version: 1, }) require.NoError(t, err) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentID) require.NoError(t, err) assert.Len(t, qs, 0) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentTwoID) require.NoError(t, err) assert.Len(t, qs, 0) textTwoId, err := queries.AddDocumentText(ctx, &repository.AddDocumentTextParams{ Cleanid: cleantwoid, Bucket: "hi", Key: "hello", Part: 0, Createdat: pgtype.Timestamp{ Time: time.Now().UTC(), Valid: true, }, }) require.NoError(t, err) err = queries.AddDocumentTextEntry(ctx, &repository.AddDocumentTextEntryParams{ Version: 1, Textid: textTwoId, }) require.NoError(t, err) _, err = queries.AddResult(ctx, &repository.AddResultParams{ Queryid: contextQueryID, Value: "context_value", Textentryid: textTwoId, Queryversion: 1, }) require.NoError(t, err) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentID) require.NoError(t, err) assert.Len(t, qs, 0) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentTwoID) require.NoError(t, err) assert.Len(t, qs, 1) assert.ElementsMatch(t, []*uuid.UUID{&jsonQueryID}, qs) _, err = queries.AddResult(ctx, &repository.AddResultParams{ Queryid: jsonQueryID, Value: "context_value", Textentryid: textTwoId, Queryversion: 1, }) require.NoError(t, err) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentID) require.NoError(t, err) assert.Len(t, qs, 0) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentTwoID) require.NoError(t, err) assert.Len(t, qs, 0) _, err = queries.AddLatestQueryVersion(ctx, jsonQueryID) require.NoError(t, err) err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{ Queryid: jsonQueryID, Versionid: 2, }) require.NoError(t, err) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentID) require.NoError(t, err) assert.Len(t, qs, 1) assert.ElementsMatch(t, []*uuid.UUID{&jsonQueryID}, qs) qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, &documentTwoID) require.NoError(t, err) assert.Len(t, qs, 1) assert.ElementsMatch(t, []*uuid.UUID{&jsonQueryID}, qs) }