package queryrunner_test import ( "fmt" "regexp" "testing" "time" queryrunner "queryorchestration/api/queryRunner" "queryorchestration/internal/database/repository" "queryorchestration/internal/query" "queryorchestration/internal/query/result" resultset "queryorchestration/internal/query/result/set" resultsync "queryorchestration/internal/query/result/sync" "queryorchestration/internal/server/runner" "queryorchestration/internal/serviceconfig/objectstore" "queryorchestration/internal/serviceconfig/queue" queryc "queryorchestration/internal/serviceconfig/queue/query" "queryorchestration/internal/test" "github.com/jackc/pgx/v5/pgtype" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) type QueryConfig struct { runner.BaseConfig[queryrunner.Body] queryc.QueryConfig queue.QueueConfig objectstore.ObjectStoreConfig } func TestQueryRunner(t *testing.T) { cfg := &QueryConfig{} test.CreateDB(t, cfg) acfg := test.CreateAWSContainer(t, cfg) test.SetQueueClient(t, t.Context(), cfg, acfg.ExternalEndpoint) cfg.QueryURL = test.CreateQueue(t, cfg, test.QueryRunnerName) que := query.New(cfg) runner := queryrunner.New(&queryrunner.Services{ ResultSet: resultset.New(cfg, &resultset.Services{ Query: que, Result: result.New(cfg, &result.Services{ Query: que, }), Sync: resultsync.New(cfg), }), }) err := cfg.GetDBQueries().CreateClient(t.Context(), &repository.CreateClientParams{ Clientid: "client_one", Name: "name_one", }) require.NoError(t, err) docId, err := cfg.GetDBQueries().CreateDocument(t.Context(), &repository.CreateDocumentParams{ Clientid: "client_one", Hash: "hash", }) require.NoError(t, err) fill := "fill" cleanId, err := cfg.GetDBQueries().AddDocumentClean(t.Context(), &repository.AddDocumentCleanParams{ Documentid: docId, Bucket: &fill, Key: &fill, Hash: &fill, Mimetype: repository.NullCleanmimetype{ Valid: true, Cleanmimetype: repository.CleanmimetypeApplicationPdf, }, }) require.NoError(t, err) err = cfg.GetDBQueries().AddDocumentCleanEntry(t.Context(), &repository.AddDocumentCleanEntryParams{ Cleanid: cleanId, Version: 1, }) require.NoError(t, err) textId, err := cfg.GetDBQueries().AddDocumentText(t.Context(), &repository.AddDocumentTextParams{ Cleanid: cleanId, Bucket: fill, Key: fill, Hash: fill, Createdat: pgtype.Timestamp{ Time: time.Now(), Valid: true, }, }) require.NoError(t, err) err = cfg.GetDBQueries().AddDocumentTextEntry(t.Context(), &repository.AddDocumentTextEntryParams{ Textid: textId, Version: 1, }) require.NoError(t, err) contextQueryId, err := cfg.GetDBQueries().CreateQuery(t.Context(), repository.QuerytypeContextFull) require.NoError(t, err) version, err := cfg.GetDBQueries().AddLatestQueryVersion(t.Context(), contextQueryId) require.NoError(t, err) err = cfg.GetDBQueries().AddActiveQueryVersion(t.Context(), &repository.AddActiveQueryVersionParams{ Versionid: version, Queryid: contextQueryId, }) require.NoError(t, err) queryId, err := cfg.GetDBQueries().CreateQuery(t.Context(), repository.QuerytypeJsonExtractor) 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().AddRequiredQuery(t.Context(), &repository.AddRequiredQueryParams{ Queryid: queryId, Requiredqueryid: contextQueryId, Addedversion: version, }) 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) doc := queryrunner.Body{ DocumentID: docId, QueryID: contextQueryId, } assert.True(t, runner.Process(t.Context(), doc)) test.AssertMessageBody(t, cfg, cfg.GetQueryURL(), regexp.MustCompile(fmt.Sprintf("{\"document_id\":\"%s\",\"query_id\":\"%s\"}", docId, queryId))) }