diff --git a/.github/workflows/linux_inmemory.yml b/.github/workflows/linux_inmemory.yml index 7a875af..ab829db 100644 --- a/.github/workflows/linux_inmemory.yml +++ b/.github/workflows/linux_inmemory.yml @@ -65,7 +65,7 @@ jobs: run: | cd tests mkdir ./coverage-ci - go test -timeout 20m -v -race -cover -tags=debug -failfast -coverpkg=github.com/roadrunner-server/memory/v6/... -coverprofile=./coverage-ci/memory_kv.out -covermode=atomic kv_memory_test.go + go test -timeout 20m -v -race -cover -tags=debug -coverpkg=github.com/roadrunner-server/memory/v6/... -coverprofile=./coverage-ci/memory_kv.out -covermode=atomic -run 'TestInMemory|TestSetMany' ./... - name: Archive code coverage results uses: actions/upload-artifact@v7 @@ -100,6 +100,12 @@ jobs: } ' summary.txt > summary.filtered.txt mv summary.filtered.txt summary.txt + # a profile that maps to no plugin source uploads fine and reports 0% + blocks=$(($(wc -l < summary.txt) - 1)) + if [ "$blocks" -lt 10 ]; then + echo "::error::coverage summary holds $blocks blocks, the profile does not map to plugin sources" + exit 1 + fi - name: upload to codecov uses: codecov/codecov-action@v7 # Docs: diff --git a/.github/workflows/linux_jobs.yml b/.github/workflows/linux_jobs.yml index 5e6b90c..880ca75 100644 --- a/.github/workflows/linux_jobs.yml +++ b/.github/workflows/linux_jobs.yml @@ -65,7 +65,7 @@ jobs: run: | cd tests mkdir ./coverage-ci - go test -timeout 20m -v -race -cover -tags=debug -failfast -coverpkg=github.com/roadrunner-server/memory/v6/... -coverprofile=./coverage-ci/memory_jobs.out -covermode=atomic jobs_memory_test.go + go test -timeout 20m -v -race -cover -tags=debug -coverpkg=github.com/roadrunner-server/memory/v6/... -coverprofile=./coverage-ci/memory_jobs.out -covermode=atomic -run 'TestBoots|TestPushAndProcess|TestPauseResume|TestStatsReport|TestPrefetch|TestProtocolError|TestResponseHandler' ./... - name: Archive code coverage results uses: actions/upload-artifact@v7 @@ -100,6 +100,12 @@ jobs: } ' summary.txt > summary.filtered.txt mv summary.filtered.txt summary.txt + # a profile that maps to no plugin source uploads fine and reports 0% + blocks=$(($(wc -l < summary.txt) - 1)) + if [ "$blocks" -lt 10 ]; then + echo "::error::coverage summary holds $blocks blocks, the profile does not map to plugin sources" + exit 1 + fi - name: upload to codecov uses: codecov/codecov-action@v7 # Docs: diff --git a/tests/helpers/rr.go b/tests/helpers/rr.go new file mode 100644 index 0000000..af20d3f --- /dev/null +++ b/tests/helpers/rr.go @@ -0,0 +1,240 @@ +package helpers + +import ( + "context" + "log/slog" + "net" + "sync" + "testing" + "time" + + mocklogger "tests/mock" + + jobState "github.com/roadrunner-server/api-plugins/v6/jobs" + + "github.com/roadrunner-server/config/v6" + "github.com/roadrunner-server/endure/v2" + "github.com/roadrunner-server/logger/v6" + "github.com/stretchr/testify/require" +) + +const ( + // defaultConfigVersion is the config schema version used by the test configs. + defaultConfigVersion = "2024.2.0" + // probeTimeout caps how long Start waits for the rpc listener to answer. + probeTimeout = time.Second * 30 + probeTick = time.Millisecond * 20 + probeDial = time.Second + // logTimeout bounds WaitLog. Jobs move through the pipeline asynchronously, + // so the log record a test is after can lag the rpc call that caused it. + logTimeout = time.Second * 30 + logTick = time.Millisecond * 50 + // statsTimeout bounds WaitStats; a delayed job only becomes ready when its + // delay lapses, so this has to outlast the longest delay a test uses. + statsTimeout = time.Second * 30 + statsTick = time.Millisecond * 100 +) + +// bootCfg holds the options applied to a container before it is started. +type bootCfg struct { + version string + logLevel slog.Level + logger loggerKind + probe func(ctx context.Context) bool +} + +// loggerKind selects which logger plugin Start registers. +type loggerKind int + +const ( + realLogger loggerKind = iota + observedLogger +) + +// Option customizes the container built by Start. +type Option func(*bootCfg) + +// WithConfigVersion overrides the config schema version. +func WithConfigVersion(v string) Option { + return func(b *bootCfg) { b.version = v } +} + +// WithLogLevel sets the endure container log level (debug by default). +func WithLogLevel(l slog.Level) Option { + return func(b *bootCfg) { b.logLevel = l } +} + +// WithObservedLogger registers an in-memory logger instead of the real logger +// plugin and exposes the captured records as RR.Logs. Most jobs assertions are +// on those records, so this is the common case here. +func WithObservedLogger() Option { + return func(b *bootCfg) { b.logger = observedLogger } +} + +// WithTCPProbe makes Start return only once addr accepts a connection. The rpc +// listener binds after the driver is constructed, so dialing it proves the +// pipeline is ready to take calls. +func WithTCPProbe(addr string) Option { + return func(b *bootCfg) { + b.probe = func(ctx context.Context) bool { + d := net.Dialer{Timeout: probeDial} + conn, err := d.DialContext(ctx, "tcp", addr) + if err != nil { + return false + } + + _ = conn.Close() + return true + } + } +} + +// RR is a running container. +type RR struct { + // Logs holds the captured log records, non-nil only with WithObservedLogger. + Logs *mocklogger.ObservedLogs +} + +// WaitLog blocks until the observed log holds at least want records matching +// snippet. Polling replaces the fixed sleeps the jobs suites used between an +// rpc call and the assertion on its effect. +func (rr *RR) WaitLog(t *testing.T, snippet string, want int) { + t.Helper() + + require.Eventually(t, func() bool { + return rr.Logs.FilterMessageSnippet(snippet).Len() >= want + }, logTimeout, logTick, "expected at least %d records matching %q, saw %d", + want, snippet, rr.Logs.FilterMessageSnippet(snippet).Len()) +} + +// RequireLogCount asserts the exact number of records matching snippet, after +// waiting for them to arrive. An exact count catches a driver that redelivers. +func (rr *RR) RequireLogCount(t *testing.T, snippet string, want int) { + t.Helper() + + rr.WaitLog(t, snippet, want) + require.Equal(t, want, rr.Logs.FilterMessageSnippet(snippet).Len(), "records matching %q", snippet) +} + +// Start registers the plugins, boots the container and waits for the probe, if +// any, to answer. Errors arriving on the container channel are reported through +// t.Errorf and stop the container, but they do not abort the test. +// +// The returned stop is idempotent and also registered with t.Cleanup, so tests +// asserting on records written during shutdown can stop the container mid-test. +func Start(t *testing.T, cfgPath string, plugins []any, opts ...Option) (*RR, func()) { + t.Helper() + + cont, rr, bc := newContainer(t, cfgPath, plugins, opts) + require.NoError(t, cont.Init()) + + ch, err := cont.Serve() + require.NoError(t, err) + + stopCont := sync.OnceValue(cont.Stop) + done := make(chan struct{}) + wg := &sync.WaitGroup{} + + wg.Go(func() { + for { + select { + case res := <-ch: + if res == nil { + return + } + t.Errorf("plugin %s reported an error: %v", res.VertexID, res.Error) + if errS := stopCont(); errS != nil { + t.Errorf("container stop: %v", errS) + } + case <-done: + if errS := stopCont(); errS != nil { + t.Errorf("container stop: %v", errS) + } + return + } + } + }) + + // The drain goroutine calls t.Errorf, so it has to be joined while the test + // is still running. + stop := sync.OnceFunc(func() { + close(done) + wg.Wait() + }) + t.Cleanup(stop) + + if bc.probe != nil { + require.Eventually(t, func() bool { return bc.probe(t.Context()) }, probeTimeout, probeTick, "rpc listener did not become ready") + } + + return rr, stop +} + +// StartExpectServeError registers the plugins, requires Init to pass and Serve +// to fail, and returns the Serve error. +func StartExpectServeError(t *testing.T, cfgPath string, plugins []any, opts ...Option) error { + t.Helper() + + cont, _, _ := newContainer(t, cfgPath, plugins, opts) + require.NoError(t, cont.Init()) + + _, err := cont.Serve() + require.Error(t, err) + t.Cleanup(func() { _ = cont.Stop() }) + + return err +} + +// newContainer builds the container and registers the config, a logger and the +// caller's plugins. The container is not initialized yet. +func newContainer(t *testing.T, cfgPath string, plugins []any, opts []Option) (*endure.Endure, *RR, *bootCfg) { + t.Helper() + + bc := &bootCfg{version: defaultConfigVersion, logLevel: slog.LevelDebug} + for _, o := range opts { + o(bc) + } + + rr := &RR{} + all := make([]any, 0, 2+len(plugins)) + all = append(all, &config.Plugin{Version: bc.version, Path: cfgPath}) + + switch bc.logger { + case realLogger: + all = append(all, &logger.Plugin{}) + case observedLogger: + l, obs := mocklogger.SlogTestLogger(slog.LevelDebug) + rr.Logs = obs + all = append(all, l) + } + + cont := endure.New(bc.logLevel) + require.NoError(t, cont.RegisterAll(append(all, plugins...)...)) + + return cont, rr, bc +} + +// FetchStats returns the current pipeline state. +func FetchStats(t *testing.T, address string) *jobState.State { + t.Helper() + + state := &jobState.State{} + Stats(address, state)(t) + + return state +} + +// WaitStats polls the pipeline state until want holds, then returns it. It +// replaces the fixed sleeps the suites used to wait for delayed jobs to become +// ready or for the queue to drain. +func WaitStats(t *testing.T, address string, want func(*jobState.State) bool) *jobState.State { + t.Helper() + + var state *jobState.State + require.Eventually(t, func() bool { + state = FetchStats(t, address) + return want(state) + }, statsTimeout, statsTick, "pipeline never reached the expected state, last: %+v", state) + + return state +} diff --git a/tests/jobs_memory_test.go b/tests/jobs_memory_test.go deleted file mode 100644 index 1f624a9..0000000 --- a/tests/jobs_memory_test.go +++ /dev/null @@ -1,1023 +0,0 @@ -package memory - -import ( - "log/slog" - "maps" - "os" - "os/signal" - "slices" - "sync" - "syscall" - "testing" - "time" - - "tests/helpers" - mocklogger "tests/mock" - - jobsProto "github.com/roadrunner-server/api-go/v6/jobs/v2" - jobState "github.com/roadrunner-server/api-plugins/v6/jobs" - "github.com/roadrunner-server/config/v6" - "github.com/roadrunner-server/endure/v2" - "github.com/roadrunner-server/informer/v6" - "github.com/roadrunner-server/jobs/v6" - "github.com/roadrunner-server/memory/v6" - "github.com/roadrunner-server/resetter/v6" - rpcPlugin "github.com/roadrunner-server/rpc/v6" - "github.com/roadrunner-server/server/v6" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - sdktrace "go.opentelemetry.io/otel/sdk/trace" - "go.opentelemetry.io/otel/sdk/trace/tracetest" - _ "google.golang.org/genproto/protobuf/ptype" //nolint:revive,nolintlint -) - -type inMemoryTracer struct { - tp *sdktrace.TracerProvider - exp *tracetest.InMemoryExporter -} - -func newInMemoryTracer(t *testing.T) *inMemoryTracer { - t.Helper() - exp := tracetest.NewInMemoryExporter() - tp := sdktrace.NewTracerProvider(sdktrace.WithSyncer(exp)) - t.Cleanup(func() { _ = tp.Shutdown(t.Context()) }) - return &inMemoryTracer{tp: tp, exp: exp} -} - -func (m *inMemoryTracer) Init() error { return nil } -func (m *inMemoryTracer) Name() string { return "inMemoryTracer" } -func (m *inMemoryTracer) Tracer() *sdktrace.TracerProvider { return m.tp } - -func TestMemoryInit(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Version: "2024.1.0", - Path: "configs/.rr-memory-init.yaml", - } - - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - l, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second) - out := &jobState.State{} - t.Run("Stats", helpers.Stats("127.0.0.1:6001", out)) - - assert.Equal(t, out.Active, int64(0)) - assert.Equal(t, out.Delayed, int64(0)) - assert.Equal(t, out.Reserved, int64(0)) - assert.Equal(t, uint64(13), out.Priority) - - helpers.DestroyPipelines("127.0.0.1:6001", "test-1", "test-2") - - stopCh <- struct{}{} - wg.Wait() - - require.Equal(t, 1, oLogger.FilterMessageSnippet("plugin was started").Len()) -} - -func TestMemoryPQ(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Version: "2024.1.0", - Path: "configs/.rr-memory-pq.yaml", - } - - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - l, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second) - - for range 100 { - t.Run("PushPipeline-1", helpers.PushToPipe("test-1-pq", false, "127.0.0.1:6601")) - t.Run("PushPipeline-2", helpers.PushToPipe("test-2-pq", false, "127.0.0.1:6601")) - } - - time.Sleep(time.Second) - - helpers.DestroyPipelines("127.0.0.1:6601", "test-1-pq", "test-2-pq") - - stopCh <- struct{}{} - wg.Wait() - assert.Equal(t, 0, oLogger.FilterMessageSnippet("job was processed successfully").Len()) - assert.Equal(t, 2, oLogger.FilterMessageSnippet("pipeline was started").Len()) - assert.Equal(t, 2, oLogger.FilterMessageSnippet("pipeline was stopped").Len()) - assert.Equal(t, 200, oLogger.FilterMessageSnippet("job was pushed successfully").Len()) - // the exact number of started jobs depends on how many nack/re-dispatch cycles - // fit into the pipeline-destroy window, which varies between machines - assert.GreaterOrEqual(t, oLogger.FilterMessageSnippet("job processing was started").Len(), 2) -} - -func TestMemoryInitV27(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Path: "configs/.rr-memory-init-v27.yaml", - Version: "2024.1.0", - } - - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - l, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second * 1) - t.Run("PushPipeline", helpers.PushToPipe("test-1", false, "127.0.0.1:6001")) - t.Run("PushPipeline", helpers.PushToPipe("test-2", false, "127.0.0.1:6001")) - time.Sleep(time.Second * 1) - - helpers.DestroyPipelines("127.0.0.1:6001", "test-1", "test-2") - - stopCh <- struct{}{} - wg.Wait() - - require.Equal(t, 2, oLogger.FilterMessageSnippet("pipeline was started").Len()) - require.Equal(t, 2, oLogger.FilterMessageSnippet("pipeline was stopped").Len()) - require.Equal(t, 2, oLogger.FilterMessageSnippet("job was pushed successfully").Len()) - require.Equal(t, 2, oLogger.FilterMessageSnippet("job processing was started").Len()) -} - -func TestMemoryInitV27BadResp(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Version: "2024.1.0", - Path: "configs/.rr-memory-init-v27-br.yaml", - } - - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - l, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second * 1) - t.Run("PushPipeline", helpers.PushToPipe("test-1", false, "127.0.0.1:6001")) - t.Run("PushPipeline", helpers.PushToPipe("test-2", false, "127.0.0.1:6001")) - time.Sleep(time.Second * 1) - - helpers.DestroyPipelines("127.0.0.1:6001", "test-1", "test-2") - - stopCh <- struct{}{} - wg.Wait() - - require.Equal(t, 2, oLogger.FilterMessageSnippet("response handler error").Len()) -} - -func TestMemoryCreate(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Version: "2024.1.0", - Path: "configs/.rr-memory-create.yaml", - } - - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - l, - // &logger.Plugin{}, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second * 5) - - helpers.DestroyPipelines("127.0.0.1:6001", "local", "example") - - stopCh <- struct{}{} - wg.Wait() - - assert.Equal(t, 2, oLogger.FilterMessageSnippet("job processing was started").Len()) - assert.Equal(t, 2, oLogger.FilterMessageSnippet("message pushed to the priority queue").Len()) - assert.Equal(t, 2, oLogger.FilterMessageSnippet("job was processed successfully").Len()) - assert.Equal(t, 2, oLogger.FilterMessageSnippet("pipeline was resumed").Len()) -} - -func TestMemoryDeclare(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Version: "2024.1.0", - Path: "configs/.rr-memory-declare.yaml", - } - - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - l, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second * 3) - - t.Run("DeclarePipeline", declareMemoryPipe("10000")) - t.Run("ConsumePipeline", consumeMemoryPipe([]string{"test-3"})) - t.Run("PushPipeline", helpers.PushToPipe("test-3", false, "127.0.0.1:6001")) - time.Sleep(time.Second) - t.Run("PausePipeline", helpers.PausePipelines("127.0.0.1:6001", "test-3")) - time.Sleep(time.Second) - t.Run("DestroyPipeline", helpers.DestroyPipelines("127.0.0.1:6001", "test-3")) - - stopCh <- struct{}{} - wg.Wait() - - require.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was resumed").Len()) - require.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was stopped").Len()) - require.Equal(t, 1, oLogger.FilterMessageSnippet("job was pushed successfully").Len()) - require.Equal(t, 1, oLogger.FilterMessageSnippet("job processing was started").Len()) - require.Equal(t, 1, oLogger.FilterMessageSnippet("job was processed successfully").Len()) -} - -func TestMemoryPauseResume(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Version: "2024.1.0", - Path: "configs/.rr-memory-pause-resume.yaml", - } - - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - l, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second * 3) - - t.Run("Pause", helpers.PausePipelines("127.0.0.1:6001", "test-local")) - t.Run("pushToDisabledPipe", helpers.PushToDisabledPipe("127.0.0.1:6001", "test-local")) - t.Run("Resume", helpers.ResumePipes("127.0.0.1:6001", "test-local")) - t.Run("pushToEnabledPipe", helpers.PushToPipe("test-local", false, "127.0.0.1:6001")) - time.Sleep(time.Second * 1) - - helpers.DestroyPipelines("127.0.0.1:6001", "test-local") - - stopCh <- struct{}{} - wg.Wait() - - require.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was resumed").Len()) - require.Equal(t, 3, oLogger.FilterMessageSnippet("pipeline was stopped").Len()) - require.Equal(t, 2, oLogger.FilterMessageSnippet("job was pushed successfully").Len()) - require.Equal(t, 2, oLogger.FilterMessageSnippet("job processing was started").Len()) - require.Equal(t, 2, oLogger.FilterMessageSnippet("job was processed successfully").Len()) -} - -func TestMemoryJobsError(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Version: "2024.1.0", - Path: "configs/.rr-memory-jobs-err.yaml", - } - - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - l, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second * 3) - - t.Run("DeclarePipeline", declareMemoryPipe("10000")) - t.Run("ConsumePipeline", helpers.ResumePipes("127.0.0.1:6001", "test-3")) - t.Run("PushPipeline", helpers.PushToPipe("test-3", false, "127.0.0.1:6001")) - time.Sleep(time.Second * 25) - t.Run("PausePipeline", helpers.PausePipelines("127.0.0.1:6001", "test-3")) - time.Sleep(time.Second) - t.Run("DestroyPipeline", helpers.DestroyPipelines("127.0.0.1:6001", "test-3")) - - stopCh <- struct{}{} - wg.Wait() - - assert.Equal(t, 1, oLogger.FilterMessageSnippet("job was pushed successfully").Len()) - assert.Equal(t, 4, oLogger.FilterMessageSnippet("job processing was started").Len()) - assert.Equal(t, 1, oLogger.FilterMessageSnippet("job was processed successfully").Len()) - assert.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was paused").Len()) - assert.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was resumed").Len()) - assert.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was stopped").Len()) - assert.Equal(t, 3, oLogger.FilterMessageSnippet("jobs protocol error").Len()) -} - -func TestMemoryStats(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Version: "2024.1.0", - Path: "configs/.rr-memory-declare.yaml", - } - - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - l, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second * 3) - - t.Run("DeclarePipeline", declareMemoryPipe("10000")) - t.Run("ConsumePipeline", consumeMemoryPipe([]string{"test-3"})) - t.Run("PushPipeline", helpers.PushToPipe("test-3", false, "127.0.0.1:6001")) - time.Sleep(time.Second) - t.Run("PausePipeline", helpers.PausePipelines("127.0.0.1:6001", "test-3")) - time.Sleep(time.Second) - - t.Run("PushPipeline", helpers.PushToPipeDelayed("127.0.0.1:6001", "test-3", 5)) - t.Run("PushPipeline", helpers.PushToPipe("test-3", false, "127.0.0.1:6001")) - - time.Sleep(time.Second) - out := &jobState.State{} - t.Run("Stats", helpers.Stats("127.0.0.1:6001", out)) - - assert.Equal(t, "test-3", out.Pipeline) - assert.Equal(t, "memory", out.Driver) - assert.Equal(t, "test-3", out.Queue) - - assert.Equal(t, int64(0), out.Active) - assert.Equal(t, int64(1), out.Delayed) - assert.Equal(t, int64(0), out.Reserved) - assert.Equal(t, uint64(33), out.Priority) - - time.Sleep(time.Second) - t.Run("ConsumePipeline", consumeMemoryPipe([]string{"test-3"})) - time.Sleep(time.Second * 7) - - out = &jobState.State{} - t.Run("Stats", helpers.Stats("127.0.0.1:6001", out)) - - assert.Equal(t, out.Pipeline, "test-3") - assert.Equal(t, out.Driver, "memory") - assert.Equal(t, out.Queue, "test-3") - - assert.Equal(t, int64(0), out.Active) - assert.Equal(t, int64(0), out.Delayed) - assert.Equal(t, int64(0), out.Reserved) - assert.Equal(t, uint64(33), out.Priority) - - t.Run("DestroyPipeline", helpers.DestroyPipelines("127.0.0.1:6001", "test-3")) - - stopCh <- struct{}{} - wg.Wait() - - require.Equal(t, 3, oLogger.FilterMessageSnippet("job was pushed successfully").Len()) - require.Equal(t, 3, oLogger.FilterMessageSnippet("job processing was started").Len()) - require.Equal(t, 3, oLogger.FilterMessageSnippet("job was processed successfully").Len()) - require.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was paused").Len()) - require.Equal(t, 2, oLogger.FilterMessageSnippet("pipeline was resumed").Len()) - require.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was stopped").Len()) -} - -func TestMemoryPrefetch(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Version: "2024.1.1", - Path: "configs/.rr-memory-prefetch.yaml", - } - - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - l, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second * 3) - - t.Run("DeclarePipeline", declareMemoryPipe("1")) - t.Run("ConsumePipeline", consumeMemoryPipe([]string{"test-3"})) - for range 10 { - t.Run("PushPipeline", helpers.PushToPipe("test-3", false, "127.0.0.1:6001")) - } - - time.Sleep(time.Second * 15) - - t.Run("DestroyPipeline", helpers.DestroyPipelines("127.0.0.1:6001", "test-3")) - - stopCh <- struct{}{} - wg.Wait() - - assert.GreaterOrEqual(t, oLogger.FilterMessageSnippet("prefetch limit was reached, waiting for the jobs to be processed").Len(), 1) - assert.Equal(t, 10, oLogger.FilterMessageSnippet("job was pushed successfully").Len()) - assert.Equal(t, 10, oLogger.FilterMessageSnippet("job processing was started").Len()) - assert.Equal(t, 10, oLogger.FilterMessageSnippet("job was processed successfully").Len()) - assert.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was resumed").Len()) - assert.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was stopped").Len()) - assert.Equal(t, 1, oLogger.FilterMessageSnippet("destroy signal received").Len()) -} - -func TestMemoryTracer(t *testing.T) { - cont := endure.New(slog.LevelDebug) - - cfg := &config.Plugin{ - Version: "2023.1.0", - Path: "configs/.rr-memory-tracer.yaml", - } - - tracer := newInMemoryTracer(t) - l, oLogger := mocklogger.SlogTestLogger(slog.LevelDebug) - err := cont.RegisterAll( - cfg, - &server.Plugin{}, - &rpcPlugin.Plugin{}, - tracer, - l, - &jobs.Plugin{}, - &resetter.Plugin{}, - &informer.Plugin{}, - &memory.Plugin{}, - ) - assert.NoError(t, err) - - err = cont.Init() - if err != nil { - t.Fatal(err) - } - - ch, err := cont.Serve() - if err != nil { - t.Fatal(err) - } - - sig := make(chan os.Signal, 1) - signal.Notify(sig, os.Interrupt, syscall.SIGINT, syscall.SIGTERM) - - wg := &sync.WaitGroup{} - stopCh := make(chan struct{}, 1) - - wg.Go(func() { - for { - select { - case e := <-ch: - assert.Fail(t, "error", e.Error.Error()) - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - case <-sig: - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - case <-stopCh: - // timeout - err = cont.Stop() - if err != nil { - assert.FailNow(t, "error", err.Error()) - } - return - } - } - }) - - time.Sleep(time.Second) - t.Run("PushPipeline", helpers.PushToPipe("test-1", false, "127.0.0.1:6001")) - time.Sleep(time.Second * 2) - - t.Run("DestroyPipeline", helpers.DestroyPipelines("127.0.0.1:6001", "test-1")) - - stopCh <- struct{}{} - wg.Wait() - - spans := tracer.exp.GetSpans() - spanNames := make(map[string]struct{}, len(spans)) - for _, s := range spans { - spanNames[s.Name] = struct{}{} - } - - uniqueNames := slices.Sorted(maps.Keys(spanNames)) - - expected := []string{ - "destroy_pipeline", - "in_memory_listener", - "in_memory_push", - "in_memory_stop", - "jobs_listener", - "push", - } - - assert.Equal(t, expected, uniqueNames) - - require.Equal(t, 1, oLogger.FilterMessageSnippet("plugin was started").Len()) - require.Equal(t, 1, oLogger.FilterMessageSnippet("job was pushed successfully").Len()) - require.Equal(t, 1, oLogger.FilterMessageSnippet("job processing was started").Len()) - require.Equal(t, 1, oLogger.FilterMessageSnippet("pipeline was stopped").Len()) -} - -func declareMemoryPipe(prefetch string) func(t *testing.T) { - return func(t *testing.T) { - client := helpers.NewJobsClient(t, "127.0.0.1:6001") - req := &jobsProto.DeclareRequest{Pipeline: map[string]string{ - "driver": "memory", - "name": "test-3", - "prefetch": prefetch, - "priority": "33", - }} - err := client.Call("jobs.Declare", req, &jobsProto.JobsHandlerResponse{}) - assert.NoError(t, err) - } -} - -func consumeMemoryPipe(pipelines []string) func(t *testing.T) { - return func(t *testing.T) { - client := helpers.NewJobsClient(t, "127.0.0.1:6001") - err := client.Call("jobs.Resume", &jobsProto.Pipelines{Pipelines: slices.Clone(pipelines)}, &jobsProto.JobsHandlerResponse{}) - assert.NoError(t, err) - } -} diff --git a/tests/jobs_test.go b/tests/jobs_test.go new file mode 100644 index 0000000..2c3ed48 --- /dev/null +++ b/tests/jobs_test.go @@ -0,0 +1,220 @@ +package memory + +import ( + "testing" + + "tests/helpers" + + jobsProto "github.com/roadrunner-server/api-go/v6/jobs/v2" + jobState "github.com/roadrunner-server/api-plugins/v6/jobs" + "github.com/roadrunner-server/informer/v6" + "github.com/roadrunner-server/jobs/v6" + memoryPlugin "github.com/roadrunner-server/memory/v6" + "github.com/roadrunner-server/resetter/v6" + rpcPlugin "github.com/roadrunner-server/rpc/v6" + "github.com/roadrunner-server/server/v6" + "github.com/stretchr/testify/require" +) + +const ( + rpcAddr = "127.0.0.1:6001" + pipeline = "test-3" +) + +func jobsPlugins() []any { + return []any{ + &server.Plugin{}, + &rpcPlugin.Plugin{}, + &jobs.Plugin{}, + &resetter.Plugin{}, + &informer.Plugin{}, + &memoryPlugin.Plugin{}, + } +} + +// bootJobs starts the container with the observed logger and waits for the rpc +// listener, which is the readiness signal the fixed sleeps used to stand in for. +func bootJobs(t *testing.T, cfgPath string) (*helpers.RR, func()) { + t.Helper() + + return helpers.Start(t, cfgPath, jobsPlugins(), + helpers.WithObservedLogger(), + helpers.WithTCPProbe(rpcAddr), + ) +} + +// declarePipe declares the memory pipeline the tests push to. +func declarePipe(t *testing.T, prefetch string) { + t.Helper() + + client := helpers.NewJobsClient(t, rpcAddr) + req := &jobsProto.DeclareRequest{Pipeline: map[string]string{ + "driver": "memory", + "name": pipeline, + "prefetch": prefetch, + "priority": "33", + }} + + require.NoError(t, client.Call("jobs.Declare", req, &jobsProto.JobsHandlerResponse{})) +} + +// consumePipe resumes consumption on the declared pipeline. +func consumePipe(t *testing.T) { + t.Helper() + + client := helpers.NewJobsClient(t, rpcAddr) + require.NoError(t, client.Call("jobs.Resume", + &jobsProto.Pipelines{Pipelines: []string{pipeline}}, + &jobsProto.JobsHandlerResponse{})) +} + +// TestBoots covers the plain init config. +func TestBoots(t *testing.T) { + rr, _ := bootJobs(t, "configs/.rr-memory-init.yaml") + + rr.WaitLog(t, "plugin was started", 1) +} + +// TestPushAndProcess declares a pipeline, pushes one job and follows it through +// to completion, waiting on the records rather than sleeping between steps. +func TestPushAndProcess(t *testing.T) { + rr, _ := bootJobs(t, "configs/.rr-memory-declare.yaml") + + declarePipe(t, "10000") + consumePipe(t) + + helpers.PushToPipe(pipeline, false, rpcAddr)(t) + + rr.WaitLog(t, "job was pushed successfully", 1) + rr.WaitLog(t, "job processing was started", 1) + rr.WaitLog(t, "job was processed successfully", 1) + + helpers.PausePipelines(rpcAddr, pipeline)(t) + rr.WaitLog(t, "pipeline was paused", 1) + + helpers.DestroyPipelines(rpcAddr, pipeline)(t) + + rr.RequireLogCount(t, "job was pushed successfully", 1) + rr.RequireLogCount(t, "job was processed successfully", 1) +} + +// TestPauseResume pauses a pipeline the config consumes at startup, checks a +// push to it is rejected while paused, then resumes and pushes again. +func TestPauseResume(t *testing.T) { + const pipe = "test-local" + + rr, _ := bootJobs(t, "configs/.rr-memory-pause-resume.yaml") + + helpers.PausePipelines(rpcAddr, pipe)(t) + rr.WaitLog(t, "pipeline was paused", 1) + + helpers.PushToDisabledPipe(rpcAddr, pipe)(t) + + helpers.ResumePipes(rpcAddr, pipe)(t) + rr.WaitLog(t, "pipeline was resumed", 1) + + helpers.PushToPipe(pipe, false, rpcAddr)(t) + rr.WaitLog(t, "job was processed successfully", 1) + + helpers.DestroyPipelines(rpcAddr, pipe)(t) + + rr.RequireLogCount(t, "pipeline was resumed", 1) +} + +// TestStatsReportDelayedAndDrained pushes a delayed job and a plain one, then +// polls the pipeline state instead of sleeping out the delay. +func TestStatsReportDelayedAndDrained(t *testing.T) { + rr, _ := bootJobs(t, "configs/.rr-memory-declare.yaml") + + declarePipe(t, "10000") + consumePipe(t) + + helpers.PushToPipe(pipeline, false, rpcAddr)(t) + rr.WaitLog(t, "job was processed successfully", 1) + + helpers.PausePipelines(rpcAddr, pipeline)(t) + rr.WaitLog(t, "pipeline was paused", 1) + + // with consumption paused, a delayed job stays counted as delayed + helpers.PushToPipeDelayed(rpcAddr, pipeline, 2)(t) + helpers.PushToPipe(pipeline, false, rpcAddr)(t) + + delayed := helpers.WaitStats(t, rpcAddr, func(s *jobState.State) bool { + return s.Delayed == 1 + }) + + require.Equal(t, pipeline, delayed.Pipeline) + require.Equal(t, "memory", delayed.Driver) + require.Equal(t, pipeline, delayed.Queue) + require.Equal(t, uint64(33), delayed.Priority) + + // resuming drains both the queued and the delayed job once its delay lapses + consumePipe(t) + + drained := helpers.WaitStats(t, rpcAddr, func(s *jobState.State) bool { + return s.Delayed == 0 && s.Active == 0 && s.Reserved == 0 + }) + + require.Equal(t, pipeline, drained.Pipeline) + require.Equal(t, uint64(33), drained.Priority) + + helpers.DestroyPipelines(rpcAddr, pipeline)(t) +} + +// TestPrefetchLimit declares a pipeline with prefetch 1 and pushes ten jobs, so +// the driver has to hold jobs back until the in-flight one finishes. The old +// test waited out a flat 15s; this waits for the tenth job to be processed. +func TestPrefetchLimit(t *testing.T) { + const jobCount = 10 + + rr, stop := bootJobs(t, "configs/.rr-memory-prefetch.yaml") + + declarePipe(t, "1") + consumePipe(t) + + for range jobCount { + helpers.PushToPipe(pipeline, false, rpcAddr)(t) + } + + rr.WaitLog(t, "job was processed successfully", jobCount) + rr.WaitLog(t, "prefetch limit was reached, waiting for the jobs to be processed", 1) + + helpers.DestroyPipelines(rpcAddr, pipeline)(t) + + rr.RequireLogCount(t, "job was pushed successfully", jobCount) + rr.RequireLogCount(t, "job was processed successfully", jobCount) + + // the destroy record is written while the container shuts down + stop() + rr.WaitLog(t, "destroy signal received", 1) +} + +// TestProtocolErrorIsReported covers a worker that answers with something the +// jobs protocol cannot parse. The old test slept 25s waiting for the error. +func TestProtocolErrorIsReported(t *testing.T) { + rr, _ := bootJobs(t, "configs/.rr-memory-jobs-err.yaml") + + declarePipe(t, "10000") + helpers.ResumePipes(rpcAddr, pipeline)(t) + helpers.PushToPipe(pipeline, false, rpcAddr)(t) + + rr.WaitLog(t, "jobs protocol error", 1) + + helpers.PausePipelines(rpcAddr, pipeline)(t) + helpers.DestroyPipelines(rpcAddr, pipeline)(t) +} + +// TestResponseHandlerError pushes to two pipelines whose worker answers with a +// payload the response handler cannot parse, so each push produces one error. +func TestResponseHandlerError(t *testing.T) { + rr, _ := bootJobs(t, "configs/.rr-memory-init-v27-br.yaml") + + helpers.PushToPipe("test-1", false, rpcAddr)(t) + helpers.PushToPipe("test-2", false, rpcAddr)(t) + + rr.WaitLog(t, "response handler error", 2) + + helpers.DestroyPipelines(rpcAddr, "test-1", "test-2")(t) + + rr.RequireLogCount(t, "response handler error", 2) +}