Files
query-orchestration/internal/database/repository/sync_test.go
T
Michael McGuinness ee776d2681 Merged in feature/mockserver (pull request #135)
Single Mock Server

* mockserver

* mockserver

* reqs

* mockserver

* slowrunner

* someoptimisedqueries

* passedfullsuite

* passedfullsuite
2025-05-06 01:59:52 +00:00

1037 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.SkipNow()
}
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
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",
})
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)
}
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, &params.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, ctx context.Context, queries *repository.Queries) string {
clientId := "EXAMPLE"
err := queries.CreateClient(ctx, &repository.CreateClientParams{
Name: "example_client",
Clientid: 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.SkipNow()
}
t.Run("no collector", func(t *testing.T) {
t.Parallel()
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
clientId := "EXAMPLE"
err := queries.CreateClient(ctx, &repository.CreateClientParams{
Name: "example_client",
Clientid: clientId,
})
require.NoError(t, err)
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 := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
_, jsonQueryID := createDependentQueries(t, ctx, queries)
clientId := createClientWithCollector(t, ctx, 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("document fail clean", func(t *testing.T) {
t.Parallel()
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
_, 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)
documentID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{
Clientid: clientId,
Hash: "example_noclean",
})
require.NoError(t, err)
getDocumentSyncState(t, ctx, queries, docSyncStateParams{
clientId: clientId,
documentId: documentID,
hash: "example_noclean",
isSynced: false,
fields: `{"first_key": null}`,
})
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_noclean",
isSynced: true,
fields: `{"first_key": null}`,
})
})
t.Run("single query", func(t *testing.T) {
t.Parallel()
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
queryId := createQuery(t, queries, repository.QuerytypeContextFull)
clientId := createClientWithCollector(t, ctx, queries)
err := queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{
Clientid: clientId,
Queryid: queryId,
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)
getDocumentSyncState(t, ctx, queries, docSyncStateParams{
clientId: clientId,
documentId: documentID,
hash: "example_hash",
isSynced: false,
fields: `{"first_key": null}`,
})
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}`,
})
_, 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("multiple queries", func(t *testing.T) {
t.Parallel()
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries)
clientId := createClientWithCollector(t, ctx, queries)
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)
getDocumentSyncState(t, ctx, queries, docSyncStateParams{
clientId: clientId,
documentId: documentID,
hash: "example_hash",
isSynced: false,
fields: `{"first_key": null}`,
})
err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{
Resultid: jsonResultId,
Requiredresultid: contextResultId,
})
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 := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries)
clientId := createClientWithCollector(t, 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)
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 extry", func(t *testing.T) {
t.Parallel()
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries)
clientId := createClientWithCollector(t, 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)
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: "updated_text_context",
Textentryid: textId,
Queryversion: 1,
})
require.NoError(t, err)
jsonResultID, err := queries.AddResult(ctx, &repository.AddResultParams{
Queryid: jsonQueryID,
Value: "updated_text_json",
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_text_json"}`,
})
})
t.Run("update clean entry", func(t *testing.T) {
t.Parallel()
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries)
clientId := createClientWithCollector(t, 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)
getDocumentSyncState(t, ctx, queries, docSyncStateParams{
clientId: clientId,
documentId: documentID,
hash: "example_hash",
isSynced: false,
fields: `{"first_key": null}`,
})
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)
contextResultId, err := queries.AddResult(ctx, &repository.AddResultParams{
Queryid: contextQueryID,
Value: "update_clean_context",
Textentryid: textThreeId,
Queryversion: 1,
})
require.NoError(t, err)
jsonResultID, err := queries.AddResult(ctx, &repository.AddResultParams{
Queryid: jsonQueryID,
Value: "update_clean_json",
Textentryid: textThreeId,
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": "update_clean_json"}`,
})
})
t.Run("change collector name", func(t *testing.T) {
t.Parallel()
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries)
clientId := createClientWithCollector(t, ctx, queries)
documentID, _, _ := createDocumentWithCollectorAndResults(t, ctx, queries, clientId, contextQueryID, jsonQueryID)
latestCollectorVersion, 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"}`,
})
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"}`,
})
err = queries.RemoveCollectorQuery(ctx, &repository.RemoveCollectorQueryParams{
Clientid: clientId,
Queryid: jsonQueryID,
Removedversion: &latestCollectorVersion,
})
require.NoError(t, err)
getDocumentSyncState(t, ctx, queries, docSyncStateParams{
clientId: clientId,
documentId: documentID,
hash: "example_hash",
isSynced: true,
fields: `{}`,
})
err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{
Clientid: clientId,
Queryid: jsonQueryID,
Addedversion: latestCollectorVersion,
Name: "second_key",
})
require.NoError(t, err)
getDocumentSyncState(t, ctx, queries, docSyncStateParams{
clientId: clientId,
documentId: documentID,
hash: "example_hash",
isSynced: true,
fields: `{"second_key": "json_value"}`,
})
})
t.Run("add existing query", func(t *testing.T) {
t.Parallel()
ctx := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries)
clientId := createClientWithCollector(t, 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)
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 := context.Background()
cfg := &serviceconfig.BaseConfig{}
net := test.GetNetwork(t)
test.CreateDB(t, cfg, net, &test.CreateDatabaseConfig{})
queries := cfg.GetDBQueries()
contextQueryID, jsonQueryID := createDependentQueries(t, ctx, queries)
clientId := createClientWithCollector(t, ctx, queries)
documentID, _, textId := createDocumentWithCollectorAndResults(t, ctx, queries, clientId, contextQueryID, jsonQueryID)
superQueryID := createQuery(t, queries, repository.QuerytypeJsonExtractor)
err := queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{
Queryid: contextQueryID,
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}`,
})
superResultId, err := queries.AddResult(ctx, &repository.AddResultParams{
Queryid: superQueryID,
Value: "super_value",
Textentryid: textId,
Queryversion: 1,
})
require.NoError(t, err)
contextResultId, err := queries.AddResult(ctx, &repository.AddResultParams{
Queryid: contextQueryID,
Value: "context_with_super_value",
Textentryid: textId,
Queryversion: 1,
})
require.NoError(t, err)
jsonResultId, err := queries.AddResult(ctx, &repository.AddResultParams{
Queryid: jsonQueryID,
Value: "json_with_super",
Textentryid: textId,
Queryversion: 1,
})
require.NoError(t, err)
err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{
Resultid: contextResultId,
Requiredresultid: superResultId,
})
require.NoError(t, err)
getDocumentSyncState(t, ctx, queries, docSyncStateParams{
clientId: clientId,
documentId: documentID,
hash: "example_hash",
isSynced: false,
fields: `{"first_key": null}`,
})
err = queries.AddResultDependency(ctx, &repository.AddResultDependencyParams{
Resultid: jsonResultId,
Requiredresultid: contextResultId,
})
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 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.GetNetwork(b)
test.CreateDB(b, 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.GetNetwork(b)
test.CreateDB(b, 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)
}
}