package queryversionsyncrunner_test import ( "regexp" "testing" queryversionsyncrunner "queryorchestration/api/queryVersionSyncRunner" "queryorchestration/internal/database/repository" queryversionsync "queryorchestration/internal/query/versionsync" "queryorchestration/internal/server/runner" "queryorchestration/internal/serviceconfig/objectstore" "queryorchestration/internal/serviceconfig/queue" "queryorchestration/internal/serviceconfig/queue/clientsync" "queryorchestration/internal/test" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) type QuerySyncConfig struct { runner.BaseConfig[queryversionsyncrunner.Body] clientsync.ClientSyncConfig queue.QueueConfig objectstore.ObjectStoreConfig } func TestQueryVersionSyncRunner(t *testing.T) { cfg := &QuerySyncConfig{} test.CreateDB(t, cfg) acfg := test.CreateAWSContainer(t, cfg) test.SetQueueClient(t, t.Context(), cfg, acfg.ExternalEndpoint) cfg.ClientSyncURL = test.CreateQueue(t, cfg, test.ClientSyncRunnerName) svc := queryversionsync.New(cfg) runner := queryversionsyncrunner.New(&queryversionsyncrunner.Services{ Sync: svc, }) queryId, err := cfg.GetDBQueries().CreateQuery(t.Context(), repository.QuerytypeContextFull) require.NoError(t, err) version, err := cfg.GetDBQueries().AddLatestQueryVersion(t.Context(), queryId) require.NoError(t, err) err = cfg.GetDBQueries().AddActiveQueryVersion(t.Context(), &repository.AddActiveQueryVersionParams{ Versionid: version, Queryid: queryId, }) require.NoError(t, err) err = cfg.GetDBQueries().CreateClient(t.Context(), &repository.CreateClientParams{ Clientid: "client_one", Name: "name_one", }) require.NoError(t, err) collectorOneVersion, err := cfg.GetDBQueries().AddLatestCollectorVersion(t.Context(), "client_one") require.NoError(t, err) err = cfg.GetDBQueries().SetActiveCollectorVersion(t.Context(), &repository.SetActiveCollectorVersionParams{ Versionid: version, Clientid: "client_one", }) require.NoError(t, err) err = cfg.GetDBQueries().AddCollectorQuery(t.Context(), &repository.AddCollectorQueryParams{ Queryid: queryId, Clientid: "client_one", Name: "EXAMPLE_ONE", Addedversion: collectorOneVersion, }) require.NoError(t, err) err = cfg.GetDBQueries().CreateClient(t.Context(), &repository.CreateClientParams{ Clientid: "client_two", Name: "name_two", }) require.NoError(t, err) collectorTwoVersion, err := cfg.GetDBQueries().AddLatestCollectorVersion(t.Context(), "client_two") require.NoError(t, err) err = cfg.GetDBQueries().SetActiveCollectorVersion(t.Context(), &repository.SetActiveCollectorVersionParams{ Versionid: version, Clientid: "client_two", }) require.NoError(t, err) err = cfg.GetDBQueries().AddCollectorQuery(t.Context(), &repository.AddCollectorQueryParams{ Queryid: queryId, Clientid: "client_two", Name: "EXAMPLE_TWO", Addedversion: collectorTwoVersion, }) require.NoError(t, err) doc := queryversionsyncrunner.Body{ QueryID: queryId, } assert.True(t, runner.Process(t.Context(), doc)) test.AssertMessageBodies(t, cfg, cfg.GetClientSyncURL(), []*regexp.Regexp{ regexp.MustCompile(`^{"id":"client_one"}$`), regexp.MustCompile(`^{"id":"client_two"}$`), }) }