Files
query-orchestration/internal/database/repository/sync_test.go
T
Michael McGuinness bbe6f4188e Merged in feature/tests (pull request #117)
Feature/tests

* improvetests

* generation

* simplifiedtesting

* simplifiedtesting

* longfile
2025-04-23 17:51:44 +00:00

916 lines
28 KiB
Go

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{
RunMigrations: true,
})
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{
RunMigrations: true,
})
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("valid 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{
RunMigrations: true,
})
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)
t.Run("document standard results", func(t *testing.T) {
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)
})
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)
t.Run("update upstream query", func(t *testing.T) {
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 entry", func(t *testing.T) {
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: 2,
})
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("updated clean entry", func(t *testing.T) {
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: 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: "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)
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": "update_clean_json"}`),
}, docExternal)
t.Run("update collector name", func(t *testing.T) {
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": "update_clean_json"}`),
}, 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": "update_clean_json"}`),
}, docExternal)
})
t.Run("add existing query", func(t *testing.T) {
err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{
Clientid: clientId,
Queryid: contextQueryID,
Addedversion: latestCollectorVersion,
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(`{"second_key": "update_clean_json", "example_key": "update_clean_context"}`),
}, docExternal)
})
t.Run("add super query", func(t *testing.T) {
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: textThreeId,
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(`{"second_key": null, "example_key": null}`),
}, docExternal)
contextResultId, err := queries.AddResult(ctx, &repository.AddResultParams{
Queryid: contextQueryID,
Value: "context_with_super_value",
Textentryid: textThreeId,
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(`{"second_key": null, "example_key": null}`),
}, docExternal)
jsonResultId, err := queries.AddResult(ctx, &repository.AddResultParams{
Queryid: jsonQueryID,
Value: "json_with_super",
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(`{"second_key": null, "example_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(`{"second_key": null, "example_key": "context_with_super_value"}`),
}, 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(`{"second_key": "json_with_super", "example_key": "context_with_super_value"}`),
}, docExternal)
})
})
})
}