package queryqueue import ( "context" "queryorchestration/internal/collector" "queryorchestration/internal/database" "queryorchestration/internal/database/repository" queryprocessor "queryorchestration/internal/queryProcessor" "queryorchestration/internal/result" "testing" "github.com/google/uuid" "github.com/jackc/pgx/v5/pgtype" "github.com/pashagolub/pgxmock/v3" "github.com/stretchr/testify/assert" ) func TestGetUnsyncedQueries(t *testing.T) { queryOne := &queryprocessor.Query{ ID: uuid.New(), Version: int32(1), } svc := Queue{ collectorQueries: []*queryprocessor.Query{ queryOne, }, results: []*result.Result{ {ID: uuid.New(), QueryID: queryOne.ID, QueryVersion: queryOne.Version}, }, } svc.getUnsyncedQueries() assert.EqualExportedValues(t, []*queryprocessor.Query(nil), svc.unsyncedQueue) svc.results = []*result.Result{} svc.unsyncedQueue = []*queryprocessor.Query{} svc.getUnsyncedQueries() assert.EqualExportedValues(t, []*queryprocessor.Query{queryOne}, svc.unsyncedQueue) } func TestGetCollectorQueries(t *testing.T) { ctx := context.Background() pool, err := pgxmock.NewPool() if err != nil { t.Fatalf("failed to open pgxmock database: %v", err) } queries := repository.New(pool) db := &database.Connection{ Queries: queries, Pool: pool, } svc := Queue{ db: db, collector: &collector.Collector{ ID: uuid.New(), }, } dbCollectorID := database.MustToDBUUID(svc.collector.ID) collectorQueries := []*queryprocessor.Query{ {ID: uuid.New(), Type: queryprocessor.TypeContextFull, RequiredQueryIDs: []uuid.UUID{}, Version: int32(1)}, } rows := pgxmock.NewRows([]string{"collectorId", "queryId", "type", "queryVersion", "requiredIds"}) for _, q := range collectorQueries { dbID := database.MustToDBUUID(q.ID) dbReqIDs := database.MustToDBUUIDArray(q.RequiredQueryIDs) ty, err := queryprocessor.ToDBNullQueryType(q.Type) assert.Nil(t, err) rows = rows. AddRow(dbCollectorID, dbID, ty, pgtype.Int4{Int32: int32(q.Version), Valid: true}, dbReqIDs) } pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorID).WillReturnRows(rows) err = svc.getCollectorQueries(ctx) assert.Nil(t, err) assert.EqualExportedValues(t, collectorQueries, svc.collectorQueries) } func TestIsQuerySynced(t *testing.T) { query := &queryprocessor.Query{ ID: uuid.New(), Version: int32(1), } svc := Queue{ results: []*result.Result{ {QueryID: query.ID, QueryVersion: query.Version}, }, } isSynced := svc.isQuerySynced(query) assert.True(t, isSynced) } func TestIsQuerySyncedNoResult(t *testing.T) { query := &queryprocessor.Query{ ID: uuid.New(), Version: int32(1), } svc := Queue{ results: []*result.Result{}, } isSynced := svc.isQuerySynced(query) assert.False(t, isSynced) } func TestIsQuerySyncedOldResult(t *testing.T) { query := &queryprocessor.Query{ ID: uuid.New(), Version: int32(1), } svc := Queue{ results: []*result.Result{ {QueryID: query.ID, QueryVersion: query.Version - 1}, }, } isSynced := svc.isQuerySynced(query) assert.False(t, isSynced) } func TestIsQuerySyncedNewResult(t *testing.T) { query := &queryprocessor.Query{ ID: uuid.New(), Version: int32(1), } svc := Queue{ results: []*result.Result{ {QueryID: query.ID, QueryVersion: query.Version + 1}, }, } isSynced := svc.isQuerySynced(query) assert.False(t, isSynced) }