Files
query-orchestration/internal/database/repository/sync_test.go
T
Michael McGuinness 33d68b7e04 Merged in feature/momocks (pull request #150)
Decrease Mocks

* feature/nomocks

* nonet

* assertsaws

* assert
2025-05-23 00:20:01 +00:00

1141 lines
31 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{}
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",
})
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, 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 := context.Background()
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 := context.Background()
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 := context.Background()
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",
})
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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 := context.Background()
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",
})
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 := context.Background()
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 := context.Background()
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)
}
}